From 282400b3b8c31b9a6d6e98a263e106b37db35cf1 Mon Sep 17 00:00:00 2001 From: Vincent Grobler Date: Mon, 21 Sep 2026 09:02:27 +0100 Subject: [PATCH] fix: recover stalled runners and make scheduled execution reliable Detect lost runner registrations and exit for a clean host restart instead of reporting healthy while queue claims return no work. Serialize queue polling, refill execution slots promptly, and bound database requests. Preserve unfinished work ownership on shutdown, retry recovery of already dead runners, evaluate cron schedules in UTC, and allow deployments to choose one scheduling source. Add readiness checks, actionable team plan errors, deployment configuration corrections, and an operations runbook. Add 16 regression tests covering registry loss, recovery, queue draining, HTTP readiness/error handling, and schedule evaluation. All 41 runner tests, 3 frontend tests, TypeScript builds, and lint pass (one existing lint warning). --- .env.example | 6 +- .github/workflows/ci.yml | 1 + docker-compose.yml | 5 +- docs/runner-operations.md | 166 +++++++++++++++ railway.json | 4 +- render.yaml | 2 +- supabase/functions/cron-evaluate/index.ts | 38 ++-- task-runner/Dockerfile | 4 +- task-runner/src/index.ts | 245 +++++++++++----------- task-runner/src/runnerRegistry.test.ts | 98 +++++++++ task-runner/src/runnerRegistry.ts | 96 ++++++--- task-runner/src/runnerRuntime.test.ts | 119 +++++++++++ task-runner/src/supabase.ts | 16 ++ task-runner/src/triggerScheduler.test.ts | 41 ++++ task-runner/src/triggerScheduler.ts | 39 ++-- 15 files changed, 687 insertions(+), 193 deletions(-) create mode 100644 docs/runner-operations.md create mode 100644 task-runner/src/runnerRegistry.test.ts create mode 100644 task-runner/src/runnerRuntime.test.ts create mode 100644 task-runner/src/triggerScheduler.test.ts 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[] = [];