diff --git a/app/api/workflows/sync/route.ts b/app/api/workflows/sync/route.ts new file mode 100644 index 0000000000..323afca10c --- /dev/null +++ b/app/api/workflows/sync/route.ts @@ -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 }) + } +} diff --git a/app/w/[id]/workflow.tsx b/app/w/[id]/workflow.tsx index 224a2534e8..5944c2ad02 100644 --- a/app/w/[id]/workflow.tsx +++ b/app/w/[id]/workflow.tsx @@ -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 ( - + + + ) diff --git a/app/w/components/workflows/sync-wrapper.tsx b/app/w/components/workflows/sync-wrapper.tsx new file mode 100644 index 0000000000..0428922e52 --- /dev/null +++ b/app/w/components/workflows/sync-wrapper.tsx @@ -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} +} diff --git a/db/migrations/0001_foamy_dakota_north.sql b/db/migrations/0001_foamy_dakota_north.sql new file mode 100644 index 0000000000..b374a0a5cc --- /dev/null +++ b/db/migrations/0001_foamy_dakota_north.sql @@ -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; \ No newline at end of file diff --git a/db/migrations/meta/0001_snapshot.json b/db/migrations/meta/0001_snapshot.json new file mode 100644 index 0000000000..ae2e592318 --- /dev/null +++ b/db/migrations/meta/0001_snapshot.json @@ -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": {} + } +} diff --git a/db/migrations/meta/_journal.json b/db/migrations/meta/_journal.json index 019d5f47f3..1971faacb9 100644 --- a/db/migrations/meta/_journal.json +++ b/db/migrations/meta/_journal.json @@ -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 } ] } diff --git a/db/schema.ts b/db/schema.ts index eb8d6ffb2b..731d9e4895 100644 --- a/db/schema.ts +++ b/db/schema.ts @@ -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(), +}) diff --git a/middleware.ts b/middleware.ts index 23153a1606..b1ef96d4d8 100644 --- a/middleware.ts +++ b/middleware.ts @@ -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*'], } diff --git a/package-lock.json b/package-lock.json index df967315cc..0e1a21d1af 100644 --- a/package-lock.json +++ b/package-lock.json @@ -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", diff --git a/package.json b/package.json index 291f34c796..4b3f5ef1a0 100644 --- a/package.json +++ b/package.json @@ -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", diff --git a/stores/workflow/sync/hooks.ts b/stores/workflow/sync/hooks.ts new file mode 100644 index 0000000000..197ea12f72 --- /dev/null +++ b/stores/workflow/sync/hooks.ts @@ -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 { + 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 | 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]) +}