Feat/db sync (#94)

* feat(db-sync): added general sync file and implemented environment sync

* improvement(workflows-store): structured workflows store system better and added getter for values across stores

* fix(stores): deleted workflows/types since unused

* improvement(db-sync): added workflow event syncs and debounce

* improvement(db-sync): clean and upgraded db-sync system; environment sync implemented

* improvement(db-sync): added batch sync with registry; init bug needs fixing

* improvement(db-sync): finalized sync system and implemented for workflow

* fix(db-sync): fixed client-side rendering

* improvement(db-sync): created backwards sync system; environment implemented

* improvement(db-sync): added colors to db

* fix(db-sync): color sync with db

* improvement(db-sync): added workflow backwards sync; fixing color bug and race condition

* fix(db-stores): color sync

* feature(db-sync): db-sync complete; need to sync history

* improvement(db-sync): added scheduling

* fix(db-sync): environment sync to db
This commit is contained in:
Emir Karabeg
2025-03-03 19:43:39 -08:00
committed by GitHub
parent 7b5168fdd6
commit 76dbc4a52f
56 changed files with 2448 additions and 667 deletions
+106
View File
@@ -0,0 +1,106 @@
import { NextRequest, NextResponse } from 'next/server'
import { eq } from 'drizzle-orm'
import { z } from 'zod'
import { getSession } from '@/lib/auth'
import { decryptSecret, encryptSecret } from '@/lib/utils'
import { EnvironmentVariable } from '@/stores/settings/environment/types'
import { db } from '@/db'
import { environment } from '@/db/schema'
// Schema for environment variable updates
const EnvVarSchema = z.object({
variables: z.record(z.string()),
})
export async function POST(req: NextRequest) {
try {
const session = await getSession()
if (!session?.user?.id) {
return NextResponse.json({ error: 'Unauthorized' }, { status: 401 })
}
const body = await req.json()
const { variables } = EnvVarSchema.parse(body)
// Encrypt all variables
const encryptedVariables = await Object.entries(variables).reduce(
async (accPromise, [key, value]) => {
const acc = await accPromise
const { encrypted } = await encryptSecret(value)
return { ...acc, [key]: encrypted }
},
Promise.resolve({})
)
// Replace all environment variables for user
await db
.insert(environment)
.values({
id: crypto.randomUUID(),
userId: session.user.id,
variables: encryptedVariables,
updatedAt: new Date(),
})
.onConflictDoUpdate({
target: [environment.userId],
set: {
variables: encryptedVariables,
updatedAt: new Date(),
},
})
return NextResponse.json({ success: true })
} catch (error) {
console.error('Error updating environment variables:', error)
if (error instanceof z.ZodError) {
return NextResponse.json(
{ error: 'Invalid request data', details: error.errors },
{ status: 400 }
)
}
return NextResponse.json({ error: 'Failed to update environment variables' }, { status: 500 })
}
}
export async function GET(request: Request) {
try {
// Get the session directly in the API route
const session = await getSession()
if (!session?.user?.id) {
return NextResponse.json({ error: 'Unauthorized' }, { status: 401 })
}
const userId = session.user.id
const result = await db
.select()
.from(environment)
.where(eq(environment.userId, userId))
.limit(1)
if (!result.length || !result[0].variables) {
return NextResponse.json({ data: {} }, { status: 200 })
}
// Decrypt the variables for client-side use
const encryptedVariables = result[0].variables as Record<string, string>
const decryptedVariables: Record<string, EnvironmentVariable> = {}
// Decrypt each variable
for (const [key, encryptedValue] of Object.entries(encryptedVariables)) {
try {
const { decrypted } = await decryptSecret(encryptedValue)
decryptedVariables[key] = { key, value: decrypted }
} catch (error) {
console.error(`Error decrypting variable ${key}:`, error)
// If decryption fails, provide a placeholder
decryptedVariables[key] = { key, value: '' }
}
}
return NextResponse.json({ data: decryptedVariables }, { status: 200 })
} catch (error: any) {
console.error('Environment fetch error:', error)
return NextResponse.json({ error: error.message }, { status: 500 })
}
}
+77
View File
@@ -0,0 +1,77 @@
import { NextResponse } from 'next/server'
import { eq } from 'drizzle-orm'
import { nanoid } from 'nanoid'
import { z } from 'zod'
import { db } from '@/db'
import { settings } from '@/db/schema'
const SettingsSchema = z.object({
userId: z.string(),
isAutoConnectEnabled: z.boolean().default(true),
})
export async function POST(request: Request) {
try {
const body = await request.json()
const { userId, isAutoConnectEnabled } = SettingsSchema.parse(body)
// Store the settings
await db
.insert(settings)
.values({
id: nanoid(),
userId,
general: { isAutoConnectEnabled },
updatedAt: new Date(),
})
.onConflictDoUpdate({
target: [settings.userId],
set: {
general: { isAutoConnectEnabled },
updatedAt: new Date(),
},
})
return NextResponse.json({ success: true }, { status: 200 })
} catch (error: any) {
console.error('Settings update error:', error)
return NextResponse.json({ error: error.message }, { status: 500 })
}
}
export async function GET(request: Request) {
try {
const { searchParams } = new URL(request.url)
const userId = searchParams.get('userId')
if (!userId) {
return NextResponse.json({ error: 'userId is required' }, { status: 400 })
}
const result = await db.select().from(settings).where(eq(settings.userId, userId)).limit(1)
if (!result.length) {
return NextResponse.json(
{
data: {
isAutoConnectEnabled: true, // Return default values
},
},
{ status: 200 }
)
}
const generalSettings = result[0].general as { isAutoConnectEnabled: boolean }
return NextResponse.json(
{
data: {
isAutoConnectEnabled: generalSettings.isAutoConnectEnabled,
},
},
{ status: 200 }
)
} catch (error: any) {
console.error('Settings fetch error:', error)
return NextResponse.json({ error: error.message }, { status: 500 })
}
}
-83
View File
@@ -1,83 +0,0 @@
import { NextResponse } from 'next/server'
import { eq, sql } 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.record(z.any()),
})
// Define the schema for batch sync
const BatchSyncSchema = z.object({
workflows: z.array(WorkflowSchema),
deletedWorkflowIds: z.array(z.string()).optional(),
})
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, deletedWorkflowIds } = BatchSyncSchema.parse(body)
const now = new Date()
// Process all operations in a single transaction
await db.transaction(async (tx) => {
// Handle deletions first
if (deletedWorkflowIds?.length) {
await tx
.delete(workflow)
.where(
sql`${workflow.id} IN ${deletedWorkflowIds} AND ${workflow.userId} = ${session.user.id}`
)
}
// Handle updates/inserts
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 })
}
}
+138
View File
@@ -0,0 +1,138 @@
import { NextRequest, 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'
// Schema for workflow data
const WorkflowStateSchema = z.object({
blocks: z.record(z.any()),
edges: z.array(z.any()),
loops: z.record(z.any()),
lastSaved: z.number().optional(),
isDeployed: z.boolean().optional(),
deployedAt: z.date().optional(),
})
const WorkflowSchema = z.object({
id: z.string(),
name: z.string(),
description: z.string().optional(),
color: z.string().optional(),
state: WorkflowStateSchema,
})
const SyncPayloadSchema = z.object({
workflows: z.record(z.string(), WorkflowSchema),
})
export async function GET(request: Request) {
try {
// Get the session directly in the API route
const session = await getSession()
if (!session?.user?.id) {
return NextResponse.json({ error: 'Unauthorized' }, { status: 401 })
}
const userId = session.user.id
// Fetch all workflows for the user
const workflows = await db.select().from(workflow).where(eq(workflow.userId, userId))
// Return the workflows
return NextResponse.json({ data: workflows }, { status: 200 })
} catch (error: any) {
console.error('Workflow fetch error:', error)
return NextResponse.json({ error: error.message }, { status: 500 })
}
}
export async function POST(req: NextRequest) {
try {
const session = await getSession()
if (!session?.user?.id) {
return NextResponse.json({ error: 'Unauthorized' }, { status: 401 })
}
const body = await req.json()
const { workflows: clientWorkflows } = SyncPayloadSchema.parse(body)
// Get all workflows for the user from the database
const dbWorkflows = await db.select().from(workflow).where(eq(workflow.userId, session.user.id))
const now = new Date()
const operations: Promise<any>[] = []
// Create a map of DB workflows for easier lookup
const dbWorkflowMap = new Map(dbWorkflows.map((w) => [w.id, w]))
const processedIds = new Set<string>()
// Process client workflows
for (const [id, clientWorkflow] of Object.entries(clientWorkflows)) {
processedIds.add(id)
const dbWorkflow = dbWorkflowMap.get(id)
if (!dbWorkflow) {
// New workflow - create
operations.push(
db.insert(workflow).values({
id: clientWorkflow.id,
userId: session.user.id,
name: clientWorkflow.name,
description: clientWorkflow.description,
color: clientWorkflow.color,
state: clientWorkflow.state,
lastSynced: now,
createdAt: now,
updatedAt: now,
})
)
} else {
// Existing workflow - update if needed
const needsUpdate =
JSON.stringify(dbWorkflow.state) !== JSON.stringify(clientWorkflow.state) ||
dbWorkflow.name !== clientWorkflow.name ||
dbWorkflow.description !== clientWorkflow.description ||
dbWorkflow.color !== clientWorkflow.color
if (needsUpdate) {
operations.push(
db
.update(workflow)
.set({
name: clientWorkflow.name,
description: clientWorkflow.description,
color: clientWorkflow.color,
state: clientWorkflow.state,
lastSynced: now,
updatedAt: now,
})
.where(eq(workflow.id, id))
)
}
}
}
// Handle deletions - workflows in DB but not in client
for (const dbWorkflow of dbWorkflows) {
if (!processedIds.has(dbWorkflow.id)) {
operations.push(db.delete(workflow).where(eq(workflow.id, dbWorkflow.id)))
}
}
// Execute all operations in parallel
await Promise.all(operations)
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: 'Workflow sync failed' }, { status: 500 })
}
}