From c24d53156d9b20056eb65a69935568b963c4cf7a Mon Sep 17 00:00:00 2001 From: brokemac79 Date: Wed, 15 Jul 2026 18:05:01 +0100 Subject: [PATCH] fix(queue): unblock exact review publication --- .github/workflows/sweep.yml | 2 +- dashboard/exact-review-health.ts | 7 +- dashboard/worker.ts | 169 +++++++++++++++++++++++-------- test/dashboard-worker.test.ts | 95 ++++++++++++++++- test/exact-review-health.test.ts | 19 ++++ test/sweep-workflow.test.ts | 5 +- 6 files changed, 252 insertions(+), 45 deletions(-) diff --git a/.github/workflows/sweep.yml b/.github/workflows/sweep.yml index 838aa04976..02a772c8e5 100644 --- a/.github/workflows/sweep.yml +++ b/.github/workflows/sweep.yml @@ -1131,7 +1131,7 @@ jobs: runs-on: ubuntu-latest timeout-minutes: 60 concurrency: - group: clawsweeper-state-publisher + group: clawsweeper-exact-review-publisher cancel-in-progress: false queue: max permissions: diff --git a/dashboard/exact-review-health.ts b/dashboard/exact-review-health.ts index c83eb8f1b6..c5daa584c1 100644 --- a/dashboard/exact-review-health.ts +++ b/dashboard/exact-review-health.ts @@ -163,7 +163,7 @@ function exactReviewPhaseStartedAt( const dispatchedAt = validTimestamp(item.dispatchedAt); const leaseExpiresAt = validTimestamp(item.leaseExpiresAt); const leaseStartedAt = - leaseExpiresAt === null ? null : validTimestamp(leaseExpiresAt - dispatchLeaseMs); + leaseExpiresAt === null ? null : timestampAtOrBefore(leaseExpiresAt - dispatchLeaseMs, now); // Rolling deploys can expose rows created before dispatchedAt existed, while a rollback can // leave an old dispatchedAt behind. The current lease start is the reliable compatibility // marker; prefer the newest plausible transition and keep an unknown age non-alarming. @@ -194,6 +194,11 @@ function validTimestamp(value: unknown): number | null { return Number.isFinite(number) && number > 0 && number <= 8_640_000_000_000_000 ? number : null; } +function timestampAtOrBefore(value: unknown, maximum: number): number | null { + const timestamp = validTimestamp(value); + return timestamp !== null && timestamp <= maximum ? timestamp : null; +} + function finiteTimestamp(value: unknown, fallback: number) { return validTimestamp(value) ?? fallback; } diff --git a/dashboard/worker.ts b/dashboard/worker.ts index e4b2251a73..cefbdb0944 100644 --- a/dashboard/worker.ts +++ b/dashboard/worker.ts @@ -154,9 +154,9 @@ const HEALTH_HISTORY_KEY_PREFIX = "health-history:"; const DEFAULT_EXACT_REVIEW_QUEUE_MAX_CONCURRENT = 64; const DEFAULT_EXACT_REVIEW_TARGET_MAX_CONCURRENT = 60; const DEFAULT_EXACT_REVIEW_DISPATCH_LEASE_MS = 6 * 60 * 1000; -// Covers the global publisher lane's 100 queued 60-minute jobs plus scheduling margin. -// Terminal-run reconciliation still releases failed or cancelled dispatches early. -const DEFAULT_EXACT_REVIEW_PUBLICATION_DISPATCH_LEASE_MS = 7 * 24 * 60 * 60 * 1000; +// Exact publications have a dedicated serial lane. Bound the unclaimed handoff so a run that +// never reaches its claim step is re-dispatched; stale runs lose the lease tuple safely. +const DEFAULT_EXACT_REVIEW_PUBLICATION_DISPATCH_LEASE_MS = 15 * 60 * 1000; const DEFAULT_EXACT_REVIEW_EXECUTION_LEASE_MS = 130 * 60 * 1000; const DEFAULT_EXACT_REVIEW_RETRY_MS = 30_000; const DEFAULT_EXACT_REVIEW_WORKFLOW_PAUSED_RETRY_MS = 60_000; @@ -580,7 +580,11 @@ export class ExactReviewQueue { const state = this.readStateSync(); // A delayed or lost alarm must not let an expired one-shot recovery // suppress the next failed shard's recovery delivery. - reclaimExpiredExactReviewLeases(state, now); + reclaimExpiredExactReviewLeases( + state, + now, + exactReviewPublicationDispatchLeaseMs(this.env), + ); const key = exactReviewItemKey(decision); const current = state.items[key]; const nextAttemptAt = exactReviewQueueEnqueueAttemptAt(state, now); @@ -645,11 +649,25 @@ export class ExactReviewQueue { const now = Date.now(); const state = this.readStateSync(); const item = tupleClaim ? state.items[itemKey] : exactReviewItemForLease(state, leaseId); + if ( + item && + reclaimExpiredExactReviewLease( + state, + item.key, + item, + now, + exactReviewPublicationDispatchLeaseMs(this.env), + ) + ) { + this.writeStateSync(state); + await this.scheduleNext(state, now); + return json({ error: "lease_not_active" }, 409); + } if ( !item || item.leaseId !== leaseId || (tupleClaim && item.leaseRevision !== leaseRevision) || - !isLiveExactReviewLease(item, now) + !isLiveExactReviewLease(item, now, exactReviewPublicationDispatchLeaseMs(this.env)) ) { return json({ error: "lease_not_active" }, 409); } @@ -876,7 +894,11 @@ export class ExactReviewQueue { const current = this.readStateSync(); // Dashboard reads are also the operational heartbeat. Reclaim leases and // restore the alarm here so a deploy or lost alarm cannot strand backlog. - const changed = reclaimExpiredExactReviewLeases(current, now); + const changed = reclaimExpiredExactReviewLeases( + current, + now, + exactReviewPublicationDispatchLeaseMs(this.env), + ); if (changed) this.writeStateSync(current); else this.syncLegacyCompatibilitySync(current); return current; @@ -890,6 +912,7 @@ export class ExactReviewQueue { exactReviewTargetCapacity(this.env), exactReviewDispatchLeaseMs(this.env), exactReviewExecutionLeaseMs(this.env), + exactReviewPublicationDispatchLeaseMs(this.env), ), delivery_receipts: this.deliveryReceiptCountSync(), storage_schema_version: EXACT_REVIEW_QUEUE_STORAGE_SCHEMA_VERSION, @@ -912,7 +935,11 @@ export class ExactReviewQueue { this.syncLegacyCompatibilitySync(this.readStateSync()); }); const snapshot = this.readStateSync(); - const reclaimedSnapshot = reclaimExpiredExactReviewLeases(snapshot, startedAt); + const reclaimedSnapshot = reclaimExpiredExactReviewLeases( + snapshot, + startedAt, + exactReviewPublicationDispatchLeaseMs(this.env), + ); const expiredSnapshot = expireExactReviewPublicationItems(snapshot, startedAt); const snapshotChanged = reclaimedSnapshot || expiredSnapshot; const capacity = exactReviewQueueCapacity(this.env); @@ -943,7 +970,7 @@ export class ExactReviewQueue { // write so concurrent enqueue, claim, or complete requests cannot be lost. const now = Date.now(); const state = this.readStateSync(); - reclaimExpiredExactReviewLeases(state, now); + reclaimExpiredExactReviewLeases(state, now, exactReviewPublicationDispatchLeaseMs(this.env)); expireExactReviewPublicationItems(state, now); const admitted = exactReviewQueueAdmittedItems(state, now, capacity, targetCapacity); if (!preflight.ok) { @@ -985,10 +1012,7 @@ export class ExactReviewQueue { item.leaseExpiresAt = now + (item.decision.sourceAction === EXACT_REVIEW_ARTIFACT_PUBLISH_SOURCE_ACTION - ? Math.max( - exactReviewDispatchLeaseMs(this.env), - DEFAULT_EXACT_REVIEW_PUBLICATION_DISPATCH_LEASE_MS, - ) + ? exactReviewPublicationDispatchLeaseMs(this.env) : exactReviewDispatchLeaseMs(this.env)); item.claimedRunId = undefined; item.claimedRunAttempt = undefined; @@ -1557,6 +1581,7 @@ export class ExactReviewQueue { now, exactReviewQueueCapacity(this.env), exactReviewTargetCapacity(this.env), + exactReviewPublicationDispatchLeaseMs(this.env), ); if (next === null) { await this.storage.deleteAlarm(); @@ -1617,9 +1642,12 @@ export default { return json({ error: "not_found" }, 404); }, async scheduled(_controller, env: DashboardEnv = {}, ctx?: DashboardContext) { - const recording = recordScheduledHealthSample(env); - if (ctx?.waitUntil) ctx.waitUntil(recording); - else await recording; + const maintenance = Promise.all([ + recordScheduledHealthSample(env), + exactReviewQueueStatusSnapshot(env).catch(() => null), + ]); + if (ctx?.waitUntil) ctx.waitUntil(maintenance); + else await maintenance; }, }; @@ -2859,37 +2887,76 @@ function clearExactReviewLease(item: ExactReviewQueueItem) { item.claimedAt = undefined; } -function isLiveExactReviewLease(item: ExactReviewQueueItem, now: number) { - return Boolean(item.leaseId && item.leaseExpiresAt && item.leaseExpiresAt > now); +function exactReviewEffectiveLeaseExpiresAt( + item: ExactReviewQueueItem, + publicationDispatchLeaseMs: number, +) { + const leaseExpiresAt = Number(item.leaseExpiresAt || 0); + if ( + !leaseExpiresAt || + item.state !== "dispatching" || + !exactReviewQueueIsPublication(item) || + item.claimedRunId || + !item.dispatchedAt + ) { + return leaseExpiresAt; + } + return Math.min(leaseExpiresAt, item.dispatchedAt + publicationDispatchLeaseMs); } -function reclaimExpiredExactReviewLeases(state: ExactReviewQueueState, now: number) { +function isLiveExactReviewLease( + item: ExactReviewQueueItem, + now: number, + publicationDispatchLeaseMs = DEFAULT_EXACT_REVIEW_PUBLICATION_DISPATCH_LEASE_MS, +) { + return Boolean( + item.leaseId && exactReviewEffectiveLeaseExpiresAt(item, publicationDispatchLeaseMs) > now, + ); +} + +function reclaimExpiredExactReviewLeases( + state: ExactReviewQueueState, + now: number, + publicationDispatchLeaseMs = DEFAULT_EXACT_REVIEW_PUBLICATION_DISPATCH_LEASE_MS, +) { let changed = false; for (const [key, item] of Object.entries(state.items)) { - if ( - (item.state === "dispatching" || item.state === "leased") && - !isLiveExactReviewLease(item, now) - ) { - const oneShotRecovery = - (item.leaseDecision || item.decision).sourceAction === - FAILED_REVIEW_SHARD_RECOVERY_SOURCE_ACTION; - const hasNewerRevision = item.revision > Number(item.leaseRevision || 0); - if (oneShotRecovery && !hasNewerRevision) { - delete state.items[key]; - changed = true; - continue; - } - clearExactReviewLease(item); - item.state = "pending"; - item.nextAttemptAt = now; - if (hasNewerRevision) item.attempts = 0; - item.updatedAt = now; + if (reclaimExpiredExactReviewLease(state, key, item, now, publicationDispatchLeaseMs)) { changed = true; } } return changed; } +function reclaimExpiredExactReviewLease( + state: ExactReviewQueueState, + key: string, + item: ExactReviewQueueItem, + now: number, + publicationDispatchLeaseMs: number, +) { + if ( + (item.state !== "dispatching" && item.state !== "leased") || + isLiveExactReviewLease(item, now, publicationDispatchLeaseMs) + ) { + return false; + } + const oneShotRecovery = + (item.leaseDecision || item.decision).sourceAction === + FAILED_REVIEW_SHARD_RECOVERY_SOURCE_ACTION; + const hasNewerRevision = item.revision > Number(item.leaseRevision || 0); + if (oneShotRecovery && !hasNewerRevision) { + delete state.items[key]; + return true; + } + clearExactReviewLease(item); + item.state = "pending"; + item.nextAttemptAt = now; + if (hasNewerRevision) item.attempts = 0; + item.updatedAt = now; + return true; +} + function expireExactReviewPublicationItems(state: ExactReviewQueueState, now: number) { let changed = false; for (const [key, item] of Object.entries(state.items)) { @@ -3011,6 +3078,7 @@ function exactReviewQueueStats( targetCapacity = Number.POSITIVE_INFINITY, dispatchLeaseMs = DEFAULT_EXACT_REVIEW_DISPATCH_LEASE_MS, executionLeaseMs = DEFAULT_EXACT_REVIEW_EXECUTION_LEASE_MS, + publicationDispatchLeaseMs = DEFAULT_EXACT_REVIEW_PUBLICATION_DISPATCH_LEASE_MS, ) { const items = Object.values(state.items); const handoffHealth = summarizeExactReviewHandoff({ @@ -3068,7 +3136,13 @@ function exactReviewQueueStats( right.dispatching + right.leased - (left.dispatching + left.leased) || left.target_repo.localeCompare(right.target_repo), ); - const nextWakeAt = exactReviewQueueNextWakeAt(state, now, capacity, targetCapacity); + const nextWakeAt = exactReviewQueueNextWakeAt( + state, + now, + capacity, + targetCapacity, + publicationDispatchLeaseMs, + ); return { pending: handoffHealth.phases.pending.count, dispatching: handoffHealth.phases.dispatching.count, @@ -3099,6 +3173,7 @@ function exactReviewQueueNextWakeAt( now: number, capacity = Number.POSITIVE_INFINITY, targetCapacity = Number.POSITIVE_INFINITY, + publicationDispatchLeaseMs = DEFAULT_EXACT_REVIEW_PUBLICATION_DISPATCH_LEASE_MS, ) { const items = Object.values(state.items); if (!items.length) return null; @@ -3109,7 +3184,13 @@ function exactReviewQueueNextWakeAt( const activeItems = items.filter( (item) => item.state === "dispatching" || item.state === "leased", ); - if (activeItems.some((item) => !item.leaseExpiresAt || item.leaseExpiresAt <= now)) { + if ( + activeItems.some( + (item) => + !item.leaseExpiresAt || + exactReviewEffectiveLeaseExpiresAt(item, publicationDispatchLeaseMs) <= now, + ) + ) { return now + 1_000; } const activeReviews = activeItems.filter((item) => !exactReviewQueueIsPublication(item)); @@ -3118,7 +3199,7 @@ function exactReviewQueueNextWakeAt( .map((item) => item.leaseExpiresAt) .filter((value): value is number => Boolean(value && value > now)); const activePublisherWakeAt = activePublishers - .map((item) => item.leaseExpiresAt) + .map((item) => exactReviewEffectiveLeaseExpiresAt(item, publicationDispatchLeaseMs)) .filter((value): value is number => Boolean(value && value > now)); const activeTargetWakeAt = new Map(); const activeTargetCounts = new Map(); @@ -3159,7 +3240,8 @@ function exactReviewQueueNextWakeAt( ), ]; } - return item.leaseExpiresAt ? [item.leaseExpiresAt] : []; + const leaseExpiresAt = exactReviewEffectiveLeaseExpiresAt(item, publicationDispatchLeaseMs); + return leaseExpiresAt ? [leaseExpiresAt] : []; }); if (!times.length) return now + DEFAULT_EXACT_REVIEW_RETRY_MS; return Math.max(now + 1_000, Math.min(...times)); @@ -3195,6 +3277,13 @@ function exactReviewDispatchLeaseMs(env) { ); } +function exactReviewPublicationDispatchLeaseMs(env) { + return Math.max( + exactReviewDispatchLeaseMs(env), + DEFAULT_EXACT_REVIEW_PUBLICATION_DISPATCH_LEASE_MS, + ); +} + function exactReviewExecutionLeaseMs(env) { return Math.max( 60_000, diff --git a/test/dashboard-worker.test.ts b/test/dashboard-worker.test.ts index d19ff209d5..6978b7b7bf 100644 --- a/test/dashboard-worker.test.ts +++ b/test/dashboard-worker.test.ts @@ -2751,6 +2751,7 @@ test("exact-review queue keeps publication artifacts durable outside review capa { createdAt: number; decision: Record; + dispatchedAt?: number; leaseDecision?: Record; leaseExpiresAt?: number; leaseId?: string; @@ -2821,7 +2822,8 @@ test("exact-review queue keeps publication artifacts durable outside review capa assert.equal(afterExpiry.items["openclaw/gogcli#803@publish:102:1"], undefined); const reservedPublisher = afterExpiry.items["openclaw/gogcli#801@publish:100:1"]; assert.equal(reservedPublisher.state, "dispatching"); - assert.ok((reservedPublisher.leaseExpiresAt ?? 0) - Date.now() > 6 * 24 * 60 * 60_000); + assert.ok((reservedPublisher.leaseExpiresAt ?? 0) - Date.now() > 14 * 60_000); + assert.ok((reservedPublisher.leaseExpiresAt ?? 0) - Date.now() <= 15 * 60_000); const publicationPayload = dispatched.find( (payload) => (payload.client_payload as Record).source_action === @@ -2839,6 +2841,97 @@ test("exact-review queue keeps publication artifacts durable outside review capa > ).producerDecision, ); + + const firstPublicationLease = { + leaseId: reservedPublisher.leaseId, + leaseRevision: reservedPublisher.leaseRevision, + }; + reservedPublisher.dispatchedAt = Date.now() - 16 * 60_000; + reservedPublisher.leaseExpiresAt = Date.now() + 7 * 24 * 60 * 60_000; + await storage.put("exact-review-queue", afterExpiry); + + let queueMaintenance: Promise | undefined; + await worker.scheduled( + {}, + { EXACT_REVIEW_QUEUE: new MemoryDurableNamespace(queue) }, + { waitUntil: (promise) => (queueMaintenance = promise) }, + ); + await queueMaintenance; + const afterScheduledMaintenance = (await storage.get("exact-review-queue")) as typeof state; + const scheduledPublisher = afterScheduledMaintenance.items["openclaw/gogcli#801@publish:100:1"]; + assert.equal(scheduledPublisher.state, "pending"); + assert.equal(scheduledPublisher.leaseId, undefined); + assert.ok(((await storage.getAlarm()) ?? Number.POSITIVE_INFINITY) <= Date.now() + 1_000); + + await queue.alarm(); + const afterScheduledRedispatch = (await storage.get("exact-review-queue")) as typeof state; + const scheduledRedispatch = afterScheduledRedispatch.items["openclaw/gogcli#801@publish:100:1"]; + assert.equal(scheduledRedispatch.state, "dispatching"); + assert.notEqual(scheduledRedispatch.leaseId, firstPublicationLease.leaseId); + assert.equal( + dispatched.filter( + (payload) => + (payload.client_payload as Record).source_action === + "exact_review_artifact_publish", + ).length, + 2, + ); + + const scheduledPublicationLease = { + leaseId: scheduledRedispatch.leaseId, + leaseRevision: scheduledRedispatch.leaseRevision, + }; + scheduledRedispatch.dispatchedAt = Date.now() - 16 * 60_000; + scheduledRedispatch.leaseExpiresAt = Date.now() + 7 * 24 * 60 * 60_000; + await storage.put("exact-review-queue", afterScheduledRedispatch); + + const expiredLegacyClaim = await queue.fetch( + new Request("https://clawsweeper-exact-review-queue/claim", { + method: "POST", + body: JSON.stringify({ + lease_id: scheduledPublicationLease.leaseId, + item_key: "openclaw/gogcli#801@publish:100:1", + lease_revision: scheduledPublicationLease.leaseRevision, + run_id: "999998", + run_attempt: 1, + }), + }), + ); + assert.equal(expiredLegacyClaim.status, 409); + assert.deepEqual(await expiredLegacyClaim.json(), { error: "lease_not_active" }); + const afterExpiredClaim = (await storage.get("exact-review-queue")) as typeof state; + const reclaimedPublisher = afterExpiredClaim.items["openclaw/gogcli#801@publish:100:1"]; + assert.equal(reclaimedPublisher.state, "pending"); + assert.equal(reclaimedPublisher.leaseId, undefined); + assert.ok(((await storage.getAlarm()) ?? Number.POSITIVE_INFINITY) <= Date.now() + 1_000); + + await queue.alarm(); + + const publicationDispatches = dispatched.filter( + (payload) => + (payload.client_payload as Record).source_action === + "exact_review_artifact_publish", + ); + assert.equal(publicationDispatches.length, 3); + const afterRedispatch = (await storage.get("exact-review-queue")) as typeof state; + const redispatchedPublisher = afterRedispatch.items["openclaw/gogcli#801@publish:100:1"]; + assert.equal(redispatchedPublisher.state, "dispatching"); + assert.notEqual(redispatchedPublisher.leaseId, scheduledPublicationLease.leaseId); + assert.ok((redispatchedPublisher.leaseExpiresAt ?? 0) - Date.now() > 14 * 60_000); + const staleClaim = await queue.fetch( + new Request("https://clawsweeper-exact-review-queue/claim", { + method: "POST", + body: JSON.stringify({ + lease_id: firstPublicationLease.leaseId, + item_key: "openclaw/gogcli#801@publish:100:1", + lease_revision: firstPublicationLease.leaseRevision, + run_id: "999999", + run_attempt: 1, + }), + }), + ); + assert.equal(staleClaim.status, 409); + assert.deepEqual(await staleClaim.json(), { error: "lease_not_active" }); } finally { globalThis.fetch = originalFetch; } diff --git a/test/exact-review-health.test.ts b/test/exact-review-health.test.ts index f185cb0be3..8889a87787 100644 --- a/test/exact-review-health.test.ts +++ b/test/exact-review-health.test.ts @@ -125,6 +125,25 @@ test("exact-review handoff health derives legacy dispatch age from its active le assert.equal(health.phases.dispatching.oldest_age_seconds, 20); }); +test("exact-review handoff health uses dispatch time when a longer lease has a future derived start", () => { + const health = summarize({ + items: [ + { + state: "dispatching", + createdAt: NOW - 10 * 60_000, + updatedAt: NOW - 6 * 60_000, + dispatchedAt: NOW - 6 * 60_000, + leaseExpiresAt: NOW + 15 * 60_000, + }, + ], + }); + + assert.equal(health.status, "stalled"); + assert.equal(health.reason, "claim_stalled"); + assert.equal(health.phases.dispatching.oldest_at, "2026-07-13T01:54:00.000Z"); + assert.equal(health.phases.dispatching.oldest_age_seconds, 360); +}); + test("exact-review handoff health ignores stale rollback telemetry and unknown legacy ages", () => { const rolledBack = summarize({ items: [ diff --git a/test/sweep-workflow.test.ts b/test/sweep-workflow.test.ts index 5eee7ff4b8..3536c8823d 100644 --- a/test/sweep-workflow.test.ts +++ b/test/sweep-workflow.test.ts @@ -323,7 +323,7 @@ test("scheduled review shards receive the compiler-backed runtime artifact", () assert.doesNotMatch(reviewJob, /npm pack "@typescript/); }); -test("exact event review hands immutable artifacts to one state publisher", () => { +test("exact event review hands immutable artifacts to one dedicated publisher", () => { type Step = { "continue-on-error"?: boolean; name?: string; @@ -401,11 +401,12 @@ test("exact event review hands immutable artifacts to one state publisher", () = step(publisher, "Claim durable exact review publication").run ?? "", /internal\/exact-review\/claim/, ); - assert.equal(publisher.concurrency?.group, "clawsweeper-state-publisher"); + assert.equal(publisher.concurrency?.group, "clawsweeper-exact-review-publisher"); assert.equal(publisher.concurrency?.["cancel-in-progress"], false); assert.equal(publisher.concurrency?.queue, "max"); assert.equal(publisher.permissions?.actions, "write"); assert.equal(batchPublisher.concurrency?.group, "clawsweeper-state-publisher"); + assert.notEqual(publisher.concurrency?.group, batchPublisher.concurrency?.group); const publicationContext = step(publisher, "Claim durable exact review publication"); assert.match( publicationContext.run ?? "",