From 4dd96eff91e8c0d813a3beaee208de48c30b9305 Mon Sep 17 00:00:00 2001 From: Vincent Grobler Date: Sun, 2 Aug 2026 15:43:07 +0100 Subject: [PATCH] Fix cron trigger execution path --- supabase/functions/cron-evaluate/index.ts | 190 +++++++++++++----- supabase/functions/webhook-trigger/index.ts | 5 +- .../migrations/079_cron_evaluate_schedule.sql | 6 +- .../083_fix_cron_trigger_execution.sql | 39 ++++ task-runner/src/triggerScheduler.ts | 134 ++++++++---- 5 files changed, 270 insertions(+), 104 deletions(-) create mode 100644 supabase/migrations/083_fix_cron_trigger_execution.sql diff --git a/supabase/functions/cron-evaluate/index.ts b/supabase/functions/cron-evaluate/index.ts index 2aeb6e94..2cf6c80e 100644 --- a/supabase/functions/cron-evaluate/index.ts +++ b/supabase/functions/cron-evaluate/index.ts @@ -4,8 +4,8 @@ /** * Cron Evaluate — serverless cron trigger evaluation. * - * Called by Supabase pg_cron on a schedule (e.g. every 30 min or hourly). - * Evaluates all enabled CRON triggers, creates pending tasks for any that + * Called by Supabase pg_cron every minute. + * Evaluates all enabled CRON triggers, creates executable work for any that * are due, and optionally pings the task runner to wake it up. * * This replaces the need for the Railway task runner to run 24/7 just for @@ -21,13 +21,15 @@ import { corsHeaders } from '../_shared/cors.ts'; interface CronTriggerRow { id: string; - agent_id: string; + agent_id: string | null; + team_id: string | null; workspace_id: string; cron_expression: string; task_title_template: string; task_description_template: string; context_options: string[]; last_fired_at: string | null; + created_at: string; } // ─── CRON Parser (mirrored from triggerScheduler.ts) ──────────────────────── @@ -91,7 +93,7 @@ function cronMatchesDate(expression: string, date: Date): boolean { const MAX_CATCHUP_MS = 48 * 60 * 60 * 1000; -function isTriggerDue(cronExpression: string, lastFiredAt: string | null): boolean { +function isTriggerDue(cronExpression: string, lastFiredAt: string | null, createdAt: string): boolean { const now = new Date(); // Current-minute match @@ -111,15 +113,13 @@ function isTriggerDue(cronExpression: string, lastFiredAt: string | null): boole return true; } - // Catch-up: check for missed firings - if (!lastFiredAt) return false; - - const lastFired = new Date(lastFiredAt); - const gapMs = now.getTime() - lastFired.getTime(); + // Catch-up: check for missed firings since last fire, or trigger creation. + const baseline = new Date(lastFiredAt ?? createdAt); + const gapMs = now.getTime() - baseline.getTime(); if (gapMs < 2 * 60 * 1000) return false; - const lookbackStart = new Date(Math.max(lastFired.getTime(), now.getTime() - MAX_CATCHUP_MS)); + const lookbackStart = new Date(Math.max(baseline.getTime(), now.getTime() - MAX_CATCHUP_MS)); const scanTime = new Date(lookbackStart); scanTime.setSeconds(0, 0); scanTime.setMinutes(scanTime.getMinutes() + 1); @@ -128,7 +128,7 @@ function isTriggerDue(cronExpression: string, lastFiredAt: string | null): boole if (cronMatchesDate(cronExpression, scanTime)) { console.log( `[CronEvaluate] Catch-up: missed firing at ${scanTime.toISOString()} ` + - `(last fired: ${lastFiredAt}, now: ${now.toISOString()})`, + `(baseline: ${lastFiredAt ?? createdAt}, now: ${now.toISOString()})`, ); return true; } @@ -138,6 +138,23 @@ function isTriggerDue(cronExpression: string, lastFiredAt: string | null): boole return false; } +async function resolveWorkspaceOwner( + db: ReturnType, + workspaceId: string, +): Promise { + const { data, error } = await db + .from('workspaces') + .select('owner_id') + .eq('id', workspaceId) + .single(); + + if (error || !data) { + throw new Error(`Could not resolve workspace owner: ${error?.message ?? 'workspace not found'}`); + } + + return (data as { owner_id: string }).owner_id; +} + // ─── Template Rendering ───────────────────────────────────────────────────── function renderTemplate(template: string): string { @@ -191,7 +208,7 @@ Deno.serve(async (req: Request) => { // Fetch all enabled CRON triggers const { data: triggers, error: fetchError } = await db .from('agent_triggers') - .select('id, agent_id, workspace_id, cron_expression, task_title_template, task_description_template, context_options, last_fired_at') + .select('id, agent_id, team_id, workspace_id, cron_expression, task_title_template, task_description_template, context_options, last_fired_at, created_at') .eq('trigger_type', 'cron') .eq('enabled', true) .not('cron_expression', 'is', null); @@ -213,45 +230,90 @@ Deno.serve(async (req: Request) => { } let fired = 0; + let firedTasks = 0; + let firedTeamRuns = 0; const firedTriggerIds: string[] = []; + const taskIds: string[] = []; + const teamRunIds: string[] = []; for (const trigger of rows) { try { - if (!isTriggerDue(trigger.cron_expression, trigger.last_fired_at)) continue; + if (!isTriggerDue(trigger.cron_expression, trigger.last_fired_at, trigger.created_at)) continue; - console.log(`[CronEvaluate] Firing trigger ${trigger.id} for agent ${trigger.agent_id}`); + const target = trigger.team_id ? `team ${trigger.team_id}` : `agent ${trigger.agent_id}`; + console.log(`[CronEvaluate] Firing trigger ${trigger.id} for ${target}`); const title = renderTemplate(trigger.task_title_template); const description = renderTemplate(trigger.task_description_template); - - // Create the task as pending - const taskResult = await db - .from('tasks') - .insert({ - workspace_id: trigger.workspace_id, - title, - description, - assigned_agent_id: trigger.agent_id, - status: 'pending', - priority: 'medium', - created_by: trigger.agent_id, - scheduled_for: new Date().toISOString(), - }) - .select('id') - .single(); - - if (taskResult.error) { - await db.from('trigger_log').insert({ - trigger_id: trigger.id, - status: 'failed', - error: taskResult.error.message, - }); - console.error(`[CronEvaluate] Failed to create task for trigger ${trigger.id}:`, taskResult.error.message); - continue; + const ownerId = await resolveWorkspaceOwner(db, trigger.workspace_id); + + let taskId: string | null = null; + let teamRunId: string | null = null; + + if (trigger.team_id) { + const inputTask = description ? `${title}\n\n${description}` : title; + const runResult = await db + .from('team_runs') + .insert({ + workspace_id: trigger.workspace_id, + team_id: trigger.team_id, + input_task: inputTask, + status: 'pending', + created_by: ownerId, + }) + .select('id') + .single(); + + if (runResult.error) { + await db.from('trigger_log').insert({ + trigger_id: trigger.id, + status: 'failed', + error: runResult.error.message, + }); + console.error(`[CronEvaluate] Failed to create team run for trigger ${trigger.id}:`, runResult.error.message); + continue; + } + + teamRunId = (runResult.data as { id: string }).id; + firedTeamRuns++; + teamRunIds.push(teamRunId); + } else if (trigger.agent_id) { + const taskResult = await db + .from('tasks') + .insert({ + workspace_id: trigger.workspace_id, + title, + description, + assigned_agent_id: trigger.agent_id, + status: 'dispatched', + priority: 'medium', + created_by: ownerId, + scheduled_for: new Date().toISOString(), + metadata: { + source: 'cron_trigger', + trigger_id: trigger.id, + }, + }) + .select('id') + .single(); + + if (taskResult.error) { + await db.from('trigger_log').insert({ + trigger_id: trigger.id, + status: 'failed', + error: taskResult.error.message, + }); + console.error(`[CronEvaluate] Failed to create task for trigger ${trigger.id}:`, taskResult.error.message); + continue; + } + + taskId = (taskResult.data as { id: string }).id; + firedTasks++; + taskIds.push(taskId); + } else { + throw new Error('Trigger has no agent_id or team_id'); } - const taskId = (taskResult.data as { id: string }).id; - // Update last_fired_at await db .from('agent_triggers') @@ -267,32 +329,46 @@ Deno.serve(async (req: Request) => { fired++; firedTriggerIds.push(trigger.id); - console.log(`[CronEvaluate] Created task ${taskId} from trigger ${trigger.id}`); + console.log(`[CronEvaluate] Created ${taskId ? `task ${taskId}` : `team run ${teamRunId}`} from trigger ${trigger.id}`); } catch (err) { const errMsg = err instanceof Error ? err.message : String(err); console.error(`[CronEvaluate] Error processing trigger ${trigger.id}: ${errMsg}`); } } - // If tasks were created, ping the task runner to wake it up + // If work was created, ping the task runner to wake it up if (fired > 0) { const taskRunnerUrl = Deno.env.get('TASK_RUNNER_URL'); if (taskRunnerUrl) { try { const webhookSecret = Deno.env.get('WEBHOOK_SECRET') ?? ''; - await fetch(`${taskRunnerUrl}/webhook/task`, { - method: 'POST', - headers: { - 'Content-Type': 'application/json', - ...(webhookSecret ? { 'x-webhook-secret': webhookSecret } : {}), - }, - body: JSON.stringify({ source: 'cron-evaluate', fired }), - signal: AbortSignal.timeout(10000), - }); - console.log(`[CronEvaluate] Pinged task runner to pick up ${fired} task(s)`); + const headers = { + 'Content-Type': 'application/json', + ...(webhookSecret ? { 'x-webhook-secret': webhookSecret } : {}), + }; + + if (firedTasks > 0) { + await fetch(`${taskRunnerUrl}/webhook/task`, { + method: 'POST', + headers, + body: JSON.stringify({ source: 'cron-evaluate', fired: firedTasks }), + signal: AbortSignal.timeout(10000), + }); + } + + if (firedTeamRuns > 0) { + await fetch(`${taskRunnerUrl}/webhook/team-run`, { + method: 'POST', + headers, + body: JSON.stringify({ source: 'cron-evaluate', fired: firedTeamRuns }), + signal: AbortSignal.timeout(10000), + }); + } + + console.log(`[CronEvaluate] Pinged task runner to pick up ${firedTasks} task(s) and ${firedTeamRuns} team run(s)`); } catch { - // Non-fatal — tasks will be picked up on next poll/startup - console.warn('[CronEvaluate] Could not reach task runner — tasks will be picked up on next poll'); + // Non-fatal — work will be picked up on next poll/startup + console.warn('[CronEvaluate] Could not reach task runner — work will be picked up on next poll'); } } } @@ -300,7 +376,11 @@ Deno.serve(async (req: Request) => { return new Response(JSON.stringify({ evaluated: rows.length, fired, + fired_tasks: firedTasks, + fired_team_runs: firedTeamRuns, triggered_ids: firedTriggerIds, + task_ids: taskIds, + team_run_ids: teamRunIds, timestamp: new Date().toISOString(), }), { status: 200, diff --git a/supabase/functions/webhook-trigger/index.ts b/supabase/functions/webhook-trigger/index.ts index f5c41815..adc06475 100644 --- a/supabase/functions/webhook-trigger/index.ts +++ b/supabase/functions/webhook-trigger/index.ts @@ -203,7 +203,7 @@ Deno.serve(async (req: Request) => { message: 'Team run created and queued for execution.', }); } else { - // Agent trigger → create a tasks row (existing behaviour) + // Agent trigger -> create an executable tasks row. const { data: task, error: taskError } = await supabase .from('tasks') .insert({ @@ -212,7 +212,7 @@ Deno.serve(async (req: Request) => { workspace_id: triggerRecord.workspace_id, assigned_agent_id: triggerRecord.agent_id, created_by: ownerId, - status: 'pending', + status: 'dispatched', priority: 'medium', metadata: { source: 'webhook_trigger', @@ -258,4 +258,3 @@ Deno.serve(async (req: Request) => { return serverError(message); } }); - diff --git a/supabase/migrations/079_cron_evaluate_schedule.sql b/supabase/migrations/079_cron_evaluate_schedule.sql index 362163be..8975677b 100644 --- a/supabase/migrations/079_cron_evaluate_schedule.sql +++ b/supabase/migrations/079_cron_evaluate_schedule.sql @@ -3,7 +3,7 @@ -- -- 079_cron_evaluate_schedule.sql — pg_cron job to evaluate triggers serverlessly -- --- Calls the cron-evaluate Edge Function every 30 minutes via pg_net. +-- Calls the cron-evaluate Edge Function every minute via pg_net. -- This removes the need for the task runner to be always-on for cron evaluation. -- @@ -12,7 +12,7 @@ CREATE EXTENSION IF NOT EXISTS pg_cron WITH SCHEMA pg_catalog; CREATE EXTENSION IF NOT EXISTS pg_net WITH SCHEMA extensions; -- ──────────────────────────────────────────────────────────────────────────── --- Schedule: run every 30 minutes +-- Schedule: run every minute -- ──────────────────────────────────────────────────────────────────────────── -- Remove existing schedule if present (idempotent redeploy) @@ -23,7 +23,7 @@ WHERE EXISTS ( SELECT cron.schedule( 'evaluate-cron-triggers', - '*/30 * * * *', + '* * * * *', $$ SELECT extensions.http_post( url := current_setting('app.settings.supabase_url') || '/functions/v1/cron-evaluate', diff --git a/supabase/migrations/083_fix_cron_trigger_execution.sql b/supabase/migrations/083_fix_cron_trigger_execution.sql new file mode 100644 index 00000000..c326fa9c --- /dev/null +++ b/supabase/migrations/083_fix_cron_trigger_execution.sql @@ -0,0 +1,39 @@ +-- SPDX-License-Identifier: AGPL-3.0-or-later +-- Copyright (C) 2026 CrewForm +-- +-- 083_fix_cron_trigger_execution.sql +-- +-- Cron trigger execution fixes for existing deployments: +-- - evaluate cron triggers every minute, so minute-level expressions can match +-- - recover trigger-created agent tasks that were queued as pending + +CREATE EXTENSION IF NOT EXISTS pg_cron WITH SCHEMA pg_catalog; +CREATE EXTENSION IF NOT EXISTS pg_net WITH SCHEMA extensions; + +SELECT cron.unschedule('evaluate-cron-triggers') +WHERE EXISTS ( + SELECT 1 FROM cron.job WHERE jobname = 'evaluate-cron-triggers' +); + +SELECT cron.schedule( + 'evaluate-cron-triggers', + '* * * * *', + $$ + SELECT extensions.http_post( + url := current_setting('app.settings.supabase_url') || '/functions/v1/cron-evaluate', + headers := jsonb_build_object( + 'Content-Type', 'application/json', + 'Authorization', 'Bearer ' || current_setting('app.settings.service_role_key') + ), + body := '{}'::jsonb + ); + $$ +); + +UPDATE public.tasks + SET status = 'dispatched', + updated_at = now() + WHERE status = 'pending' + AND assigned_agent_id IS NOT NULL + AND assigned_team_id IS NULL + AND metadata->>'source' IN ('cron_trigger', 'webhook_trigger'); diff --git a/task-runner/src/triggerScheduler.ts b/task-runner/src/triggerScheduler.ts index b9f39281..7dbc2a73 100644 --- a/task-runner/src/triggerScheduler.ts +++ b/task-runner/src/triggerScheduler.ts @@ -11,13 +11,15 @@ import { supabase } from './supabase'; interface CronTriggerRow { id: string; - agent_id: string; + agent_id: string | null; + team_id: string | null; workspace_id: string; cron_expression: string; task_title_template: string; task_description_template: string; context_options: string[]; last_fired_at: string | null; + created_at: string; } // ─── CRON Parser ──────────────────────────────────────────────────────────── @@ -102,11 +104,12 @@ const MAX_CATCHUP_MS = 48 * 60 * 60 * 1000; * 1. The current time matches the CRON expression, OR * 2. A firing was MISSED while the runner was offline (catch-up) * - * Catch-up: scans each minute between last_fired_at and now (capped at 48h). + * Catch-up: scans each minute between last_fired_at (or created_at for first + * execution) and now, capped at 48h. * If any minute matched the cron expression and the trigger didn't fire, it fires now. * This ensures daily/weekly triggers work even when the runner isn't always-on. */ -function isTriggerDue(cronExpression: string, lastFiredAt: string | null): boolean { +function isTriggerDue(cronExpression: string, lastFiredAt: string | null, createdAt: string): boolean { const now = new Date(); // ── 1. Current-minute match (existing real-time check) ── @@ -128,19 +131,14 @@ function isTriggerDue(cronExpression: string, lastFiredAt: string | null): boole } // ── 2. Catch-up: check for missed firings while runner was offline ── - if (!lastFiredAt) { - // Never fired — don't spam on first boot. Only fire on current-minute match. - return false; - } - - const lastFired = new Date(lastFiredAt); - const gapMs = now.getTime() - lastFired.getTime(); + const baseline = new Date(lastFiredAt ?? createdAt); + const gapMs = now.getTime() - baseline.getTime(); // Only catch up if there's a meaningful gap (> 2 minutes, since we eval every 60s) if (gapMs < 2 * 60 * 1000) return false; // Cap lookback to prevent flooding after long outages - const lookbackStart = new Date(Math.max(lastFired.getTime(), now.getTime() - MAX_CATCHUP_MS)); + const lookbackStart = new Date(Math.max(baseline.getTime(), now.getTime() - MAX_CATCHUP_MS)); // Scan each minute from lookback start to now, looking for a missed cron match // For daily triggers with a 48h lookback, this is at most 2,880 iterations — trivial. @@ -153,7 +151,7 @@ function isTriggerDue(cronExpression: string, lastFiredAt: string | null): boole if (cronMatchesDate(cronExpression, scanTime)) { console.log( `[TriggerScheduler] Catch-up: missed firing at ${scanTime.toISOString()} ` + - `(last fired: ${lastFiredAt}, now: ${now.toISOString()})`, + `(baseline: ${lastFiredAt ?? createdAt}, now: ${now.toISOString()})`, ); return true; // Missed this one — fire now } @@ -163,6 +161,20 @@ function isTriggerDue(cronExpression: string, lastFiredAt: string | null): boole return false; } +async function resolveWorkspaceOwner(workspaceId: string): Promise { + const { data, error } = await supabase + .from('workspaces') + .select('owner_id') + .eq('id', workspaceId) + .single(); + + if (error || !data) { + throw new Error(`Could not resolve workspace owner: ${error?.message ?? 'workspace not found'}`); + } + + return (data as { owner_id: string }).owner_id; +} + // ─── Template Rendering ───────────────────────────────────────────────────── function renderTemplate(template: string): string { @@ -298,7 +310,7 @@ export async function evaluateTriggers(): Promise { // Fetch all enabled CRON triggers const result = await supabase .from('agent_triggers') - .select('id, agent_id, workspace_id, cron_expression, task_title_template, task_description_template, context_options, last_fired_at') + .select('id, agent_id, team_id, workspace_id, cron_expression, task_title_template, task_description_template, context_options, last_fired_at, created_at') .eq('trigger_type', 'cron') .eq('enabled', true) .not('cron_expression', 'is', null); @@ -312,11 +324,11 @@ export async function evaluateTriggers(): Promise { for (const trigger of triggers) { try { - if (!isTriggerDue(trigger.cron_expression, trigger.last_fired_at)) continue; + if (!isTriggerDue(trigger.cron_expression, trigger.last_fired_at, trigger.created_at)) continue; - console.log(`[TriggerScheduler] Firing trigger ${trigger.id} for agent ${trigger.agent_id}`); + const target = trigger.team_id ? `team ${trigger.team_id}` : `agent ${trigger.agent_id}`; + console.log(`[TriggerScheduler] Firing trigger ${trigger.id} for ${target}`); - // Create the task const title = renderTemplate(trigger.task_title_template); let description = renderTemplate(trigger.task_description_template); @@ -332,33 +344,69 @@ export async function evaluateTriggers(): Promise { } } - const taskResult = await supabase - .from('tasks') - .insert({ - workspace_id: trigger.workspace_id, - title, - description, - assigned_agent_id: trigger.agent_id, - status: 'pending', - priority: 'medium', - created_by: trigger.agent_id, // Agent self-creates - scheduled_for: new Date().toISOString(), - }) - .select('id') - .single(); - - if (taskResult.error) { - // Log failure - await supabase.from('trigger_log').insert({ - trigger_id: trigger.id, - status: 'failed', - error: taskResult.error.message, - }); - console.error(`[TriggerScheduler] Failed to create task for trigger ${trigger.id}:`, taskResult.error.message); - continue; - } + const ownerId = await resolveWorkspaceOwner(trigger.workspace_id); + let taskId: string | null = null; + let teamRunId: string | null = null; + + if (trigger.team_id) { + const inputTask = description ? `${title}\n\n${description}` : title; + const runResult = await supabase + .from('team_runs') + .insert({ + workspace_id: trigger.workspace_id, + team_id: trigger.team_id, + input_task: inputTask, + status: 'pending', + created_by: ownerId, + }) + .select('id') + .single(); + + if (runResult.error) { + await supabase.from('trigger_log').insert({ + trigger_id: trigger.id, + status: 'failed', + error: runResult.error.message, + }); + console.error(`[TriggerScheduler] Failed to create team run for trigger ${trigger.id}:`, runResult.error.message); + continue; + } - const taskId = (taskResult.data as { id: string }).id; + teamRunId = (runResult.data as { id: string }).id; + } else if (trigger.agent_id) { + const taskResult = await supabase + .from('tasks') + .insert({ + workspace_id: trigger.workspace_id, + title, + description, + assigned_agent_id: trigger.agent_id, + status: 'dispatched', + priority: 'medium', + created_by: ownerId, + scheduled_for: new Date().toISOString(), + metadata: { + source: 'cron_trigger', + trigger_id: trigger.id, + }, + }) + .select('id') + .single(); + + if (taskResult.error) { + await supabase.from('trigger_log').insert({ + trigger_id: trigger.id, + status: 'failed', + error: taskResult.error.message, + }); + console.error(`[TriggerScheduler] Failed to create task for trigger ${trigger.id}:`, taskResult.error.message); + continue; + } + + taskId = (taskResult.data as { id: string }).id; + } else { + throw new Error('Trigger has no agent_id or team_id'); + } // Update last_fired_at await supabase @@ -373,7 +421,7 @@ export async function evaluateTriggers(): Promise { status: 'fired', }); - console.log(`[TriggerScheduler] Created task ${taskId} from trigger ${trigger.id}`); + console.log(`[TriggerScheduler] Created ${taskId ? `task ${taskId}` : `team run ${teamRunId}`} from trigger ${trigger.id}`); } catch (err: unknown) { const errMsg = err instanceof Error ? err.message : String(err); console.error(`[TriggerScheduler] Error processing trigger ${trigger.id}: ${errMsg}`);