mirror of
https://github.com/simstudioai/sim.git
synced 2026-09-24 15:45:35 +08:00
feat(sync): added batch sync for registry (#66)
This commit is contained in:
@@ -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 })
|
||||
}
|
||||
}
|
||||
@@ -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 })
|
||||
}
|
||||
}
|
||||
+3
-3
@@ -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 = () => {
|
||||
|
||||
+48
-33
@@ -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<boolean> {
|
||||
async function syncWorkflowsToServer(payloads: WorkflowSyncPayload[]): Promise<boolean> {
|
||||
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<boolean> {
|
||||
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)
|
||||
|
||||
@@ -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)
|
||||
|
||||
|
||||
Reference in New Issue
Block a user