diff --git a/.env.example b/.env.example index 8604aac5..20768353 100644 --- a/.env.example +++ b/.env.example @@ -11,7 +11,8 @@ POSTGRES_PORT=5432 # ── Supabase ────────────────────────────────────────────────────────────── # For hosted Supabase: use your project URL and keys -# For self-hosting without Supabase: leave blank (direct Postgres mode) +# A complete self-hosted Supabase stack is also supported; plain PostgreSQL alone +# does not supply the Auth, REST, Realtime, Storage, and Functions APIs. VITE_SUPABASE_URL= VITE_SUPABASE_ANON_KEY= SUPABASE_SERVICE_ROLE_KEY= @@ -23,6 +24,9 @@ FRONTEND_PORT=3000 # ── Task Runner ─────────────────────────────────────────────────────── # URL of the task runner (used by MCP Server Publishing config snippet) VITE_TASK_RUNNER_URL=http://localhost:3001 +# Leave true for the always-on worker; set false only after configuring pg_cron. +# Do not run both schedulers concurrently. See docs/runner-operations.md. +TRIGGER_SCHEDULER_ENABLED=true # ── Encryption ──────────────────────────────────────────────────────────── # 32-byte hex string for AES-256-GCM API key encryption diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 77226a5a..f8def60b 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -7,6 +7,7 @@ name: CI on: + workflow_dispatch: push: branches: [main] pull_request: diff --git a/docker-compose.yml b/docker-compose.yml index 4ab3c522..952baf36 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -68,8 +68,8 @@ services: # ── Task Runner ───────────────────────────────────────────────────────── task-runner: build: - context: ./task-runner - dockerfile: Dockerfile + context: . + dockerfile: task-runner/Dockerfile restart: unless-stopped environment: VITE_SUPABASE_URL: ${VITE_SUPABASE_URL:-http://localhost:8000} @@ -77,6 +77,7 @@ services: API_KEY_ENCRYPTION_KEY: ${API_KEY_ENCRYPTION_KEY:?API_KEY_ENCRYPTION_KEY is required} ALLOW_PRIVATE_PROVIDER_URLS: ${ALLOW_PRIVATE_PROVIDER_URLS:-false} WEBHOOK_SECRET: ${WEBHOOK_SECRET:?WEBHOOK_SECRET is required} + TRIGGER_SCHEDULER_ENABLED: ${TRIGGER_SCHEDULER_ENABLED:-true} OPENAI_API_KEY: ${OPENAI_API_KEY:-} ANTHROPIC_API_KEY: ${ANTHROPIC_API_KEY:-} GOOGLE_GENERATIVE_AI_API_KEY: ${GOOGLE_GENERATIVE_AI_API_KEY:-} diff --git a/docs/runner-operations.md b/docs/runner-operations.md new file mode 100644 index 00000000..89e0a94d --- /dev/null +++ b/docs/runner-operations.md @@ -0,0 +1,166 @@ +# Runner and schedule operations + +The frontend runs on Vercel, and Supabase stores users, configuration, queues, +and results. The Node runner executes model/tool calls and serves MCP, A2A, +AG-UI, chat-widget, and knowledge-search endpoints. Changing the runner host +does not require moving the frontend or database. + +## Diagnose before restarting + +`GET /health` must return 200 only when the runner has a live registration and +a heartbeat confirmed within the last minute. Compare its `runnerId` with +`public.task_runners`. Older runners returned 200 even after that row was +deleted, leaving a process that could schedule work but could never claim it. + +Run these read-only queries in the Supabase SQL editor: + +```sql +select id, status, current_load, max_concurrency, last_heartbeat, started_at +from public.task_runners; + +select 'tasks' as queue, status, count(*), min(created_at) as oldest +from public.tasks group by status +union all +select 'team_runs', status, count(*), min(created_at) +from public.team_runs group by status; + +select jobid, jobname, schedule, active from cron.job; +select jobid, status, return_message, start_time +from cron.job_run_details order by start_time desc limit 10; + +select status_code, timed_out, error_msg, created +from net._http_response order by created desc limit 10; +``` + +A successful pg_cron job only means PostgreSQL queued the HTTP request. +Inspect the HTTP response separately: 401 is not a successful evaluation. +Do not print cron commands or HTTP request headers; they may contain secrets. + +Before restoring a stuck worker, decide which queued jobs remain relevant. +Restarting immediately can execute old prompts, spend model credits, and send +historical outputs to configured destinations. Review task IDs, destinations, +and dates before deciding to run or cancel them. Do not bulk-convert all +`pending` tasks to `dispatched`: pending tasks may be intentional drafts. + +## Choose one scheduler + +The current runner and Edge Function each evaluate triggers independently. +They do not share an atomic claim for a scheduled firing, so running both can +create duplicates. Use one until that transaction is implemented. + +### Always-on runner (smallest change) + +Keep `TRIGGER_SCHEDULER_ENABLED=true` (the default), use one replica, and disable +the redundant pg_cron job after confirming the local scheduler is working: + +```sql +select cron.alter_job(jobid, active := false) +from cron.job where jobname = 'evaluate-cron-triggers'; +``` + +This SQL changes production configuration; it is a deployment step, not a +diagnostic query. A restarted runner catches up missed schedules within a +48-hour window, producing one current job rather than replaying every interval. +Cron expressions are evaluated in UTC. The custom parser is not a complete +cron implementation; prefer simple expressions until parser validation and +standard day-of-month/day-of-week semantics are added. + +### Supabase pg_cron (for an external scheduler) + +Before switching, store two secrets in Supabase Vault through its dashboard: + +- `crewform_project_url`: the Supabase API origin, not the frontend URL. +- `crewform_cron_secret`: exactly the same value as the `CRON_SECRET` Edge + Function secret. Do not use the public anon key as a scheduler credential. + +Deploy `cron-evaluate` with JWT verification disabled because it authenticates +the `x-cron-secret` header itself. Keep that secret authentication enabled. +Set its `TASK_RUNNER_URL` and `WEBHOOK_SECRET` to match the worker. The function +only enqueues work; a worker must still execute it. + +After verifying the configuration, set the worker's +`TRIGGER_SCHEDULER_ENABLED=false` and replace the database job: + +```sql +select cron.unschedule(jobid) +from cron.job where jobname = 'evaluate-cron-triggers'; + +select cron.schedule('evaluate-cron-triggers', '* * * * *', $job$ + select net.http_post( + url := (select decrypted_secret from vault.decrypted_secrets + where name = 'crewform_project_url') || '/functions/v1/cron-evaluate', + headers := jsonb_build_object( + 'Content-Type', 'application/json', + 'x-cron-secret', (select decrypted_secret from vault.decrypted_secrets + where name = 'crewform_cron_secret') + ), + body := '{}'::jsonb, + timeout_milliseconds := 30000 + ); +$job$); +``` + +The correct pg_net function is `net.http_post`, not the +`extensions.http_post` call in historical migrations 079 and 083. Existing +deployments may have a manually configured job that differs from those files. +This runbook intentionally does not replay historical migrations against a +database with uncertain migration history. + +The Edge evaluator currently does not add the local scheduler's workspace +context enrichment. Do not switch enriched schedules without implementing +parity. Verify `net._http_response`, `trigger_log`, and actual task completion; +checking just the cron job or the enqueue log is insufficient. + +## Safe rollout of runner reliability changes + +1. Review the historical queue and choose what to keep. Resolve invalid team + configuration or plan requirements before reenabling recurring jobs. +2. Verify `SUPABASE_URL` is the API origin and `API_KEY_ENCRYPTION_KEY` matches + the existing Edge Function encryption key. Do not generate a replacement + key for already-encrypted credentials. `ENCRYPTION_KEY` is not the current + variable name. +3. Keep one runner replica and one scheduler. Deploy the runner changes with + the existing credentials and endpoints. +4. Confirm the new runner registration has advancing heartbeats and `/health` + returns 200. A lost registration makes the process exit with code 1 so the + hosting restart policy can obtain a fresh registration. Railway's configured + healthcheck is a deployment readiness check, not continuous auto-healing. +5. Run one explicitly selected task, then a scheduled task. Confirm completion + and output delivery. Observe multiple schedule cycles and a restart. +6. Keep the previous image available for rollback. Reverting restores the old + registry bug, so continue monitoring the registration if rolling back. + +Shutdown preserves dead runner rows until recovery, rather than deleting the +ownership link on unfinished jobs. Recovery is still at-least-once execution: +external writes need idempotency, and approvals need durable continuation. +The fixes do not establish exactly-once delivery or safe resumption of every +interrupted workflow. + +## Hosting choices + +At Railway's $5/month Hobby minimum, lower memory consumption does not reduce +the subscription below $5. Heartbeats, polling, and Realtime keep this worker +active, so Railway serverless sleep is not a solution for the current design. + +Vercel's free cron offering and Supabase's free Edge Function duration limit +do not replace a persistent worker with potentially long agent runs. An +event-driven Cloud Run service is a possible future migration: keep the queue +and scheduler in Supabase, dispatch authenticated requests through a reliable +delivery mechanism, and execute work within the request's lifetime. Disable +permanent polling and persist continuation state for long runs and approvals. +The current webhook responds before execution finishes, so the existing image +is not a drop-in request-billed Cloud Run worker. + +Cloud Run's compute free allowance may cover a small workload, but billing +setup, network traffic, container storage/builds, and usage above allowances +mean it is not a guaranteed zero-cost service. Choose a region near Supabase, +keep minimum instances at zero, and measure actual runtime before estimating +savings. A VPS adds another fixed bill and maintenance responsibility. + +Sources checked 20 September 2026: [Railway plans](https://docs.railway.com/pricing/plans), +[Railway serverless](https://docs.railway.com/deployments/serverless), +[Supabase scheduling](https://supabase.com/docs/guides/functions/schedule-functions), +[Supabase function limits](https://supabase.com/docs/guides/functions/limits), +[Vercel cron limits](https://vercel.com/docs/cron-jobs/usage-and-pricing), +[Cloud Run pricing](https://cloud.google.com/run/pricing), +[Cloud Run runtime](https://docs.cloud.google.com/run/docs/container-contract). diff --git a/railway.json b/railway.json index 11546ddb..5cdfc924 100644 --- a/railway.json +++ b/railway.json @@ -9,8 +9,10 @@ ] }, "deploy": { + "healthcheckPath": "/health", + "healthcheckTimeout": 120, "restartPolicyType": "ON_FAILURE", "restartPolicyMaxRetries": 10, "numReplicas": 1 } -} \ No newline at end of file +} diff --git a/render.yaml b/render.yaml index 6529ba7b..f0a1c178 100644 --- a/render.yaml +++ b/render.yaml @@ -18,7 +18,7 @@ services: sync: false - key: SUPABASE_SERVICE_ROLE_KEY sync: false - - key: ENCRYPTION_KEY + - key: API_KEY_ENCRYPTION_KEY sync: false - key: WEBHOOK_SECRET sync: false diff --git a/supabase/functions/cron-evaluate/index.ts b/supabase/functions/cron-evaluate/index.ts index 29e6a438..c38861dd 100644 --- a/supabase/functions/cron-evaluate/index.ts +++ b/supabase/functions/cron-evaluate/index.ts @@ -83,11 +83,11 @@ function cronMatchesDate(expression: string, date: Date): boolean { const [minute, hour, dayOfMonth, month, dayOfWeek] = parts; return ( - matchesCronField(minute, date.getMinutes(), 59) && - matchesCronField(hour, date.getHours(), 23) && - matchesCronField(dayOfMonth, date.getDate(), 31) && - matchesCronField(month, date.getMonth() + 1, 12) && - matchesCronField(dayOfWeek, date.getDay(), 6) + matchesCronField(minute, date.getUTCMinutes(), 59) && + matchesCronField(hour, date.getUTCHours(), 23) && + matchesCronField(dayOfMonth, date.getUTCDate(), 31) && + matchesCronField(month, date.getUTCMonth() + 1, 12) && + matchesCronField(dayOfWeek, date.getUTCDay(), 6) ); } @@ -101,11 +101,11 @@ function isTriggerDue(cronExpression: string, lastFiredAt: string | null, create if (lastFiredAt) { const last = new Date(lastFiredAt); if ( - last.getFullYear() === now.getFullYear() && - last.getMonth() === now.getMonth() && - last.getDate() === now.getDate() && - last.getHours() === now.getHours() && - last.getMinutes() === now.getMinutes() + last.getUTCFullYear() === now.getUTCFullYear() && + last.getUTCMonth() === now.getUTCMonth() && + last.getUTCDate() === now.getUTCDate() && + last.getUTCHours() === now.getUTCHours() && + last.getUTCMinutes() === now.getUTCMinutes() ) { return false; } @@ -121,8 +121,8 @@ function isTriggerDue(cronExpression: string, lastFiredAt: string | null, create 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); + scanTime.setUTCSeconds(0, 0); + scanTime.setUTCMinutes(scanTime.getUTCMinutes() + 1); while (scanTime < now) { if (cronMatchesDate(cronExpression, scanTime)) { @@ -132,7 +132,7 @@ function isTriggerDue(cronExpression: string, lastFiredAt: string | null, create ); return true; } - scanTime.setMinutes(scanTime.getMinutes() + 1); + scanTime.setUTCMinutes(scanTime.getUTCMinutes() + 1); } return false; @@ -161,7 +161,7 @@ function renderTemplate(template: string): string { const now = new Date(); return template .replace(/\{\{date\}\}/g, now.toISOString().split('T')[0]) - .replace(/\{\{time\}\}/g, now.toTimeString().split(' ')[0]) + .replace(/\{\{time\}\}/g, now.toISOString().slice(11, 19)) .replace(/\{\{datetime\}\}/g, now.toISOString()); } @@ -333,27 +333,29 @@ Deno.serve(async (req: Request) => { }; if (firedTasks > 0) { - await fetch(`${taskRunnerUrl}/webhook/task`, { + const response = await fetch(`${taskRunnerUrl}/webhook/task`, { method: 'POST', headers, body: JSON.stringify({ source: 'cron-evaluate', fired: firedTasks }), signal: AbortSignal.timeout(10000), }); + if (!response.ok) throw new Error(`Task runner returned HTTP ${response.status}`); } if (firedTeamRuns > 0) { - await fetch(`${taskRunnerUrl}/webhook/team-run`, { + const response = await fetch(`${taskRunnerUrl}/webhook/team-run`, { method: 'POST', headers, body: JSON.stringify({ source: 'cron-evaluate', fired: firedTeamRuns }), signal: AbortSignal.timeout(10000), }); + if (!response.ok) throw new Error(`Task runner returned HTTP ${response.status}`); } console.log(`[CronEvaluate] Pinged task runner to pick up ${firedTasks} task(s) and ${firedTeamRuns} team run(s)`); - } catch { + } catch (error) { // 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'); + console.warn('[CronEvaluate] Runner wake-up failed; work remains queued:', error instanceof Error ? error.message : 'unknown error'); } } } diff --git a/task-runner/Dockerfile b/task-runner/Dockerfile index 77cebd15..138c1a1f 100644 --- a/task-runner/Dockerfile +++ b/task-runner/Dockerfile @@ -26,9 +26,9 @@ COPY task-runner/src ./src # Build TypeScript to dist/ RUN npx tsc -p tsconfig.json -# Health check — verify the process is running +# Readiness includes the runner's database registration and recent heartbeat. HEALTHCHECK --interval=30s --timeout=10s --start-period=15s --retries=3 \ - CMD pgrep -f "node" > /dev/null || exit 1 + CMD node -e "fetch('http://127.0.0.1:'+(process.env.PORT||3001)+'/health',{signal:AbortSignal.timeout(5000)}).then(r=>process.exit(r.ok?0:1)).catch(()=>process.exit(1))" # Set explicit Node.js heap limit. # Railway typically provides 512MB–1GB RAM. diff --git a/task-runner/src/index.ts b/task-runner/src/index.ts index f134f528..3fe8c563 100644 --- a/task-runner/src/index.ts +++ b/task-runner/src/index.ts @@ -13,7 +13,7 @@ import { handleChatRequest } from './chatServer'; import { handleKbSearchRequest } from './kbSearchEndpoint'; import { registerRunner, deregisterRunner, getRunnerId, getInstanceName, - runRecoverySweep, RECOVERY_INTERVAL_MS, MAX_CONCURRENT, decrementLoad, + runRecoverySweep, RECOVERY_INTERVAL_MS, MAX_CONCURRENT, decrementLoad, isRunnerHealthy, } from './runnerRegistry'; import { evaluateTriggers, TRIGGER_EVAL_INTERVAL_MS } from './triggerScheduler'; import { initTracing, isTracingEnabled, startTrace, startSpan, endSpan, endTrace, flushTraces } from './tracing'; @@ -34,6 +34,11 @@ let pollTimer: ReturnType | null = null; /** Active slots — how many tasks/runs are currently being processed. */ let activeSlots = 0; +let stopping = false; +let pollInFlight = false; +let pollRequested = false; +let server: http.Server | undefined; +const maintenanceTimers: ReturnType[] = []; function log(msg: string) { const name = getInstanceName(); @@ -51,6 +56,7 @@ function logError(msg: string, err?: unknown) { /** Schedule the next poll with the current adaptive interval. */ function scheduleNextPoll() { + if (stopping) return; if (pollTimer) clearTimeout(pollTimer); pollTimer = setTimeout(() => { void poll(); }, pollIntervalMs); } @@ -70,9 +76,15 @@ function backOffPollInterval() { * Decrements the active slot count, updates DB load, and triggers a new poll. */ function onSlotFreed() { - activeSlots = Math.max(activeSlots - 1, 0); - log(`Slot freed — active: ${activeSlots}/${MAX_CONCURRENT}`); - void decrementLoad(); + void decrementLoad().catch((err: unknown) => { + logError('Failed to release runner capacity:', err); + }).finally(() => { + activeSlots = Math.max(activeSlots - 1, 0); + log(`Slot freed — active: ${activeSlots}/${MAX_CONCURRENT}`); + // A burst can exceed capacity. Drain queued work as each slot frees up, + // rather than waiting up to five minutes for the next fallback poll. + if (!stopping) void poll(); + }); } // ─── Team Run Executor Router ──────────────────────────────────────────────── @@ -105,9 +117,11 @@ async function executeTeamRun(run: TeamRun): Promise { if (teamMode === 'orchestrator') { const allowed = await isFeatureEnabled(run.workspace_id, 'orchestrator_mode'); if (!allowed) { - log(`Orchestrator mode requires an Enterprise license — failing run ${run.id}`); + log(`Orchestrator mode requires a Pro plan or above — failing run ${run.id}`); await supabase.from('team_runs').update({ status: 'failed', + error_message: 'Orchestrator mode requires a Pro plan or above.', + completed_at: new Date().toISOString(), output: 'Orchestrator mode requires a Pro plan or above. Please upgrade at crewform.tech/pricing.', }).eq('id', run.id); if (traceCtx) endTrace(traceCtx, 'error', 'License check failed: orchestrator mode'); @@ -119,9 +133,11 @@ async function executeTeamRun(run: TeamRun): Promise { } else if (teamMode === 'collaboration') { const allowed = await isFeatureEnabled(run.workspace_id, 'collaboration_mode'); if (!allowed) { - log(`Collaboration mode requires an Enterprise license — failing run ${run.id}`); + log(`Collaboration mode requires a Team plan or above — failing run ${run.id}`); await supabase.from('team_runs').update({ status: 'failed', + error_message: 'Collaboration mode requires a Team plan or above.', + completed_at: new Date().toISOString(), output: 'Collaboration mode requires a Team plan or above. Please upgrade at crewform.tech/pricing.', }).eq('id', run.id); if (traceCtx) endTrace(traceCtx, 'error', 'License check failed: collaboration mode'); @@ -155,7 +171,7 @@ async function executeTeamRun(run: TeamRun): Promise { * Returns true if work was claimed. */ async function tryClaimTask(): Promise { - if (activeSlots >= MAX_CONCURRENT) return false; + if (stopping || !isRunnerHealthy() || activeSlots >= MAX_CONCURRENT) return false; const runnerId = getRunnerId(); const rpcResponse = await supabase.rpc('claim_next_task', { @@ -190,7 +206,7 @@ async function tryClaimTask(): Promise { * Returns true if work was claimed. */ async function tryClaimTeamRun(): Promise { - if (activeSlots >= MAX_CONCURRENT) return false; + if (stopping || !isRunnerHealthy() || activeSlots >= MAX_CONCURRENT) return false; const runnerId = getRunnerId(); const teamRunResponse = await supabase.rpc('claim_next_team_run', { @@ -226,121 +242,96 @@ async function tryClaimTeamRun(): Promise { // ─── Poll Loop (Slow Fallback) ─────────────────────────────────────────────── async function poll() { - if (pollTimer) clearTimeout(pollTimer); - - if (activeSlots >= MAX_CONCURRENT) { - scheduleNextPoll(); + if (stopping) return; + if (pollInFlight) { + pollRequested = true; return; } - + pollInFlight = true; + if (pollTimer) clearTimeout(pollTimer); let foundWork = false; - try { - if (await tryClaimTask()) foundWork = true; - if (await tryClaimTeamRun()) foundWork = true; + do { + pollRequested = false; + while (!stopping && isRunnerHealthy() && activeSlots < MAX_CONCURRENT) { + const claimedTask = await tryClaimTask(); + const claimedRun = await tryClaimTeamRun(); + if (!claimedTask && !claimedRun) break; + foundWork = true; + } + } while (pollRequested && !stopping && activeSlots < MAX_CONCURRENT); } catch (err: unknown) { - const errMsg = err instanceof Error ? err.message : String(err); - logError(`Unexpected error in polling loop: ${errMsg}`); - } - - if (foundWork) { - resetPollInterval(); - } else { - backOffPollInterval(); + logError('Unexpected error in polling loop:', err); + } finally { + pollInFlight = false; + if (foundWork) resetPollInterval(); + else backOffPollInterval(); + scheduleNextPoll(); } - - scheduleNextPoll(); } // ─── Webhook Server ────────────────────────────────────────────────────────── function createWebhookServer(): http.Server { - const server = http.createServer((req, res) => { - // A2A protocol endpoints (Agent Card + JSON-RPC) - void handleA2ARequest(req, res).then((a2aHandled) => { - if (a2aHandled) return; - - // AG-UI protocol endpoints (SSE streaming) - void handleAgUiRequest(req, res).then((agUiHandled) => { - if (agUiHandled) return; - - // MCP Server endpoint (expose agents as MCP tools) - void handleMcpServerRequest(req, res).then((mcpHandled) => { - if (mcpHandled) return; - - // Chat Widget endpoints - void handleChatRequest(req, res).then((chatHandled) => { - if (chatHandled) return; - - // KB Search endpoint - void handleKbSearchRequest(req, res).then((kbHandled) => { - if (kbHandled) return; - - // Health check - if (req.method === 'GET' && req.url === '/health') { - res.writeHead(200, { 'Content-Type': 'application/json' }); - res.end(JSON.stringify({ - status: 'ok', - activeSlots, - maxConcurrent: MAX_CONCURRENT, - runnerId: getRunnerId(), - })); - return; - } - - // Webhook endpoints - if (req.method === 'POST' && (req.url === '/webhook/task' || req.url === '/webhook/team-run')) { - // Validate webhook secret - if (!WEBHOOK_SECRET || req.headers['x-webhook-secret'] !== WEBHOOK_SECRET) { - log('Webhook rejected — invalid secret'); - res.writeHead(401, { 'Content-Type': 'application/json' }); - res.end(JSON.stringify({ error: 'Unauthorized' })); - return; - } - - // Read body (we don't actually need the payload — we use claim_next RPC) - let body = ''; - req.on('data', (chunk: Buffer) => { body += chunk.toString(); }); - req.on('end', () => { - const endpoint = req.url; - log(`Webhook received: ${endpoint} (${body.length} bytes)`); - - // Respond immediately — processing is async - res.writeHead(200, { 'Content-Type': 'application/json' }); - res.end(JSON.stringify({ accepted: true })); - - // Attempt to claim and process - if (endpoint === '/webhook/task') { - void tryClaimTask().catch((err: unknown) => { - logError('Webhook task claim failed:', err); - }); - } else { - void tryClaimTeamRun().catch((err: unknown) => { - logError('Webhook team-run claim failed:', err); - }); - } - }); - return; - } - - // 404 for everything else - res.writeHead(404, { 'Content-Type': 'application/json' }); - res.end(JSON.stringify({ error: 'Not found' })); - }); // end handleKbSearchRequest.then - }); // end handleChatRequest.then - }); // end handleMcpServerRequest.then - }); // end handleAgUiRequest.then - }); // end handleA2ARequest.then + return http.createServer((req, res) => { + void routeRequest(req, res).catch((err: unknown) => { + logError('HTTP request failed:', err); + if (res.headersSent) res.destroy(); + else { + res.writeHead(500, { 'Content-Type': 'application/json' }); + res.end(JSON.stringify({ error: 'Internal server error' })); + } + }); }); +} - return server; +async function routeRequest(req: http.IncomingMessage, res: http.ServerResponse): Promise { + if (req.method === 'GET' && req.url === '/health') { + const ready = !stopping && isRunnerHealthy(); + res.writeHead(ready ? 200 : 503, { 'Content-Type': 'application/json' }); + res.end(JSON.stringify({ + status: ready ? 'ok' : 'unavailable', + activeSlots, + maxConcurrent: MAX_CONCURRENT, + runnerId: getRunnerId(), + })); + return; + } + if (stopping) { + res.writeHead(503, { 'Content-Type': 'application/json' }); + res.end(JSON.stringify({ error: 'Runner is shutting down' })); + return; + } + for (const handler of [handleA2ARequest, handleAgUiRequest, handleMcpServerRequest, handleChatRequest, handleKbSearchRequest]) { + if (await handler(req, res)) return; + } + if (req.method === 'POST' && (req.url === '/webhook/task' || req.url === '/webhook/team-run')) { + if (!WEBHOOK_SECRET || req.headers['x-webhook-secret'] !== WEBHOOK_SECRET) { + res.writeHead(401, { 'Content-Type': 'application/json' }); + res.end(JSON.stringify({ error: 'Unauthorized' })); + return; + } + // Payload contents are unnecessary: the database claim is authoritative. + req.resume(); + res.writeHead(200, { 'Content-Type': 'application/json' }); + res.end(JSON.stringify({ accepted: true })); + void poll(); + return; + } + res.writeHead(404, { 'Content-Type': 'application/json' }); + res.end(JSON.stringify({ error: 'Not found' })); } // ─── Startup ───────────────────────────────────────────────────────────────── async function start() { try { - const id = await registerRunner(); + const id = await registerRunner(() => { + // Recovery may have reassigned in-flight work. Stop this process; + // the host restart policy will acquire a new registration. + stopping = true; + process.exit(1); + }); log(`Registered with ID ${id}`); // Initialize OpenTelemetry tracing (no-op if no env vars set) @@ -357,7 +348,7 @@ async function start() { if (!WEBHOOK_SECRET) log('⚠️ WEBHOOK_SECRET is not set — webhook endpoints are disabled'); // Start HTTP webhook server - const server = createWebhookServer(); + server = createWebhookServer(); server.listen(PORT, '0.0.0.0', () => { log(`Webhook server listening on 0.0.0.0:${PORT}`); }); @@ -389,7 +380,7 @@ async function start() { }, (payload) => { log(`Realtime: task ${(payload.new as { id: string }).id} dispatched — claiming`); - void tryClaimTask().catch((err: unknown) => { + void poll().catch((err: unknown) => { logError('Realtime task claim failed:', err); }); }, @@ -404,7 +395,7 @@ async function start() { }, (payload) => { log(`Realtime: new task ${(payload.new as { id: string }).id} — claiming`); - void tryClaimTask().catch((err: unknown) => { + void poll().catch((err: unknown) => { logError('Realtime task claim failed:', err); }); }, @@ -419,7 +410,7 @@ async function start() { }, (payload) => { log(`Realtime: new team run ${(payload.new as { id: string }).id} — claiming`); - void tryClaimTeamRun().catch((err: unknown) => { + void poll().catch((err: unknown) => { logError('Realtime team-run claim failed:', err); }); }, @@ -444,7 +435,7 @@ async function start() { } async function reconnectRealtime() { - if (isReconnecting || realtimeDisabled) return; + if (stopping || isReconnecting || realtimeDisabled) return; isReconnecting = true; realtimeReconnectAttempt++; @@ -469,6 +460,8 @@ async function start() { log(`Realtime reconnect attempt ${realtimeReconnectAttempt}/${REALTIME_GIVE_UP_AFTER} in ${delay}ms`); await new Promise(resolve => setTimeout(resolve, delay)); + if (stopping) return; + // Remove ALL channels — ensures no leaked subscriptions try { await supabase.removeAllChannels(); @@ -489,8 +482,8 @@ async function start() { } // Health check: detect silent disconnects - setInterval(() => { - if (isReconnecting || realtimeDisabled) return; + maintenanceTimers.push(setInterval(() => { + if (stopping || isReconnecting || realtimeDisabled) return; const ch = (globalThis as Record).__realtimeChannel as ReturnType | undefined; @@ -501,18 +494,25 @@ async function start() { log(`Realtime health check: channel state "${state}" — triggering reconnect`); void reconnectRealtime(); } - }, REALTIME_HEALTH_CHECK_MS); + }, REALTIME_HEALTH_CHECK_MS)); const channel = createRealtimeChannel(); (globalThis as Record).__realtimeChannel = channel; // Start recovery sweep and trigger evaluation on fixed intervals - setInterval(() => { void runRecoverySweep(); }, RECOVERY_INTERVAL_MS); - setInterval(() => { void evaluateTriggers(); }, TRIGGER_EVAL_INTERVAL_MS); - - // Immediate trigger catch-up: fire any crons missed while runner was offline - log('Running startup trigger catch-up sweep...'); - void evaluateTriggers(); + maintenanceTimers.push(setInterval(() => { + void runRecoverySweep().then((recovered) => { + if (recovered > 0) void poll(); + }); + }, RECOVERY_INTERVAL_MS)); + if (process.env.TRIGGER_SCHEDULER_ENABLED !== 'false') { + maintenanceTimers.push(setInterval(() => { void evaluateTriggers().then(() => poll()); }, TRIGGER_EVAL_INTERVAL_MS)); + // Use one scheduler: disable this when pg_cron owns scheduling. + log('Running startup trigger catch-up sweep...'); + void evaluateTriggers().then(() => poll()); + } else { + log('Local trigger scheduler disabled; an external scheduler must enqueue scheduled work.'); + } // Initial poll (slow fallback chain starts here) void poll(); @@ -526,6 +526,10 @@ async function start() { // ─── Graceful Shutdown ─────────────────────────────────────────────────────── async function shutdown(signal: string) { + if (stopping) return; + stopping = true; + for (const timer of maintenanceTimers) clearInterval(timer); + server?.close(); log(`Received ${signal}, shutting down gracefully...`); if (pollTimer) clearTimeout(pollTimer); // Unsubscribe from Realtime @@ -537,7 +541,10 @@ async function shutdown(signal: string) { // Wait briefly for in-flight tasks to complete (best effort) if (activeSlots > 0) { log(`Waiting for ${activeSlots} active task(s) to finish...`); - await new Promise(resolve => setTimeout(resolve, 5000)); + const deadline = Date.now() + 30_000; + while (activeSlots > 0 && Date.now() < deadline) { + await new Promise(resolve => setTimeout(resolve, 250)); + } } await deregisterRunner(); process.exit(0); diff --git a/task-runner/src/runnerRegistry.test.ts b/task-runner/src/runnerRegistry.test.ts new file mode 100644 index 00000000..53f2dc18 --- /dev/null +++ b/task-runner/src/runnerRegistry.test.ts @@ -0,0 +1,98 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later +// Copyright (C) 2026 CrewForm + +import { beforeEach, afterEach, describe, expect, it, vi } from 'vitest'; + +const db = vi.hoisted(() => ({ from: vi.fn(), rpc: vi.fn() })); +vi.mock('./supabase', () => ({ supabase: db })); + +let registry: typeof import('./runnerRegistry'); +let chain: Record>; + +beforeEach(async () => { + vi.resetModules(); + vi.useFakeTimers(); + vi.clearAllMocks(); + chain = {}; + for (const name of ['insert', 'update', 'eq', 'select']) chain[name] = vi.fn(() => chain); + chain.single = vi.fn().mockResolvedValue({ data: { id: 'runner-1' }, error: null }); + chain.maybeSingle = vi.fn().mockResolvedValue({ data: { id: 'runner-1' }, error: null }); + db.from.mockReturnValue(chain); + db.rpc.mockResolvedValue({ data: 0, error: null }); + registry = await import('./runnerRegistry'); +}); + +afterEach(() => { + vi.clearAllTimers(); + vi.useRealTimers(); +}); + +describe('runner registration and recovery', () => { + it('is not ready until registered, and becomes unhealthy after missed heartbeats', async () => { + expect(registry.isRunnerHealthy()).toBe(false); + await registry.registerRunner(vi.fn()); + expect(registry.isRunnerHealthy()).toBe(true); + vi.setSystemTime(Date.now() + 60_001); + expect(registry.isRunnerHealthy()).toBe(false); + await registry.sendHeartbeat(); + expect(registry.isRunnerHealthy()).toBe(true); + }); + + it('detects a deleted registration instead of silently updating zero rows forever', async () => { + const lost = vi.fn(); + await registry.registerRunner(lost); + chain.maybeSingle.mockResolvedValue({ data: null, error: null }); + await registry.sendHeartbeat(); + expect(lost).toHaveBeenCalledTimes(1); + expect(registry.isRunnerHealthy()).toBe(false); + expect(chain.eq).toHaveBeenCalledWith('status', 'active'); + await registry.sendHeartbeat(); + expect(lost).toHaveBeenCalledTimes(1); + }); + + it('does not treat a transient database error as proof of a lost lease', async () => { + const lost = vi.fn(); + await registry.registerRunner(lost); + chain.maybeSingle.mockResolvedValue({ data: null, error: { message: 'network error' } }); + await registry.sendHeartbeat(); + expect(lost).not.toHaveBeenCalled(); + vi.setSystemTime(Date.now() + 60_001); + expect(registry.isRunnerHealthy()).toBe(false); + }); + + it('recovers already-dead runners even when no new runner was marked stale', async () => { + db.rpc.mockResolvedValueOnce({ data: 0, error: null }) + .mockResolvedValueOnce({ data: 2, error: null }); + expect(await registry.runRecoverySweep()).toBe(2); + expect(db.rpc).toHaveBeenNthCalledWith(1, 'mark_stale_runners', { stale_threshold: '2 minutes' }); + expect(db.rpc).toHaveBeenNthCalledWith(2, 'recover_stale_tasks'); + }); + + it('retries recovery on the next sweep after the first recovery request fails', async () => { + db.rpc.mockResolvedValueOnce({ data: 1, error: null }) + .mockResolvedValueOnce({ data: null, error: { message: 'temporary failure' } }) + .mockResolvedValueOnce({ data: 0, error: null }) + .mockResolvedValueOnce({ data: 1, error: null }); + expect(await registry.runRecoverySweep()).toBe(0); + expect(await registry.runRecoverySweep()).toBe(1); + }); + + it('preserves runner ownership on shutdown so unfinished tasks remain recoverable', async () => { + await registry.registerRunner(vi.fn()); + await registry.deregisterRunner(); + expect(chain.update).toHaveBeenCalledWith({ status: 'dead' }); + expect(registry.getRunnerId()).toBeNull(); + expect(registry.isRunnerHealthy()).toBe(false); + }); + + it('does not overlap heartbeat requests', async () => { + await registry.registerRunner(vi.fn()); + let complete!: (value: unknown) => void; + chain.maybeSingle.mockImplementation(() => new Promise(resolve => { complete = resolve; })); + const first = registry.sendHeartbeat(); + await registry.sendHeartbeat(); + expect(chain.maybeSingle).toHaveBeenCalledTimes(1); + complete({ data: { id: 'runner-1' }, error: null }); + await first; + }); +}); diff --git a/task-runner/src/runnerRegistry.ts b/task-runner/src/runnerRegistry.ts index e5790da0..5cc94e58 100644 --- a/task-runner/src/runnerRegistry.ts +++ b/task-runner/src/runnerRegistry.ts @@ -3,19 +3,26 @@ import { supabase } from './supabase'; let runnerId: string | null = null; let heartbeatInterval: ReturnType | null = null; +let lastHeartbeatAt = 0; +let heartbeatInFlight = false; +let recoveryInFlight = false; +let leaseLost = false; +let onLeaseLost: () => void = () => { process.exit(1); }; const HEARTBEAT_INTERVAL_MS = 10_000; export const RECOVERY_INTERVAL_MS = 30_000; const INSTANCE_NAME = `${os.hostname()}-${process.pid}`; /** Max concurrent tasks this runner can handle. */ -export const MAX_CONCURRENT = Math.max(1, parseInt(process.env.MAX_CONCURRENT ?? '3', 10)); +const configuredConcurrency = Number(process.env.MAX_CONCURRENT ?? '3'); +export const MAX_CONCURRENT = Number.isInteger(configuredConcurrency) && configuredConcurrency > 0 + ? configuredConcurrency : 3; /** * Register this task runner instance in the database. * Returns the assigned runner UUID. */ -export async function registerRunner(): Promise { +export async function registerRunner(handleLeaseLost?: () => void): Promise { const { data, error } = await supabase .from('task_runners') .insert({ @@ -32,10 +39,15 @@ export async function registerRunner(): Promise { } runnerId = data.id as string; + lastHeartbeatAt = Date.now(); + leaseLost = false; + if (handleLeaseLost) onLeaseLost = handleLeaseLost; // Start heartbeat loop heartbeatInterval = setInterval(() => { - void sendHeartbeat(); + void sendHeartbeat().catch((err: unknown) => { + console.error(`[Runner ${INSTANCE_NAME}] Heartbeat failed:`, err); + }); }, HEARTBEAT_INTERVAL_MS); return runnerId; @@ -44,19 +56,41 @@ export async function registerRunner(): Promise { /** * Send a heartbeat to update last_heartbeat timestamp. */ -async function sendHeartbeat(): Promise { - if (!runnerId) return; - - const { error } = await supabase - .from('task_runners') - .update({ last_heartbeat: new Date().toISOString() }) - .eq('id', runnerId); - - if (error) { - console.error(`[Runner ${INSTANCE_NAME}] Heartbeat failed:`, error.message); +export async function sendHeartbeat(): Promise { + if (!runnerId || heartbeatInFlight || leaseLost) return; + heartbeatInFlight = true; + try { + const { data, error } = await supabase + .from('task_runners') + .update({ last_heartbeat: new Date().toISOString() }) + .eq('id', runnerId) + .eq('status', 'active') + .select('id') + .maybeSingle(); + + if (error) { + console.error(`[Runner ${INSTANCE_NAME}] Heartbeat failed:`, error.message); + return; + } + if (!data) { + // An UPDATE of a deleted row is a successful request with zero rows. + // Never revive this lease: recovery may already have reassigned its work. + leaseLost = true; + console.error(`[Runner ${INSTANCE_NAME}] Runner registration lost; exiting for a clean restart.`); + onLeaseLost(); + return; + } + lastHeartbeatAt = Date.now(); + } finally { + heartbeatInFlight = false; } } +/** Readiness, rather than merely whether the Node process is alive. */ +export function isRunnerHealthy(): boolean { + return !!runnerId && !leaseLost && Date.now() - lastHeartbeatAt < 60_000; +} + /** * Decrement runner load in the database after a task/run completes. */ @@ -77,9 +111,11 @@ export async function decrementLoad(): Promise { * Returns the number of recovered tasks/runs. */ export async function runRecoverySweep(): Promise { + if (recoveryInFlight) return 0; + recoveryInFlight = true; try { // 1. Mark stale runners as dead - const markResult = await supabase.rpc('mark_stale_runners'); + const markResult = await supabase.rpc('mark_stale_runners', { stale_threshold: '2 minutes' }); const staleCount = (markResult.data as number | null) ?? 0; if (markResult.error) { @@ -89,33 +125,34 @@ export async function runRecoverySweep(): Promise { if (staleCount > 0) { console.warn(`[Runner ${INSTANCE_NAME}] Marked ${staleCount} stale runner(s) as dead.`); + } - // 2. Recover orphaned tasks from dead runners - const recoverResult = await supabase.rpc('recover_stale_tasks'); - const recoveredCount = (recoverResult.data as number | null) ?? 0; - - if (recoverResult.error) { - console.error(`[Runner ${INSTANCE_NAME}] recover_stale_tasks failed:`, recoverResult.error.message); - return 0; - } + // Always recover: a previous sweep may have marked runners dead and + // then failed, or a shutdown may have marked one dead explicitly. + const recoverResult = await supabase.rpc('recover_stale_tasks'); + const recoveredCount = (recoverResult.data as number | null) ?? 0; - if (recoveredCount > 0) { - console.warn(`[Runner ${INSTANCE_NAME}] Recovered ${recoveredCount} orphaned task(s)/run(s).`); - } + if (recoverResult.error) { + console.error(`[Runner ${INSTANCE_NAME}] recover_stale_tasks failed:`, recoverResult.error.message); + return 0; + } - return recoveredCount; + if (recoveredCount > 0) { + console.warn(`[Runner ${INSTANCE_NAME}] Recovered ${recoveredCount} orphaned task(s)/run(s).`); } - return 0; + return recoveredCount; } catch (err: unknown) { const errMsg = err instanceof Error ? err.message : String(err); console.error(`[Runner ${INSTANCE_NAME}] Recovery sweep error: ${errMsg}`); return 0; + } finally { + recoveryInFlight = false; } } /** - * Deregister this runner (delete its row) on graceful shutdown. + * Retire this runner without losing ownership of unfinished work. */ export async function deregisterRunner(): Promise { if (heartbeatInterval) { @@ -127,7 +164,7 @@ export async function deregisterRunner(): Promise { const { error } = await supabase .from('task_runners') - .delete() + .update({ status: 'dead' }) .eq('id', runnerId); if (error) { @@ -137,6 +174,7 @@ export async function deregisterRunner(): Promise { } runnerId = null; + lastHeartbeatAt = 0; } /** diff --git a/task-runner/src/runnerRuntime.test.ts b/task-runner/src/runnerRuntime.test.ts new file mode 100644 index 00000000..3e4ed683 --- /dev/null +++ b/task-runner/src/runnerRuntime.test.ts @@ -0,0 +1,119 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later +// Copyright (C) 2026 CrewForm + +import type { IncomingMessage, ServerResponse } from 'http'; +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; + +const state = vi.hoisted(() => ({ + healthy: true, + tasks: [] as Array<{ id: string }>, + finish: [] as Array<() => void>, + httpHandler: undefined as ((req: IncomingMessage, res: ServerResponse) => void) | undefined, + processTask: vi.fn(), + a2a: vi.fn(), + release: vi.fn(), +})); + +vi.mock('http', () => ({ default: { + createServer: (handler: typeof state.httpHandler) => { + state.httpHandler = handler; + return { listen: (_port: number, _host: string, ready: () => void) => ready(), close: vi.fn() }; + }, +} })); +vi.mock('./supabase', () => ({ supabase: { + rpc: async (name: string) => ({ + data: name === 'claim_next_task' && state.tasks.length ? [state.tasks.shift()] : [], error: null, + }), + channel: () => { + const channel = { on: () => channel, subscribe: () => channel }; + return channel; + }, +} })); +vi.mock('./executor', () => ({ processTask: state.processTask })); +vi.mock('./pipelineExecutor', () => ({ processPipelineRun: vi.fn() })); +vi.mock('./orchestratorExecutor', () => ({ processOrchestratorRun: vi.fn() })); +vi.mock('./collaborationExecutor', () => ({ processCollaborationRun: vi.fn() })); +vi.mock('./auditWriter', () => ({ writeTeamRunAudit: vi.fn() })); +vi.mock('./license', () => ({ isFeatureEnabled: vi.fn(), validateLicensesOnStartup: async () => {} })); +vi.mock('./a2aServer', () => ({ handleA2ARequest: state.a2a })); +vi.mock('./agUiServer', () => ({ handleAgUiRequest: async () => false })); +vi.mock('./mcpServer', () => ({ handleMcpServerRequest: async () => false })); +vi.mock('./chatServer', () => ({ handleChatRequest: async () => false })); +vi.mock('./kbSearchEndpoint', () => ({ handleKbSearchRequest: async () => false })); +vi.mock('./runnerRegistry', () => ({ + registerRunner: async () => 'runner-test', deregisterRunner: vi.fn(), + getRunnerId: () => 'runner-test', getInstanceName: () => 'test', + runRecoverySweep: async () => 0, RECOVERY_INTERVAL_MS: 30_000, + MAX_CONCURRENT: 2, decrementLoad: state.release, isRunnerHealthy: () => state.healthy, +})); +vi.mock('./triggerScheduler', () => ({ evaluateTriggers: async () => {}, TRIGGER_EVAL_INTERVAL_MS: 60_000 })); +vi.mock('./tracing', () => ({ initTracing: async () => {}, isTracingEnabled: () => false })); + +const listeners = new Map<'SIGINT' | 'SIGTERM', Set>(); +async function flush() { + // Drain asynchronous startup, claim, and completion continuations without + // advancing the fallback timer (which used to hide delayed queue pickup). + for (let i = 0; i < 40; i++) await Promise.resolve(); +} + +beforeEach(() => { + vi.resetModules(); + vi.clearAllMocks(); + vi.useFakeTimers(); + vi.stubEnv('TRIGGER_SCHEDULER_ENABLED', 'false'); + state.healthy = true; + state.tasks = []; + state.finish = []; + state.httpHandler = undefined; + state.a2a.mockResolvedValue(false); + state.release.mockResolvedValue(undefined); + state.processTask.mockImplementation(() => new Promise(resolve => state.finish.push(resolve))); + for (const signal of ['SIGINT', 'SIGTERM'] as const) listeners.set(signal, new Set(process.listeners(signal))); +}); + +afterEach(() => { + vi.clearAllTimers(); + vi.useRealTimers(); + vi.unstubAllEnvs(); + for (const signal of ['SIGINT', 'SIGTERM'] as const) { + for (const listener of process.listeners(signal)) { + if (!listeners.get(signal)?.has(listener)) process.removeListener(signal, listener); + } + } + delete (globalThis as Record).__realtimeChannel; +}); + +describe('runner runtime', () => { + it('fills capacity and starts queued work on completion without waiting for the polling timer', async () => { + state.tasks = [{ id: 'first' }, { id: 'second' }, { id: 'third' }]; + await import('./index'); + await flush(); + expect(state.processTask).toHaveBeenCalledTimes(2); + state.finish[0](); + await flush(); + expect(state.release).toHaveBeenCalledTimes(1); + expect(state.processTask).toHaveBeenCalledTimes(3); + expect(state.processTask).toHaveBeenLastCalledWith({ id: 'third' }); + }); + + it('reports 503 when registration is unhealthy even though the HTTP process is alive', async () => { + await import('./index'); + await flush(); + state.healthy = false; + const res = { writeHead: vi.fn(), end: vi.fn() }; + state.httpHandler!({ method: 'GET', url: '/health' } as IncomingMessage, res as unknown as ServerResponse); + await flush(); + expect(res.writeHead).toHaveBeenCalledWith(503, { 'Content-Type': 'application/json' }); + expect(JSON.parse(res.end.mock.calls[0][0]).status).toBe('unavailable'); + }); + + it('turns a rejected protocol handler into an HTTP error instead of an unhandled rejection', async () => { + await import('./index'); + await flush(); + state.a2a.mockRejectedValueOnce(new Error('upstream unavailable')); + const res = { writeHead: vi.fn(), end: vi.fn(), headersSent: false }; + state.httpHandler!({ method: 'GET', url: '/a2a/test' } as IncomingMessage, res as unknown as ServerResponse); + await flush(); + expect(res.writeHead).toHaveBeenCalledWith(500, { 'Content-Type': 'application/json' }); + }); +}); diff --git a/task-runner/src/supabase.ts b/task-runner/src/supabase.ts index 185c17e3..c69e1da8 100644 --- a/task-runner/src/supabase.ts +++ b/task-runner/src/supabase.ts @@ -14,6 +14,22 @@ if (!SUPABASE_URL || !SUPABASE_SERVICE_ROLE_KEY) { } export const supabase = createClient(SUPABASE_URL, SUPABASE_SERVICE_ROLE_KEY, { + // Bound database requests so a stalled fetch cannot wedge scheduling or + // heartbeats indefinitely. Preserve cancellation supplied by the caller. + global: { + fetch: (input, init) => { + const url = input instanceof Request ? input.url : String(input); + // Large storage uploads/downloads have different duration needs. + if (!new URL(url).pathname.startsWith('/rest/v1/')) return fetch(input, init); + const callerSignal = init?.signal ?? (input instanceof Request ? input.signal : undefined); + return fetch(input, { + ...init, + signal: callerSignal + ? AbortSignal.any([callerSignal, AbortSignal.timeout(15_000)]) + : AbortSignal.timeout(15_000), + }); + }, + }, auth: { persistSession: false, autoRefreshToken: false, diff --git a/task-runner/src/triggerScheduler.test.ts b/task-runner/src/triggerScheduler.test.ts new file mode 100644 index 00000000..1bd388a8 --- /dev/null +++ b/task-runner/src/triggerScheduler.test.ts @@ -0,0 +1,41 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later +// Copyright (C) 2026 CrewForm + +import { describe, expect, it, vi } from 'vitest'; +vi.mock('./supabase', () => ({ supabase: {} })); +import { cronMatchesDate, isTriggerDue } from './triggerScheduler'; + +describe('schedule evaluation', () => { + it('evaluates UTC regardless of the host local-time getters', () => { + const date = new Date('2026-09-20T09:00:00Z'); + // Local time could be 10:00 in London; UTC remains 09:00. + vi.spyOn(date, 'getHours').mockReturnValue(10); + expect(cronMatchesDate('0 9 * * *', date)).toBe(true); + expect(cronMatchesDate('0 10 * * *', date)).toBe(false); + }); + + it('does not fire twice in the same minute', () => { + expect(isTriggerDue('0 9 * * *', '2026-09-20T09:00:02Z', + '2026-09-01T00:00:00Z', new Date('2026-09-20T09:00:55Z'))).toBe(false); + }); + + it('catches up one missed daily run after a restart', () => { + expect(isTriggerDue('0 9 * * *', '2026-09-19T09:00:00Z', + '2026-09-01T00:00:00Z', new Date('2026-09-20T09:08:00Z'))).toBe(true); + }); + + it('uses creation time when a new trigger has never fired', () => { + expect(isTriggerDue('0 9 * * *', null, + '2026-09-20T08:00:00Z', new Date('2026-09-20T09:08:00Z'))).toBe(true); + }); + + it('does not replay missed weekly runs older than the catch-up window', () => { + expect(isTriggerDue('0 9 * * 1', '2026-08-31T09:00:00Z', + '2026-08-01T00:00:00Z', new Date('2026-09-20T12:00:00Z'))).toBe(false); + }); + + it('does not replay schedules from before creation', () => { + expect(isTriggerDue('0 9 * * *', null, + '2026-09-20T10:00:00Z', new Date('2026-09-20T10:08:00Z'))).toBe(false); + }); +}); diff --git a/task-runner/src/triggerScheduler.ts b/task-runner/src/triggerScheduler.ts index 7dbc2a73..7a840a5f 100644 --- a/task-runner/src/triggerScheduler.ts +++ b/task-runner/src/triggerScheduler.ts @@ -77,18 +77,18 @@ function matchesCronField(field: string, value: number, max: number): boolean { /** * Check if a CRON expression matches a given date. */ -function cronMatchesDate(expression: string, date: Date): boolean { +export function cronMatchesDate(expression: string, date: Date): boolean { const parts = expression.trim().split(/\s+/); if (parts.length !== 5) return false; const [minute, hour, dayOfMonth, month, dayOfWeek] = parts; return ( - matchesCronField(minute, date.getMinutes(), 59) && - matchesCronField(hour, date.getHours(), 23) && - matchesCronField(dayOfMonth, date.getDate(), 31) && - matchesCronField(month, date.getMonth() + 1, 12) && - matchesCronField(dayOfWeek, date.getDay(), 6) + matchesCronField(minute, date.getUTCMinutes(), 59) && + matchesCronField(hour, date.getUTCHours(), 23) && + matchesCronField(dayOfMonth, date.getUTCDate(), 31) && + matchesCronField(month, date.getUTCMonth() + 1, 12) && + matchesCronField(dayOfWeek, date.getUTCDay(), 6) ); } @@ -109,8 +109,7 @@ const MAX_CATCHUP_MS = 48 * 60 * 60 * 1000; * 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, createdAt: string): boolean { - const now = new Date(); +export function isTriggerDue(cronExpression: string, lastFiredAt: string | null, createdAt: string, now = new Date()): boolean { // ── 1. Current-minute match (existing real-time check) ── if (cronMatchesDate(cronExpression, now)) { @@ -118,11 +117,11 @@ function isTriggerDue(cronExpression: string, lastFiredAt: string | null, create const last = new Date(lastFiredAt); // Already fired this minute — skip if ( - last.getFullYear() === now.getFullYear() && - last.getMonth() === now.getMonth() && - last.getDate() === now.getDate() && - last.getHours() === now.getHours() && - last.getMinutes() === now.getMinutes() + last.getUTCFullYear() === now.getUTCFullYear() && + last.getUTCMonth() === now.getUTCMonth() && + last.getUTCDate() === now.getUTCDate() && + last.getUTCHours() === now.getUTCHours() && + last.getUTCMinutes() === now.getUTCMinutes() ) { return false; } @@ -143,9 +142,9 @@ function isTriggerDue(cronExpression: string, lastFiredAt: string | null, create // 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. const scanTime = new Date(lookbackStart); - scanTime.setSeconds(0, 0); + scanTime.setUTCSeconds(0, 0); // Start from the minute after last fired - scanTime.setMinutes(scanTime.getMinutes() + 1); + scanTime.setUTCMinutes(scanTime.getUTCMinutes() + 1); while (scanTime < now) { if (cronMatchesDate(cronExpression, scanTime)) { @@ -155,7 +154,7 @@ function isTriggerDue(cronExpression: string, lastFiredAt: string | null, create ); return true; // Missed this one — fire now } - scanTime.setMinutes(scanTime.getMinutes() + 1); + scanTime.setUTCMinutes(scanTime.getUTCMinutes() + 1); } return false; @@ -181,7 +180,7 @@ function renderTemplate(template: string): string { const now = new Date(); return template .replace(/\{\{date\}\}/g, now.toISOString().split('T')[0]) - .replace(/\{\{time\}\}/g, now.toTimeString().split(' ')[0]) + .replace(/\{\{time\}\}/g, now.toISOString().slice(11, 19)) .replace(/\{\{datetime\}\}/g, now.toISOString()); } @@ -201,11 +200,11 @@ async function buildContextBlock( if (options.length === 0) return ''; const yesterday = new Date(); - yesterday.setDate(yesterday.getDate() - 1); + yesterday.setUTCDate(yesterday.getUTCDate() - 1); const startOfYesterday = new Date(yesterday); - startOfYesterday.setHours(0, 0, 0, 0); + startOfYesterday.setUTCHours(0, 0, 0, 0); const endOfYesterday = new Date(yesterday); - endOfYesterday.setHours(23, 59, 59, 999); + endOfYesterday.setUTCHours(23, 59, 59, 999); const sections: string[] = [];