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
190 changes: 135 additions & 55 deletions supabase/functions/cron-evaluate/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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) ────────────────────────
Expand Down Expand Up @@ -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
Expand All @@ -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);
Expand All @@ -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;
}
Expand All @@ -138,6 +138,23 @@ function isTriggerDue(cronExpression: string, lastFiredAt: string | null): boole
return false;
}

async function resolveWorkspaceOwner(
db: ReturnType<typeof createClient>,
workspaceId: string,
): Promise<string> {
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 {
Expand Down Expand Up @@ -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);
Expand All @@ -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')
Expand All @@ -267,40 +329,58 @@ 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');
}
}
}

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,
Expand Down
5 changes: 2 additions & 3 deletions supabase/functions/webhook-trigger/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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({
Expand All @@ -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',
Expand Down Expand Up @@ -258,4 +258,3 @@ Deno.serve(async (req: Request) => {
return serverError(message);
}
});

6 changes: 3 additions & 3 deletions supabase/migrations/079_cron_evaluate_schedule.sql
Original file line number Diff line number Diff line change
Expand Up @@ -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.
--

Expand All @@ -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)
Expand All @@ -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',
Expand Down
39 changes: 39 additions & 0 deletions supabase/migrations/083_fix_cron_trigger_execution.sql
Original file line number Diff line number Diff line change
@@ -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');
Loading
Loading