From 4c7630d72e4f1694617da831531c9c29ad427866 Mon Sep 17 00:00:00 2001 From: Emir Karabeg <78010029+emir-karabeg@users.noreply.github.com> Date: Tue, 18 Feb 2025 19:44:07 -0800 Subject: [PATCH] feat(sync): added batch sync for registry (#66) --- app/api/db/sync/route.ts | 72 +++++++++++++++++++++++++++++ app/api/workflows/sync/route.ts | 67 --------------------------- stores/index.ts | 6 +-- stores/sync-manager.ts | 81 +++++++++++++++++++-------------- stores/workflow/utils.ts | 10 ++++ 5 files changed, 133 insertions(+), 103 deletions(-) create mode 100644 app/api/db/sync/route.ts delete mode 100644 app/api/workflows/sync/route.ts diff --git a/app/api/db/sync/route.ts b/app/api/db/sync/route.ts new file mode 100644 index 0000000000..2bcac86ef8 --- /dev/null +++ b/app/api/db/sync/route.ts @@ -0,0 +1,72 @@ +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 a single workflow +const WorkflowSchema = z.object({ + id: z.string(), + name: z.string(), + description: z.string().optional(), + state: z.string(), // JSON stringified workflow state +}) + +// Define the schema for batch sync +const BatchSyncSchema = z.object({ + workflows: z.array(WorkflowSchema), +}) + +export async function POST(request: Request) { + try { + const session = await getSession() + if (!session?.user?.id) { + return NextResponse.json({ error: 'Unauthorized' }, { status: 401 }) + } + + const body = await request.json() + const { workflows } = BatchSyncSchema.parse(body) + const now = new Date() + + // Process all workflows in a single transaction + await db.transaction(async (tx) => { + for (const workflowData of workflows) { + await tx + .insert(workflow) + .values({ + id: workflowData.id, + userId: session.user.id, + name: workflowData.name, + description: workflowData.description, + state: workflowData.state, + lastSynced: now, + createdAt: now, + updatedAt: now, + }) + .onConflictDoUpdate({ + target: [workflow.id], + set: { + name: workflowData.name, + description: workflowData.description, + state: workflowData.state, + lastSynced: now, + updatedAt: now, + }, + where: eq(workflow.userId, session.user.id), + }) + } + }) + + return NextResponse.json({ success: true }) + } catch (error) { + console.error('Batch sync error:', error) + if (error instanceof z.ZodError) { + return NextResponse.json( + { error: 'Invalid request data', details: error.errors }, + { status: 400 } + ) + } + return NextResponse.json({ error: 'Batch sync failed' }, { status: 500 }) + } +} diff --git a/app/api/workflows/sync/route.ts b/app/api/workflows/sync/route.ts deleted file mode 100644 index 323afca10c..0000000000 --- a/app/api/workflows/sync/route.ts +++ /dev/null @@ -1,67 +0,0 @@ -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/stores/index.ts b/stores/index.ts index 0db589e878..085a707604 100644 --- a/stores/index.ts +++ b/stores/index.ts @@ -9,9 +9,9 @@ import { useWorkflowRegistry } from './workflow/registry/store' import { useWorkflowStore } from './workflow/store' // Initialize sync manager when the store is first imported -// if (typeof window !== 'undefined') { -// initializeSyncManager() -// } +if (typeof window !== 'undefined') { + initializeSyncManager() +} // Reset all application stores to their initial state export const resetAllStores = () => { diff --git a/stores/sync-manager.ts b/stores/sync-manager.ts index 5f44f4563b..59bf5f8659 100644 --- a/stores/sync-manager.ts +++ b/stores/sync-manager.ts @@ -1,20 +1,21 @@ import { useWorkflowRegistry } from './workflow/registry/store' import { useWorkflowStore } from './workflow/store' +import { mergeSubblockState } from './workflow/utils' -interface SyncPayload { +interface WorkflowSyncPayload { id: string name: string - description?: string + description?: string | undefined state: string } -async function syncWorkflowToServer(payload: SyncPayload): Promise { +async function syncWorkflowsToServer(payloads: WorkflowSyncPayload[]): Promise { try { - const response = await fetch('/api/workflows/sync', { + const response = await fetch('/api/db/sync', { method: 'POST', headers: { 'Content-Type': 'application/json' }, - body: JSON.stringify(payload), - keepalive: true, // Ensures request completes even during page unload + body: JSON.stringify({ workflows: payloads }), + keepalive: true, }) if (!response.ok) { @@ -22,13 +23,13 @@ async function syncWorkflowToServer(payload: SyncPayload): Promise { window.location.href = '/login' return false } - throw new Error(`Sync failed: ${response.statusText}`) + throw new Error(`Batch sync failed: ${response.statusText}`) } - console.log('Workflow synced successfully') + console.log('Workflows synced successfully') return true } catch (error) { - console.error('Error syncing workflow:', error) + console.error('Error syncing workflows:', error) return false } } @@ -37,32 +38,46 @@ export function initializeSyncManager() { if (typeof window === 'undefined') return const handleBeforeUnload = async (event: BeforeUnloadEvent) => { - const { activeWorkflowId, workflows } = useWorkflowRegistry.getState() - const workflowState = useWorkflowStore.getState() + const { workflows } = useWorkflowRegistry.getState() - if (!activeWorkflowId || !workflows[activeWorkflowId]) { - return + // Prepare sync payloads for all workflows + const syncPayloads: (WorkflowSyncPayload | null)[] = await Promise.all( + Object.entries(workflows).map(async ([id, metadata]) => { + // Get workflow state from localStorage + const savedState = localStorage.getItem(`workflow-${id}`) + if (!savedState) return null + + const state = JSON.parse(savedState) + // Merge subblock states for all blocks in the workflow + const mergedBlocks = mergeSubblockState(state.blocks) + + return { + id, + name: metadata.name, + description: metadata.description, + state: JSON.stringify({ + blocks: mergedBlocks, + edges: state.edges, + loops: state.loops, + lastSaved: state.lastSaved, + }), + } + }) + ) + + // Filter out null values and sync if there are workflows to sync + const validPayloads = syncPayloads.filter( + (payload): payload is WorkflowSyncPayload => payload !== null + ) + + if (validPayloads.length > 0) { + // Show confirmation dialog + event.preventDefault() + event.returnValue = '' + + // Attempt to sync + await syncWorkflowsToServer(validPayloads) } - - const activeWorkflow = workflows[activeWorkflowId] - const payload: SyncPayload = { - id: activeWorkflowId, - name: activeWorkflow.name, - description: activeWorkflow.description, - state: JSON.stringify({ - blocks: workflowState.blocks, - edges: workflowState.edges, - loops: workflowState.loops, - lastSaved: workflowState.lastSaved, - }), - } - - // Show confirmation dialog - event.preventDefault() - event.returnValue = '' - - // Attempt to sync - await syncWorkflowToServer(payload) } window.addEventListener('beforeunload', handleBeforeUnload) diff --git a/stores/workflow/utils.ts b/stores/workflow/utils.ts index b6005b2a0e..2e61b3a92c 100644 --- a/stores/workflow/utils.ts +++ b/stores/workflow/utils.ts @@ -16,9 +16,19 @@ export function mergeSubblockState( return Object.entries(blocksToProcess).reduce( (acc, [id, block]) => { + // Skip if block is undefined or doesn't have subBlocks + if (!block || !block.subBlocks) { + return acc + } + // Create a deep copy of the block's subBlocks to maintain structure const mergedSubBlocks = Object.entries(block.subBlocks).reduce( (subAcc, [subBlockId, subBlock]) => { + // Skip if subBlock is undefined + if (!subBlock) { + return subAcc + } + // Get the stored value for this subblock const storedValue = useSubBlockStore.getState().getValue(id, subBlockId)