Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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({
Expand Down Expand Up @@ -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({
Expand Down Expand Up @@ -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' },
Expand All @@ -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 () => {
Expand Down Expand Up @@ -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,
Expand Down
46 changes: 11 additions & 35 deletions apps/web/src/app/api/workflows/[workflowId]/route.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,33 +7,27 @@ 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
.select()
.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 }) };
Expand Down Expand Up @@ -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);
Expand All @@ -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)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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();
});

Expand Down
5 changes: 3 additions & 2 deletions apps/web/src/app/api/workflows/[workflowId]/run/route.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand All @@ -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 });
}

Expand Down Expand Up @@ -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;

Expand Down
17 changes: 17 additions & 0 deletions apps/web/src/app/api/workflows/__tests__/route.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 }));

Expand Down
48 changes: 17 additions & 31 deletions apps/web/src/app/api/workflows/route.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6,32 +6,20 @@ 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),
name: z.string().min(1).max(200),
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;
Expand All @@ -55,15 +43,15 @@ 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 } });

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;
Expand Down Expand Up @@ -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,
Expand All @@ -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 });
}
Loading
Loading