fix(020): don't burn jobs_started on pre-dispatch failures
Root cause: ConsolidatorWorker.tryDispatch / launchOne incremented
the per-(owner, day) `jobs_started` counter immediately after
JobsStore.Create, BEFORE the dispatch flip succeeded. When the
subsequent steps fail (token issue, agent lookup, dispatch flip
race) the job row is Completed as `failed` but `jobs_started`
remains incremented — and the error paths never call
RecordCompletion, so `jobs_failed` stays flat while `jobs_started`
drifts upward.
Over hours/days, owners whose dispatches fail consistently (e.g.
owner_id=2 in the kubic deployment, hitting one of the harness
failure modes from commit bfb2551) accumulate phantom
`jobs_started` until the default DreamDailyJobLimit=100 trips. From
that point every tick logs `circuit broken … reason=jobs_exceeded`
for all four job types, even though no real jobs ran — and the
counter never decays until midnight UTC.
Fix: move `usage.RecordStart(...)` to AFTER a successful
`jobs.Dispatch(...)` in both tryDispatch (consolidator.go:518)
and launchOne (consolidator.go:455). Now only dispatches that
actually transitioned a row to `dispatched` count against the
daily-job-limit gate.
Test: TestConsolidator_PreDispatchFailureDoesNotBurnJobsStarted
seeds DreamDailyJobLimit=2, makes the agent lookup fail, calls
ForceRun three times, asserts jobs_started stays 0 and the gate
still allows. Verified to fail without the fix
(jobs_started=2 / reason=jobs_exceeded) and pass with it.
Counterpart TestConsolidator_DispatchSuccessIncrementsJobsStarted
asserts jobs_started=1 on a real successful dispatch so the
counter still feeds the gate correctly.
Operational note: this prevents future inflation. Existing stuck
rows for owner_id=2 in today's `memory_dream_usage` bucket need a
one-shot SQL fix —
UPDATE memory_dream_usage
SET jobs_started = jobs_succeeded + jobs_failed
WHERE date = date('now')
AND owner_id = '2';
or simply wait for the next UTC-midnight reset.
Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 4.7
parent
205093e654
commit
0bb2b6500a
@@ -427,10 +427,13 @@ func (w *ConsolidatorWorker) ForceRunN(ctx context.Context, ownerID, jobType str
|
||||
// launchOne wires up the per-job tokens / harness dispatch for an
|
||||
// already-Created job row. Shared by ForceRunN and the worker's
|
||||
// internal tryDispatch path.
|
||||
//
|
||||
// jobs_started is only incremented after a successful Dispatch flip so
|
||||
// pre-dispatch failures (token issue, agent lookup, dispatch race) do
|
||||
// not consume daily-job-limit slots — they are surfaced through the
|
||||
// job row's status=failed and the jobs_failed counter via the caller
|
||||
// of Complete, not through the gate's jobs_started counter.
|
||||
func (w *ConsolidatorWorker) launchOne(ctx context.Context, ownerID, jobType string, jobID int64) error {
|
||||
if w.usage != nil {
|
||||
_ = w.usage.RecordStart(ctx, ownerID)
|
||||
}
|
||||
tok, _, err := w.tokens.Issue(ctx, ownerID, jobID)
|
||||
if err != nil {
|
||||
_ = w.jobs.Complete(ctx, jobID, JobStatusFailed, "", "token issue: "+err.Error())
|
||||
@@ -446,6 +449,9 @@ func (w *ConsolidatorWorker) launchOne(ctx context.Context, ownerID, jobType str
|
||||
_ = w.jobs.Complete(ctx, jobID, JobStatusFailed, "", "dispatch flip: "+err.Error())
|
||||
return fmt.Errorf("dispatch flip: %w", err)
|
||||
}
|
||||
if w.usage != nil {
|
||||
_ = w.usage.RecordStart(ctx, ownerID)
|
||||
}
|
||||
w.wg.Add(1)
|
||||
go func() {
|
||||
defer w.wg.Done()
|
||||
@@ -512,9 +518,6 @@ func (w *ConsolidatorWorker) tryDispatch(ctx context.Context, ownerID, jobType,
|
||||
)
|
||||
return
|
||||
}
|
||||
if w.usage != nil {
|
||||
_ = w.usage.RecordStart(ctx, ownerID)
|
||||
}
|
||||
tok, _, err := w.tokens.Issue(ctx, ownerID, jobID)
|
||||
if err != nil {
|
||||
w.logger.Warn("issue token failed", "job_id", jobID, "error", err)
|
||||
@@ -539,6 +542,12 @@ func (w *ConsolidatorWorker) tryDispatch(ctx context.Context, ownerID, jobType,
|
||||
_ = w.jobs.Complete(ctx, jobID, JobStatusFailed, "", "dispatch flip: "+err.Error())
|
||||
return
|
||||
}
|
||||
// Only count jobs_started after a successful Dispatch flip so
|
||||
// pre-dispatch failures (token issue, agent lookup, race) don't
|
||||
// burn daily-job-limit slots without actually running anything.
|
||||
if w.usage != nil {
|
||||
_ = w.usage.RecordStart(ctx, ownerID)
|
||||
}
|
||||
|
||||
w.wg.Add(1)
|
||||
go func() {
|
||||
|
||||
@@ -180,6 +180,96 @@ func TestConsolidator_WallclockTerminatesRunaway(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// TestConsolidator_PreDispatchFailureDoesNotBurnJobsStarted verifies
|
||||
// that when launchOne fails before the Dispatch flip (e.g. agent
|
||||
// lookup fails), the per-(owner, day) jobs_started counter is NOT
|
||||
// incremented. Otherwise, repeated pre-dispatch failures inflate the
|
||||
// counter without any actual runs and eventually trip the
|
||||
// jobs_exceeded circuit breaker, wedging the worker for the rest of
|
||||
// the UTC day.
|
||||
func TestConsolidator_PreDispatchFailureDoesNotBurnJobsStarted(t *testing.T) {
|
||||
h := &stubHarness{}
|
||||
// Pass nil agent → stubAgentLookup.GetAgent returns "not found".
|
||||
// launchOne will Complete the job as failed before reaching Dispatch.
|
||||
cfg := MemoryConfig{
|
||||
DreamEnabled: true,
|
||||
DreamWatermark: 1,
|
||||
DreamMaxConcurrent: 1,
|
||||
DreamParallel: 1,
|
||||
DreamWallclockBudget: 100 * time.Millisecond,
|
||||
DreamAgent: "claude-code",
|
||||
DreamDailyJobLimit: 2,
|
||||
}
|
||||
w, db := newWorkerForTest(t, h, nil, cfg)
|
||||
usage := NewDreamUsageStore(db)
|
||||
gate := NewUsageGate(cfg, usage)
|
||||
w.SetUsageGate(usage, gate)
|
||||
|
||||
ctx := context.Background()
|
||||
owner := "1"
|
||||
|
||||
// Three attempts to ForceRun. Each one fails inside launchOne at the
|
||||
// agent-lookup step. None of them actually dispatch, so none should
|
||||
// count against DreamDailyJobLimit.
|
||||
for i := 0; i < 3; i++ {
|
||||
_, _ = w.ForceRun(ctx, owner, JobTypeReflection)
|
||||
}
|
||||
|
||||
u, err := usage.Today(ctx, owner)
|
||||
if err != nil {
|
||||
t.Fatalf("Today: %v", err)
|
||||
}
|
||||
if u.JobsStarted != 0 {
|
||||
t.Errorf("jobs_started must stay 0 when no dispatch flip succeeded; got %d", u.JobsStarted)
|
||||
}
|
||||
|
||||
allowed, reason, _ := gate.Allow(ctx, owner)
|
||||
if !allowed {
|
||||
t.Errorf("gate should still allow after pre-dispatch failures; got denied (%s)", reason)
|
||||
}
|
||||
}
|
||||
|
||||
// TestConsolidator_DispatchSuccessIncrementsJobsStarted is the positive
|
||||
// counterpart: when launchOne succeeds end-to-end (token issue + agent
|
||||
// lookup + dispatch flip), jobs_started IS incremented so the daily
|
||||
// limit is enforced correctly.
|
||||
func TestConsolidator_DispatchSuccessIncrementsJobsStarted(t *testing.T) {
|
||||
h := &stubHarness{}
|
||||
agent := DreamAgentNamed{Name: "claude-code"}
|
||||
cfg := MemoryConfig{
|
||||
DreamEnabled: true,
|
||||
DreamWatermark: 1,
|
||||
DreamMaxConcurrent: 1,
|
||||
DreamParallel: 1,
|
||||
DreamWallclockBudget: 500 * time.Millisecond,
|
||||
DreamAgent: "claude-code",
|
||||
DreamDailyJobLimit: 10,
|
||||
}
|
||||
w, db := newWorkerForTest(t, h, agent, cfg)
|
||||
usage := NewDreamUsageStore(db)
|
||||
gate := NewUsageGate(cfg, usage)
|
||||
w.SetUsageGate(usage, gate)
|
||||
|
||||
ctx := context.Background()
|
||||
owner := "1"
|
||||
|
||||
ids, err := w.ForceRun(ctx, owner, JobTypeReflection)
|
||||
if err != nil {
|
||||
t.Fatalf("ForceRun: %v", err)
|
||||
}
|
||||
if ids == 0 {
|
||||
t.Fatalf("ForceRun returned zero job id")
|
||||
}
|
||||
// Wait for the async runJob goroutine to settle so the counter
|
||||
// snapshot is stable. The stub harness returns immediately.
|
||||
w.Stop()
|
||||
|
||||
u, _ := usage.Today(ctx, owner)
|
||||
if u.JobsStarted != 1 {
|
||||
t.Errorf("jobs_started: want 1 after a successful dispatch, got %d", u.JobsStarted)
|
||||
}
|
||||
}
|
||||
|
||||
// seedMemoryWithChannel creates the open-brain channel + an agent
|
||||
// owned by owner_id=1 + one message.
|
||||
func seedMemoryWithChannel(t *testing.T, db *sql.DB, agentName string, channelID int64, body string) int64 {
|
||||
|
||||
Reference in New Issue
Block a user