diff --git a/server/src/__tests__/heartbeat-issue-liveness-escalation.test.ts b/server/src/__tests__/heartbeat-issue-liveness-escalation.test.ts index 60a0203276..bcfadd6b6e 100644 --- a/server/src/__tests__/heartbeat-issue-liveness-escalation.test.ts +++ b/server/src/__tests__/heartbeat-issue-liveness-escalation.test.ts @@ -66,6 +66,7 @@ import { heartbeatService } from "../services/heartbeat.ts"; import { instanceSettingsService } from "../services/instance-settings.ts"; import { issueService } from "../services/issues.ts"; import { runningProcesses } from "../adapters/index.ts"; +import { DEFAULT_LIVENESS_REESCALATION_COOLDOWN_MS } from "../services/recovery/service.ts"; const embeddedPostgresSupport = await getEmbeddedPostgresTestSupport(); const describeEmbeddedPostgres = embeddedPostgresSupport.supported ? describe : describe.skip; @@ -1153,10 +1154,11 @@ describeEmbeddedPostgres("heartbeat issue graph liveness escalation", () => { ); }); - it("creates a fresh escalation when the previous matching escalation is terminal", async () => { + it("holds a recently closed matching escalation, then re-escalates after the cooldown", async () => { await enableAutoRecovery(); const { companyId, managerId, blockedIssueId, blockerIssueId } = await seedBlockedChain(); const heartbeat = heartbeatService(db); + const now = new Date(); const incidentKey = [ "harness_liveness", companyId, @@ -1178,9 +1180,18 @@ describeEmbeddedPostgres("heartbeat issue graph liveness escalation", () => { identifier: "CLOSED-3", originKind: "harness_liveness_escalation", originId: incidentKey, + createdAt: new Date(now.getTime() - 30 * 60 * 1000), + updatedAt: now, }); - const result = await heartbeat.reconcileIssueGraphLiveness(); + const held = await heartbeat.reconcileIssueGraphLiveness({ now }); + + expect(held.escalationsCreated).toBe(0); + expect(held.skippedReescalationCooldown).toBe(1); + + const result = await heartbeat.reconcileIssueGraphLiveness({ + now: new Date(now.getTime() + DEFAULT_LIVENESS_REESCALATION_COOLDOWN_MS + 1), + }); expect(result.escalationsCreated).toBe(1); expect(result.existingEscalations).toBe(0); @@ -1211,6 +1222,41 @@ describeEmbeddedPostgres("heartbeat issue graph liveness escalation", () => { expect(blockers.some((row) => row.blockerIssueId === freshEscalation?.id)).toBe(true); }); + it("re-escalates immediately after a matching escalation is cancelled", async () => { + await enableAutoRecovery(); + const { companyId, managerId, blockedIssueId, blockerIssueId } = await seedBlockedChain(); + const heartbeat = heartbeatService(db); + const now = new Date(); + const incidentKey = [ + "harness_liveness", + companyId, + blockedIssueId, + "blocked_by_unassigned_issue", + blockerIssueId, + ].join(":"); + + await db.insert(issues).values({ + id: randomUUID(), + companyId, + title: "Cancelled escalation", + status: "cancelled", + priority: "high", + parentId: blockedIssueId, + assigneeAgentId: managerId, + issueNumber: 3, + identifier: "CANCELLED-3", + originKind: "harness_liveness_escalation", + originId: incidentKey, + createdAt: new Date(now.getTime() - 30 * 60 * 1000), + updatedAt: now, + }); + + const result = await heartbeat.reconcileIssueGraphLiveness({ now }); + + expect(result.escalationsCreated).toBe(1); + expect(result.skippedReescalationCooldown).toBe(0); + }); + it("removes closed liveness escalations from blocker relations during reconciliation", async () => { await enableAutoRecovery(); const { companyId, blockedIssueId, blockerIssueId } = await seedBlockedChain(); diff --git a/server/src/__tests__/productivity-review-service.test.ts b/server/src/__tests__/productivity-review-service.test.ts index d58d2aac1b..9dbf784923 100644 --- a/server/src/__tests__/productivity-review-service.test.ts +++ b/server/src/__tests__/productivity-review-service.test.ts @@ -255,7 +255,44 @@ describeEmbeddedPostgres("productivity review service", () => { expect(await listRefreshComments(review!.id)).toHaveLength(DEFAULT_PRODUCTIVITY_REVIEW_MAX_REFRESH_COMMENTS); }); - it("caps productivity review creation per source issue in the rolling creation window", async () => { + it("allows only one productivity review per source issue in 24 hours", async () => { + const now = new Date("2026-04-28T12:00:00.000Z"); + const seeded = await seedAssignedIssue(); + await insertRuns({ + companyId: seeded.companyId, + agentId: seeded.coderId, + issueId: seeded.issueId, + count: DEFAULT_PRODUCTIVITY_REVIEW_NO_COMMENT_STREAK_RUNS, + now, + }); + const createdAt = new Date(now.getTime() - 8 * 60 * 60 * 1000); + await db.insert(issues).values({ + id: randomUUID(), + companyId: seeded.companyId, + title: "Completed productivity review", + status: "done", + priority: "high", + originKind: PRODUCTIVITY_REVIEW_ORIGIN_KIND, + originId: seeded.issueId, + originFingerprint: `productivity-review:${seeded.issueId}`, + parentId: seeded.issueId, + issueNumber: 2, + identifier: `${seeded.issuePrefix}-2`, + createdAt, + updatedAt: createdAt, + }); + + const result = await productivityReviewService(db).reconcileProductivityReviews({ + now, + companyId: seeded.companyId, + }); + + expect(result.created).toBe(0); + expect(result.creationCapped).toBe(1); + expect(await listProductivityReviews(seeded.companyId)).toHaveLength(1); + }); + + it("suppresses creation after three consecutive completed reviews with no source action", async () => { const now = new Date("2026-04-28T12:00:00.000Z"); const seeded = await seedAssignedIssue(); await insertRuns({ @@ -266,12 +303,12 @@ describeEmbeddedPostgres("productivity review service", () => { now, }); await db.insert(issues).values( - [8, 9, 10].map((hoursAgo, index) => { + [96, 72, 48].map((hoursAgo, index) => { const createdAt = new Date(now.getTime() - hoursAgo * 60 * 60 * 1000); return { id: randomUUID(), companyId: seeded.companyId, - title: `Completed productivity review ${index + 1}`, + title: `No-action productivity review ${index + 1}`, status: "done", priority: "high", originKind: PRODUCTIVITY_REVIEW_ORIGIN_KIND, @@ -281,7 +318,7 @@ describeEmbeddedPostgres("productivity review service", () => { issueNumber: index + 2, identifier: `${seeded.issuePrefix}-${index + 2}`, createdAt, - updatedAt: createdAt, + updatedAt: new Date(createdAt.getTime() + 60 * 60 * 1000), }; }), ); @@ -292,10 +329,62 @@ describeEmbeddedPostgres("productivity review service", () => { }); expect(result.created).toBe(0); - expect(result.creationCapped).toBe(1); + expect(result.noActionSuppressed).toBe(1); expect(await listProductivityReviews(seeded.companyId)).toHaveLength(3); }); + it("resets no-action suppression for source action after a zero-duration review", async () => { + const now = new Date("2026-04-28T12:00:00.000Z"); + const seeded = await seedAssignedIssue(); + await insertRuns({ + companyId: seeded.companyId, + agentId: seeded.coderId, + issueId: seeded.issueId, + count: DEFAULT_PRODUCTIVITY_REVIEW_NO_COMMENT_STREAK_RUNS, + now, + }); + const reviewWindows = [96, 72, 48].map((hoursAgo, index) => { + const createdAt = new Date(now.getTime() - hoursAgo * 60 * 60 * 1000); + return { + id: randomUUID(), + companyId: seeded.companyId, + title: `Productivity review ${index + 1}`, + status: "done" as const, + priority: "high" as const, + originKind: PRODUCTIVITY_REVIEW_ORIGIN_KIND, + originId: seeded.issueId, + originFingerprint: `productivity-review:${seeded.issueId}`, + parentId: seeded.issueId, + issueNumber: index + 2, + identifier: `${seeded.issuePrefix}-${index + 2}`, + createdAt, + updatedAt: new Date(createdAt.getTime() + 60 * 60 * 1000), + }; + }); + const actedReview = reviewWindows[1]!; + actedReview.updatedAt = actedReview.createdAt; + await db.insert(issues).values(reviewWindows); + await db.insert(activityLog).values({ + companyId: seeded.companyId, + actorType: "agent", + actorId: seeded.coderId, + agentId: seeded.coderId, + action: "issue.updated", + entityType: "issue", + entityId: seeded.issueId, + createdAt: new Date(actedReview.createdAt.getTime() + 2 * 60 * 60 * 1000), + }); + + const result = await productivityReviewService(db).reconcileProductivityReviews({ + now, + companyId: seeded.companyId, + }); + + expect(result.created).toBe(1); + expect(result.noActionSuppressed).toBe(0); + expect(await listProductivityReviews(seeded.companyId)).toHaveLength(4); + }); + it("does not count cancelled productivity reviews toward the creation cap", async () => { const now = new Date("2026-04-28T12:00:00.000Z"); const seeded = await seedAssignedIssue(); diff --git a/server/src/services/heartbeat.ts b/server/src/services/heartbeat.ts index a11d13783e..e7ca5cd2e5 100644 --- a/server/src/services/heartbeat.ts +++ b/server/src/services/heartbeat.ts @@ -11540,6 +11540,8 @@ export function heartbeatService(db: Db, options: HeartbeatServiceOptions = {}) runId?: string | null; force?: boolean; lookbackHours?: number; + now?: Date; + reescalationCooldownMs?: number; }) { return recovery.reconcileIssueGraphLiveness({ ...opts, issueCreatedAtGte: await getWorktreeExecutionCutoff() }); } diff --git a/server/src/services/productivity-review.ts b/server/src/services/productivity-review.ts index 8441c5c675..8301ea6e08 100644 --- a/server/src/services/productivity-review.ts +++ b/server/src/services/productivity-review.ts @@ -2,6 +2,7 @@ import { and, asc, desc, eq, gt, gte, inArray, isNull, notInArray, sql } from "d import type { Db } from "@paperclipai/db"; import { clampIssueRequestDepth } from "@paperclipai/shared"; import { + activityLog, agents, companies, costEvents, @@ -30,7 +31,8 @@ export const DEFAULT_PRODUCTIVITY_REVIEW_RESOLVED_SNOOZE_MS = 6 * 60 * 60 * 1000 export const DEFAULT_PRODUCTIVITY_REVIEW_REFRESH_INTERVAL_MS = 60 * 60 * 1000; export const DEFAULT_PRODUCTIVITY_REVIEW_MAX_REFRESH_COMMENTS = 3; export const DEFAULT_PRODUCTIVITY_REVIEW_CREATION_WINDOW_MS = 24 * 60 * 60 * 1000; -export const DEFAULT_PRODUCTIVITY_REVIEW_MAX_CREATIONS_PER_WINDOW = 3; +export const DEFAULT_PRODUCTIVITY_REVIEW_MAX_CREATIONS_PER_WINDOW = 1; +export const DEFAULT_PRODUCTIVITY_REVIEW_MAX_CONSECUTIVE_NO_ACTION_REVIEWS = 3; const TERMINAL_RUN_STATUSES = ["succeeded", "interrupted", "failed", "cancelled", "timed_out"] as const; const ACTIVE_RUN_STATUSES = ["queued", "running", "scheduled_retry"] as const; @@ -54,6 +56,7 @@ type ProductivityReviewThresholds = { maxRefreshComments: number; creationWindowMs: number; maxCreationsPerWindow: number; + maxConsecutiveNoActionReviews: number; }; type ProductivityReviewEvidence = { @@ -177,6 +180,10 @@ function buildThresholds(overrides?: Partial): Pro overrides?.maxCreationsPerWindow ?? DEFAULT_PRODUCTIVITY_REVIEW_MAX_CREATIONS_PER_WINDOW, DEFAULT_PRODUCTIVITY_REVIEW_MAX_CREATIONS_PER_WINDOW, ), + maxConsecutiveNoActionReviews: readPositiveInteger( + overrides?.maxConsecutiveNoActionReviews ?? DEFAULT_PRODUCTIVITY_REVIEW_MAX_CONSECUTIVE_NO_ACTION_REVIEWS, + DEFAULT_PRODUCTIVITY_REVIEW_MAX_CONSECUTIVE_NO_ACTION_REVIEWS, + ), }; } @@ -307,6 +314,55 @@ export function productivityReviewService(db: Db, deps?: { enqueueWakeup?: Enque .then((rows) => Number(rows[0]?.count ?? 0)); } + async function countConsecutiveNoActionProductivityReviews( + companyId: string, + sourceIssueId: string, + thresholds: ProductivityReviewThresholds, + ) { + const completedReviews = await db + .select({ + createdAt: issues.createdAt, + }) + .from(issues) + .where( + and( + eq(issues.companyId, companyId), + eq(issues.originKind, PRODUCTIVITY_REVIEW_ORIGIN_KIND), + eq(issues.originId, sourceIssueId), + eq(issues.status, "done"), + visibleIssueCondition(), + ), + ) + .orderBy(desc(issues.updatedAt), desc(issues.id)) + .limit(thresholds.maxConsecutiveNoActionReviews); + + const earliestReviewCreatedAt = completedReviews.at(-1)?.createdAt; + if (!earliestReviewCreatedAt) return 0; + const sourceActions = await db + .select({ createdAt: activityLog.createdAt }) + .from(activityLog) + .where( + and( + eq(activityLog.companyId, companyId), + eq(activityLog.entityType, "issue"), + eq(activityLog.entityId, sourceIssueId), + gte(activityLog.createdAt, earliestReviewCreatedAt), + ), + ); + + let streak = 0; + for (const [index, review] of completedReviews.entries()) { + const nextNewerReviewCreatedAt = completedReviews[index - 1]?.createdAt ?? null; + const sourceAction = sourceActions.some((activity) => { + if (activity.createdAt < review.createdAt) return false; + return !nextNewerReviewCreatedAt || activity.createdAt < nextNewerReviewCreatedAt; + }); + if (sourceAction) break; + streak += 1; + } + return streak; + } + async function getRefreshCommentState(companyId: string, reviewIssueId: string) { return db .select({ @@ -679,6 +735,15 @@ export function productivityReviewService(db: Db, deps?: { enqueueWakeup?: Enque return { kind: "creation_capped" as const, reviewIssueId: null }; } + const consecutiveNoActionReviews = await countConsecutiveNoActionProductivityReviews( + evidence.sourceIssue.companyId, + evidence.sourceIssue.id, + opts.thresholds, + ); + if (consecutiveNoActionReviews >= opts.thresholds.maxConsecutiveNoActionReviews) { + return { kind: "no_action_suppressed" as const, reviewIssueId: null }; + } + const ownerAgentId = await resolveReviewOwnerAgentId(evidence.sourceIssue, evidence.sourceAgent); let review: Awaited>; try { @@ -791,6 +856,7 @@ export function productivityReviewService(db: Db, deps?: { enqueueWakeup?: Enque existing: 0, snoozed: 0, creationCapped: 0, + noActionSuppressed: 0, skipped: 0, failed: 0, reviewIssueIds: [] as string[], @@ -831,6 +897,7 @@ export function productivityReviewService(db: Db, deps?: { enqueueWakeup?: Enque if (outcome.kind === "created") result.created += 1; else if (outcome.kind === "updated") result.updated += 1; else if (outcome.kind === "creation_capped") result.creationCapped += 1; + else if (outcome.kind === "no_action_suppressed") result.noActionSuppressed += 1; else result.existing += 1; if (outcome.reviewIssueId) result.reviewIssueIds.push(outcome.reviewIssueId); } catch (err) { diff --git a/server/src/services/recovery/service.ts b/server/src/services/recovery/service.ts index b6a335ccc2..6d5b0d7da0 100644 --- a/server/src/services/recovery/service.ts +++ b/server/src/services/recovery/service.ts @@ -75,6 +75,7 @@ const UNSUCCESSFUL_HEARTBEAT_RUN_TERMINAL_STATUSES = ["interrupted", "failed", " export const ACTIVE_RUN_OUTPUT_SUSPICION_THRESHOLD_MS = 60 * 60 * 1000; export const ACTIVE_RUN_OUTPUT_CRITICAL_THRESHOLD_MS = 4 * 60 * 60 * 1000; export const ACTIVE_RUN_OUTPUT_CONTINUE_REARM_MS = 30 * 60 * 1000; +export const DEFAULT_LIVENESS_REESCALATION_COOLDOWN_MS = 60 * 60 * 1000; const ACTIVE_RUN_OUTPUT_EVIDENCE_TAIL_BYTES = 8 * 1024; const STRANDED_ISSUE_RECOVERY_ORIGIN_KIND = RECOVERY_ORIGIN_KINDS.strandedIssueRecovery; const STALE_ACTIVE_RUN_EVALUATION_ORIGIN_KIND = RECOVERY_ORIGIN_KINDS.staleActiveRunEvaluation; @@ -4019,6 +4020,34 @@ export function recoveryService(db: Db, deps: { enqueueWakeup: RecoveryWakeup }) }) ?? null; } + async function findRecentCompletedLivenessRecoveryIssue( + finding: IssueLivenessFinding, + now: Date, + cooldownMs: number, + ) { + if (cooldownMs <= 0) return null; + const cutoff = new Date(now.getTime() - cooldownMs); + return db + .select({ id: issues.id }) + .from(issues) + .where( + and( + eq(issues.companyId, finding.companyId), + eq(issues.originKind, RECOVERY_ORIGIN_KINDS.issueGraphLivenessEscalation), + or( + eq(issues.originId, finding.incidentKey), + eq(issues.originFingerprint, livenessRecoveryLeafFingerprint(finding)), + ), + visibleIssueCondition(), + eq(issues.status, "done"), + gte(issues.updatedAt, cutoff), + ), + ) + .orderBy(desc(issues.updatedAt), desc(issues.id)) + .limit(1) + .then((rows) => rows[0] ?? null); + } + async function removeRecoveryBlockerFromSource(recovery: typeof issues.$inferSelect) { const parsed = parseLivenessIncidentKey(recovery.originId); if (!parsed) return false; @@ -4374,6 +4403,8 @@ export function recoveryService(db: Db, deps: { enqueueWakeup: RecoveryWakeup }) async function createIssueGraphLivenessEscalation(input: { finding: IssueLivenessFinding; runId?: string | null; + now: Date; + reescalationCooldownMs: number; }) { const issue = await db .select() @@ -4404,6 +4435,13 @@ export function recoveryService(db: Db, deps: { enqueueWakeup: RecoveryWakeup }) }); return { kind: "existing" as const, escalationIssueId: existing.id }; } + if (await findRecentCompletedLivenessRecoveryIssue( + input.finding, + input.now, + input.reescalationCooldownMs, + )) { + return { kind: "cooldown" as const }; + } const ownerSelection = await resolveEscalationOwnerAgentId(input.finding, recoveryIssue); if (!ownerSelection) return { kind: "skipped" as const }; @@ -4778,6 +4816,8 @@ export function recoveryService(db: Db, deps: { enqueueWakeup: RecoveryWakeup }) force?: boolean; lookbackHours?: number; issueCreatedAtGte?: Date | null; + now?: Date; + reescalationCooldownMs?: number; }) { let findings = await collectIssueGraphLivenessFindings(); if (opts?.issueCreatedAtGte) { @@ -4804,7 +4844,11 @@ export function recoveryService(db: Db, deps: { enqueueWakeup: RecoveryWakeup }) const lookbackHours = normalizeIssueGraphLivenessAutoRecoveryLookbackHours( opts?.lookbackHours ?? experimentalSettings.issueGraphLivenessAutoRecoveryLookbackHours, ); - const now = new Date(); + const now = opts?.now ?? new Date(); + const reescalationCooldownMs = Math.max( + 0, + Math.floor(asNumber(opts?.reescalationCooldownMs, DEFAULT_LIVENESS_REESCALATION_COOLDOWN_MS)), + ); const cutoff = new Date(now.getTime() - lookbackHours * 60 * 60 * 1000); const obsoleteRecoveryCleanup = await retireObsoleteLivenessRecoveryIssues(findings); const doneRecoveryBlockerCleanup = await retireDoneLivenessRecoveryBlockers(); @@ -4819,6 +4863,7 @@ export function recoveryService(db: Db, deps: { enqueueWakeup: RecoveryWakeup }) skipped: 0, skippedAutoRecoveryDisabled: 0, skippedOutsideLookback: 0, + skippedReescalationCooldown: 0, obsoleteRecoveriesRetired: obsoleteRecoveryCleanup.retired, obsoleteRecoveriesActiveSkipped: obsoleteRecoveryCleanup.activeSkipped, obsoleteRecoveryBlockerRelationsRemoved: obsoleteRecoveryCleanup.blockerRelationsRemoved, @@ -4868,6 +4913,8 @@ export function recoveryService(db: Db, deps: { enqueueWakeup: RecoveryWakeup }) const escalation = await createIssueGraphLivenessEscalation({ finding, runId: opts?.runId ?? null, + now, + reescalationCooldownMs, }); if (escalation.kind === "created") { result.escalationsCreated += 1; @@ -4877,6 +4924,9 @@ export function recoveryService(db: Db, deps: { enqueueWakeup: RecoveryWakeup }) result.existingEscalations += 1; result.issueIds.push(finding.issueId); result.escalationIssueIds.push(escalation.escalationIssueId); + } else if (escalation.kind === "cooldown") { + result.skippedReescalationCooldown += 1; + result.skipped += 1; } else { result.skipped += 1; }