Merge pull request #55 from simstudioai/feature/dbsync

Feature/dbsync
This commit is contained in:
waleedlatif1
2025-02-16 23:04:17 -08:00
committed by GitHub
11 changed files with 701 additions and 2 deletions
+67
View File
@@ -0,0 +1,67 @@
import { NextResponse } from 'next/server'
import { eq } from 'drizzle-orm'
import { z } from 'zod'
import { getSession } from '@/lib/auth'
import { db } from '@/db'
import { workflow } from '@/db/schema'
// Define the schema for incoming data
const WorkflowSyncSchema = z.object({
id: z.string(),
name: z.string(),
description: z.string().optional(),
state: z.string(), // JSON stringified workflow state
})
export async function POST(request: Request) {
try {
// Get the authenticated user
const session = await getSession()
if (!session?.user?.id) {
return NextResponse.json({ error: 'Unauthorized' }, { status: 401 })
}
// Parse and validate the request body
const body = await request.json()
const { id, name, description, state } = WorkflowSyncSchema.parse(body)
// Get the current timestamp
const now = new Date()
// Upsert the workflow
await db
.insert(workflow)
.values({
id,
userId: session.user.id,
name,
description,
state,
lastSynced: now,
createdAt: now,
updatedAt: now,
})
.onConflictDoUpdate({
target: [workflow.id],
set: {
name,
description,
state,
lastSynced: now,
updatedAt: now,
},
where: eq(workflow.userId, session.user.id), // Only update if the workflow belongs to the user
})
return NextResponse.json({ success: true })
} catch (error) {
console.error('Workflow sync error:', error)
if (error instanceof z.ZodError) {
return NextResponse.json(
{ error: 'Invalid request data', details: error.errors },
{ status: 400 }
)
}
return NextResponse.json({ error: 'Sync failed' }, { status: 500 })
}
}
+4 -1
View File
@@ -19,6 +19,7 @@ import { initializeStateLogger } from '@/stores/workflow/logger'
import { useWorkflowRegistry } from '@/stores/workflow/registry/store'
import { useWorkflowStore } from '@/stores/workflow/store'
import { NotificationList } from '@/app/w/components/notifications/notifications'
import { WorkflowSyncWrapper } from '@/app/w/components/workflows/sync-wrapper'
import { getBlock } from '../../../blocks'
import { ErrorBoundary } from '../components/error-boundary/error-boundary'
import { CustomEdge } from './components/custom-edge/custom-edge'
@@ -378,7 +379,9 @@ export default function Workflow() {
return (
<ReactFlowProvider>
<ErrorBoundary>
<WorkflowContent />
<WorkflowSyncWrapper>
<WorkflowContent />
</WorkflowSyncWrapper>
</ErrorBoundary>
</ReactFlowProvider>
)
@@ -0,0 +1,18 @@
import { ReactNode } from 'react'
import {
useDebouncedWorkflowSync,
usePeriodicWorkflowSync,
useSyncOnUnload,
} from '@/stores/workflow/sync/hooks'
interface WorkflowSyncWrapperProps {
children: ReactNode
}
export function WorkflowSyncWrapper({ children }: WorkflowSyncWrapperProps) {
useDebouncedWorkflowSync()
usePeriodicWorkflowSync()
useSyncOnUnload()
return <>{children}</>
}
+12
View File
@@ -0,0 +1,12 @@
CREATE TABLE "workflow" (
"id" text PRIMARY KEY NOT NULL,
"user_id" text NOT NULL,
"name" text NOT NULL,
"description" text,
"state" text NOT NULL,
"last_synced" timestamp NOT NULL,
"created_at" timestamp NOT NULL,
"updated_at" timestamp NOT NULL
);
--> statement-breakpoint
ALTER TABLE "workflow" ADD CONSTRAINT "workflow_user_id_user_id_fk" FOREIGN KEY ("user_id") REFERENCES "public"."user"("id") ON DELETE cascade ON UPDATE no action;
+378
View File
@@ -0,0 +1,378 @@
{
"id": "a17085d8-94da-40fe-b71c-e143059e3f80",
"prevId": "00000000-0000-0000-0000-000000000000",
"version": "7",
"dialect": "postgresql",
"tables": {
"public.account": {
"name": "account",
"schema": "",
"columns": {
"id": {
"name": "id",
"type": "text",
"primaryKey": true,
"notNull": true
},
"account_id": {
"name": "account_id",
"type": "text",
"primaryKey": false,
"notNull": true
},
"provider_id": {
"name": "provider_id",
"type": "text",
"primaryKey": false,
"notNull": true
},
"user_id": {
"name": "user_id",
"type": "text",
"primaryKey": false,
"notNull": true
},
"access_token": {
"name": "access_token",
"type": "text",
"primaryKey": false,
"notNull": false
},
"refresh_token": {
"name": "refresh_token",
"type": "text",
"primaryKey": false,
"notNull": false
},
"id_token": {
"name": "id_token",
"type": "text",
"primaryKey": false,
"notNull": false
},
"access_token_expires_at": {
"name": "access_token_expires_at",
"type": "timestamp",
"primaryKey": false,
"notNull": false
},
"refresh_token_expires_at": {
"name": "refresh_token_expires_at",
"type": "timestamp",
"primaryKey": false,
"notNull": false
},
"scope": {
"name": "scope",
"type": "text",
"primaryKey": false,
"notNull": false
},
"password": {
"name": "password",
"type": "text",
"primaryKey": false,
"notNull": false
},
"created_at": {
"name": "created_at",
"type": "timestamp",
"primaryKey": false,
"notNull": true
},
"updated_at": {
"name": "updated_at",
"type": "timestamp",
"primaryKey": false,
"notNull": true
}
},
"indexes": {},
"foreignKeys": {
"account_user_id_user_id_fk": {
"name": "account_user_id_user_id_fk",
"tableFrom": "account",
"tableTo": "user",
"columnsFrom": ["user_id"],
"columnsTo": ["id"],
"onDelete": "cascade",
"onUpdate": "no action"
}
},
"compositePrimaryKeys": {},
"uniqueConstraints": {},
"policies": {},
"checkConstraints": {},
"isRLSEnabled": false
},
"public.session": {
"name": "session",
"schema": "",
"columns": {
"id": {
"name": "id",
"type": "text",
"primaryKey": true,
"notNull": true
},
"expires_at": {
"name": "expires_at",
"type": "timestamp",
"primaryKey": false,
"notNull": true
},
"token": {
"name": "token",
"type": "text",
"primaryKey": false,
"notNull": true
},
"created_at": {
"name": "created_at",
"type": "timestamp",
"primaryKey": false,
"notNull": true
},
"updated_at": {
"name": "updated_at",
"type": "timestamp",
"primaryKey": false,
"notNull": true
},
"ip_address": {
"name": "ip_address",
"type": "text",
"primaryKey": false,
"notNull": false
},
"user_agent": {
"name": "user_agent",
"type": "text",
"primaryKey": false,
"notNull": false
},
"user_id": {
"name": "user_id",
"type": "text",
"primaryKey": false,
"notNull": true
}
},
"indexes": {},
"foreignKeys": {
"session_user_id_user_id_fk": {
"name": "session_user_id_user_id_fk",
"tableFrom": "session",
"tableTo": "user",
"columnsFrom": ["user_id"],
"columnsTo": ["id"],
"onDelete": "cascade",
"onUpdate": "no action"
}
},
"compositePrimaryKeys": {},
"uniqueConstraints": {
"session_token_unique": {
"name": "session_token_unique",
"nullsNotDistinct": false,
"columns": ["token"]
}
},
"policies": {},
"checkConstraints": {},
"isRLSEnabled": false
},
"public.user": {
"name": "user",
"schema": "",
"columns": {
"id": {
"name": "id",
"type": "text",
"primaryKey": true,
"notNull": true
},
"name": {
"name": "name",
"type": "text",
"primaryKey": false,
"notNull": true
},
"email": {
"name": "email",
"type": "text",
"primaryKey": false,
"notNull": true
},
"email_verified": {
"name": "email_verified",
"type": "boolean",
"primaryKey": false,
"notNull": true
},
"image": {
"name": "image",
"type": "text",
"primaryKey": false,
"notNull": false
},
"created_at": {
"name": "created_at",
"type": "timestamp",
"primaryKey": false,
"notNull": true
},
"updated_at": {
"name": "updated_at",
"type": "timestamp",
"primaryKey": false,
"notNull": true
}
},
"indexes": {},
"foreignKeys": {},
"compositePrimaryKeys": {},
"uniqueConstraints": {
"user_email_unique": {
"name": "user_email_unique",
"nullsNotDistinct": false,
"columns": ["email"]
}
},
"policies": {},
"checkConstraints": {},
"isRLSEnabled": false
},
"public.verification": {
"name": "verification",
"schema": "",
"columns": {
"id": {
"name": "id",
"type": "text",
"primaryKey": true,
"notNull": true
},
"identifier": {
"name": "identifier",
"type": "text",
"primaryKey": false,
"notNull": true
},
"value": {
"name": "value",
"type": "text",
"primaryKey": false,
"notNull": true
},
"expires_at": {
"name": "expires_at",
"type": "timestamp",
"primaryKey": false,
"notNull": true
},
"created_at": {
"name": "created_at",
"type": "timestamp",
"primaryKey": false,
"notNull": false
},
"updated_at": {
"name": "updated_at",
"type": "timestamp",
"primaryKey": false,
"notNull": false
}
},
"indexes": {},
"foreignKeys": {},
"compositePrimaryKeys": {},
"uniqueConstraints": {},
"policies": {},
"checkConstraints": {},
"isRLSEnabled": false
},
"public.workflow": {
"name": "workflow",
"schema": "",
"columns": {
"id": {
"name": "id",
"type": "text",
"primaryKey": true,
"notNull": true
},
"user_id": {
"name": "user_id",
"type": "text",
"primaryKey": false,
"notNull": true
},
"name": {
"name": "name",
"type": "text",
"primaryKey": false,
"notNull": true
},
"description": {
"name": "description",
"type": "text",
"primaryKey": false,
"notNull": false
},
"state": {
"name": "state",
"type": "text",
"primaryKey": false,
"notNull": true
},
"last_synced": {
"name": "last_synced",
"type": "timestamp",
"primaryKey": false,
"notNull": true
},
"created_at": {
"name": "created_at",
"type": "timestamp",
"primaryKey": false,
"notNull": true
},
"updated_at": {
"name": "updated_at",
"type": "timestamp",
"primaryKey": false,
"notNull": true
}
},
"indexes": {},
"foreignKeys": {
"workflow_user_id_user_id_fk": {
"name": "workflow_user_id_user_id_fk",
"tableFrom": "workflow",
"tableTo": "user",
"columnsFrom": ["user_id"],
"columnsTo": ["id"],
"onDelete": "cascade",
"onUpdate": "no action"
}
},
"compositePrimaryKeys": {},
"uniqueConstraints": {},
"policies": {},
"checkConstraints": {},
"isRLSEnabled": false
}
},
"enums": {},
"schemas": {},
"sequences": {},
"roles": {},
"policies": {},
"views": {},
"_meta": {
"columns": {},
"schemas": {},
"tables": {}
}
}
+7
View File
@@ -8,6 +8,13 @@
"when": 1739697832964,
"tag": "0000_careless_black_knight",
"breakpoints": true
},
{
"idx": 1,
"version": "7",
"when": 1739773751302,
"tag": "0001_foamy_dakota_north",
"breakpoints": true
}
]
}
+13
View File
@@ -49,3 +49,16 @@ export const verification = pgTable('verification', {
createdAt: timestamp('created_at'),
updatedAt: timestamp('updated_at'),
})
export const workflow = pgTable('workflow', {
id: text('id').primaryKey(),
userId: text('user_id')
.notNull()
.references(() => user.id, { onDelete: 'cascade' }),
name: text('name').notNull(),
description: text('description'),
state: text('state').notNull(), // JSON stringified workflow state
lastSynced: timestamp('last_synced').notNull(),
createdAt: timestamp('created_at').notNull(),
updatedAt: timestamp('updated_at').notNull(),
})
+1 -1
View File
@@ -11,5 +11,5 @@ export async function middleware(request: NextRequest) {
// TODO: Add protected routes
export const config = {
matcher: ['/dashboard/:path*'],
matcher: ['/dashboard/:path*', '/w/:path*'],
}
+23
View File
@@ -21,12 +21,14 @@
"@radix-ui/react-switch": "^1.1.2",
"@radix-ui/react-tabs": "^1.1.2",
"@radix-ui/react-tooltip": "^1.1.6",
"@types/lodash.debounce": "^4.0.9",
"better-auth": "^1.1.18",
"class-variance-authority": "^0.7.1",
"clsx": "^2.1.1",
"cmdk": "^1.0.0",
"date-fns": "^4.1.0",
"drizzle-orm": "^0.39.3",
"lodash.debounce": "^4.0.8",
"lucide-react": "^0.469.0",
"next": "15.1.3",
"openai": "^4.83.0",
@@ -4512,6 +4514,21 @@
"pretty-format": "^29.0.0"
}
},
"node_modules/@types/lodash": {
"version": "4.17.15",
"resolved": "https://registry.npmjs.org/@types/lodash/-/lodash-4.17.15.tgz",
"integrity": "sha512-w/P33JFeySuhN6JLkysYUK2gEmy9kHHFN7E8ro0tkfmlDOgxBDzWEZ/J8cWA+fHqFevpswDTFZnDx+R9lbL6xw==",
"license": "MIT"
},
"node_modules/@types/lodash.debounce": {
"version": "4.0.9",
"resolved": "https://registry.npmjs.org/@types/lodash.debounce/-/lodash.debounce-4.0.9.tgz",
"integrity": "sha512-Ma5JcgTREwpLRwMM+XwBR7DaWe96nC38uCBDFKZWbNKD+osjVzdpnUSwBcqCptrp16sSOLBAUb50Car5I0TCsQ==",
"license": "MIT",
"dependencies": {
"@types/lodash": "*"
}
},
"node_modules/@types/node": {
"version": "20.17.11",
"resolved": "https://registry.npmjs.org/@types/node/-/node-20.17.11.tgz",
@@ -8513,6 +8530,12 @@
"dev": true,
"license": "MIT"
},
"node_modules/lodash.debounce": {
"version": "4.0.8",
"resolved": "https://registry.npmjs.org/lodash.debounce/-/lodash.debounce-4.0.8.tgz",
"integrity": "sha512-FT1yDzDYEoYWhnSGnpE/4Kj1fLZkDFyqRb7fNt6FdYOSxlUWAtp42Eh6Wb0rGIv/m9Bgo7x4GhQbm5Ys4SG5ow==",
"license": "MIT"
},
"node_modules/lodash.memoize": {
"version": "4.1.2",
"resolved": "https://registry.npmjs.org/lodash.memoize/-/lodash.memoize-4.1.2.tgz",
+2
View File
@@ -29,12 +29,14 @@
"@radix-ui/react-switch": "^1.1.2",
"@radix-ui/react-tabs": "^1.1.2",
"@radix-ui/react-tooltip": "^1.1.6",
"@types/lodash.debounce": "^4.0.9",
"better-auth": "^1.1.18",
"class-variance-authority": "^0.7.1",
"clsx": "^2.1.1",
"cmdk": "^1.0.0",
"date-fns": "^4.1.0",
"drizzle-orm": "^0.39.3",
"lodash.debounce": "^4.0.8",
"lucide-react": "^0.469.0",
"next": "15.1.3",
"openai": "^4.83.0",
+176
View File
@@ -0,0 +1,176 @@
import { useEffect, useRef } from 'react'
import { useRouter } from 'next/navigation'
import debounce from 'lodash.debounce'
import { useNotificationStore } from '@/stores/notifications/store'
import { useWorkflowRegistry } from '../registry/store'
import { useWorkflowStore } from '../store'
const SYNC_DEBOUNCE_MS = 2000 // 2 seconds
const PERIODIC_SYNC_MS = 30000 // 30 seconds
interface SyncPayload {
id: string
name: string
description?: string
state: string
}
async function syncWorkflowToServer(payload: SyncPayload): Promise<boolean> {
try {
const response = await fetch('/api/workflows/sync', {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify(payload),
})
if (!response.ok) {
if (response.status === 401) {
// Auth error - will be handled by the middleware
window.location.href = '/login'
return false
}
throw new Error(`Sync failed: ${response.statusText}`)
}
return true
} catch (error) {
console.error('Error syncing workflow:', error)
return false
}
}
export function useDebouncedWorkflowSync() {
const router = useRouter()
const { addNotification } = useNotificationStore()
const workflowState = useWorkflowStore((state) => ({
blocks: state.blocks,
edges: state.edges,
loops: state.loops,
lastSaved: state.lastSaved,
}))
const { activeWorkflowId, workflows } = useWorkflowRegistry()
const debouncedSyncRef = useRef<ReturnType<typeof debounce> | null>(null)
useEffect(() => {
if (!activeWorkflowId || !workflows[activeWorkflowId]) return
const syncWorkflow = async () => {
const activeWorkflow = workflows[activeWorkflowId]
const payload: SyncPayload = {
id: activeWorkflowId,
name: activeWorkflow.name,
description: activeWorkflow.description,
state: JSON.stringify(workflowState),
}
const success = await syncWorkflowToServer(payload)
if (!success) {
addNotification(
'error',
'Failed to save workflow changes. Please try again.',
activeWorkflowId
)
}
}
// Create a debounced version of syncWorkflow
if (!debouncedSyncRef.current) {
debouncedSyncRef.current = debounce(syncWorkflow, SYNC_DEBOUNCE_MS)
}
// Call the debounced sync
debouncedSyncRef.current()
// Cleanup
return () => {
debouncedSyncRef.current?.cancel()
}
}, [activeWorkflowId, workflows, workflowState, addNotification])
}
export function usePeriodicWorkflowSync() {
const { addNotification } = useNotificationStore()
const workflowState = useWorkflowStore((state) => ({
blocks: state.blocks,
edges: state.edges,
loops: state.loops,
lastSaved: state.lastSaved,
}))
const { activeWorkflowId, workflows } = useWorkflowRegistry()
useEffect(() => {
if (!activeWorkflowId || !workflows[activeWorkflowId]) return
const syncWorkflow = async () => {
const activeWorkflow = workflows[activeWorkflowId]
const payload: SyncPayload = {
id: activeWorkflowId,
name: activeWorkflow.name,
description: activeWorkflow.description,
state: JSON.stringify(workflowState),
}
const success = await syncWorkflowToServer(payload)
if (!success) {
addNotification(
'error',
'Failed to auto-save workflow changes. Please save manually.',
activeWorkflowId
)
}
}
const intervalId = setInterval(syncWorkflow, PERIODIC_SYNC_MS)
return () => clearInterval(intervalId)
}, [activeWorkflowId, workflows, workflowState, addNotification])
}
export function useSyncOnUnload() {
const { addNotification } = useNotificationStore()
const workflowState = useWorkflowStore((state) => ({
blocks: state.blocks,
edges: state.edges,
loops: state.loops,
lastSaved: state.lastSaved,
}))
const { activeWorkflowId, workflows } = useWorkflowRegistry()
useEffect(() => {
if (!activeWorkflowId || !workflows[activeWorkflowId]) return
const handleBeforeUnload = async (event: BeforeUnloadEvent) => {
const activeWorkflow = workflows[activeWorkflowId]
const payload: SyncPayload = {
id: activeWorkflowId,
name: activeWorkflow.name,
description: activeWorkflow.description,
state: JSON.stringify(workflowState),
}
// Use the keepalive option to try to complete the request even during unload
const response = await fetch('/api/workflows/sync', {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify(payload),
keepalive: true,
})
if (!response.ok) {
addNotification(
'error',
'Failed to save workflow changes before closing.',
activeWorkflowId
)
}
// Show a confirmation dialog
event.preventDefault()
event.returnValue = ''
}
window.addEventListener('beforeunload', handleBeforeUnload)
return () => window.removeEventListener('beforeunload', handleBeforeUnload)
}, [activeWorkflowId, workflows, workflowState, addNotification])
}