fix(sweeper): give a queued task its own reconnect grace before expiring it

Review catch: keying expiry on runtime liveness alone was too aggressive in
one direction. Enqueue binds a task to agent.runtime_id without checking
that the runtime is up (task.go CreateAgentTask call sites), so for a
runtime that has already been dark for longer than the grace, the liveness
clause is satisfied the moment a new task is created — assigning an issue
to a laptop closed overnight would fail it inside one 30s sweep tick
instead of waiting for the machine to come back. That is the opposite of
what "reconnect grace" promises.

Require the row's own age as well: a queued task now waits a full grace,
counted from when it started waiting, before it can be given up on. A
heartbeating runtime still never expires its backlog, so the MUL-6558 fix
is unchanged. The age bound also stops the sweep re-evaluating the runtime
subquery against every queued row on every tick.

Also remove the deploy surfaces the previous commit missed — the knob was
deleted from the server but helm (values + configmap), the self-host
compose file, and the CI-run helm render assertions still advertised it, so
operators would have kept setting a variable nothing reads. Same for five
orphaned comment lines in .env.example describing the removed variable.

Correct the fail-closed comment on the unbound-runtime arms: they are not
reachable for "rows predating migration 251" — runtime_id was NOT NULL from
migration 004 until 251 replaced it with a CHECK that every insert and
update is verified against, and the migration-004 FK is ON DELETE RESTRICT.
The arms stay as defence against a future schema change, which is what the
comment now says.

Tests: a third phase covers the reported case — a task enqueued against an
already-dead runtime survives, then expires once it has waited a full grace
of its own. Phase 2's expectation is corrected accordingly: the row created
at now() no longer expires when its runtime dies, because it has not waited
its own grace yet.

Co-authored-by: multica-agent <github@multica.ai>
This commit is contained in:
J
2026-08-28 17:32:50 +08:00
co-authored by multica-agent
parent 786f43d7c4
commit b567511eef
8 changed files with 109 additions and 33 deletions
-5
View File
@@ -88,11 +88,6 @@ MULTICA_SHUTDOWN_HOLD_DURATION=
# Accepts a positive Go duration. The default is 3h; values below 150s are
# clamped to the runtime heartbeat freshness window.
MULTICA_RUNTIME_RECONNECT_GRACE=3h
# How long a task may sit in 'queued' without being claimed before failing
# with 'queued_expired'. Accepts a positive Go duration; empty, unset, or
# invalid values keep the default of 2h. Raise this when a runtime has low
# task concurrency and tasks can legitimately wait behind long-running work
# — with a concurrency-1 runtime a backlog can exceed the default 2h.
# Comma-separated CIDR list of reverse proxies whose X-Forwarded-For /
# X-Real-IP headers the per-IP webhook rate limiter is allowed to trust.
# Empty (the default) means "trust no headers" — the limiter uses
@@ -28,7 +28,6 @@ data:
CLOUDFRONT_DOMAIN: {{ .Values.backend.config.cloudfrontDomain | quote }}
CLOUDFRONT_KEY_PAIR_ID: {{ .Values.backend.config.cloudfrontKeyPairId | quote }}
LOCAL_UPLOAD_BASE_URL: {{ .Values.backend.config.localUploadBaseUrl | quote }}
MULTICA_TASK_QUEUED_TTL: {{ .Values.backend.config.taskQueuedTTL | quote }}
{{- if not .Values.postgres.external.enabled }}
# --- PostgreSQL (consumed by the backend to build DATABASE_URL) ---
-5
View File
@@ -146,11 +146,6 @@ backend:
cloudfrontDomain: ""
cloudfrontKeyPairId: ""
localUploadBaseUrl: http://api.multica.dev.lan
# How long a task may sit in the queue unclaimed before failing with
# `queued_expired`; raise it when a runtime's low task concurrency lets
# tasks legitimately wait behind long-running work.
# Accepts a positive duration and empty/unset falls back to the server default of 2h.
taskQueuedTTL: 2h
resources:
requests:
cpu: 100m
-1
View File
@@ -100,7 +100,6 @@ services:
MULTICA_DATABASE_CONNECT_TIMEOUT: ${MULTICA_DATABASE_CONNECT_TIMEOUT:-5s}
MULTICA_SHUTDOWN_HOLD_DURATION: ${MULTICA_SHUTDOWN_HOLD_DURATION:-}
MULTICA_RUNTIME_RECONNECT_GRACE: ${MULTICA_RUNTIME_RECONNECT_GRACE:-}
MULTICA_TASK_QUEUED_TTL: ${MULTICA_TASK_QUEUED_TTL:-}
ALLOW_SIGNUP: ${ALLOW_SIGNUP:-true}
ALLOWED_EMAILS: ${ALLOWED_EMAILS:-}
ALLOWED_EMAIL_DOMAINS: ${ALLOWED_EMAIL_DOMAINS:-}
-8
View File
@@ -34,17 +34,9 @@ default_config="$(
)"
require_rendered_value "$default_config" 'MULTICA_VCS_INTEGRATION_ENABLED: "true"'
require_rendered_value "$default_config" 'MULTICA_CLOUD_URL: ""'
require_rendered_value "$default_config" 'MULTICA_TASK_QUEUED_TTL: "2h"'
require_rendered_value "$default_config" 'MULTICA_DATABASE_STARTUP_TIMEOUT: "3m"'
require_rendered_value "$default_config" 'MULTICA_DATABASE_CONNECT_TIMEOUT: "5s"'
queued_ttl_config="$(
helm template multica "$CHART_DIR" \
--show-only templates/configmap.yaml \
--set-string backend.config.taskQueuedTTL=12h
)"
require_rendered_value "$queued_ttl_config" 'MULTICA_TASK_QUEUED_TTL: "12h"'
default_backend="$(
helm template multica "$CHART_DIR" \
--show-only templates/backend.yaml
+63 -5
View File
@@ -1039,7 +1039,12 @@ func TestSweepDoesNotResetIssueAlreadyInReview(t *testing.T) {
// than queue age (MUL-6558). The same ancient queued task must survive while
// its runtime is heartbeating — a busy runtime is not a dead one — and only
// become expirable once that runtime has been silent past the reconnect grace.
// A runtime_offline retry stays exempt in both phases; FailExpiredRuntimeReconnectRetries owns its exit.
// A third phase covers the other direction: liveness alone is not enough
// either, because enqueue binds a task to its agent's runtime without checking
// that the runtime is up. A task assigned to an already-dead runtime must still
// get its own full grace to wait, or assigning an issue to a laptop that is
// closed overnight fails inside one sweep tick.
// A runtime_offline retry stays exempt throughout; FailExpiredRuntimeReconnectRetries owns its exit.
func TestExpireStaleQueuedTasks(t *testing.T) {
if testPool == nil {
t.Skip("no database connection")
@@ -1159,8 +1164,14 @@ func TestExpireStaleQueuedTasks(t *testing.T) {
t.Fatalf("ExpireStaleQueuedTasks (dead runtime) failed: %v", err)
}
expired := expiredIDs(failed)
if !expired[parseUUIDBytes(oldTaskID)] || !expired[parseUUIDBytes(freshTaskID)] {
t.Fatalf("dead runtime: expected both non-exempt queued tasks to expire, got %v", expired)
if !expired[parseUUIDBytes(oldTaskID)] {
t.Fatal("dead runtime: the 5h-old queued task should expire once its runtime is gone")
}
// The fresh row is the same case as phase 3 seen from the other side: its
// runtime is dead, but it has not waited a grace of its own yet, so both
// conditions are required and it survives this sweep.
if expired[parseUUIDBytes(freshTaskID)] {
t.Fatal("dead runtime: a task queued moments ago must still get its own full grace before failing")
}
if expired[parseUUIDBytes(recoveryTaskID)] {
t.Fatal("runtime_offline retry must stay exempt from the queued sweep")
@@ -1190,8 +1201,8 @@ func TestExpireStaleQueuedTasks(t *testing.T) {
`, freshTaskID).Scan(&freshStatus); err != nil {
t.Fatalf("failed to read fresh task: %v", err)
}
if freshStatus != "failed" {
t.Fatalf("fresh task: expected status=failed once its runtime went away, got %q", freshStatus)
if freshStatus != "queued" {
t.Fatalf("fresh task: expected status=queued (it has not waited a full grace yet), got %q", freshStatus)
}
var recoveryStatus string
@@ -1201,6 +1212,53 @@ func TestExpireStaleQueuedTasks(t *testing.T) {
if recoveryStatus != "queued" {
t.Fatalf("runtime recovery retry: expected status=queued, got %q", recoveryStatus)
}
// Phase 3 — a task enqueued AFTER the runtime went dark. The runtime has
// been silent for hours, so the liveness clause is already satisfied the
// moment this row is created; only the row's own age keeps it alive. It
// must survive until it has waited a full grace of its own, otherwise
// assigning an issue to a machine that is merely asleep fails in ~30s.
lateIssueID := mkIssue("Queued TTL test (enqueued after runtime died)")
t.Cleanup(func() {
testPool.Exec(ctx, `DELETE FROM agent_task_queue WHERE issue_id = $1`, lateIssueID)
testPool.Exec(ctx, `DELETE FROM issue WHERE id = $1`, lateIssueID)
})
var lateTaskID string
if err := testPool.QueryRow(ctx, `
INSERT INTO agent_task_queue (agent_id, runtime_id, issue_id, status, priority, created_at)
VALUES ($1, $2, $3, 'queued', 0, now())
RETURNING id
`, agentID, runtimeID, lateIssueID).Scan(&lateTaskID); err != nil {
t.Fatalf("failed to insert late queued task: %v", err)
}
lateSweep, err := queries.ExpireStaleQueuedTasks(ctx, db.ExpireStaleQueuedTasksParams{
ReconnectGraceSecs: graceSecs,
MaxPerTick: 100,
})
if err != nil {
t.Fatalf("ExpireStaleQueuedTasks (late enqueue) failed: %v", err)
}
if expiredIDs(lateSweep)[parseUUIDBytes(lateTaskID)] {
t.Fatal("a task enqueued against an already-offline runtime was failed immediately; it must get a full reconnect grace of its own to wait")
}
// Once it HAS waited a full grace, it expires like any other.
if _, err := testPool.Exec(ctx, `
UPDATE agent_task_queue SET created_at = now() - make_interval(secs => $1) WHERE id = $2
`, graceSecs+60, lateTaskID); err != nil {
t.Fatalf("failed to age the late task: %v", err)
}
agedSweep, err := queries.ExpireStaleQueuedTasks(ctx, db.ExpireStaleQueuedTasksParams{
ReconnectGraceSecs: graceSecs,
MaxPerTick: 100,
})
if err != nil {
t.Fatalf("ExpireStaleQueuedTasks (aged late enqueue) failed: %v", err)
}
if !expiredIDs(agedSweep)[parseUUIDBytes(lateTaskID)] {
t.Fatal("a task that waited a full grace against a dead runtime should expire")
}
}
// TestExpireStaleQueuedTasksRespectsBatchLimit verifies the per-tick cap so
+23 -4
View File
@@ -3298,6 +3298,7 @@ const expireStaleQueuedTasks = `-- name: ExpireStaleQueuedTasks :many
WITH victims AS (
SELECT id FROM agent_task_queue
WHERE status = 'queued'
AND created_at < now() - make_interval(secs => $1::double precision)
AND (
runtime_id IS NULL
OR NOT EXISTS (
@@ -3328,6 +3329,7 @@ SET status = 'failed',
FROM victims v
WHERE t.id = v.id
AND t.status = 'queued'
AND t.created_at < now() - make_interval(secs => $1::double precision)
AND (
t.runtime_id IS NULL
OR NOT EXISTS (SELECT 1 FROM agent_runtime r WHERE r.id = t.runtime_id)
@@ -3368,12 +3370,29 @@ type ExpireStaleQueuedTasksParams struct {
// dispatched/running rows, so a daemon going down now retires its queued and
// its in-flight work on one clock instead of two.
//
// The row must ALSO have been queued for a full grace of its own. Enqueue binds
// a task to agent.runtime_id without checking that the runtime is up, so
// runtime liveness alone would fail a task the instant it is assigned to a
// runtime that has been offline a while — a laptop closed overnight would turn
// "assign this issue" into a failure inside one 30s sweep tick instead of
// waiting for the machine to come back. Requiring the task's own age keeps the
// promise the name makes: a queued task gets one full reconnect grace before it
// is given up on, counted from when it started waiting. It also bounds the
// scan, which would otherwise re-evaluate the runtime subquery against every
// queued row on every tick.
//
// Heartbeat age is read directly rather than gated on runtime.status='online',
// so a row stuck at 'online' with a long-dead heartbeat still releases its
// queue. Rows with no usable runtime binding at all — runtime_id IS NULL
// (possible for pre-migration-251 rows, whose CHECK landed NOT VALID) or a
// runtime row that no longer exists — can never acquire a liveness signal, so
// they are expired on sight rather than left to sit forever.
// queue.
//
// The runtime_id IS NULL / missing-runtime arms are fail-closed defence, not
// live paths: the schema already excludes both. runtime_id was NOT NULL from
// migration 004 until 251 replaced it with CHECK (runtime_id IS NOT NULL OR
// completed_at IS NOT NULL), which every insert and update is checked against
// (NOT VALID only skips the backfill scan), so no queued row can be unbound;
// and the migration-004 FK is ON DELETE RESTRICT, so runtime_id cannot point at
// a deleted runtime. They are kept so a future schema change cannot silently
// strand rows that have no liveness signal at all.
//
// A retry created by runtime_offline is exempt: it deliberately waits for that
// runtime to reconnect, and FailExpiredRuntimeReconnectRetries owns its exit.
+23 -4
View File
@@ -1411,12 +1411,29 @@ RETURNING *;
-- dispatched/running rows, so a daemon going down now retires its queued and
-- its in-flight work on one clock instead of two.
--
-- The row must ALSO have been queued for a full grace of its own. Enqueue binds
-- a task to agent.runtime_id without checking that the runtime is up, so
-- runtime liveness alone would fail a task the instant it is assigned to a
-- runtime that has been offline a while — a laptop closed overnight would turn
-- "assign this issue" into a failure inside one 30s sweep tick instead of
-- waiting for the machine to come back. Requiring the task's own age keeps the
-- promise the name makes: a queued task gets one full reconnect grace before it
-- is given up on, counted from when it started waiting. It also bounds the
-- scan, which would otherwise re-evaluate the runtime subquery against every
-- queued row on every tick.
--
-- Heartbeat age is read directly rather than gated on runtime.status='online',
-- so a row stuck at 'online' with a long-dead heartbeat still releases its
-- queue. Rows with no usable runtime binding at all — runtime_id IS NULL
-- (possible for pre-migration-251 rows, whose CHECK landed NOT VALID) or a
-- runtime row that no longer exists — can never acquire a liveness signal, so
-- they are expired on sight rather than left to sit forever.
-- queue.
--
-- The runtime_id IS NULL / missing-runtime arms are fail-closed defence, not
-- live paths: the schema already excludes both. runtime_id was NOT NULL from
-- migration 004 until 251 replaced it with CHECK (runtime_id IS NOT NULL OR
-- completed_at IS NOT NULL), which every insert and update is checked against
-- (NOT VALID only skips the backfill scan), so no queued row can be unbound;
-- and the migration-004 FK is ON DELETE RESTRICT, so runtime_id cannot point at
-- a deleted runtime. They are kept so a future schema change cannot silently
-- strand rows that have no liveness signal at all.
--
-- A retry created by runtime_offline is exempt: it deliberately waits for that
-- runtime to reconnect, and FailExpiredRuntimeReconnectRetries owns its exit.
@@ -1438,6 +1455,7 @@ RETURNING *;
WITH victims AS (
SELECT id FROM agent_task_queue
WHERE status = 'queued'
AND created_at < now() - make_interval(secs => @reconnect_grace_secs::double precision)
AND (
runtime_id IS NULL
OR NOT EXISTS (
@@ -1468,6 +1486,7 @@ SET status = 'failed',
FROM victims v
WHERE t.id = v.id
AND t.status = 'queued'
AND t.created_at < now() - make_interval(secs => @reconnect_grace_secs::double precision)
AND (
t.runtime_id IS NULL
OR NOT EXISTS (SELECT 1 FROM agent_runtime r WHERE r.id = t.runtime_id)