1605 lines
51 KiB
TypeScript
1605 lines
51 KiB
TypeScript
import { randomUUID } from "node:crypto";
|
|
import { eq, sql } from "drizzle-orm";
|
|
import { afterAll, afterEach, beforeAll, describe, expect, it, vi } from "vitest";
|
|
import {
|
|
agents,
|
|
agentWakeupRequests,
|
|
companies,
|
|
costEvents,
|
|
createDb,
|
|
documentRevisions,
|
|
documents,
|
|
heartbeatRuns,
|
|
issueComments,
|
|
issueDocuments,
|
|
issues,
|
|
} from "@paperclipai/db";
|
|
import { ISSUE_CONTINUATION_SUMMARY_DOCUMENT_KEY } from "@paperclipai/shared";
|
|
import {
|
|
getEmbeddedPostgresTestSupport,
|
|
startEmbeddedPostgresTestDatabase,
|
|
} from "./helpers/embedded-postgres.js";
|
|
import {
|
|
MAX_TURN_CONTINUATION_RETRY_REASON,
|
|
MAX_TURN_CONTINUATION_WAKE_REASON,
|
|
heartbeatService,
|
|
} from "../services/heartbeat.ts";
|
|
import { runningProcesses } from "../adapters/index.ts";
|
|
import { recoveryService } from "../services/recovery/service.ts";
|
|
|
|
const mockAdapterExecute = vi.hoisted(() =>
|
|
vi.fn(async () => ({
|
|
exitCode: 0,
|
|
signal: null,
|
|
timedOut: false,
|
|
errorMessage: null,
|
|
summary: "Stale-queue invalidation test run.",
|
|
provider: "test",
|
|
model: "test-model",
|
|
})),
|
|
);
|
|
|
|
vi.mock("../adapters/index.ts", async () => {
|
|
const actual = await vi.importActual<typeof import("../adapters/index.ts")>("../adapters/index.ts");
|
|
return {
|
|
...actual,
|
|
getServerAdapter: vi.fn(() => ({
|
|
supportsLocalAgentJwt: false,
|
|
execute: mockAdapterExecute,
|
|
})),
|
|
};
|
|
});
|
|
|
|
const embeddedPostgresSupport = await getEmbeddedPostgresTestSupport();
|
|
const describeEmbeddedPostgres = embeddedPostgresSupport.supported ? describe : describe.skip;
|
|
|
|
if (!embeddedPostgresSupport.supported) {
|
|
console.warn(
|
|
`Skipping embedded Postgres heartbeat stale-queue invalidation tests on this host: ${embeddedPostgresSupport.reason ?? "unsupported environment"}`,
|
|
);
|
|
}
|
|
|
|
async function ensureIssueRelationsTable(db: ReturnType<typeof createDb>) {
|
|
await db.execute(sql.raw(`
|
|
CREATE TABLE IF NOT EXISTS "issue_relations" (
|
|
"id" uuid PRIMARY KEY DEFAULT gen_random_uuid(),
|
|
"company_id" uuid NOT NULL,
|
|
"issue_id" uuid NOT NULL,
|
|
"related_issue_id" uuid NOT NULL,
|
|
"type" text NOT NULL,
|
|
"created_by_agent_id" uuid,
|
|
"created_by_user_id" text,
|
|
"created_at" timestamptz NOT NULL DEFAULT now(),
|
|
"updated_at" timestamptz NOT NULL DEFAULT now()
|
|
);
|
|
`));
|
|
}
|
|
|
|
async function waitForCondition(fn: () => Promise<boolean>, timeoutMs = 3_000) {
|
|
const deadline = Date.now() + timeoutMs;
|
|
while (Date.now() < deadline) {
|
|
if (await fn()) return true;
|
|
await new Promise((resolve) => setTimeout(resolve, 50));
|
|
}
|
|
return fn();
|
|
}
|
|
|
|
async function cleanupHeartbeatInvalidationFixture(db: ReturnType<typeof createDb>) {
|
|
for (let attempt = 0; attempt < 10; attempt += 1) {
|
|
try {
|
|
await db.execute(sql.raw(`
|
|
TRUNCATE TABLE
|
|
"company_skills",
|
|
"issue_comments",
|
|
"issue_documents",
|
|
"document_revisions",
|
|
"documents",
|
|
"issue_relations",
|
|
"issue_tree_holds",
|
|
"issues",
|
|
"heartbeat_run_events",
|
|
"cost_events",
|
|
"activity_log",
|
|
"heartbeat_runs",
|
|
"agent_wakeup_requests",
|
|
"agent_runtime_state",
|
|
"agents",
|
|
"companies"
|
|
RESTART IDENTITY CASCADE
|
|
`));
|
|
return;
|
|
} catch (error) {
|
|
const isLateCommentRace =
|
|
error instanceof Error &&
|
|
error.message.includes("issue_comments_issue_id_issues_id_fk");
|
|
if (!isLateCommentRace || attempt === 9) {
|
|
throw error;
|
|
}
|
|
|
|
// Heartbeat completion can write issue-thread comments shortly after the
|
|
// run leaves queued/running. Retry the dependent deletes once those land.
|
|
await new Promise((resolve) => setTimeout(resolve, 100));
|
|
}
|
|
}
|
|
}
|
|
|
|
type SeedOptions = {
|
|
agentName?: string;
|
|
agentRole?: string;
|
|
maxConcurrentRuns?: number;
|
|
heartbeatConfig?: Record<string, unknown>;
|
|
};
|
|
|
|
type SeedResult = {
|
|
companyId: string;
|
|
agentId: string;
|
|
};
|
|
|
|
describeEmbeddedPostgres("heartbeat stale queued-run invalidation", () => {
|
|
let db!: ReturnType<typeof createDb>;
|
|
let heartbeat!: ReturnType<typeof heartbeatService>;
|
|
let tempDb: Awaited<ReturnType<typeof startEmbeddedPostgresTestDatabase>> | null = null;
|
|
let beforeContinuationDispatchCheck:
|
|
| ((input: { runId: string; issueId: string }) => Promise<void>)
|
|
| null = null;
|
|
let afterContinuationDispatchCheck:
|
|
| ((input: { runId: string; issueId: string }) => Promise<void>)
|
|
| null = null;
|
|
|
|
const countExecuteCallsForRun = (runId: string) =>
|
|
mockAdapterExecute.mock.calls.filter(([context]) => context?.runId === runId).length;
|
|
|
|
beforeAll(async () => {
|
|
tempDb = await startEmbeddedPostgresTestDatabase("paperclip-heartbeat-stale-queue-");
|
|
db = createDb(tempDb.connectionString);
|
|
heartbeat = heartbeatService(db, {
|
|
beforeResolvedInteractionContinuationDispatchCheck: async (input) => {
|
|
await beforeContinuationDispatchCheck?.(input);
|
|
},
|
|
afterResolvedInteractionContinuationDispatchCheck: async (input) => {
|
|
await afterContinuationDispatchCheck?.(input);
|
|
},
|
|
});
|
|
await ensureIssueRelationsTable(db);
|
|
}, 20_000);
|
|
|
|
afterEach(async () => {
|
|
beforeContinuationDispatchCheck = null;
|
|
afterContinuationDispatchCheck = null;
|
|
mockAdapterExecute.mockReset();
|
|
mockAdapterExecute.mockImplementation(async () => ({
|
|
exitCode: 0,
|
|
signal: null,
|
|
timedOut: false,
|
|
errorMessage: null,
|
|
summary: "Stale-queue invalidation test run.",
|
|
provider: "test",
|
|
model: "test-model",
|
|
}));
|
|
runningProcesses.clear();
|
|
let idlePolls = 0;
|
|
for (let attempt = 0; attempt < 100; attempt += 1) {
|
|
const runs = await db
|
|
.select({ status: heartbeatRuns.status })
|
|
.from(heartbeatRuns);
|
|
const hasActiveRun = runs.some((run) => run.status === "queued" || run.status === "running");
|
|
if (!hasActiveRun) {
|
|
idlePolls += 1;
|
|
if (idlePolls >= 3) break;
|
|
} else {
|
|
idlePolls = 0;
|
|
}
|
|
await new Promise((resolve) => setTimeout(resolve, 50));
|
|
}
|
|
await new Promise((resolve) => setTimeout(resolve, 50));
|
|
await cleanupHeartbeatInvalidationFixture(db);
|
|
});
|
|
|
|
afterAll(async () => {
|
|
await tempDb?.cleanup();
|
|
});
|
|
|
|
async function seedCompanyAndAgent(opts: SeedOptions = {}): Promise<SeedResult> {
|
|
const companyId = randomUUID();
|
|
const agentId = randomUUID();
|
|
await db.insert(companies).values({
|
|
id: companyId,
|
|
name: "Paperclip",
|
|
issuePrefix: `T${companyId.replace(/-/g, "").slice(0, 6).toUpperCase()}`,
|
|
defaultResponsibleUserId: "responsible-user",
|
|
requireBoardApprovalForNewAgents: false,
|
|
});
|
|
await db.insert(agents).values({
|
|
id: agentId,
|
|
companyId,
|
|
name: opts.agentName ?? "ClaudeCoder",
|
|
role: opts.agentRole ?? "engineer",
|
|
status: "active",
|
|
adapterType: "codex_local",
|
|
adapterConfig: {},
|
|
runtimeConfig: {
|
|
heartbeat: {
|
|
wakeOnDemand: true,
|
|
maxConcurrentRuns: opts.maxConcurrentRuns ?? 1,
|
|
...(opts.heartbeatConfig ?? {}),
|
|
},
|
|
},
|
|
permissions: {},
|
|
});
|
|
return { companyId, agentId };
|
|
}
|
|
|
|
async function seedQueuedRun(input: {
|
|
companyId: string;
|
|
agentId: string;
|
|
issueId: string;
|
|
wakeReason: string;
|
|
contextExtras?: Record<string, unknown>;
|
|
invocationSource?: "assignment" | "automation";
|
|
scheduledRetryReason?: string | null;
|
|
}) {
|
|
const wakeupRequestId = randomUUID();
|
|
const runId = randomUUID();
|
|
await db.insert(agentWakeupRequests).values({
|
|
id: wakeupRequestId,
|
|
companyId: input.companyId,
|
|
agentId: input.agentId,
|
|
source: input.invocationSource ?? "assignment",
|
|
triggerDetail: "system",
|
|
reason: input.wakeReason,
|
|
payload: { issueId: input.issueId },
|
|
status: "queued",
|
|
});
|
|
await db.insert(heartbeatRuns).values({
|
|
id: runId,
|
|
companyId: input.companyId,
|
|
agentId: input.agentId,
|
|
invocationSource: input.invocationSource ?? "assignment",
|
|
triggerDetail: "system",
|
|
status: "queued",
|
|
wakeupRequestId,
|
|
scheduledRetryReason: input.scheduledRetryReason ?? null,
|
|
contextSnapshot: {
|
|
issueId: input.issueId,
|
|
wakeReason: input.wakeReason,
|
|
...(input.contextExtras ?? {}),
|
|
},
|
|
});
|
|
await db
|
|
.update(agentWakeupRequests)
|
|
.set({ runId })
|
|
.where(eq(agentWakeupRequests.id, wakeupRequestId));
|
|
return { runId, wakeupRequestId };
|
|
}
|
|
|
|
async function seedContinuationSummary(input: {
|
|
companyId: string;
|
|
issueId: string;
|
|
agentId: string;
|
|
body: string;
|
|
}) {
|
|
const documentId = randomUUID();
|
|
const revisionId = randomUUID();
|
|
await db.insert(documents).values({
|
|
id: documentId,
|
|
companyId: input.companyId,
|
|
title: "Continuation Summary",
|
|
format: "markdown",
|
|
latestBody: input.body,
|
|
latestRevisionId: revisionId,
|
|
latestRevisionNumber: 1,
|
|
createdByAgentId: input.agentId,
|
|
updatedByAgentId: input.agentId,
|
|
});
|
|
await db.insert(documentRevisions).values({
|
|
id: revisionId,
|
|
companyId: input.companyId,
|
|
documentId,
|
|
revisionNumber: 1,
|
|
title: "Continuation Summary",
|
|
format: "markdown",
|
|
body: input.body,
|
|
createdByAgentId: input.agentId,
|
|
});
|
|
await db.insert(issueDocuments).values({
|
|
companyId: input.companyId,
|
|
issueId: input.issueId,
|
|
documentId,
|
|
key: ISSUE_CONTINUATION_SUMMARY_DOCUMENT_KEY,
|
|
});
|
|
}
|
|
|
|
it("skips generic timer wakes with no actionable assigned work before adapter execution", async () => {
|
|
const { agentId } = await seedCompanyAndAgent({
|
|
heartbeatConfig: {
|
|
enabled: true,
|
|
skipTimerWhenNoActionableWork: true,
|
|
},
|
|
});
|
|
|
|
const run = await heartbeat.wakeup(agentId, {
|
|
source: "timer",
|
|
triggerDetail: "schedule",
|
|
});
|
|
|
|
expect(run).toBeNull();
|
|
expect(mockAdapterExecute).not.toHaveBeenCalled();
|
|
|
|
const [wakeup] = await db
|
|
.select({
|
|
status: agentWakeupRequests.status,
|
|
reason: agentWakeupRequests.reason,
|
|
payload: agentWakeupRequests.payload,
|
|
})
|
|
.from(agentWakeupRequests)
|
|
.where(eq(agentWakeupRequests.agentId, agentId));
|
|
const runRows = await db.select({ id: heartbeatRuns.id }).from(heartbeatRuns);
|
|
|
|
expect(wakeup).toMatchObject({
|
|
status: "skipped",
|
|
reason: "heartbeat.timer.no_actionable_work",
|
|
});
|
|
expect(wakeup?.payload).toMatchObject({
|
|
heartbeatSkip: {
|
|
reason: expect.stringContaining("No assigned todo or in_progress issue"),
|
|
},
|
|
});
|
|
expect(runRows).toHaveLength(0);
|
|
});
|
|
|
|
it("checks guarded issue status and assignee under the enqueue lock", async () => {
|
|
const { companyId, agentId } = await seedCompanyAndAgent();
|
|
const parkedIssueId = randomUUID();
|
|
const reassignedIssueId = randomUUID();
|
|
await db.insert(issues).values([
|
|
{
|
|
id: parkedIssueId,
|
|
companyId,
|
|
title: "Parked connection intent",
|
|
status: "backlog" as const,
|
|
priority: "medium" as const,
|
|
assigneeAgentId: agentId,
|
|
},
|
|
{
|
|
id: reassignedIssueId,
|
|
companyId,
|
|
title: "Reassigned connection intent",
|
|
status: "in_progress" as const,
|
|
priority: "medium" as const,
|
|
assigneeAgentId: null,
|
|
},
|
|
]);
|
|
|
|
for (const issueId of [parkedIssueId, reassignedIssueId]) {
|
|
const run = await heartbeat.wakeup(agentId, {
|
|
source: "automation",
|
|
triggerDetail: "system",
|
|
reason: "issue_commented",
|
|
payload: { issueId, interactionId: randomUUID() },
|
|
contextSnapshot: { issueId, wakeReason: "issue_commented" },
|
|
requestedByActorType: "user",
|
|
requestedByActorId: "responsible-user",
|
|
issueStateGuard: {
|
|
statuses: ["in_progress"],
|
|
assigneeAgentId: agentId,
|
|
},
|
|
});
|
|
expect(run).toBeNull();
|
|
}
|
|
|
|
expect(mockAdapterExecute).not.toHaveBeenCalled();
|
|
expect(await db.select().from(heartbeatRuns)).toHaveLength(0);
|
|
expect(await db.select({ status: agentWakeupRequests.status, reason: agentWakeupRequests.reason })
|
|
.from(agentWakeupRequests)).toEqual([
|
|
{ status: "skipped", reason: "issue_state_guard_mismatch" },
|
|
{ status: "skipped", reason: "issue_state_guard_mismatch" },
|
|
]);
|
|
});
|
|
|
|
it.each([
|
|
{ runtimeMode: "native", status: "done", reassigned: false },
|
|
{ runtimeMode: "legacy", status: "done", reassigned: false },
|
|
{ runtimeMode: "native", status: "cancelled", reassigned: false },
|
|
{ runtimeMode: "native", status: "backlog", reassigned: false },
|
|
{ runtimeMode: "native", status: "in_review", reassigned: false },
|
|
{ runtimeMode: "native", status: "blocked", reassigned: false },
|
|
{ runtimeMode: "native", status: "in_progress", reassigned: true },
|
|
] as const)("skips stale $runtimeMode productive recovery after status=$status reassigned=$reassigned commits under the enqueue lock", async ({ runtimeMode, status, reassigned }) => {
|
|
const { companyId, agentId } = await seedCompanyAndAgent();
|
|
const issueId = randomUUID();
|
|
const runId = randomUUID();
|
|
await db.insert(issues).values({
|
|
id: issueId,
|
|
companyId,
|
|
title: "Completion racing with productive recovery",
|
|
status: "in_progress",
|
|
assigneeAgentId: agentId,
|
|
});
|
|
await db.insert(heartbeatRuns).values({
|
|
id: runId,
|
|
companyId,
|
|
agentId,
|
|
invocationSource: "assignment",
|
|
runtimeMode,
|
|
status: "succeeded",
|
|
livenessState: "completed",
|
|
contextSnapshot: { issueId, wakeReason: "issue_assigned" },
|
|
startedAt: new Date(),
|
|
finishedAt: new Date(),
|
|
});
|
|
|
|
// Let the real sweep select its in-progress snapshot, then hold the issue
|
|
// lock until the real enqueue transaction is waiting on the newer state.
|
|
const enqueueWakeup = vi.fn(async (...[targetAgentId, options]: Parameters<typeof heartbeat.wakeup>) => {
|
|
let pendingWake!: ReturnType<typeof heartbeat.wakeup>;
|
|
await db.transaction(async (tx) => {
|
|
await tx.update(issues).set({
|
|
status,
|
|
...(reassigned ? { assigneeAgentId: null, assigneeUserId: "responsible-user" } : {}),
|
|
}).where(eq(issues.id, issueId));
|
|
const [{ pid }] = await tx.execute<{ pid: number }>(sql`select pg_backend_pid() as pid`);
|
|
pendingWake = heartbeat.wakeup(targetAgentId, options);
|
|
try {
|
|
expect(await waitForCondition(async () => {
|
|
const [{ waiting }] = await db.execute<{ waiting: boolean }>(sql`
|
|
select exists (
|
|
select 1 from pg_stat_activity
|
|
where ${pid} = any(pg_blocking_pids(pid))
|
|
) as waiting
|
|
`);
|
|
return waiting;
|
|
})).toBe(true);
|
|
} catch (error) {
|
|
// Observe a pending rejection even if the lock assertion fails.
|
|
void pendingWake.catch(() => {});
|
|
throw error;
|
|
}
|
|
});
|
|
return pendingWake;
|
|
});
|
|
const recovery = recoveryService(db, { enqueueWakeup });
|
|
const result = await recovery.reconcileStrandedAssignedIssues();
|
|
|
|
expect(enqueueWakeup).toHaveBeenCalledOnce();
|
|
expect(result).toMatchObject({ continuationRequeued: 0, escalated: 0, skipped: 1, issueIds: [] });
|
|
expect(await db.select({ id: heartbeatRuns.id }).from(heartbeatRuns)).toEqual([{ id: runId }]);
|
|
expect(await db.select().from(issueComments)).toHaveLength(0);
|
|
expect(mockAdapterExecute).not.toHaveBeenCalled();
|
|
const [wakeup] = await db.select().from(agentWakeupRequests);
|
|
expect(wakeup).toMatchObject({
|
|
status: "skipped",
|
|
reason: "issue_state_guard_mismatch",
|
|
runId: null,
|
|
payload: {
|
|
heartbeatSkip: {
|
|
expectedStatuses: ["in_progress"],
|
|
actualStatus: status,
|
|
expectedAssigneeAgentId: agentId,
|
|
actualAssigneeAgentId: reassigned ? null : agentId,
|
|
},
|
|
},
|
|
});
|
|
const [issue] = await db.select().from(issues).where(eq(issues.id, issueId));
|
|
expect(issue).toMatchObject({ status, assigneeAgentId: reassigned ? null : agentId });
|
|
});
|
|
|
|
it("cancels a resolved connection-intent wake parked before queued-run claim", async () => {
|
|
const { companyId, agentId } = await seedCompanyAndAgent();
|
|
const issueId = randomUUID();
|
|
await db.insert(issues).values({
|
|
id: issueId,
|
|
companyId,
|
|
title: "Connection intent parked after enqueue",
|
|
status: "in_progress",
|
|
priority: "medium",
|
|
assigneeAgentId: agentId,
|
|
});
|
|
const { runId, wakeupRequestId } = await seedQueuedRun({
|
|
companyId,
|
|
agentId,
|
|
issueId,
|
|
wakeReason: "issue_commented",
|
|
invocationSource: "automation",
|
|
contextExtras: {
|
|
interactionId: randomUUID(),
|
|
interactionKind: "connection_intent",
|
|
interactionStatus: "accepted",
|
|
interactionResolvedAt: "2026-08-28T13:30:00.000Z",
|
|
mutation: "interaction",
|
|
source: "connection_intent.resolved",
|
|
},
|
|
});
|
|
|
|
await db.update(issues).set({ status: "backlog" }).where(eq(issues.id, issueId));
|
|
await heartbeat.resumeQueuedRuns();
|
|
|
|
const [run, wakeup, issue] = await Promise.all([
|
|
db.select({ status: heartbeatRuns.status, errorCode: heartbeatRuns.errorCode })
|
|
.from(heartbeatRuns)
|
|
.where(eq(heartbeatRuns.id, runId))
|
|
.then((rows) => rows[0] ?? null),
|
|
db.select({ status: agentWakeupRequests.status, error: agentWakeupRequests.error })
|
|
.from(agentWakeupRequests)
|
|
.where(eq(agentWakeupRequests.id, wakeupRequestId))
|
|
.then((rows) => rows[0] ?? null),
|
|
db.select({ status: issues.status }).from(issues).where(eq(issues.id, issueId))
|
|
.then((rows) => rows[0] ?? null),
|
|
]);
|
|
expect(run).toMatchObject({ status: "cancelled", errorCode: "issue_not_in_progress" });
|
|
expect(wakeup).toMatchObject({ status: "skipped", error: expect.stringContaining("no longer in_progress") });
|
|
expect(issue?.status).toBe("backlog");
|
|
expect(countExecuteCallsForRun(runId)).toBe(0);
|
|
});
|
|
|
|
it("does not re-open a resolved connection-intent issue parked after claim", async () => {
|
|
const { companyId, agentId } = await seedCompanyAndAgent();
|
|
const issueId = randomUUID();
|
|
await db.insert(issues).values({
|
|
id: issueId,
|
|
companyId,
|
|
title: "Connection intent parked between claim and checkout",
|
|
status: "in_progress",
|
|
priority: "medium",
|
|
assigneeAgentId: agentId,
|
|
});
|
|
const { runId, wakeupRequestId } = await seedQueuedRun({
|
|
companyId,
|
|
agentId,
|
|
issueId,
|
|
wakeReason: "issue_commented",
|
|
invocationSource: "automation",
|
|
contextExtras: {
|
|
interactionId: randomUUID(),
|
|
interactionKind: "connection_intent",
|
|
interactionStatus: "accepted",
|
|
interactionResolvedAt: "2026-08-28T13:30:00.000Z",
|
|
mutation: "interaction",
|
|
source: "connection_intent.resolved",
|
|
},
|
|
});
|
|
|
|
await db.execute(sql.raw(`
|
|
CREATE OR REPLACE FUNCTION park_connection_intent_after_claim()
|
|
RETURNS trigger AS $trigger$
|
|
BEGIN
|
|
IF NEW.id = '${runId}'::uuid AND NEW.status = 'running' THEN
|
|
UPDATE issues SET status = 'backlog' WHERE id = '${issueId}'::uuid;
|
|
END IF;
|
|
RETURN NEW;
|
|
END;
|
|
$trigger$ LANGUAGE plpgsql;
|
|
|
|
CREATE TRIGGER park_connection_intent_after_claim
|
|
AFTER UPDATE OF status ON heartbeat_runs
|
|
FOR EACH ROW EXECUTE FUNCTION park_connection_intent_after_claim();
|
|
`));
|
|
|
|
await heartbeat.resumeQueuedRuns();
|
|
await waitForCondition(async () => {
|
|
const [run, wakeup] = await Promise.all([
|
|
db.select({ status: heartbeatRuns.status })
|
|
.from(heartbeatRuns)
|
|
.where(eq(heartbeatRuns.id, runId))
|
|
.then((rows) => rows[0] ?? null),
|
|
db.select({ status: agentWakeupRequests.status })
|
|
.from(agentWakeupRequests)
|
|
.where(eq(agentWakeupRequests.id, wakeupRequestId))
|
|
.then((rows) => rows[0] ?? null),
|
|
]);
|
|
return run?.status === "cancelled" && wakeup?.status === "skipped";
|
|
});
|
|
|
|
const [run, wakeup, issue] = await Promise.all([
|
|
db.select({ status: heartbeatRuns.status, errorCode: heartbeatRuns.errorCode })
|
|
.from(heartbeatRuns)
|
|
.where(eq(heartbeatRuns.id, runId))
|
|
.then((rows) => rows[0] ?? null),
|
|
db.select({ status: agentWakeupRequests.status, error: agentWakeupRequests.error })
|
|
.from(agentWakeupRequests)
|
|
.where(eq(agentWakeupRequests.id, wakeupRequestId))
|
|
.then((rows) => rows[0] ?? null),
|
|
db.select({ status: issues.status }).from(issues).where(eq(issues.id, issueId))
|
|
.then((rows) => rows[0] ?? null),
|
|
]);
|
|
expect(run).toMatchObject({ status: "cancelled", errorCode: "issue_not_in_progress" });
|
|
expect(wakeup).toMatchObject({ status: "skipped", error: expect.stringContaining("no longer in_progress") });
|
|
expect(issue?.status).toBe("backlog");
|
|
expect(countExecuteCallsForRun(runId)).toBe(0);
|
|
});
|
|
|
|
it.each([
|
|
{
|
|
mutation: "parked",
|
|
expectedErrorCode: "issue_not_in_progress",
|
|
expectedError: "no longer in_progress",
|
|
},
|
|
{
|
|
mutation: "reassigned",
|
|
expectedErrorCode: "issue_assignee_changed",
|
|
expectedError: "changed assignee",
|
|
},
|
|
])(
|
|
"cancels a resolved connection-intent wake $mutation after checkout but before adapter dispatch",
|
|
async ({ mutation, expectedErrorCode, expectedError }) => {
|
|
const { companyId, agentId } = await seedCompanyAndAgent();
|
|
const replacementAgentId = randomUUID();
|
|
if (mutation === "reassigned") {
|
|
await db.insert(agents).values({
|
|
id: replacementAgentId,
|
|
companyId,
|
|
name: "ReplacementCoder",
|
|
role: "engineer",
|
|
status: "active",
|
|
adapterType: "codex_local",
|
|
adapterConfig: {},
|
|
runtimeConfig: { heartbeat: { wakeOnDemand: true, maxConcurrentRuns: 1 } },
|
|
permissions: {},
|
|
});
|
|
}
|
|
const issueId = randomUUID();
|
|
await db.insert(issues).values({
|
|
id: issueId,
|
|
companyId,
|
|
title: `Connection intent ${mutation} at final dispatch`,
|
|
status: "in_progress",
|
|
priority: "medium",
|
|
assigneeAgentId: agentId,
|
|
});
|
|
const { runId, wakeupRequestId } = await seedQueuedRun({
|
|
companyId,
|
|
agentId,
|
|
issueId,
|
|
wakeReason: "issue_commented",
|
|
invocationSource: "automation",
|
|
contextExtras: {
|
|
interactionId: randomUUID(),
|
|
interactionKind: "connection_intent",
|
|
interactionStatus: "accepted",
|
|
interactionResolvedAt: "2026-08-28T13:30:00.000Z",
|
|
mutation: "interaction",
|
|
source: "connection_intent.resolved",
|
|
},
|
|
});
|
|
beforeContinuationDispatchCheck = async ({ runId: guardedRunId, issueId: guardedIssueId }) => {
|
|
expect(guardedRunId).toBe(runId);
|
|
expect(guardedIssueId).toBe(issueId);
|
|
await db
|
|
.update(issues)
|
|
.set(mutation === "parked"
|
|
? { status: "backlog", updatedAt: new Date() }
|
|
: { assigneeAgentId: replacementAgentId, updatedAt: new Date() })
|
|
.where(eq(issues.id, issueId));
|
|
};
|
|
|
|
await heartbeat.resumeQueuedRuns();
|
|
await waitForCondition(async () => {
|
|
const [run, wakeup] = await Promise.all([
|
|
db.select({ status: heartbeatRuns.status })
|
|
.from(heartbeatRuns)
|
|
.where(eq(heartbeatRuns.id, runId))
|
|
.then((rows) => rows[0] ?? null),
|
|
db.select({ status: agentWakeupRequests.status })
|
|
.from(agentWakeupRequests)
|
|
.where(eq(agentWakeupRequests.id, wakeupRequestId))
|
|
.then((rows) => rows[0] ?? null),
|
|
]);
|
|
return run?.status === "cancelled" && wakeup?.status === "skipped";
|
|
});
|
|
|
|
const [run, wakeup, issue] = await Promise.all([
|
|
db.select({ status: heartbeatRuns.status, errorCode: heartbeatRuns.errorCode })
|
|
.from(heartbeatRuns)
|
|
.where(eq(heartbeatRuns.id, runId))
|
|
.then((rows) => rows[0] ?? null),
|
|
db.select({ status: agentWakeupRequests.status, error: agentWakeupRequests.error })
|
|
.from(agentWakeupRequests)
|
|
.where(eq(agentWakeupRequests.id, wakeupRequestId))
|
|
.then((rows) => rows[0] ?? null),
|
|
db.select({ status: issues.status, assigneeAgentId: issues.assigneeAgentId })
|
|
.from(issues)
|
|
.where(eq(issues.id, issueId))
|
|
.then((rows) => rows[0] ?? null),
|
|
]);
|
|
expect(run).toMatchObject({ status: "cancelled", errorCode: expectedErrorCode });
|
|
expect(wakeup).toMatchObject({ status: "skipped", error: expect.stringContaining(expectedError) });
|
|
expect(issue).toMatchObject(mutation === "parked"
|
|
? { status: "backlog", assigneeAgentId: agentId }
|
|
: { status: "in_progress", assigneeAgentId: replacementAgentId });
|
|
expect(countExecuteCallsForRun(runId)).toBe(0);
|
|
},
|
|
);
|
|
|
|
it("rejects ownership changes immediately before the final continuation handoff", async () => {
|
|
const { companyId, agentId } = await seedCompanyAndAgent();
|
|
const issueId = randomUUID();
|
|
await db.insert(issues).values({
|
|
id: issueId,
|
|
companyId,
|
|
title: "Connection intent parked at the atomic dispatch gate",
|
|
status: "in_progress",
|
|
priority: "medium",
|
|
assigneeAgentId: agentId,
|
|
});
|
|
const { runId } = await seedQueuedRun({
|
|
companyId,
|
|
agentId,
|
|
issueId,
|
|
wakeReason: "issue_commented",
|
|
invocationSource: "automation",
|
|
contextExtras: {
|
|
interactionId: randomUUID(),
|
|
interactionKind: "connection_intent",
|
|
interactionStatus: "accepted",
|
|
interactionResolvedAt: "2026-08-28T13:30:00.000Z",
|
|
mutation: "interaction",
|
|
source: "connection_intent.resolved",
|
|
},
|
|
});
|
|
|
|
const ordering: string[] = [];
|
|
let parkPromise: Promise<unknown> | null = null;
|
|
afterContinuationDispatchCheck = async ({ runId: guardedRunId, issueId: guardedIssueId }) => {
|
|
expect(guardedRunId).toBe(runId);
|
|
expect(guardedIssueId).toBe(issueId);
|
|
ordering.push("validated");
|
|
parkPromise = Promise.resolve(
|
|
db
|
|
.update(issues)
|
|
.set({
|
|
status: "backlog",
|
|
checkoutRunId: null,
|
|
executionRunId: null,
|
|
executionAgentNameKey: null,
|
|
executionLockedAt: null,
|
|
updatedAt: new Date(),
|
|
})
|
|
.where(eq(issues.id, issueId))
|
|
.returning({ id: issues.id }),
|
|
).then((rows) => {
|
|
expect(rows).toHaveLength(1);
|
|
ordering.push("parked");
|
|
});
|
|
// Admission is committed before adapter-owned setup. This concurrent
|
|
// update must not wait on a lock held by the adapter callback.
|
|
await new Promise((resolve) => setTimeout(resolve, 25));
|
|
await parkPromise;
|
|
expect(ordering).toEqual(["validated", "parked"]);
|
|
};
|
|
await heartbeat.resumeQueuedRuns();
|
|
await waitForCondition(async () => {
|
|
const run = await db
|
|
.select({ status: heartbeatRuns.status })
|
|
.from(heartbeatRuns)
|
|
.where(eq(heartbeatRuns.id, runId))
|
|
.then((rows) => rows[0] ?? null);
|
|
return run?.status === "cancelled";
|
|
});
|
|
await parkPromise;
|
|
|
|
const issue = await db
|
|
.select({ status: issues.status })
|
|
.from(issues)
|
|
.where(eq(issues.id, issueId))
|
|
.then((rows) => rows[0] ?? null);
|
|
expect(issue?.status).toBe("backlog");
|
|
expect(ordering).toEqual(["validated", "parked"]);
|
|
expect(countExecuteCallsForRun(runId)).toBe(0);
|
|
});
|
|
|
|
it("rate-limits skipped generic timer wakes by advancing the timer baseline", async () => {
|
|
const { agentId } = await seedCompanyAndAgent({
|
|
heartbeatConfig: {
|
|
enabled: true,
|
|
intervalSec: 60,
|
|
skipTimerWhenNoActionableWork: true,
|
|
},
|
|
});
|
|
const now = new Date();
|
|
await db
|
|
.update(agents)
|
|
.set({ lastHeartbeatAt: new Date(now.getTime() - 120_000) })
|
|
.where(eq(agents.id, agentId));
|
|
|
|
const firstTick = await heartbeat.tickTimers(now);
|
|
const secondTick = await heartbeat.tickTimers(now);
|
|
|
|
expect(firstTick.skipped).toBe(1);
|
|
expect(secondTick.skipped).toBe(0);
|
|
expect(mockAdapterExecute).not.toHaveBeenCalled();
|
|
|
|
const wakeups = await db
|
|
.select({ reason: agentWakeupRequests.reason })
|
|
.from(agentWakeupRequests)
|
|
.where(eq(agentWakeupRequests.agentId, agentId));
|
|
const [agent] = await db
|
|
.select({ lastHeartbeatAt: agents.lastHeartbeatAt })
|
|
.from(agents)
|
|
.where(eq(agents.id, agentId));
|
|
|
|
expect(wakeups).toHaveLength(1);
|
|
expect(wakeups[0]?.reason).toBe("heartbeat.timer.no_actionable_work");
|
|
expect(agent?.lastHeartbeatAt).toBeInstanceOf(Date);
|
|
expect(agent?.lastHeartbeatAt?.getTime()).toBeGreaterThan(now.getTime() - 120_000);
|
|
});
|
|
|
|
it("atomically claims a due timer interval across overlapping scheduler ticks", async () => {
|
|
const { agentId } = await seedCompanyAndAgent({
|
|
heartbeatConfig: {
|
|
enabled: true,
|
|
intervalSec: 60,
|
|
},
|
|
});
|
|
const now = new Date();
|
|
await db
|
|
.update(agents)
|
|
.set({
|
|
createdAt: new Date(now.getTime() - 120_000),
|
|
lastHeartbeatAt: null,
|
|
})
|
|
.where(eq(agents.id, agentId));
|
|
|
|
const results = await Promise.all([
|
|
heartbeat.tickTimers(now),
|
|
heartbeat.tickTimers(now),
|
|
]);
|
|
|
|
expect(results.reduce((total, result) => total + result.enqueued, 0)).toBe(1);
|
|
|
|
const runs = await db
|
|
.select({
|
|
id: heartbeatRuns.id,
|
|
contextSnapshot: heartbeatRuns.contextSnapshot,
|
|
})
|
|
.from(heartbeatRuns)
|
|
.where(eq(heartbeatRuns.agentId, agentId));
|
|
const [agent] = await db
|
|
.select({ lastHeartbeatAt: agents.lastHeartbeatAt })
|
|
.from(agents)
|
|
.where(eq(agents.id, agentId));
|
|
|
|
expect(runs).toHaveLength(1);
|
|
expect(runs[0]?.contextSnapshot).toMatchObject({
|
|
timerClaimWasFirstHeartbeat: true,
|
|
});
|
|
expect(agent?.lastHeartbeatAt?.getTime()).toBeGreaterThanOrEqual(now.getTime());
|
|
});
|
|
|
|
it("allows generic timer wakes when the agent has assigned todo work", async () => {
|
|
const { companyId, agentId } = await seedCompanyAndAgent({
|
|
heartbeatConfig: {
|
|
enabled: true,
|
|
skipTimerWhenNoActionableWork: true,
|
|
},
|
|
});
|
|
await db.insert(issues).values({
|
|
id: randomUUID(),
|
|
companyId,
|
|
title: "Assigned work",
|
|
status: "todo",
|
|
priority: "high",
|
|
assigneeAgentId: agentId,
|
|
});
|
|
|
|
const run = await heartbeat.wakeup(agentId, {
|
|
source: "timer",
|
|
triggerDetail: "schedule",
|
|
});
|
|
|
|
expect(run).not.toBeNull();
|
|
await waitForCondition(async () => countExecuteCallsForRun(run!.id) > 0);
|
|
|
|
expect(countExecuteCallsForRun(run!.id)).toBe(1);
|
|
});
|
|
|
|
it("allows legacy generic timer wakes by default when no skip policy is set", async () => {
|
|
const { agentId } = await seedCompanyAndAgent({
|
|
heartbeatConfig: {
|
|
enabled: true,
|
|
},
|
|
});
|
|
|
|
const run = await heartbeat.wakeup(agentId, {
|
|
source: "timer",
|
|
triggerDetail: "schedule",
|
|
});
|
|
|
|
expect(run).not.toBeNull();
|
|
await waitForCondition(async () => countExecuteCallsForRun(run!.id) > 0);
|
|
expect(countExecuteCallsForRun(run!.id)).toBe(1);
|
|
});
|
|
|
|
it("allows explicit proactive generic timer wakes without assigned issue work", async () => {
|
|
const { agentId } = await seedCompanyAndAgent({
|
|
heartbeatConfig: {
|
|
enabled: true,
|
|
skipTimerWhenNoActionableWork: false,
|
|
},
|
|
});
|
|
|
|
const run = await heartbeat.wakeup(agentId, {
|
|
source: "timer",
|
|
triggerDetail: "schedule",
|
|
});
|
|
|
|
expect(run).not.toBeNull();
|
|
await waitForCondition(async () => countExecuteCallsForRun(run!.id) > 0);
|
|
expect(countExecuteCallsForRun(run!.id)).toBe(1);
|
|
});
|
|
|
|
it("skips wakes before queueing when per-agent daily run cap is reached", async () => {
|
|
const { companyId, agentId } = await seedCompanyAndAgent({
|
|
heartbeatConfig: {
|
|
maxDailyRuns: 1,
|
|
},
|
|
});
|
|
await db.insert(heartbeatRuns).values({
|
|
id: randomUUID(),
|
|
companyId,
|
|
agentId,
|
|
invocationSource: "on_demand",
|
|
triggerDetail: "manual",
|
|
status: "succeeded",
|
|
createdAt: new Date(),
|
|
startedAt: new Date(),
|
|
finishedAt: new Date(),
|
|
contextSnapshot: {},
|
|
});
|
|
|
|
const run = await heartbeat.wakeup(agentId, {
|
|
source: "on_demand",
|
|
triggerDetail: "manual",
|
|
});
|
|
|
|
expect(run).toBeNull();
|
|
expect(mockAdapterExecute).not.toHaveBeenCalled();
|
|
|
|
const [wakeup] = await db
|
|
.select({
|
|
status: agentWakeupRequests.status,
|
|
reason: agentWakeupRequests.reason,
|
|
payload: agentWakeupRequests.payload,
|
|
})
|
|
.from(agentWakeupRequests)
|
|
.where(eq(agentWakeupRequests.agentId, agentId));
|
|
|
|
expect(wakeup).toMatchObject({
|
|
status: "skipped",
|
|
reason: "heartbeat.daily_run_limit",
|
|
});
|
|
expect(wakeup?.payload).toMatchObject({
|
|
heartbeatSkip: {
|
|
observed: 1,
|
|
limit: 1,
|
|
},
|
|
});
|
|
});
|
|
|
|
it("treats zero daily run cap as a hard stop", async () => {
|
|
const { agentId } = await seedCompanyAndAgent({
|
|
heartbeatConfig: {
|
|
maxDailyRuns: 0,
|
|
},
|
|
});
|
|
|
|
const run = await heartbeat.wakeup(agentId, {
|
|
source: "on_demand",
|
|
triggerDetail: "manual",
|
|
});
|
|
|
|
expect(run).toBeNull();
|
|
expect(mockAdapterExecute).not.toHaveBeenCalled();
|
|
|
|
const [wakeup] = await db
|
|
.select({
|
|
status: agentWakeupRequests.status,
|
|
reason: agentWakeupRequests.reason,
|
|
payload: agentWakeupRequests.payload,
|
|
})
|
|
.from(agentWakeupRequests)
|
|
.where(eq(agentWakeupRequests.agentId, agentId));
|
|
|
|
expect(wakeup).toMatchObject({
|
|
status: "skipped",
|
|
reason: "heartbeat.daily_run_limit",
|
|
});
|
|
expect(wakeup?.payload).toMatchObject({
|
|
heartbeatSkip: {
|
|
observed: 0,
|
|
limit: 0,
|
|
},
|
|
});
|
|
});
|
|
|
|
it("counts started cancelled runs toward the per-agent daily run cap", async () => {
|
|
const { companyId, agentId } = await seedCompanyAndAgent({
|
|
heartbeatConfig: {
|
|
maxDailyRuns: 1,
|
|
},
|
|
});
|
|
await db.insert(heartbeatRuns).values({
|
|
id: randomUUID(),
|
|
companyId,
|
|
agentId,
|
|
invocationSource: "on_demand",
|
|
triggerDetail: "manual",
|
|
status: "cancelled",
|
|
createdAt: new Date(),
|
|
startedAt: new Date(),
|
|
finishedAt: new Date(),
|
|
contextSnapshot: {},
|
|
});
|
|
|
|
const run = await heartbeat.wakeup(agentId, {
|
|
source: "on_demand",
|
|
triggerDetail: "manual",
|
|
});
|
|
|
|
expect(run).toBeNull();
|
|
expect(mockAdapterExecute).not.toHaveBeenCalled();
|
|
|
|
const [wakeup] = await db
|
|
.select({
|
|
status: agentWakeupRequests.status,
|
|
reason: agentWakeupRequests.reason,
|
|
payload: agentWakeupRequests.payload,
|
|
})
|
|
.from(agentWakeupRequests)
|
|
.where(eq(agentWakeupRequests.agentId, agentId));
|
|
|
|
expect(wakeup).toMatchObject({
|
|
status: "skipped",
|
|
reason: "heartbeat.daily_run_limit",
|
|
});
|
|
expect(wakeup?.payload).toMatchObject({
|
|
heartbeatSkip: {
|
|
observed: 1,
|
|
limit: 1,
|
|
},
|
|
});
|
|
});
|
|
|
|
it("coalesces same-issue wakes before enforcing the daily run cap", async () => {
|
|
const { companyId, agentId } = await seedCompanyAndAgent({
|
|
heartbeatConfig: {
|
|
maxDailyRuns: 1,
|
|
},
|
|
});
|
|
const issueId = randomUUID();
|
|
const wakeupRequestId = randomUUID();
|
|
const queuedRunId = randomUUID();
|
|
await db.insert(heartbeatRuns).values({
|
|
id: randomUUID(),
|
|
companyId,
|
|
agentId,
|
|
invocationSource: "on_demand",
|
|
triggerDetail: "manual",
|
|
status: "succeeded",
|
|
createdAt: new Date(),
|
|
startedAt: new Date(),
|
|
finishedAt: new Date(),
|
|
contextSnapshot: {},
|
|
});
|
|
await db.insert(agentWakeupRequests).values({
|
|
id: wakeupRequestId,
|
|
companyId,
|
|
agentId,
|
|
source: "on_demand",
|
|
triggerDetail: "manual",
|
|
reason: "manual",
|
|
payload: { issueId },
|
|
status: "queued",
|
|
});
|
|
await db.insert(heartbeatRuns).values({
|
|
id: queuedRunId,
|
|
companyId,
|
|
agentId,
|
|
invocationSource: "on_demand",
|
|
triggerDetail: "manual",
|
|
status: "queued",
|
|
wakeupRequestId,
|
|
contextSnapshot: { issueId },
|
|
});
|
|
await db.insert(issues).values({
|
|
id: issueId,
|
|
companyId,
|
|
title: "Queued issue work",
|
|
status: "in_progress",
|
|
priority: "high",
|
|
assigneeAgentId: agentId,
|
|
executionRunId: queuedRunId,
|
|
});
|
|
await db
|
|
.update(agentWakeupRequests)
|
|
.set({ runId: queuedRunId })
|
|
.where(eq(agentWakeupRequests.id, wakeupRequestId));
|
|
|
|
const run = await heartbeat.wakeup(agentId, {
|
|
source: "on_demand",
|
|
triggerDetail: "manual",
|
|
payload: { issueId },
|
|
});
|
|
|
|
expect(run?.id).toBe(queuedRunId);
|
|
expect(mockAdapterExecute).not.toHaveBeenCalled();
|
|
|
|
const wakeups = await db
|
|
.select({
|
|
status: agentWakeupRequests.status,
|
|
reason: agentWakeupRequests.reason,
|
|
runId: agentWakeupRequests.runId,
|
|
})
|
|
.from(agentWakeupRequests)
|
|
.where(eq(agentWakeupRequests.agentId, agentId));
|
|
|
|
expect(wakeups).toEqual(
|
|
expect.arrayContaining([
|
|
expect.objectContaining({
|
|
status: "coalesced",
|
|
reason: "issue_execution_same_name",
|
|
runId: queuedRunId,
|
|
}),
|
|
]),
|
|
);
|
|
});
|
|
|
|
it("skips wakes before queueing when per-agent daily cost cap is reached", async () => {
|
|
const { companyId, agentId } = await seedCompanyAndAgent({
|
|
heartbeatConfig: {
|
|
maxDailyCostCents: 75,
|
|
},
|
|
});
|
|
await db.insert(costEvents).values({
|
|
companyId,
|
|
agentId,
|
|
provider: "test",
|
|
biller: "test",
|
|
billingType: "metered_api",
|
|
model: "test-model",
|
|
inputTokens: 100,
|
|
outputTokens: 50,
|
|
costCents: 75,
|
|
occurredAt: new Date(),
|
|
});
|
|
|
|
const run = await heartbeat.wakeup(agentId, {
|
|
source: "on_demand",
|
|
triggerDetail: "manual",
|
|
});
|
|
|
|
expect(run).toBeNull();
|
|
expect(mockAdapterExecute).not.toHaveBeenCalled();
|
|
|
|
const [wakeup] = await db
|
|
.select({
|
|
status: agentWakeupRequests.status,
|
|
reason: agentWakeupRequests.reason,
|
|
payload: agentWakeupRequests.payload,
|
|
})
|
|
.from(agentWakeupRequests)
|
|
.where(eq(agentWakeupRequests.agentId, agentId));
|
|
|
|
expect(wakeup).toMatchObject({
|
|
status: "skipped",
|
|
reason: "heartbeat.daily_cost_limit",
|
|
});
|
|
expect(wakeup?.payload).toMatchObject({
|
|
heartbeatSkip: {
|
|
observed: 75,
|
|
limit: 75,
|
|
},
|
|
});
|
|
});
|
|
|
|
it("treats zero daily cost cap as a hard stop", async () => {
|
|
const { agentId } = await seedCompanyAndAgent({
|
|
heartbeatConfig: {
|
|
maxDailyCostCents: 0,
|
|
},
|
|
});
|
|
|
|
const run = await heartbeat.wakeup(agentId, {
|
|
source: "on_demand",
|
|
triggerDetail: "manual",
|
|
});
|
|
|
|
expect(run).toBeNull();
|
|
expect(mockAdapterExecute).not.toHaveBeenCalled();
|
|
|
|
const [wakeup] = await db
|
|
.select({
|
|
status: agentWakeupRequests.status,
|
|
reason: agentWakeupRequests.reason,
|
|
payload: agentWakeupRequests.payload,
|
|
})
|
|
.from(agentWakeupRequests)
|
|
.where(eq(agentWakeupRequests.agentId, agentId));
|
|
|
|
expect(wakeup).toMatchObject({
|
|
status: "skipped",
|
|
reason: "heartbeat.daily_cost_limit",
|
|
});
|
|
expect(wakeup?.payload).toMatchObject({
|
|
heartbeatSkip: {
|
|
observed: 0,
|
|
limit: 0,
|
|
},
|
|
});
|
|
});
|
|
|
|
it("skips already queued runs before adapter execution when the daily cost cap is reached", async () => {
|
|
const { companyId, agentId } = await seedCompanyAndAgent({
|
|
heartbeatConfig: {
|
|
maxDailyCostCents: 75,
|
|
},
|
|
});
|
|
const wakeupRequestId = randomUUID();
|
|
const runId = randomUUID();
|
|
await db.insert(agentWakeupRequests).values({
|
|
id: wakeupRequestId,
|
|
companyId,
|
|
agentId,
|
|
source: "on_demand",
|
|
triggerDetail: "manual",
|
|
reason: "manual",
|
|
payload: {},
|
|
status: "queued",
|
|
});
|
|
await db.insert(heartbeatRuns).values({
|
|
id: runId,
|
|
companyId,
|
|
agentId,
|
|
invocationSource: "on_demand",
|
|
triggerDetail: "manual",
|
|
status: "queued",
|
|
wakeupRequestId,
|
|
contextSnapshot: {},
|
|
});
|
|
await db
|
|
.update(agentWakeupRequests)
|
|
.set({ runId })
|
|
.where(eq(agentWakeupRequests.id, wakeupRequestId));
|
|
await db.insert(costEvents).values({
|
|
companyId,
|
|
agentId,
|
|
provider: "test",
|
|
biller: "test",
|
|
billingType: "metered_api",
|
|
model: "test-model",
|
|
inputTokens: 100,
|
|
outputTokens: 50,
|
|
costCents: 75,
|
|
occurredAt: new Date(),
|
|
});
|
|
|
|
await heartbeat.resumeQueuedRuns();
|
|
|
|
expect(mockAdapterExecute).not.toHaveBeenCalled();
|
|
|
|
const [run] = await db
|
|
.select({
|
|
status: heartbeatRuns.status,
|
|
errorCode: heartbeatRuns.errorCode,
|
|
resultJson: heartbeatRuns.resultJson,
|
|
})
|
|
.from(heartbeatRuns)
|
|
.where(eq(heartbeatRuns.id, runId));
|
|
const [wakeup] = await db
|
|
.select({
|
|
status: agentWakeupRequests.status,
|
|
error: agentWakeupRequests.error,
|
|
})
|
|
.from(agentWakeupRequests)
|
|
.where(eq(agentWakeupRequests.id, wakeupRequestId));
|
|
|
|
expect(run).toMatchObject({
|
|
status: "cancelled",
|
|
errorCode: "heartbeat.daily_cost_limit",
|
|
});
|
|
expect(run?.resultJson).toMatchObject({
|
|
stopReason: "heartbeat.daily_cost_limit",
|
|
observed: 75,
|
|
limit: 75,
|
|
});
|
|
expect(wakeup).toMatchObject({
|
|
status: "skipped",
|
|
error: expect.stringContaining("per-day heartbeat budget cap"),
|
|
});
|
|
});
|
|
|
|
it("skips already queued issue runs at the daily run cap and releases the execution lock", async () => {
|
|
const { companyId, agentId } = await seedCompanyAndAgent({
|
|
heartbeatConfig: {
|
|
maxDailyRuns: 1,
|
|
},
|
|
});
|
|
const issueId = randomUUID();
|
|
const wakeupRequestId = randomUUID();
|
|
const queuedRunId = randomUUID();
|
|
await db.insert(heartbeatRuns).values({
|
|
id: randomUUID(),
|
|
companyId,
|
|
agentId,
|
|
invocationSource: "on_demand",
|
|
triggerDetail: "manual",
|
|
status: "succeeded",
|
|
createdAt: new Date(),
|
|
startedAt: new Date(),
|
|
finishedAt: new Date(),
|
|
contextSnapshot: {},
|
|
});
|
|
await db.insert(agentWakeupRequests).values({
|
|
id: wakeupRequestId,
|
|
companyId,
|
|
agentId,
|
|
source: "on_demand",
|
|
triggerDetail: "manual",
|
|
reason: "manual",
|
|
payload: { issueId },
|
|
status: "queued",
|
|
});
|
|
await db.insert(heartbeatRuns).values({
|
|
id: queuedRunId,
|
|
companyId,
|
|
agentId,
|
|
invocationSource: "on_demand",
|
|
triggerDetail: "manual",
|
|
status: "queued",
|
|
wakeupRequestId,
|
|
contextSnapshot: { issueId },
|
|
});
|
|
await db.insert(issues).values({
|
|
id: issueId,
|
|
companyId,
|
|
title: "Queued issue work",
|
|
status: "in_progress",
|
|
priority: "high",
|
|
assigneeAgentId: agentId,
|
|
executionRunId: queuedRunId,
|
|
});
|
|
await db
|
|
.update(agentWakeupRequests)
|
|
.set({ runId: queuedRunId })
|
|
.where(eq(agentWakeupRequests.id, wakeupRequestId));
|
|
|
|
await heartbeat.resumeQueuedRuns();
|
|
|
|
expect(mockAdapterExecute).not.toHaveBeenCalled();
|
|
|
|
const [run] = await db
|
|
.select({
|
|
status: heartbeatRuns.status,
|
|
errorCode: heartbeatRuns.errorCode,
|
|
})
|
|
.from(heartbeatRuns)
|
|
.where(eq(heartbeatRuns.id, queuedRunId));
|
|
const [wakeup] = await db
|
|
.select({ status: agentWakeupRequests.status })
|
|
.from(agentWakeupRequests)
|
|
.where(eq(agentWakeupRequests.id, wakeupRequestId));
|
|
const [issue] = await db
|
|
.select({ executionRunId: issues.executionRunId })
|
|
.from(issues)
|
|
.where(eq(issues.id, issueId));
|
|
|
|
expect(run).toMatchObject({
|
|
status: "cancelled",
|
|
errorCode: "heartbeat.daily_run_limit",
|
|
});
|
|
expect(wakeup).toMatchObject({ status: "skipped" });
|
|
expect(issue?.executionRunId).toBeNull();
|
|
});
|
|
|
|
it("promotes deferred issue wakes when a queued holder is cancelled by the daily run cap", async () => {
|
|
const { companyId, agentId } = await seedCompanyAndAgent({
|
|
heartbeatConfig: {
|
|
maxDailyRuns: 1,
|
|
},
|
|
});
|
|
const peerAgentId = randomUUID();
|
|
const issueId = randomUUID();
|
|
const wakeupRequestId = randomUUID();
|
|
const queuedRunId = randomUUID();
|
|
const deferredWakeupId = randomUUID();
|
|
await db.insert(agents).values({
|
|
id: peerAgentId,
|
|
companyId,
|
|
name: "PeerAgent",
|
|
role: "engineer",
|
|
status: "active",
|
|
adapterType: "codex_local",
|
|
adapterConfig: {},
|
|
runtimeConfig: {
|
|
heartbeat: {
|
|
wakeOnDemand: true,
|
|
maxConcurrentRuns: 1,
|
|
},
|
|
},
|
|
permissions: {},
|
|
});
|
|
await db.insert(heartbeatRuns).values({
|
|
id: randomUUID(),
|
|
companyId,
|
|
agentId,
|
|
invocationSource: "on_demand",
|
|
triggerDetail: "manual",
|
|
status: "succeeded",
|
|
createdAt: new Date(),
|
|
startedAt: new Date(),
|
|
finishedAt: new Date(),
|
|
contextSnapshot: {},
|
|
});
|
|
await db.insert(agentWakeupRequests).values({
|
|
id: wakeupRequestId,
|
|
companyId,
|
|
agentId,
|
|
source: "on_demand",
|
|
triggerDetail: "manual",
|
|
reason: "manual",
|
|
payload: { issueId },
|
|
status: "queued",
|
|
});
|
|
await db.insert(heartbeatRuns).values({
|
|
id: queuedRunId,
|
|
companyId,
|
|
agentId,
|
|
invocationSource: "on_demand",
|
|
triggerDetail: "manual",
|
|
status: "queued",
|
|
wakeupRequestId,
|
|
contextSnapshot: { issueId },
|
|
});
|
|
await db.insert(issues).values({
|
|
id: issueId,
|
|
companyId,
|
|
title: "Queued issue work",
|
|
status: "in_progress",
|
|
priority: "high",
|
|
assigneeAgentId: agentId,
|
|
executionRunId: queuedRunId,
|
|
});
|
|
await db.insert(agentWakeupRequests).values({
|
|
id: deferredWakeupId,
|
|
companyId,
|
|
agentId: peerAgentId,
|
|
source: "comment",
|
|
triggerDetail: "mention",
|
|
reason: "issue_execution_deferred",
|
|
payload: {
|
|
issueId,
|
|
_paperclipWakeContext: {
|
|
issueId,
|
|
wakeReason: "issue_mention",
|
|
},
|
|
},
|
|
status: "deferred_issue_execution",
|
|
});
|
|
await db
|
|
.update(agentWakeupRequests)
|
|
.set({ runId: queuedRunId })
|
|
.where(eq(agentWakeupRequests.id, wakeupRequestId));
|
|
|
|
await heartbeat.resumeQueuedRuns();
|
|
await waitForCondition(async () => {
|
|
const [deferred] = await db
|
|
.select({ status: agentWakeupRequests.status, runId: agentWakeupRequests.runId })
|
|
.from(agentWakeupRequests)
|
|
.where(eq(agentWakeupRequests.id, deferredWakeupId));
|
|
return Boolean(deferred?.runId) && deferred?.status !== "deferred_issue_execution";
|
|
});
|
|
|
|
const [deferred] = await db
|
|
.select({ status: agentWakeupRequests.status, runId: agentWakeupRequests.runId })
|
|
.from(agentWakeupRequests)
|
|
.where(eq(agentWakeupRequests.id, deferredWakeupId));
|
|
const [promotedRun] = deferred?.runId
|
|
? await db
|
|
.select({ agentId: heartbeatRuns.agentId })
|
|
.from(heartbeatRuns)
|
|
.where(eq(heartbeatRuns.id, deferred.runId))
|
|
: [];
|
|
|
|
expect(deferred?.status).not.toBe("deferred_issue_execution");
|
|
expect(promotedRun?.agentId).toBe(peerAgentId);
|
|
});
|
|
|
|
it("cancels queued max-turn continuations when another continuation owns the issue lock", async () => {
|
|
const { companyId, agentId } = await seedCompanyAndAgent();
|
|
const issueId = randomUUID();
|
|
const lockOwnerRunId = randomUUID();
|
|
|
|
await db.insert(heartbeatRuns).values({
|
|
id: lockOwnerRunId,
|
|
companyId,
|
|
agentId,
|
|
invocationSource: "automation",
|
|
triggerDetail: "system",
|
|
status: "scheduled_retry",
|
|
scheduledRetryReason: MAX_TURN_CONTINUATION_RETRY_REASON,
|
|
scheduledRetryAttempt: 1,
|
|
scheduledRetryAt: new Date("2026-04-20T12:00:00.000Z"),
|
|
contextSnapshot: {
|
|
issueId,
|
|
wakeReason: MAX_TURN_CONTINUATION_WAKE_REASON,
|
|
retryReason: MAX_TURN_CONTINUATION_RETRY_REASON,
|
|
},
|
|
});
|
|
|
|
await db.insert(issues).values({
|
|
id: issueId,
|
|
companyId,
|
|
title: "Duplicate max-turn continuation",
|
|
status: "in_progress",
|
|
priority: "medium",
|
|
assigneeAgentId: agentId,
|
|
executionRunId: lockOwnerRunId,
|
|
executionAgentNameKey: "claudecoder",
|
|
executionLockedAt: new Date("2026-04-20T11:59:00.000Z"),
|
|
});
|
|
|
|
const { runId, wakeupRequestId } = await seedQueuedRun({
|
|
companyId,
|
|
agentId,
|
|
issueId,
|
|
wakeReason: MAX_TURN_CONTINUATION_WAKE_REASON,
|
|
invocationSource: "automation",
|
|
scheduledRetryReason: MAX_TURN_CONTINUATION_RETRY_REASON,
|
|
contextExtras: {
|
|
retryReason: MAX_TURN_CONTINUATION_RETRY_REASON,
|
|
},
|
|
});
|
|
|
|
await heartbeat.resumeQueuedRuns();
|
|
|
|
await waitForCondition(async () => {
|
|
const run = await db
|
|
.select({ status: heartbeatRuns.status })
|
|
.from(heartbeatRuns)
|
|
.where(eq(heartbeatRuns.id, runId))
|
|
.then((rows) => rows[0] ?? null);
|
|
return run?.status === "cancelled";
|
|
});
|
|
|
|
const [run, wakeup, issue] = await Promise.all([
|
|
db
|
|
.select({
|
|
status: heartbeatRuns.status,
|
|
errorCode: heartbeatRuns.errorCode,
|
|
resultJson: heartbeatRuns.resultJson,
|
|
})
|
|
.from(heartbeatRuns)
|
|
.where(eq(heartbeatRuns.id, runId))
|
|
.then((rows) => rows[0] ?? null),
|
|
db
|
|
.select({ status: agentWakeupRequests.status, error: agentWakeupRequests.error })
|
|
.from(agentWakeupRequests)
|
|
.where(eq(agentWakeupRequests.id, wakeupRequestId))
|
|
.then((rows) => rows[0] ?? null),
|
|
db
|
|
.select({ executionRunId: issues.executionRunId })
|
|
.from(issues)
|
|
.where(eq(issues.id, issueId))
|
|
.then((rows) => rows[0] ?? null),
|
|
]);
|
|
|
|
expect(run?.status).toBe("cancelled");
|
|
expect(run?.errorCode).toBe("issue_execution_lock_changed");
|
|
expect(run?.resultJson).toMatchObject({ stopReason: "issue_execution_lock_changed" });
|
|
expect(wakeup?.status).toBe("skipped");
|
|
expect(wakeup?.error).toContain("execution lock");
|
|
expect(issue?.executionRunId).toBe(lockOwnerRunId);
|
|
expect(countExecuteCallsForRun(runId)).toBe(0);
|
|
});
|
|
|
|
it.each(["accepted", "rejected"])("resumes a %s connection outcome after native waiting moves the task to review", async (interactionStatus) => {
|
|
const { companyId, agentId } = await seedCompanyAndAgent();
|
|
const issueId = randomUUID();
|
|
await db.insert(issues).values({ id: issueId, companyId, title: "Waiting for connection", status: "in_review", priority: "medium", assigneeAgentId: agentId });
|
|
const { runId } = await seedQueuedRun({ companyId, agentId, issueId, wakeReason: "issue_commented", invocationSource: "automation",
|
|
contextExtras: { interactionId: randomUUID(), interactionKind: "connection_intent", interactionStatus,
|
|
interactionResolvedAt: new Date().toISOString(), mutation: "interaction", source: "connection_intent.resolved", forceFreshSession: true } });
|
|
await heartbeat.resumeQueuedRuns();
|
|
await waitForCondition(async () => (await db.select({ status: heartbeatRuns.status }).from(heartbeatRuns).where(eq(heartbeatRuns.id, runId)))[0]?.status === "succeeded");
|
|
expect(countExecuteCallsForRun(runId)).toBe(1);
|
|
});
|
|
});
|