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

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 5 additions & 1 deletion .env.example
Original file line number Diff line number Diff line change
Expand Up @@ -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=
Expand All @@ -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
Expand Down
1 change: 1 addition & 0 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@
name: CI

on:
workflow_dispatch:
push:
branches: [main]
pull_request:
Expand Down
5 changes: 3 additions & 2 deletions docker-compose.yml
Original file line number Diff line number Diff line change
Expand Up @@ -68,15 +68,16 @@ 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}
SUPABASE_SERVICE_ROLE_KEY: ${SUPABASE_SERVICE_ROLE_KEY:-}
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:-}
Expand Down
166 changes: 166 additions & 0 deletions docs/runner-operations.md
Original file line number Diff line number Diff line change
@@ -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).
4 changes: 3 additions & 1 deletion railway.json
Original file line number Diff line number Diff line change
Expand Up @@ -9,8 +9,10 @@
]
},
"deploy": {
"healthcheckPath": "/health",
"healthcheckTimeout": 120,
"restartPolicyType": "ON_FAILURE",
"restartPolicyMaxRetries": 10,
"numReplicas": 1
}
}
}
2 changes: 1 addition & 1 deletion render.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
38 changes: 20 additions & 18 deletions supabase/functions/cron-evaluate/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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)
);
}

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

Expand Down Expand Up @@ -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');
}
}
}
Expand Down
4 changes: 2 additions & 2 deletions task-runner/Dockerfile
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
Loading
Loading