diff --git a/apps/web/src/app/api/workflows/[workflowId]/__tests__/route.test.ts b/apps/web/src/app/api/workflows/[workflowId]/__tests__/route.test.ts index 7b0c5ba945..314336dcb9 100644 --- a/apps/web/src/app/api/workflows/[workflowId]/__tests__/route.test.ts +++ b/apps/web/src/app/api/workflows/[workflowId]/__tests__/route.test.ts @@ -162,6 +162,15 @@ describe('GET /api/workflows/[workflowId]', () => { expect(response!.status).toBe(404); }); + it('should hide non-scheduled workflows', async () => { + mockSelectWhere.mockResolvedValue([{ ...mockWorkflow, triggerType: 'event' as const }]); + + const request = new Request('https://example.com/api/workflows/wf_legacy'); + const response = await GET(request, createContext('wf_legacy')); + + expect(response!.status).toBe(404); + }); + it('should return 403 when user is not owner or admin', async () => { mockSelectWhere.mockResolvedValue([mockWorkflow]); vi.mocked(checkDriveAccess).mockResolvedValue(createAccessFixture({ @@ -245,6 +254,19 @@ describe('PATCH /api/workflows/[workflowId]', () => { expect(response!.status).toBe(404); }); + it('should hide non-scheduled workflows', async () => { + mockSelectWhere.mockResolvedValue([{ ...mockWorkflow, triggerType: 'event' as const }]); + + const request = new Request('https://example.com/api/workflows/wf_legacy', { + method: 'PATCH', + headers: { 'Content-Type': 'application/json' }, + body: JSON.stringify({ name: 'New Name' }), + }); + const response = await PATCH(request, createContext('wf_legacy')); + + expect(response!.status).toBe(404); + }); + it('should return 403 when user is not owner or admin', async () => { mockSelectWhere.mockResolvedValue([mockWorkflow]); vi.mocked(checkDriveAccess).mockResolvedValue(createAccessFixture({ @@ -275,10 +297,7 @@ describe('PATCH /api/workflows/[workflowId]', () => { expect(body.error).toContain('cron expression'); }); - it('should return 400 when setting eventTriggers to null on an event workflow', async () => { - const eventWorkflow = { ...mockWorkflow, triggerType: 'event' as const, eventTriggers: [{ operation: 'update', resourceType: 'page' }] }; - mockSelectWhere.mockResolvedValue([eventWorkflow]); - + it('should reject legacy event workflow fields', async () => { const request = new Request('https://example.com/api/workflows/wf_1', { method: 'PATCH', headers: { 'Content-Type': 'application/json' }, @@ -288,7 +307,7 @@ describe('PATCH /api/workflows/[workflowId]', () => { expect(response!.status).toBe(400); const body = await response!.json(); - expect(body.error).toContain('event trigger'); + expect(body.error).toBe('Invalid input'); }); it('should return updated workflow on success', async () => { @@ -345,6 +364,15 @@ describe('DELETE /api/workflows/[workflowId]', () => { expect(response!.status).toBe(404); }); + it('should hide non-scheduled workflows', async () => { + mockSelectWhere.mockResolvedValue([{ ...mockWorkflow, triggerType: 'event' as const }]); + + const request = new Request('https://example.com/api/workflows/wf_legacy', { method: 'DELETE' }); + const response = await DELETE(request, createContext('wf_legacy')); + + expect(response!.status).toBe(404); + }); + it('should return 403 when user is not owner or admin', async () => { vi.mocked(checkDriveAccess).mockResolvedValue(createAccessFixture({ isMember: true, diff --git a/apps/web/src/app/api/workflows/[workflowId]/route.ts b/apps/web/src/app/api/workflows/[workflowId]/route.ts index 0b1004bc7b..76769c7adf 100644 --- a/apps/web/src/app/api/workflows/[workflowId]/route.ts +++ b/apps/web/src/app/api/workflows/[workflowId]/route.ts @@ -7,25 +7,17 @@ import { validateCronExpression, validateTimezone, getNextRunDate } from '@/lib/ const AUTH_OPTIONS_READ = { allow: ['session'] as const, requireCSRF: false }; const AUTH_OPTIONS_WRITE = { allow: ['session'] as const, requireCSRF: true }; - -const eventTriggerSchema = z.object({ - operation: z.string().min(1), - resourceType: z.string().min(1), -}); +const MANAGEABLE_TRIGGER_TYPE = 'cron' as const; const updateWorkflowSchema = z.object({ name: z.string().min(1).max(200).optional(), agentPageId: z.string().min(1).optional(), prompt: z.string().min(1).optional(), contextPageIds: z.array(z.string()).optional(), - triggerType: z.enum(['cron', 'event']).optional(), cronExpression: z.string().min(1).optional().nullable(), timezone: z.string().optional(), isEnabled: z.boolean().optional(), - eventTriggers: z.array(eventTriggerSchema).optional().nullable(), - watchedFolderIds: z.array(z.string()).optional().nullable(), - eventDebounceSecs: z.number().int().min(5).max(3600).optional().nullable(), -}); +}).strict(); async function getWorkflowWithAuth(workflowId: string, userId: string) { const [workflow] = await db @@ -33,7 +25,9 @@ async function getWorkflowWithAuth(workflowId: string, userId: string) { .from(workflows) .where(eq(workflows.id, workflowId)); - if (!workflow) return { error: NextResponse.json({ error: 'Workflow not found' }, { status: 404 }) }; + if (!workflow || workflow.triggerType !== MANAGEABLE_TRIGGER_TYPE) { + return { error: NextResponse.json({ error: 'Workflow not found' }, { status: 404 }) }; + } const access = await checkDriveAccess(workflow.driveId, userId); if (!access.drive) return { error: NextResponse.json({ error: 'Drive not found' }, { status: 404 }) }; @@ -95,9 +89,6 @@ export async function PATCH( } } - // Determine effective trigger type - const triggerType = data.triggerType ?? workflow.triggerType; - // Validate timezone const effectiveTimezone = data.timezone ?? workflow.timezone; const tzValidation = validateTimezone(effectiveTimezone); @@ -113,31 +104,16 @@ export async function PATCH( } } - // Validate: cron workflows need a cron expression, event workflows need triggers - if (triggerType === 'cron') { - // Resolve effective cronExpression: explicit null from payload means "clear it" - const cronExpr = data.cronExpression !== undefined ? data.cronExpression : workflow.cronExpression; - if (!cronExpr) { - return NextResponse.json({ error: 'Cron workflows require a cron expression' }, { status: 400 }); - } - } - if (triggerType === 'event') { - const triggers = data.eventTriggers !== undefined - ? data.eventTriggers - : (workflow.eventTriggers as Array<{ operation: string; resourceType: string }> | null); - if (!triggers || triggers.length === 0) { - return NextResponse.json({ error: 'Event workflows require at least one event trigger' }, { status: 400 }); - } + // Resolve effective cronExpression: explicit null from payload means "clear it" + const cronExpr = data.cronExpression !== undefined ? data.cronExpression : workflow.cronExpression; + if (!cronExpr) { + return NextResponse.json({ error: 'Cron workflows require a cron expression' }, { status: 400 }); } // Compute nextRunAt based on updated fields (only for cron workflows) const isEnabled = data.isEnabled ?? workflow.isEnabled; - let nextRunAt: Date | null = null; - if (triggerType === 'cron') { - const cronExpression = data.cronExpression ?? workflow.cronExpression; - const timezone = data.timezone ?? workflow.timezone; - nextRunAt = isEnabled && cronExpression ? getNextRunDate(cronExpression, timezone) : null; - } + const timezone = data.timezone ?? workflow.timezone; + const nextRunAt = isEnabled ? getNextRunDate(cronExpr, timezone) : null; const [updated] = await db .update(workflows) diff --git a/apps/web/src/app/api/workflows/[workflowId]/run/__tests__/route.test.ts b/apps/web/src/app/api/workflows/[workflowId]/run/__tests__/route.test.ts index 7419d813e1..4febdd5376 100644 --- a/apps/web/src/app/api/workflows/[workflowId]/run/__tests__/route.test.ts +++ b/apps/web/src/app/api/workflows/[workflowId]/run/__tests__/route.test.ts @@ -222,23 +222,19 @@ describe('POST /api/workflows/[workflowId]/run', () => { expect(executeWorkflow).toHaveBeenCalledWith(mockWorkflow); }); - test('event workflow with stale cronExpression does not compute nextRunAt', async () => { + test('non-scheduled workflow is treated as not found', async () => { const eventWorkflow = { ...mockWorkflow, triggerType: 'event' as const, cronExpression: '0 9 * * 1-5', // stale leftover }; mockSelectWhere.mockResolvedValue([eventWorkflow]); - vi.mocked(executeWorkflow).mockResolvedValue({ - success: true, - responseText: 'Done', - toolCallCount: 0, - durationMs: 100, - }); const request = new Request('https://example.com/api/workflows/wf_1/run', { method: 'POST' }); - await POST(request, createContext('wf_1')); + const response = await POST(request, createContext('wf_1')); + expect(response.status).toBe(404); + expect(executeWorkflow).not.toHaveBeenCalled(); expect(getNextRunDate).not.toHaveBeenCalled(); }); diff --git a/apps/web/src/app/api/workflows/[workflowId]/run/route.ts b/apps/web/src/app/api/workflows/[workflowId]/run/route.ts index baea7e9b5d..94efe8a35f 100644 --- a/apps/web/src/app/api/workflows/[workflowId]/run/route.ts +++ b/apps/web/src/app/api/workflows/[workflowId]/run/route.ts @@ -6,6 +6,7 @@ import { executeWorkflow } from '@/lib/workflows/workflow-executor'; import { getNextRunDate } from '@/lib/workflows/cron-utils'; const AUTH_OPTIONS = { allow: ['session'] as const, requireCSRF: true }; +const MANAGEABLE_TRIGGER_TYPE = 'cron' as const; // POST /api/workflows/[workflowId]/run - Manual trigger export async function POST( @@ -22,7 +23,7 @@ export async function POST( .from(workflows) .where(eq(workflows.id, workflowId)); - if (!workflow) { + if (!workflow || workflow.triggerType !== MANAGEABLE_TRIGGER_TYPE) { return NextResponse.json({ error: 'Workflow not found' }, { status: 404 }); } @@ -59,7 +60,7 @@ export async function POST( } // Update status — only compute nextRunAt for cron workflows - const nextRunAt = (workflow.triggerType === 'cron' && workflow.isEnabled && workflow.cronExpression) + const nextRunAt = (workflow.isEnabled && workflow.cronExpression) ? getNextRunDate(workflow.cronExpression, workflow.timezone) : null; diff --git a/apps/web/src/app/api/workflows/__tests__/route.test.ts b/apps/web/src/app/api/workflows/__tests__/route.test.ts index a3cf95dd1d..6135f7af1c 100644 --- a/apps/web/src/app/api/workflows/__tests__/route.test.ts +++ b/apps/web/src/app/api/workflows/__tests__/route.test.ts @@ -239,6 +239,23 @@ describe('POST /api/workflows', () => { expect(body.error).toBe('Invalid input'); }); + it('should reject legacy event workflow fields', async () => { + const request = new Request('https://example.com/api/workflows', { + method: 'POST', + headers: { 'Content-Type': 'application/json' }, + body: JSON.stringify({ + ...validBody, + triggerType: 'event', + eventTriggers: [{ operation: 'upload', resourceType: 'file' }], + }), + }); + const response = await POST(request); + + expect(response.status).toBe(400); + const body = await response.json(); + expect(body.error).toBe('Invalid input'); + }); + it('should return 404 when drive not found', async () => { vi.mocked(checkDriveAccess).mockResolvedValue(createAccessFixture({ drive: null })); diff --git a/apps/web/src/app/api/workflows/route.ts b/apps/web/src/app/api/workflows/route.ts index ccfa048fc2..c2e9f53a77 100644 --- a/apps/web/src/app/api/workflows/route.ts +++ b/apps/web/src/app/api/workflows/route.ts @@ -6,11 +6,7 @@ import { db, workflows, pages, eq, and } from '@pagespace/db'; import { validateCronExpression, validateTimezone, getNextRunDate } from '@/lib/workflows/cron-utils'; const AUTH_OPTIONS = { allow: ['session'] as const, requireCSRF: true }; - -const eventTriggerSchema = z.object({ - operation: z.string().min(1), - resourceType: z.string().min(1), -}); +const MANAGEABLE_TRIGGER_TYPE = 'cron' as const; const createWorkflowSchema = z.object({ driveId: z.string().min(1), @@ -18,20 +14,12 @@ const createWorkflowSchema = z.object({ agentPageId: z.string().min(1), prompt: z.string().min(1), contextPageIds: z.array(z.string()).default([]), - triggerType: z.enum(['cron', 'event']).default('cron'), - cronExpression: z.string().min(1).optional(), + cronExpression: z.string().min(1), timezone: z.string().default('UTC'), isEnabled: z.boolean().default(true), - eventTriggers: z.array(eventTriggerSchema).optional(), - watchedFolderIds: z.array(z.string()).optional(), - eventDebounceSecs: z.number().int().min(5).max(3600).default(30), -}).refine(data => { - if (data.triggerType === 'cron') return !!data.cronExpression; - if (data.triggerType === 'event') return data.eventTriggers && data.eventTriggers.length > 0; - return true; -}, { message: 'Cron workflows need cronExpression; event workflows need eventTriggers' }); - -// GET /api/workflows?driveId=xxx - List workflows for a drive +}).strict(); + +// GET /api/workflows?driveId=xxx - List scheduled workflows for a drive export async function GET(request: Request) { const auth = await authenticateRequestWithOptions(request, { allow: ['session'] as const, requireCSRF: false }); if (isAuthError(auth)) return auth.error; @@ -55,7 +43,7 @@ export async function GET(request: Request) { const results = await db .select() .from(workflows) - .where(eq(workflows.driveId, driveId)) + .where(and(eq(workflows.driveId, driveId), eq(workflows.triggerType, MANAGEABLE_TRIGGER_TYPE))) .orderBy(workflows.createdAt); auditRequest(request, { eventType: 'data.read', userId, resourceType: 'workflow', resourceId: driveId, details: { count: results.length } }); @@ -63,7 +51,7 @@ export async function GET(request: Request) { return NextResponse.json(results); } -// POST /api/workflows - Create a new workflow +// POST /api/workflows - Create a new scheduled workflow export async function POST(request: Request) { const auth = await authenticateRequestWithOptions(request, AUTH_OPTIONS); if (isAuthError(auth)) return auth.error; @@ -108,13 +96,11 @@ export async function POST(request: Request) { // Validate cron expression for cron-type workflows let nextRunAt: Date | null = null; - if (data.triggerType === 'cron') { - const cronValidation = validateCronExpression(data.cronExpression!); - if (!cronValidation.valid) { - return NextResponse.json({ error: `Invalid cron expression: ${cronValidation.error}` }, { status: 400 }); - } - nextRunAt = data.isEnabled ? getNextRunDate(data.cronExpression!, data.timezone) : null; + const cronValidation = validateCronExpression(data.cronExpression); + if (!cronValidation.valid) { + return NextResponse.json({ error: `Invalid cron expression: ${cronValidation.error}` }, { status: 400 }); } + nextRunAt = data.isEnabled ? getNextRunDate(data.cronExpression, data.timezone) : null; const [workflow] = await db.insert(workflows).values({ driveId: data.driveId, @@ -123,18 +109,18 @@ export async function POST(request: Request) { agentPageId: data.agentPageId, prompt: data.prompt, contextPageIds: data.contextPageIds, - triggerType: data.triggerType, - cronExpression: data.triggerType === 'cron' ? data.cronExpression! : null, + triggerType: MANAGEABLE_TRIGGER_TYPE, + cronExpression: data.cronExpression, timezone: data.timezone, isEnabled: data.isEnabled, - eventTriggers: data.triggerType === 'event' ? data.eventTriggers : null, - watchedFolderIds: data.triggerType === 'event' ? (data.watchedFolderIds ?? null) : null, - eventDebounceSecs: data.triggerType === 'event' ? data.eventDebounceSecs : null, + eventTriggers: null, + watchedFolderIds: null, + eventDebounceSecs: null, nextRunAt, updatedAt: new Date(), }).returning(); - auditRequest(request, { eventType: 'data.write', userId, resourceType: 'workflow', resourceId: workflow.id, details: { driveId: data.driveId, triggerType: data.triggerType } }); + auditRequest(request, { eventType: 'data.write', userId, resourceType: 'workflow', resourceId: workflow.id, details: { driveId: data.driveId, triggerType: MANAGEABLE_TRIGGER_TYPE } }); return NextResponse.json(workflow, { status: 201 }); } diff --git a/apps/web/src/components/workflows/WorkflowForm.tsx b/apps/web/src/components/workflows/WorkflowForm.tsx index ee99c45e1f..79359ba8ce 100644 --- a/apps/web/src/components/workflows/WorkflowForm.tsx +++ b/apps/web/src/components/workflows/WorkflowForm.tsx @@ -22,20 +22,15 @@ import { } from '@/components/ui/select'; import { fetchJSON } from '@/lib/auth/auth-fetch'; import { getHumanReadableCron } from '@/lib/workflows/cron-utils'; -import type { EventTrigger } from './types'; interface WorkflowFormData { name: string; agentPageId: string; prompt: string; contextPageIds: string[]; - triggerType: 'cron' | 'event'; - cronExpression?: string; + cronExpression: string; timezone: string; isEnabled: boolean; - eventTriggers?: EventTrigger[]; - watchedFolderIds?: string[]; - eventDebounceSecs?: number; } interface AgentPage { @@ -61,25 +56,15 @@ const CRON_PRESETS = [ { label: 'Monthly on 1st', value: '0 9 1 * *' }, ]; -const EVENT_TRIGGER_PRESETS: { label: string; description: string; operation: string; resourceType: string }[] = [ - { label: 'Page created', description: 'When a new page is created', operation: 'create', resourceType: 'page' }, - { label: 'File uploaded', description: 'When a file is uploaded', operation: 'upload', resourceType: 'file' }, - { label: 'Page moved', description: 'When a page is moved to a folder', operation: 'move', resourceType: 'page' }, - { label: 'Member added', description: 'When a new member joins the drive', operation: 'member_add', resourceType: 'member' }, -]; - const fetcher = (url: string) => fetchJSON(url); export function WorkflowForm({ open, onOpenChange, driveId, initialData, onSubmit }: WorkflowFormProps) { const [name, setName] = useState(initialData?.name ?? ''); const [agentPageId, setAgentPageId] = useState(initialData?.agentPageId ?? ''); const [prompt, setPrompt] = useState(initialData?.prompt ?? ''); - const [triggerType, setTriggerType] = useState<'cron' | 'event'>(initialData?.triggerType ?? 'cron'); const [cronExpression, setCronExpression] = useState(initialData?.cronExpression ?? '0 9 * * 1-5'); const [timezone, setTimezone] = useState(initialData?.timezone ?? Intl.DateTimeFormat().resolvedOptions().timeZone); const [isEnabled, setIsEnabled] = useState(initialData?.isEnabled ?? true); - const [eventTriggers, setEventTriggers] = useState(initialData?.eventTriggers ?? []); - const [eventDebounceSecs, setEventDebounceSecs] = useState(initialData?.eventDebounceSecs ?? 30); const [cronPreview, setCronPreview] = useState(''); const [isSubmitting, setIsSubmitting] = useState(false); const [error, setError] = useState(''); @@ -107,30 +92,13 @@ export function WorkflowForm({ open, onOpenChange, driveId, initialData, onSubmi setName(initialData?.name ?? ''); setAgentPageId(initialData?.agentPageId ?? ''); setPrompt(initialData?.prompt ?? ''); - setTriggerType(initialData?.triggerType ?? 'cron'); setCronExpression(initialData?.cronExpression ?? '0 9 * * 1-5'); setTimezone(initialData?.timezone ?? Intl.DateTimeFormat().resolvedOptions().timeZone); setIsEnabled(initialData?.isEnabled ?? true); - setEventTriggers(initialData?.eventTriggers ?? []); - setEventDebounceSecs(initialData?.eventDebounceSecs ?? 30); setError(''); } }, [open, initialData]); - const toggleEventTrigger = (operation: string, resourceType: string) => { - setEventTriggers(prev => { - const exists = prev.some(t => t.operation === operation && t.resourceType === resourceType); - if (exists) { - return prev.filter(t => !(t.operation === operation && t.resourceType === resourceType)); - } - return [...prev, { operation, resourceType }]; - }); - }; - - const isEventTriggerSelected = (operation: string, resourceType: string) => { - return eventTriggers.some(t => t.operation === operation && t.resourceType === resourceType); - }; - const handleSubmit = async (e: FormEvent) => { e.preventDefault(); setError(''); @@ -142,12 +110,9 @@ export function WorkflowForm({ open, onOpenChange, driveId, initialData, onSubmi agentPageId, prompt, contextPageIds: [], - triggerType, - cronExpression: triggerType === 'cron' ? cronExpression : undefined, + cronExpression, timezone, isEnabled, - eventTriggers: triggerType === 'event' ? eventTriggers : undefined, - eventDebounceSecs: triggerType === 'event' ? eventDebounceSecs : undefined, }); onOpenChange(false); } catch (err) { @@ -166,9 +131,7 @@ export function WorkflowForm({ open, onOpenChange, driveId, initialData, onSubmi } })(); - const isValid = name && agentPageId && prompt && isTimezoneValid && ( - triggerType === 'cron' ? !!cronExpression : eventTriggers.length > 0 - ); + const isValid = name && agentPageId && prompt && isTimezoneValid && !!cronExpression; return ( @@ -222,112 +185,34 @@ export function WorkflowForm({ open, onOpenChange, driveId, initialData, onSubmi /> - {/* Trigger Type Toggle */}
- -
- - + + setCronExpression(e.target.value)} + placeholder="0 9 * * 1-5" + required + /> + {cronPreview && ( +

{cronPreview}

+ )} +
+ {CRON_PRESETS.map(preset => ( + + ))}
- {/* Cron Schedule Fields */} - {triggerType === 'cron' && ( -
- - setCronExpression(e.target.value)} - placeholder="0 9 * * 1-5" - required - /> - {cronPreview && ( -

{cronPreview}

- )} -
- {CRON_PRESETS.map(preset => ( - - ))} -
-
- )} - - {/* Event Trigger Fields */} - {triggerType === 'event' && ( - <> -
- -
- {EVENT_TRIGGER_PRESETS.map(preset => ( - - ))} -
- {eventTriggers.length === 0 && ( -

Select at least one event trigger

- )} -
- -
- - setEventDebounceSecs(Number(e.target.value))} - /> -

- Wait this long after the first event before running. Coalesces rapid events (e.g., bulk uploads) into a single run. -

-
- - )} -
{workflow.name} - {workflow.triggerType === 'event' ? ( -
- - - {(workflow.eventTriggers ?? []).map(t => `${t.operation}:${t.resourceType}`).join(', ') || 'Event'} - -
- ) : ( -
- - {workflow.cronExpression} -
- )} +
+ + {workflow.cronExpression ?? '-'} +
diff --git a/apps/web/src/components/workflows/WorkflowsDashboard.tsx b/apps/web/src/components/workflows/WorkflowsDashboard.tsx index 236aeef4f2..b0f9fdc34e 100644 --- a/apps/web/src/components/workflows/WorkflowsDashboard.tsx +++ b/apps/web/src/components/workflows/WorkflowsDashboard.tsx @@ -16,43 +16,29 @@ interface WorkflowsDashboardProps { driveName: string; } +interface WorkflowFormData { + name: string; + agentPageId: string; + prompt: string; + contextPageIds: string[]; + cronExpression: string; + timezone: string; + isEnabled: boolean; +} + export function WorkflowsDashboard({ driveId, driveName }: WorkflowsDashboardProps) { const { workflows, isLoading, mutate, runWorkflow, toggleWorkflow, deleteWorkflow } = useWorkflows(driveId); const [formOpen, setFormOpen] = useState(false); const [editingWorkflow, setEditingWorkflow] = useState(null); const [deleteTarget, setDeleteTarget] = useState(null); - const handleCreate = async (data: { - name: string; - agentPageId: string; - prompt: string; - contextPageIds: string[]; - triggerType: 'cron' | 'event'; - cronExpression?: string; - timezone: string; - isEnabled: boolean; - eventTriggers?: { operation: string; resourceType: string }[]; - watchedFolderIds?: string[]; - eventDebounceSecs?: number; - }) => { + const handleCreate = async (data: WorkflowFormData) => { await post('/api/workflows', { ...data, driveId }); mutate(); toast.success('Workflow created'); }; - const handleUpdate = async (data: { - name: string; - agentPageId: string; - prompt: string; - contextPageIds: string[]; - triggerType: 'cron' | 'event'; - cronExpression?: string; - timezone: string; - isEnabled: boolean; - eventTriggers?: { operation: string; resourceType: string }[]; - watchedFolderIds?: string[]; - eventDebounceSecs?: number; - }) => { + const handleUpdate = async (data: WorkflowFormData) => { if (!editingWorkflow) return; await patch(`/api/workflows/${editingWorkflow.id}`, data); mutate(); @@ -123,7 +109,7 @@ export function WorkflowsDashboard({ driveId, driveName }: WorkflowsDashboardPro

- Automate AI agents with scheduled cron jobs or event triggers. Each workflow executes an agent with a prompt — on a schedule or when something happens in your drive. + Automate AI agents with scheduled cron jobs. Event-driven workflows have been removed while folder-centric automation is redesigned.

{isLoading ? ( @@ -150,13 +136,9 @@ export function WorkflowsDashboard({ driveId, driveName }: WorkflowsDashboardPro agentPageId: editingWorkflow.agentPageId, prompt: editingWorkflow.prompt, contextPageIds: editingWorkflow.contextPageIds ?? [], - triggerType: editingWorkflow.triggerType ?? 'cron', - cronExpression: editingWorkflow.cronExpression ?? undefined, + cronExpression: editingWorkflow.cronExpression ?? '0 9 * * 1-5', timezone: editingWorkflow.timezone, isEnabled: editingWorkflow.isEnabled, - eventTriggers: editingWorkflow.eventTriggers ?? undefined, - watchedFolderIds: editingWorkflow.watchedFolderIds ?? undefined, - eventDebounceSecs: editingWorkflow.eventDebounceSecs ?? undefined, } : undefined} onSubmit={editingWorkflow ? handleUpdate : handleCreate} /> diff --git a/apps/web/src/components/workflows/types.ts b/apps/web/src/components/workflows/types.ts index ebb4d30136..6cd01d1da7 100644 --- a/apps/web/src/components/workflows/types.ts +++ b/apps/web/src/components/workflows/types.ts @@ -1,6 +1,3 @@ -import type { EventTrigger } from '@pagespace/db'; -export type { EventTrigger }; - /** JSON-serialized workflow from the API (dates are strings, not Date objects). */ export interface Workflow { id: string; @@ -10,13 +7,10 @@ export interface Workflow { agentPageId: string; prompt: string; contextPageIds: string[]; - triggerType: 'cron' | 'event'; + triggerType: 'cron'; cronExpression: string | null; timezone: string; isEnabled: boolean; - eventTriggers?: EventTrigger[] | null; - watchedFolderIds?: string[] | null; - eventDebounceSecs?: number | null; lastRunAt: string | null; nextRunAt: string | null; lastRunStatus: 'never_run' | 'success' | 'error' | 'running'; diff --git a/apps/web/src/instrumentation.ts b/apps/web/src/instrumentation.ts index 172f0ce1b4..2b852593d6 100644 --- a/apps/web/src/instrumentation.ts +++ b/apps/web/src/instrumentation.ts @@ -14,14 +14,11 @@ export async function register() { console.log('[Instrumentation] Environment validation passed'); // Initialize activity broadcast hook for real-time updates - const { setActivityBroadcastHook, setWorkflowTriggerHook } = await import('@pagespace/lib'); + const { setActivityBroadcastHook } = await import('@pagespace/lib'); const { broadcastActivityEvent } = await import('@/lib/websocket/socket-utils'); - const { emitWorkflowEvent } = await import('@/lib/workflows/event-trigger'); setActivityBroadcastHook(broadcastActivityEvent); - setWorkflowTriggerHook(emitWorkflowEvent); console.log('[Instrumentation] Activity broadcast hook initialized'); - console.log('[Instrumentation] Workflow trigger hook initialized'); } } diff --git a/apps/web/src/lib/workflows/__tests__/event-trigger.test.ts b/apps/web/src/lib/workflows/__tests__/event-trigger.test.ts deleted file mode 100644 index fbce48586c..0000000000 --- a/apps/web/src/lib/workflows/__tests__/event-trigger.test.ts +++ /dev/null @@ -1,495 +0,0 @@ -import { describe, test, expect, beforeEach, afterEach, vi } from 'vitest'; -import type { WorkflowEvent } from '../event-trigger'; - -// ============================================================================ -// Tests for event-trigger.ts -// ============================================================================ - -const { - mockReturning, - mockUpdateWhere, - mockUpdateSet, - mockUpdate, - mockSelectWhere, - mockSelectFrom, - mockSelect, -} = vi.hoisted(() => ({ - mockReturning: vi.fn().mockResolvedValue([]), - mockUpdateWhere: vi.fn(), - mockUpdateSet: vi.fn(), - mockUpdate: vi.fn(), - mockSelectWhere: vi.fn().mockResolvedValue([]), - mockSelectFrom: vi.fn(), - mockSelect: vi.fn(), -})); - -vi.mock('@pagespace/db', () => ({ - db: { - select: mockSelect, - update: mockUpdate, - }, - workflows: { - isEnabled: 'isEnabled', - triggerType: 'triggerType', - driveId: 'driveId', - id: 'id', - lastRunStatus: 'lastRunStatus', - lastRunAt: 'lastRunAt', - }, - pages: { id: 'id', parentId: 'parentId' }, - eq: vi.fn(), - and: vi.fn(), - ne: vi.fn(), -})); - -vi.mock('../workflow-executor', () => ({ - executeWorkflow: vi.fn(), -})); - -vi.mock('@pagespace/lib/server', () => ({ - loggers: { - api: { info: vi.fn(), error: vi.fn(), warn: vi.fn(), debug: vi.fn() }, - }, -})); - -import { emitWorkflowEvent } from '../event-trigger'; -import { executeWorkflow } from '../workflow-executor'; - -// ============================================================================ -// Fixtures -// ============================================================================ - -const createWorkflow = (overrides: Record = {}) => ({ - id: 'wf_1', - driveId: 'drive_abc', - name: 'Test Workflow', - triggerType: 'event' as const, - isEnabled: true, - agentPageId: 'page_1', - prompt: 'Process this event', - contextPageIds: [], - cronExpression: null, - timezone: 'UTC', - eventTriggers: [{ operation: 'create', resourceType: 'page' }], - watchedFolderIds: null, - eventDebounceSecs: 5, - lastRunStatus: 'never_run', - lastRunAt: null, - lastRunError: null, - lastRunDurationMs: null, - nextRunAt: null, - createdBy: 'user_123', - createdAt: new Date('2024-01-01'), - updatedAt: new Date('2024-01-01'), - ...overrides, -}); - -const createEvent = (overrides: Partial = {}): WorkflowEvent => ({ - operation: 'create', - resourceType: 'page', - resourceId: 'res_1', - driveId: 'drive_abc', - pageId: null, - userId: 'user_456', - ...overrides, -}); - -// ============================================================================ -// Setup -// ============================================================================ - -describe('emitWorkflowEvent', () => { - beforeEach(() => { - vi.resetAllMocks(); - vi.useFakeTimers(); - - // Default DB chain: select().from().where() - mockSelect.mockReturnValue({ from: mockSelectFrom }); - mockSelectFrom.mockReturnValue({ where: mockSelectWhere }); - mockSelectWhere.mockResolvedValue([]); - - // Default update chain: update().set().where().returning() - // mockUpdateWhere must serve as both thenable (for post-exec updates) - // and have .returning() (for atomic claim) - mockUpdate.mockReturnValue({ set: mockUpdateSet }); - mockUpdateSet.mockReturnValue({ where: mockUpdateWhere }); - mockUpdateWhere.mockImplementation(() => { - const p = Promise.resolve(undefined) as Promise & { returning: typeof mockReturning }; - p.returning = mockReturning; - return p; - }); - // Default: atomic claim succeeds (returns claimed row) - mockReturning.mockResolvedValue([{ id: 'wf_1' }]); - - vi.mocked(executeWorkflow).mockResolvedValue({ - success: true, - responseText: 'Done', - toolCallCount: 0, - durationMs: 100, - }); - }); - - afterEach(() => { - vi.useRealTimers(); - }); - - // -------------------------------------------------------------------------- - // Early exits - // -------------------------------------------------------------------------- - - test('no driveId', async () => { - const event = createEvent({ driveId: null }); - - await emitWorkflowEvent(event); - - expect(mockSelect).not.toHaveBeenCalled(); - }); - - test('recursive trigger prevention', async () => { - const event = createEvent({ - isAiGenerated: true, - aiConversationId: 'workflow-wf_1-1234', - }); - - await emitWorkflowEvent(event); - - expect(mockSelect).not.toHaveBeenCalled(); - }); - - test('no enabled event workflows in drive', async () => { - mockSelectWhere.mockResolvedValue([]); - - await emitWorkflowEvent(createEvent()); - - expect(executeWorkflow).not.toHaveBeenCalled(); - }); - - // -------------------------------------------------------------------------- - // Event matching - // -------------------------------------------------------------------------- - - test('matching event triggers execution after debounce', async () => { - const workflow = createWorkflow(); - mockSelectWhere.mockResolvedValue([workflow]); - - await emitWorkflowEvent(createEvent()); - - // Before debounce fires - expect(executeWorkflow).not.toHaveBeenCalled(); - - // Fire the debounce timer - await vi.advanceTimersByTimeAsync(5000); - - expect(executeWorkflow).toHaveBeenCalledTimes(1); - // Verify the prompt has event context prepended - const calledArg = vi.mocked(executeWorkflow).mock.calls[0][0]; - expect(calledArg.prompt).toContain(''); - expect(calledArg.prompt).toContain('Event: create on page'); - expect(calledArg.prompt).toContain('Process this event'); - }); - - test('non-matching event skips execution', async () => { - const workflow = createWorkflow({ - eventTriggers: [{ operation: 'delete', resourceType: 'page' }], - }); - mockSelectWhere.mockResolvedValue([workflow]); - - await emitWorkflowEvent(createEvent({ operation: 'create', resourceType: 'page' })); - - await vi.advanceTimersByTimeAsync(10000); - - expect(executeWorkflow).not.toHaveBeenCalled(); - }); - - // -------------------------------------------------------------------------- - // Debounce coalescing - // -------------------------------------------------------------------------- - - test('rapid events coalesce into single execution', async () => { - const workflow = createWorkflow({ eventDebounceSecs: 10 }); - mockSelectWhere.mockResolvedValue([workflow]); - - // Fire 3 events rapidly - await emitWorkflowEvent(createEvent({ resourceId: 'res_1' })); - await emitWorkflowEvent(createEvent({ resourceId: 'res_2' })); - await emitWorkflowEvent(createEvent({ resourceId: 'res_3' })); - - await vi.advanceTimersByTimeAsync(10000); - - // Only one execution despite 3 events - expect(executeWorkflow).toHaveBeenCalledTimes(1); - // Should use the LATEST event context - const calledArg = vi.mocked(executeWorkflow).mock.calls[0][0]; - expect(calledArg.prompt).toContain('res_3'); - }); - - // -------------------------------------------------------------------------- - // Folder scoping - // -------------------------------------------------------------------------- - - test('watched folder matches via event pageId', async () => { - const workflow = createWorkflow({ - watchedFolderIds: ['folder_a'], - }); - mockSelectWhere.mockResolvedValue([workflow]); - - await emitWorkflowEvent(createEvent({ pageId: 'folder_a' })); - - await vi.advanceTimersByTimeAsync(5000); - - expect(executeWorkflow).toHaveBeenCalledTimes(1); - }); - - test('watched folder matches via resource parentId', async () => { - const workflow = createWorkflow({ - watchedFolderIds: ['folder_b'], - }); - - // Call 1: find matching workflows - // Call 2: resource parent lookup - // Call 3: re-validation during debounce (returns the workflow as enabled) - let callCount = 0; - mockSelectWhere.mockImplementation(async () => { - callCount++; - if (callCount === 1) return [workflow]; - if (callCount === 2) return [{ parentId: 'folder_b' }]; - return [workflow]; // re-validation - }); - mockSelectFrom.mockReturnValue({ where: mockSelectWhere }); - mockSelect.mockReturnValue({ from: mockSelectFrom }); - - await emitWorkflowEvent(createEvent({ pageId: 'unrelated_page' })); - - await vi.advanceTimersByTimeAsync(5000); - - expect(executeWorkflow).toHaveBeenCalledTimes(1); - }); - - test('event outside watched folders skips execution', async () => { - const workflow = createWorkflow({ - watchedFolderIds: ['folder_a'], - }); - - let callCount = 0; - mockSelectWhere.mockImplementation(async () => { - callCount++; - if (callCount === 1) return [workflow]; - // Resource is in a different folder - return [{ parentId: 'folder_z' }]; - }); - - await emitWorkflowEvent(createEvent({ pageId: 'other_page' })); - - await vi.advanceTimersByTimeAsync(5000); - - expect(executeWorkflow).not.toHaveBeenCalled(); - }); - - // -------------------------------------------------------------------------- - // Debounce re-validation - // -------------------------------------------------------------------------- - - test('workflow disabled during debounce is not executed', async () => { - const workflow = createWorkflow({ eventDebounceSecs: 5 }); - - let callCount = 0; - mockSelectWhere.mockImplementation(async () => { - callCount++; - // First call: find matching workflows - if (callCount === 1) return [workflow]; - // Re-fetch during debounce: workflow is now disabled - return [{ ...workflow, isEnabled: false }]; - }); - - await emitWorkflowEvent(createEvent()); - - await vi.advanceTimersByTimeAsync(5000); - - expect(executeWorkflow).not.toHaveBeenCalled(); - }); - - test('workflow deleted during debounce is not executed', async () => { - const workflow = createWorkflow({ eventDebounceSecs: 5 }); - - let callCount = 0; - mockSelectWhere.mockImplementation(async () => { - callCount++; - if (callCount === 1) return [workflow]; - // Re-fetch during debounce: workflow deleted - return []; - }); - - await emitWorkflowEvent(createEvent()); - - await vi.advanceTimersByTimeAsync(5000); - - expect(executeWorkflow).not.toHaveBeenCalled(); - }); - - test('workflow edited during debounce executes with fresh prompt', async () => { - const workflow = createWorkflow({ eventDebounceSecs: 5, prompt: 'Old prompt' }); - - let callCount = 0; - mockSelectWhere.mockImplementation(async () => { - callCount++; - if (callCount === 1) return [workflow]; - // Re-fetch during debounce: prompt was edited - return [{ ...workflow, prompt: 'Updated prompt' }]; - }); - - await emitWorkflowEvent(createEvent()); - - await vi.advanceTimersByTimeAsync(5000); - - expect(executeWorkflow).toHaveBeenCalledTimes(1); - const calledArg = vi.mocked(executeWorkflow).mock.calls[0][0]; - expect(calledArg.prompt).toContain('Updated prompt'); - expect(calledArg.prompt).not.toContain('Old prompt'); - }); - - test('workflow trigger type changed to cron during debounce is not executed', async () => { - const workflow = createWorkflow({ eventDebounceSecs: 5 }); - - let callCount = 0; - mockSelectWhere.mockImplementation(async () => { - callCount++; - if (callCount === 1) return [workflow]; - // Re-fetch during debounce: trigger type was switched to cron - return [{ ...workflow, triggerType: 'cron' }]; - }); - - await emitWorkflowEvent(createEvent()); - - await vi.advanceTimersByTimeAsync(5000); - - expect(executeWorkflow).not.toHaveBeenCalled(); - }); - - test('workflow already running during debounce is not executed', async () => { - const workflow = createWorkflow({ eventDebounceSecs: 5 }); - - let callCount = 0; - mockSelectWhere.mockImplementation(async () => { - callCount++; - if (callCount === 1) return [workflow]; - // Re-fetch during debounce: workflow is currently running - return [{ ...workflow, lastRunStatus: 'running' }]; - }); - - await emitWorkflowEvent(createEvent()); - - await vi.advanceTimersByTimeAsync(5000); - - expect(executeWorkflow).not.toHaveBeenCalled(); - }); - - // -------------------------------------------------------------------------- - // Workflow status updates - // -------------------------------------------------------------------------- - - test('marks workflow running then updates status on success', async () => { - const workflow = createWorkflow(); - mockSelectWhere.mockResolvedValue([workflow]); - - await emitWorkflowEvent(createEvent()); - await vi.advanceTimersByTimeAsync(5000); - - // update() called for: mark running + update status - expect(mockUpdate).toHaveBeenCalledTimes(2); - }); - - test('marks workflow as error when execution fails', async () => { - const workflow = createWorkflow(); - mockSelectWhere.mockResolvedValue([workflow]); - vi.mocked(executeWorkflow).mockResolvedValue({ - success: false, - durationMs: 50, - error: 'Agent failed', - }); - - await emitWorkflowEvent(createEvent()); - await vi.advanceTimersByTimeAsync(5000); - - expect(mockUpdateSet).toHaveBeenCalledWith( - expect.objectContaining({ lastRunStatus: 'error', lastRunError: 'Agent failed' }) - ); - }); -}); - -// ============================================================================ -// buildEventContext (tested indirectly via prompt injection) -// ============================================================================ - -describe('buildEventContext', () => { - beforeEach(() => { - vi.resetAllMocks(); - vi.useFakeTimers(); - - mockSelect.mockReturnValue({ from: mockSelectFrom }); - mockSelectFrom.mockReturnValue({ where: mockSelectWhere }); - mockUpdate.mockReturnValue({ set: mockUpdateSet }); - mockUpdateSet.mockReturnValue({ where: mockUpdateWhere }); - mockUpdateWhere.mockImplementation(() => { - const p = Promise.resolve(undefined) as Promise & { returning: typeof mockReturning }; - p.returning = mockReturning; - return p; - }); - mockReturning.mockResolvedValue([{ id: 'wf_1' }]); - - vi.mocked(executeWorkflow).mockResolvedValue({ - success: true, - responseText: 'Done', - toolCallCount: 0, - durationMs: 100, - }); - }); - - afterEach(() => { - vi.useRealTimers(); - }); - - test('includes metadata title and folder in context', async () => { - const workflow = createWorkflow({ eventDebounceSecs: 1 }); - mockSelectWhere.mockResolvedValue([workflow]); - - await emitWorkflowEvent(createEvent({ - metadata: { resourceTitle: 'My Page', folderName: 'Reports' }, - })); - - await vi.advanceTimersByTimeAsync(1000); - - const calledArg = vi.mocked(executeWorkflow).mock.calls[0][0]; - expect(calledArg.prompt).toContain('Name: My Page'); - expect(calledArg.prompt).toContain('Folder: Reports'); - }); - - test('truncates excessively long metadata values', async () => { - const workflow = createWorkflow({ eventDebounceSecs: 1 }); - mockSelectWhere.mockResolvedValue([workflow]); - - const longTitle = 'A'.repeat(300); - await emitWorkflowEvent(createEvent({ - metadata: { resourceTitle: longTitle }, - })); - - await vi.advanceTimersByTimeAsync(1000); - - const calledArg = vi.mocked(executeWorkflow).mock.calls[0][0]; - // Should be truncated to 200 chars + "..." - expect(calledArg.prompt).not.toContain(longTitle); - expect(calledArg.prompt).toContain('A'.repeat(200) + '...'); - }); - - test('wraps context in structured delimiters', async () => { - const workflow = createWorkflow({ eventDebounceSecs: 1 }); - mockSelectWhere.mockResolvedValue([workflow]); - - await emitWorkflowEvent(createEvent()); - - await vi.advanceTimersByTimeAsync(1000); - - const calledArg = vi.mocked(executeWorkflow).mock.calls[0][0]; - expect(calledArg.prompt).toContain(''); - expect(calledArg.prompt).toContain(''); - }); -}); diff --git a/apps/web/src/lib/workflows/event-trigger.ts b/apps/web/src/lib/workflows/event-trigger.ts deleted file mode 100644 index 472666a52c..0000000000 --- a/apps/web/src/lib/workflows/event-trigger.ts +++ /dev/null @@ -1,230 +0,0 @@ -import { db, workflows, pages, eq, and, ne } from '@pagespace/db'; -import type { EventTrigger } from '@pagespace/db'; -import { executeWorkflow } from './workflow-executor'; -import { loggers } from '@pagespace/lib/server'; - -export interface WorkflowEvent { - operation: string; - resourceType: string; - resourceId: string; - driveId: string | null; - pageId: string | null; - userId: string; - isAiGenerated?: boolean; - aiConversationId?: string | null; - metadata?: Record; -} - -// In-memory debounce buffer: workflowId -> timeout handle -// NOTE: These Maps are process-local. In a multi-instance deployment, events -// hitting different instances won't debounce against each other. Acceptable for -// single-instance Docker; move to Redis if horizontal scaling is needed. -const debounceTimers = new Map>(); - -// Track pending event context per workflow for injection into the prompt -const pendingEventContexts = new Map(); - -/** - * Called from the activity logger hook. Finds matching event-triggered - * workflows and executes them (debounced). - */ -export async function emitWorkflowEvent(event: WorkflowEvent): Promise { - // Skip if no driveId - event-triggered workflows are scoped to drives - if (!event.driveId) return; - - // Recursive trigger prevention: skip if event was generated by a workflow - if (event.isAiGenerated && event.aiConversationId?.startsWith('workflow-')) { - return; - } - - try { - // Find enabled event-triggered workflows in this drive - const matchingWorkflows = await db - .select() - .from(workflows) - .where( - and( - eq(workflows.isEnabled, true), - eq(workflows.triggerType, 'event'), - eq(workflows.driveId, event.driveId) - ) - ); - - if (matchingWorkflows.length === 0) return; - - for (const workflow of matchingWorkflows) { - const eventTriggers = (workflow.eventTriggers as EventTrigger[] | null) ?? []; - - // Check if any trigger matches this event - const matches = eventTriggers.some( - trigger => trigger.operation === event.operation && trigger.resourceType === event.resourceType - ); - if (!matches) continue; - - // Check folder scoping - const watchedFolderIds = (workflow.watchedFolderIds as string[] | null) ?? []; - if (watchedFolderIds.length > 0) { - // The event's pageId is the parent folder. If the resource IS a page, - // we need to check its parentId as well. - let folderMatch = false; - - if (event.pageId && watchedFolderIds.includes(event.pageId)) { - folderMatch = true; - } - - // Also check if the resource itself is inside a watched folder - if (!folderMatch && event.resourceId) { - try { - const [resource] = await db - .select({ parentId: pages.parentId }) - .from(pages) - .where(eq(pages.id, event.resourceId)); - - if (resource?.parentId && watchedFolderIds.includes(resource.parentId)) { - folderMatch = true; - } - } catch { - // Resource might not be a page (could be a file, member, etc.) - } - } - - if (!folderMatch) continue; - } - - // Build event context string for prompt injection - const contextStr = buildEventContext(event); - - // Apply debounce - const debounceSecs = workflow.eventDebounceSecs ?? 30; - const existingTimer = debounceTimers.get(workflow.id); - - if (existingTimer) { - // Update pending context (latest event info) - pendingEventContexts.set(workflow.id, contextStr); - // Timer already running - events are coalescing - continue; - } - - // Store context and set debounce timer - pendingEventContexts.set(workflow.id, contextStr); - - const timer = setTimeout(() => { - void (async () => { - try { - debounceTimers.delete(workflow.id); - const eventContext = pendingEventContexts.get(workflow.id); - pendingEventContexts.delete(workflow.id); - - // Re-validate: workflow may have been disabled or deleted during debounce - const [current] = await db - .select() - .from(workflows) - .where(eq(workflows.id, workflow.id)); - - if (!current || !current.isEnabled) return; - - // Skip if trigger type was changed away from event during debounce - if (current.triggerType !== 'event') return; - - // Skip if another execution is already in progress - if (current.lastRunStatus === 'running') return; - - await executeEventWorkflow(current, eventContext ?? contextStr); - } catch (error) { - loggers.api.error('Error in debounced workflow execution', { - workflowId: workflow.id, - error: error instanceof Error ? error.message : String(error), - }); - } - })(); - }, debounceSecs * 1000); - - debounceTimers.set(workflow.id, timer); - } - } catch (error) { - loggers.api.error('Error in emitWorkflowEvent', { - error: error instanceof Error ? error.message : String(error), - operation: event.operation, - resourceType: event.resourceType, - driveId: event.driveId, - }); - } -} - -const maxContextFieldLength = 200; - -const truncate = (value: unknown, max = maxContextFieldLength): string => { - const str = String(value ?? ''); - return str.length > max ? `${str.slice(0, max)}...` : str; -}; - -function buildEventContext(event: WorkflowEvent): string { - const parts = [`Event: ${event.operation} on ${event.resourceType}`]; - if (event.resourceId) parts.push(`Resource ID: ${event.resourceId}`); - if (event.metadata) { - const title = event.metadata.resourceTitle || event.metadata.title; - if (title) parts.push(`Name: ${truncate(title)}`); - const folder = event.metadata.folderName || event.metadata.parentTitle; - if (folder) parts.push(`Folder: ${truncate(folder)}`); - } - return `\n${parts.join('\n')}\n`; -} - -async function executeEventWorkflow( - workflow: typeof workflows.$inferSelect, - eventContext: string -): Promise { - try { - // Atomically claim: only proceed if not already running (prevents double-execution) - const [claimed] = await db - .update(workflows) - .set({ lastRunStatus: 'running', lastRunAt: new Date() }) - .where(and(eq(workflows.id, workflow.id), ne(workflows.lastRunStatus, 'running'))) - .returning(); - - if (!claimed) { - loggers.api.info('Event workflow already running, skipping', { workflowId: workflow.id }); - return; - } - - // Inject event context into prompt - const enhancedWorkflow = { - ...workflow, - prompt: `${eventContext}\n\n${workflow.prompt}`, - }; - - const result = await executeWorkflow(enhancedWorkflow); - - // Update status (event workflows don't have nextRunAt from cron) - await db - .update(workflows) - .set({ - lastRunAt: new Date(), - lastRunStatus: result.success ? 'success' : 'error', - lastRunError: result.error || null, - lastRunDurationMs: result.durationMs, - }) - .where(eq(workflows.id, workflow.id)); - - loggers.api.info('Event-triggered workflow executed', { - workflowId: workflow.id, - workflowName: workflow.name, - success: result.success, - durationMs: result.durationMs, - }); - } catch (error) { - const errorMsg = error instanceof Error ? error.message : String(error); - loggers.api.error('Event-triggered workflow failed', { - workflowId: workflow.id, - error: errorMsg, - }); - - await db - .update(workflows) - .set({ - lastRunStatus: 'error', - lastRunError: errorMsg, - }) - .where(eq(workflows.id, workflow.id)); - } -}