mirror of
https://github.com/paperclipai/paperclip.git
synced 2026-10-06 10:48:12 +02:00
test(heartbeat): pin 8 uncovered branches of the deferred issue-execution wake state machine (#13103)
## Thinking Path > - Paperclip is the open source app people use to manage AI agents for work > - The server manages issue execution, queued comments, and agent wake state > - Deferred wakes can pass through several failure, hold, rollback, and steering branches > - These branches had no direct tests, so a later change could fail without clear evidence > - This pull request adds characterization tests for eight uncovered branches > - The benefit is clear test evidence for future changes to deferred issue execution ## Linked Issues or Issue Description **What existing behavior does this improve?** The server behavior for deferred issue-execution wakes and queued-comment steering transactions. **Current behavior** The code handles missing agents, cross-company agents, pause holds, rollback, steering outcomes, and invalid reorder sets. These branches had no direct tests. **Proposed behavior** Keep the current behavior and test each branch against the current implementation. **Reason and benefit** The tests make silent behavior changes visible. They also protect the wake queue while later work moves this logic into a dedicated module. **Breaking changes** None. This pull request changes no production code and no runtime behavior. Related public pull requests: #12671, #11168, and #10199. ## What Changed - Add tests for missing and cross-company deferred agents. - Add tests for pause-hold promotion and cancellation. - Add a test for atomic rollback when the responsible user cannot resolve. - Add tests for successful, timed-out, and rejected native steering. - Add a test for invalid queued-comment reorder input. ## Verification - Run `server/src/__tests__/heartbeat-comment-wake-batching.test.ts` against an embedded Postgres database. - Run `server/src/__tests__/issue-queued-comments-routes.test.ts` against an embedded Postgres database. - Confirm that all eight new tests pass. - Confirm that the pull request CI checks pass. ## Risks Low risk. The diff changes test files only. It does not change production code, database schema, API behavior, or runtime behavior. ## Model Used Codex, GPT-5, tool use and code review support. The model did not author production code. ## Checklist - [x] I have included a thinking path that traces from project context to this change - [x] I have specified the model used (with version and capability details) - [x] I have checked ROADMAP.md and confirmed this PR does not duplicate planned core work - [x] I have searched GitHub for duplicate or related PRs and linked them above - [x] I have either (a) linked existing issues with `Fixes: #` / `Closes #` / `Refs #` OR (b) described the issue in-PR following the relevant issue template - [x] I have not referenced internal/instance-local Paperclip issues or links (only public GitHub `#NNN` / `github.com/paperclipai/paperclip` URLs) - [x] My branch name describes the change (e.g. `docs/...`, `fix/...`) and contains no internal Paperclip ticket id or instance-derived details - [x] I have run tests locally and they pass - [x] I have added or updated tests where applicable - [x] I have updated relevant documentation to reflect my changes - [x] I have considered and documented any risks above - [x] All Paperclip CI gates are green - [x] Greptile is 5/5 with no open P2s, recommendations, or follow-ups - [x] I will address all Greptile and reviewer comments before requesting merge --------- Co-authored-by: Paperclip <noreply@paperclip.ing>
This commit is contained in:
1 parent
2991a59b17
commit
7d84b183fb
2 files changed
+672
-2
No files matched your search
@@ -1,6 +1,6 @@
|
||||
import { randomUUID } from "node:crypto";
|
||||
import { createServer } from "node:http";
|
||||
import { and, asc, eq } from "drizzle-orm";
|
||||
import { and, asc, eq, sql } from "drizzle-orm";
|
||||
import { WebSocketServer } from "ws";
|
||||
import { afterAll, afterEach, beforeAll, describe, expect, it } from "vitest";
|
||||
import {
|
||||
@@ -12,6 +12,7 @@ import {
|
||||
issueComments,
|
||||
issueRecoveryActions,
|
||||
issues,
|
||||
issueTreeHolds,
|
||||
} from "@paperclipai/db";
|
||||
import { runningProcesses } from "../adapters/index.js";
|
||||
import { heartbeatService } from "../services/heartbeat.ts";
|
||||
@@ -2526,4 +2527,501 @@ describeEmbeddedPostgres("heartbeat comment wake batching", () => {
|
||||
await gateway.close();
|
||||
}
|
||||
}, 20_000);
|
||||
|
||||
it("fails a deferred wake whose agent no longer exists, then still promotes the next queued wake", async () => {
|
||||
const companyId = randomUUID();
|
||||
const finishingAgentId = randomUUID();
|
||||
const validAgentId = randomUUID();
|
||||
const missingAgentId = randomUUID();
|
||||
const issueId = randomUUID();
|
||||
const runId = randomUUID();
|
||||
const issuePrefix = `T${companyId.replace(/-/g, "").slice(0, 6).toUpperCase()}`;
|
||||
// Pin scheduling suppression off with the runtimeEnv test seam. Do not rely
|
||||
// on the ambient PAPERCLIP_IN_WORKTREE value: startNextQueuedRunForAgent
|
||||
// no-ops under suppression and would leave the promoted wake at "queued".
|
||||
const heartbeat = heartbeatService(db, { runtimeEnv: {} });
|
||||
|
||||
await db.insert(companies).values({
|
||||
id: companyId,
|
||||
name: "Paperclip",
|
||||
issuePrefix,
|
||||
requireBoardApprovalForNewAgents: false,
|
||||
defaultResponsibleUserId: "responsible-user",
|
||||
});
|
||||
await db.insert(agents).values([
|
||||
{
|
||||
id: finishingAgentId,
|
||||
companyId,
|
||||
name: "Finishing Agent",
|
||||
role: "engineer",
|
||||
status: "idle",
|
||||
adapterType: "process",
|
||||
adapterConfig: {},
|
||||
runtimeConfig: {},
|
||||
permissions: {},
|
||||
},
|
||||
{
|
||||
id: validAgentId,
|
||||
companyId,
|
||||
name: "Assignee Agent",
|
||||
role: "engineer",
|
||||
status: "idle",
|
||||
adapterType: "process",
|
||||
adapterConfig: {},
|
||||
runtimeConfig: {},
|
||||
permissions: {},
|
||||
},
|
||||
]);
|
||||
await db.insert(heartbeatRuns).values({
|
||||
id: runId,
|
||||
companyId,
|
||||
agentId: finishingAgentId,
|
||||
invocationSource: "assignment",
|
||||
triggerDetail: "system",
|
||||
status: "running",
|
||||
runtimeMode: "legacy",
|
||||
startedAt: new Date(),
|
||||
contextSnapshot: { issueId },
|
||||
responsibleUserId: "responsible-user",
|
||||
});
|
||||
await db.insert(issues).values({
|
||||
id: issueId,
|
||||
companyId,
|
||||
title: "Continue past a missing deferred agent",
|
||||
status: "in_progress",
|
||||
priority: "medium",
|
||||
responsibleUserId: "responsible-user",
|
||||
assigneeAgentId: validAgentId,
|
||||
issueNumber: 1,
|
||||
identifier: `${issuePrefix}-1`,
|
||||
executionRunId: runId,
|
||||
});
|
||||
|
||||
const missingAgentWakeId = randomUUID();
|
||||
// The agent_id foreign key always holds in the running system, so a wake row
|
||||
// can never outlive its agent through the application. Bypass the check for
|
||||
// this one insert to pin the release code's defensive branch for that state.
|
||||
await db.transaction(async (tx) => {
|
||||
await tx.execute(sql`set local session_replication_role = 'replica'`);
|
||||
await tx.insert(agentWakeupRequests).values({
|
||||
id: missingAgentWakeId,
|
||||
companyId,
|
||||
agentId: missingAgentId,
|
||||
source: "automation",
|
||||
reason: "issue_commented",
|
||||
status: "deferred_issue_execution",
|
||||
requestedByActorType: "system",
|
||||
requestedByActorId: "test",
|
||||
requestedAt: new Date("2026-08-22T15:00:00.000Z"),
|
||||
payload: { issueId },
|
||||
});
|
||||
});
|
||||
|
||||
const validWakeId = randomUUID();
|
||||
await db.insert(agentWakeupRequests).values({
|
||||
id: validWakeId,
|
||||
companyId,
|
||||
agentId: validAgentId,
|
||||
source: "automation",
|
||||
reason: "issue_commented",
|
||||
status: "deferred_issue_execution",
|
||||
requestedByActorType: "system",
|
||||
requestedByActorId: "test",
|
||||
requestedAt: new Date("2026-08-22T15:01:00.000Z"),
|
||||
payload: { issueId },
|
||||
});
|
||||
|
||||
// A plain legacy cancel now stops promotion early for board reconciliation
|
||||
// (see legacyExecutionNeedsReconciliation in legacy-execution-recovery.ts).
|
||||
// Cancel as an in-flight workspace wait instead. That shape still reaches
|
||||
// the deferred-wake promotion loop under test.
|
||||
await heartbeat.cancelRun(runId, undefined, {
|
||||
errorCode: "workspace_busy",
|
||||
resultJson: {
|
||||
executionRecovery: { kind: "workspace_wait", providerWorkStarted: false },
|
||||
},
|
||||
});
|
||||
|
||||
const [missingAgentWake, validWake, issueRow] = await Promise.all([
|
||||
db.select().from(agentWakeupRequests).where(eq(agentWakeupRequests.id, missingAgentWakeId)).then((rows) => rows[0]),
|
||||
db.select().from(agentWakeupRequests).where(eq(agentWakeupRequests.id, validWakeId)).then((rows) => rows[0]),
|
||||
db.select().from(issues).where(eq(issues.id, issueId)).then((rows) => rows[0]),
|
||||
]);
|
||||
|
||||
expect(missingAgentWake).toMatchObject({
|
||||
status: "failed",
|
||||
error: "Deferred wake could not be promoted: agent is not invokable",
|
||||
});
|
||||
// The promotion writes "queued", then releaseIssueExecutionAndPromote
|
||||
// immediately calls startNextQueuedRunForAgent for the idle promoted
|
||||
// agent, which claims the run in the same call. Assert the settled
|
||||
// state, not the intermediate one.
|
||||
expect(validWake?.status).toBe("claimed");
|
||||
expect(validWake?.runId).not.toBeNull();
|
||||
expect(issueRow?.executionRunId).toBe(validWake?.runId);
|
||||
});
|
||||
|
||||
it("fails a deferred wake with the same status and error text when its agent belongs to another company", async () => {
|
||||
const companyId = randomUUID();
|
||||
const otherCompanyId = randomUUID();
|
||||
const finishingAgentId = randomUUID();
|
||||
const crossCompanyAgentId = randomUUID();
|
||||
const issueId = randomUUID();
|
||||
const runId = randomUUID();
|
||||
const issuePrefix = `T${companyId.replace(/-/g, "").slice(0, 6).toUpperCase()}`;
|
||||
const otherIssuePrefix = `T${otherCompanyId.replace(/-/g, "").slice(0, 6).toUpperCase()}`;
|
||||
const heartbeat = heartbeatService(db);
|
||||
|
||||
await db.insert(companies).values([
|
||||
{
|
||||
id: companyId,
|
||||
name: "Paperclip",
|
||||
issuePrefix,
|
||||
requireBoardApprovalForNewAgents: false,
|
||||
defaultResponsibleUserId: "responsible-user",
|
||||
},
|
||||
{
|
||||
id: otherCompanyId,
|
||||
name: "Other Paperclip",
|
||||
issuePrefix: otherIssuePrefix,
|
||||
requireBoardApprovalForNewAgents: false,
|
||||
},
|
||||
]);
|
||||
await db.insert(agents).values([
|
||||
{
|
||||
id: finishingAgentId,
|
||||
companyId,
|
||||
name: "Finishing Agent",
|
||||
role: "engineer",
|
||||
status: "idle",
|
||||
adapterType: "process",
|
||||
adapterConfig: {},
|
||||
runtimeConfig: {},
|
||||
permissions: {},
|
||||
},
|
||||
{
|
||||
id: crossCompanyAgentId,
|
||||
companyId: otherCompanyId,
|
||||
name: "Other Company Agent",
|
||||
role: "engineer",
|
||||
status: "idle",
|
||||
adapterType: "process",
|
||||
adapterConfig: {},
|
||||
runtimeConfig: {},
|
||||
permissions: {},
|
||||
},
|
||||
]);
|
||||
await db.insert(heartbeatRuns).values({
|
||||
id: runId,
|
||||
companyId,
|
||||
agentId: finishingAgentId,
|
||||
invocationSource: "assignment",
|
||||
triggerDetail: "system",
|
||||
status: "running",
|
||||
runtimeMode: "legacy",
|
||||
startedAt: new Date(),
|
||||
contextSnapshot: { issueId },
|
||||
responsibleUserId: "responsible-user",
|
||||
});
|
||||
await db.insert(issues).values({
|
||||
id: issueId,
|
||||
companyId,
|
||||
title: "Cross-company deferred agent",
|
||||
status: "in_progress",
|
||||
priority: "medium",
|
||||
responsibleUserId: "responsible-user",
|
||||
assigneeAgentId: finishingAgentId,
|
||||
issueNumber: 1,
|
||||
identifier: `${issuePrefix}-1`,
|
||||
executionRunId: runId,
|
||||
});
|
||||
const wakeId = randomUUID();
|
||||
await db.insert(agentWakeupRequests).values({
|
||||
id: wakeId,
|
||||
companyId,
|
||||
agentId: crossCompanyAgentId,
|
||||
source: "automation",
|
||||
reason: "issue_commented",
|
||||
status: "deferred_issue_execution",
|
||||
requestedByActorType: "system",
|
||||
requestedByActorId: "test",
|
||||
payload: { issueId },
|
||||
});
|
||||
|
||||
// A plain legacy cancel now stops promotion early for board reconciliation
|
||||
// (see legacyExecutionNeedsReconciliation in legacy-execution-recovery.ts).
|
||||
// Cancel as an in-flight workspace wait instead. That shape still reaches
|
||||
// the deferred-wake promotion loop under test.
|
||||
await heartbeat.cancelRun(runId, undefined, {
|
||||
errorCode: "workspace_busy",
|
||||
resultJson: {
|
||||
executionRecovery: { kind: "workspace_wait", providerWorkStarted: false },
|
||||
},
|
||||
});
|
||||
|
||||
const wake = await db.select().from(agentWakeupRequests).where(eq(agentWakeupRequests.id, wakeId)).then((rows) => rows[0]);
|
||||
expect(wake).toMatchObject({
|
||||
status: "failed",
|
||||
error: "Deferred wake could not be promoted: agent is not invokable",
|
||||
});
|
||||
});
|
||||
|
||||
it("cancels a deferred wake under an active pause hold, but promotes a verified hold interaction with the hold context", async () => {
|
||||
const companyId = randomUUID();
|
||||
const finishingAgentId = randomUUID();
|
||||
const holdAgentId = randomUUID();
|
||||
const issueId = randomUUID();
|
||||
const runId = randomUUID();
|
||||
const issuePrefix = `T${companyId.replace(/-/g, "").slice(0, 6).toUpperCase()}`;
|
||||
// Pin scheduling suppression off with the runtimeEnv test seam. Do not rely
|
||||
// on the ambient PAPERCLIP_IN_WORKTREE value: startNextQueuedRunForAgent
|
||||
// no-ops under suppression and would leave the promoted wake at "queued".
|
||||
const heartbeat = heartbeatService(db, { runtimeEnv: {} });
|
||||
|
||||
await db.insert(companies).values({
|
||||
id: companyId,
|
||||
name: "Paperclip",
|
||||
issuePrefix,
|
||||
requireBoardApprovalForNewAgents: false,
|
||||
defaultResponsibleUserId: "responsible-user",
|
||||
});
|
||||
await db.insert(agents).values([
|
||||
{
|
||||
id: finishingAgentId,
|
||||
companyId,
|
||||
name: "Finishing Agent",
|
||||
role: "engineer",
|
||||
status: "idle",
|
||||
adapterType: "process",
|
||||
adapterConfig: {},
|
||||
runtimeConfig: {},
|
||||
permissions: {},
|
||||
},
|
||||
{
|
||||
id: holdAgentId,
|
||||
companyId,
|
||||
name: "Hold Interaction Agent",
|
||||
role: "engineer",
|
||||
status: "idle",
|
||||
adapterType: "process",
|
||||
adapterConfig: {},
|
||||
runtimeConfig: {},
|
||||
permissions: {},
|
||||
},
|
||||
]);
|
||||
await db.insert(heartbeatRuns).values({
|
||||
id: runId,
|
||||
companyId,
|
||||
agentId: finishingAgentId,
|
||||
invocationSource: "assignment",
|
||||
triggerDetail: "system",
|
||||
status: "running",
|
||||
runtimeMode: "legacy",
|
||||
startedAt: new Date(),
|
||||
contextSnapshot: { issueId },
|
||||
responsibleUserId: "responsible-user",
|
||||
});
|
||||
await db.insert(issues).values({
|
||||
id: issueId,
|
||||
companyId,
|
||||
title: "A pause hold gates a plain wake but not a verified one",
|
||||
status: "in_progress",
|
||||
priority: "medium",
|
||||
responsibleUserId: "responsible-user",
|
||||
assigneeAgentId: holdAgentId,
|
||||
issueNumber: 1,
|
||||
identifier: `${issuePrefix}-1`,
|
||||
executionRunId: runId,
|
||||
});
|
||||
const [hold] = await db.insert(issueTreeHolds).values({
|
||||
companyId,
|
||||
rootIssueId: issueId,
|
||||
mode: "pause",
|
||||
status: "active",
|
||||
reason: "Investigating a regression",
|
||||
}).returning();
|
||||
|
||||
const plainWakeId = randomUUID();
|
||||
await db.insert(agentWakeupRequests).values({
|
||||
id: plainWakeId,
|
||||
companyId,
|
||||
agentId: finishingAgentId,
|
||||
source: "automation",
|
||||
reason: "issue_commented",
|
||||
status: "deferred_issue_execution",
|
||||
requestedByActorType: "system",
|
||||
requestedByActorId: "test",
|
||||
requestedAt: new Date("2026-08-22T15:00:00.000Z"),
|
||||
payload: { issueId },
|
||||
});
|
||||
|
||||
const holdComment = await db.insert(issueComments).values({
|
||||
companyId,
|
||||
issueId,
|
||||
authorUserId: "hold-user",
|
||||
body: "Please continue despite the hold",
|
||||
}).returning().then((rows) => rows[0]!);
|
||||
const verifiedWakeId = randomUUID();
|
||||
await db.insert(agentWakeupRequests).values({
|
||||
id: verifiedWakeId,
|
||||
companyId,
|
||||
agentId: holdAgentId,
|
||||
source: "issue_comment",
|
||||
reason: "issue_commented",
|
||||
status: "deferred_issue_execution",
|
||||
requestedByActorType: "user",
|
||||
requestedByActorId: "hold-user",
|
||||
requestedAt: new Date("2026-08-22T15:01:00.000Z"),
|
||||
payload: {
|
||||
issueId,
|
||||
commentId: holdComment.id,
|
||||
_paperclipWakeContext: {
|
||||
wakeReason: "issue_commented",
|
||||
source: "issue.comment",
|
||||
wakeCommentIds: [holdComment.id],
|
||||
},
|
||||
},
|
||||
});
|
||||
|
||||
// A plain legacy cancel now stops promotion early for board reconciliation
|
||||
// (see legacyExecutionNeedsReconciliation in legacy-execution-recovery.ts).
|
||||
// Cancel as an in-flight workspace wait instead. That shape still reaches
|
||||
// the deferred-wake promotion loop under test.
|
||||
await heartbeat.cancelRun(runId, undefined, {
|
||||
errorCode: "workspace_busy",
|
||||
resultJson: {
|
||||
executionRecovery: { kind: "workspace_wait", providerWorkStarted: false },
|
||||
},
|
||||
});
|
||||
|
||||
const [plainWake, verifiedWake] = await Promise.all([
|
||||
db.select().from(agentWakeupRequests).where(eq(agentWakeupRequests.id, plainWakeId)).then((rows) => rows[0]),
|
||||
db.select().from(agentWakeupRequests).where(eq(agentWakeupRequests.id, verifiedWakeId)).then((rows) => rows[0]),
|
||||
]);
|
||||
expect(plainWake).toMatchObject({
|
||||
status: "cancelled",
|
||||
error: "Deferred wake suppressed by active subtree pause hold",
|
||||
});
|
||||
// Same settle-then-assert reasoning as the missing-agent test above:
|
||||
// the idle promoted agent's run is claimed synchronously.
|
||||
expect(verifiedWake?.status).toBe("claimed");
|
||||
const promotedRun = await db
|
||||
.select({ contextSnapshot: heartbeatRuns.contextSnapshot })
|
||||
.from(heartbeatRuns)
|
||||
.where(eq(heartbeatRuns.id, verifiedWake!.runId!))
|
||||
.then((rows) => rows[0]);
|
||||
expect(promotedRun?.contextSnapshot).toMatchObject({
|
||||
treeHoldInteraction: true,
|
||||
activeTreeHold: {
|
||||
holdId: hold!.id,
|
||||
rootIssueId: issueId,
|
||||
mode: "pause",
|
||||
reason: "Investigating a regression",
|
||||
interaction: true,
|
||||
},
|
||||
});
|
||||
});
|
||||
|
||||
it("rolls back the wake row, the run row, and the issue lock together when the responsible user cannot resolve", async () => {
|
||||
const companyId = randomUUID();
|
||||
const finishingAgentId = randomUUID();
|
||||
const deferredAgentId = randomUUID();
|
||||
const issueId = randomUUID();
|
||||
const runId = randomUUID();
|
||||
const issuePrefix = `T${companyId.replace(/-/g, "").slice(0, 6).toUpperCase()}`;
|
||||
const heartbeat = heartbeatService(db);
|
||||
|
||||
await db.insert(companies).values({
|
||||
id: companyId,
|
||||
name: "Paperclip",
|
||||
issuePrefix,
|
||||
requireBoardApprovalForNewAgents: false,
|
||||
// No defaultResponsibleUserId: the company default must not resolve this wake.
|
||||
});
|
||||
await db.insert(agents).values([
|
||||
{
|
||||
id: finishingAgentId,
|
||||
companyId,
|
||||
name: "Finishing Agent",
|
||||
role: "engineer",
|
||||
status: "idle",
|
||||
adapterType: "process",
|
||||
adapterConfig: {},
|
||||
runtimeConfig: {},
|
||||
permissions: {},
|
||||
},
|
||||
{
|
||||
id: deferredAgentId,
|
||||
companyId,
|
||||
name: "Deferred Agent",
|
||||
role: "engineer",
|
||||
status: "idle",
|
||||
adapterType: "process",
|
||||
adapterConfig: {},
|
||||
runtimeConfig: {},
|
||||
permissions: {},
|
||||
},
|
||||
]);
|
||||
await db.insert(heartbeatRuns).values({
|
||||
id: runId,
|
||||
companyId,
|
||||
agentId: finishingAgentId,
|
||||
invocationSource: "assignment",
|
||||
triggerDetail: "system",
|
||||
status: "running",
|
||||
runtimeMode: "legacy",
|
||||
startedAt: new Date(),
|
||||
contextSnapshot: { issueId },
|
||||
// No responsibleUserId: the finishing run itself must not resolve this wake.
|
||||
});
|
||||
await db.insert(issues).values({
|
||||
id: issueId,
|
||||
companyId,
|
||||
title: "No responsible user can be resolved",
|
||||
status: "in_progress",
|
||||
priority: "medium",
|
||||
assigneeAgentId: deferredAgentId,
|
||||
issueNumber: 1,
|
||||
identifier: `${issuePrefix}-1`,
|
||||
executionRunId: runId,
|
||||
// No responsibleUserId: the issue itself must not resolve this wake.
|
||||
});
|
||||
const wakeId = randomUUID();
|
||||
await db.insert(agentWakeupRequests).values({
|
||||
id: wakeId,
|
||||
companyId,
|
||||
agentId: deferredAgentId,
|
||||
source: "automation",
|
||||
reason: "issue_commented",
|
||||
status: "deferred_issue_execution",
|
||||
requestedByActorType: "system",
|
||||
requestedByActorId: "test",
|
||||
payload: { issueId },
|
||||
});
|
||||
|
||||
// A plain legacy cancel now stops promotion early for board reconciliation
|
||||
// (see legacyExecutionNeedsReconciliation in legacy-execution-recovery.ts).
|
||||
// Cancel as an in-flight workspace wait instead. That shape still reaches
|
||||
// the deferred-wake promotion loop under test.
|
||||
await expect(
|
||||
heartbeat.cancelRun(runId, undefined, {
|
||||
errorCode: "workspace_busy",
|
||||
resultJson: {
|
||||
executionRecovery: { kind: "workspace_wait", providerWorkStarted: false },
|
||||
},
|
||||
}),
|
||||
).rejects.toMatchObject({
|
||||
status: 422,
|
||||
details: expect.objectContaining({ code: "responsible_user_unresolved" }),
|
||||
});
|
||||
|
||||
const [wake, issueRow, runs] = await Promise.all([
|
||||
db.select().from(agentWakeupRequests).where(eq(agentWakeupRequests.id, wakeId)).then((rows) => rows[0]),
|
||||
db.select().from(issues).where(eq(issues.id, issueId)).then((rows) => rows[0]),
|
||||
db.select().from(heartbeatRuns).where(eq(heartbeatRuns.companyId, companyId)),
|
||||
]);
|
||||
expect(wake?.status).toBe("deferred_issue_execution");
|
||||
expect(issueRow?.executionRunId).toBe(runId);
|
||||
expect(runs).toHaveLength(1);
|
||||
});
|
||||
});
|
||||
@@ -2,7 +2,7 @@ import { randomUUID } from "node:crypto";
|
||||
import express from "express";
|
||||
import request from "supertest";
|
||||
import { eq } from "drizzle-orm";
|
||||
import { afterAll, afterEach, beforeAll, describe, expect, it } from "vitest";
|
||||
import { afterAll, afterEach, beforeAll, describe, expect, it, vi } from "vitest";
|
||||
import {
|
||||
agentWakeupRequests,
|
||||
agents,
|
||||
@@ -14,15 +14,25 @@ import {
|
||||
heartbeatRuns,
|
||||
issueComments,
|
||||
issues,
|
||||
runIdentityContexts,
|
||||
} from "@paperclipai/db";
|
||||
import { errorHandler } from "../middleware/index.js";
|
||||
import { issueRoutes } from "../routes/issues.js";
|
||||
import { heartbeatService } from "../services/heartbeat.js";
|
||||
import { reconcileSteeredIdentity } from "../services/run-identity.js";
|
||||
import {
|
||||
getEmbeddedPostgresTestSupport,
|
||||
startEmbeddedPostgresTestDatabase,
|
||||
} from "./helpers/embedded-postgres.js";
|
||||
|
||||
const steerNativeSessionMock = vi.hoisted(() => vi.fn());
|
||||
vi.mock("../services/native-runtime/native-session-executor.js", async (importOriginal) => {
|
||||
const actual = await importOriginal<typeof import("../services/native-runtime/native-session-executor.js")>();
|
||||
steerNativeSessionMock.mockImplementation(actual.steerNativeSession);
|
||||
return { ...actual, steerNativeSession: steerNativeSessionMock };
|
||||
});
|
||||
const { NativeSessionSteeringError } = await import("../services/native-runtime/native-session-executor.js");
|
||||
|
||||
const embeddedPostgresSupport = await getEmbeddedPostgresTestSupport();
|
||||
const describeEmbeddedPostgres = embeddedPostgresSupport.supported ? describe.sequential : describe.skip;
|
||||
|
||||
@@ -679,4 +689,166 @@ describeEmbeddedPostgres("issue queued-comment routes", () => {
|
||||
}
|
||||
await heartbeat.drainActiveRunExecutions();
|
||||
}, 30_000);
|
||||
|
||||
async function seedDispatchIdentity(seeded: Awaited<ReturnType<typeof seedQueue>>) {
|
||||
const [identity] = await db.insert(runIdentityContexts).values({
|
||||
companyId: seeded.companyId,
|
||||
runId: seeded.runId,
|
||||
revision: 1,
|
||||
cause: "dispatch",
|
||||
correlationId: "dispatch",
|
||||
status: "accepted",
|
||||
acceptedAt: new Date("2026-08-22T15:00:00.000Z"),
|
||||
}).returning();
|
||||
await db.update(heartbeatRuns)
|
||||
.set({ activeIdentityContextId: identity!.id })
|
||||
.where(eq(heartbeatRuns.id, seeded.runId));
|
||||
return identity!;
|
||||
}
|
||||
|
||||
it("writes the identity acceptance, the wake payload, the run acknowledgement, and the activity row on a successful steering transaction", async () => {
|
||||
const seeded = await seedQueue();
|
||||
await seedDispatchIdentity(seeded);
|
||||
steerNativeSessionMock.mockResolvedValueOnce({ turnId: "turn-1" });
|
||||
|
||||
const initial = await request(app(seeded.companyId))
|
||||
.get(`/api/issues/${seeded.issueId}/queued-comments`);
|
||||
const steered = await request(app(seeded.companyId))
|
||||
.post(`/api/issues/${seeded.issueId}/queued-comments/${seeded.commentIds[0]}/steer`)
|
||||
.send({ queueId: seeded.wakeId, targetRunId: seeded.runId, revision: initial.body.revision });
|
||||
|
||||
expect(steered.status, JSON.stringify(steered.body)).toBe(200);
|
||||
|
||||
const steeringIdentity = await db
|
||||
.select()
|
||||
.from(runIdentityContexts)
|
||||
.where(eq(runIdentityContexts.messageId, seeded.commentIds[0]))
|
||||
.then((rows) => rows[0]);
|
||||
expect(steeringIdentity).toMatchObject({
|
||||
status: "accepted",
|
||||
responsibleUserId: "queue-owner",
|
||||
cause: "steering",
|
||||
});
|
||||
|
||||
const wake = await db
|
||||
.select({ payload: agentWakeupRequests.payload })
|
||||
.from(agentWakeupRequests)
|
||||
.where(eq(agentWakeupRequests.id, seeded.wakeId))
|
||||
.then((rows) => rows[0]);
|
||||
expect((wake?.payload as any)?._paperclipWakeContext?.wakeCommentIds).toEqual([seeded.commentIds[1]]);
|
||||
|
||||
const run = await db
|
||||
.select({ resultJson: heartbeatRuns.resultJson, activeIdentityContextId: heartbeatRuns.activeIdentityContextId })
|
||||
.from(heartbeatRuns)
|
||||
.where(eq(heartbeatRuns.id, seeded.runId))
|
||||
.then((rows) => rows[0]);
|
||||
expect((run?.resultJson as any)?.queuedSteeringAcknowledgements?.[seeded.commentIds[0]]).toMatchObject({
|
||||
status: "acknowledged",
|
||||
queueId: seeded.wakeId,
|
||||
turnId: "turn-1",
|
||||
});
|
||||
expect(run?.activeIdentityContextId).toBe(steeringIdentity!.id);
|
||||
|
||||
const activity = await db
|
||||
.select({ details: activityLog.details })
|
||||
.from(activityLog)
|
||||
.where(eq(activityLog.action, "issue.queued_comment_steered"))
|
||||
.then((rows) => rows[0]);
|
||||
expect(activity?.details).toMatchObject({
|
||||
commentId: seeded.commentIds[0],
|
||||
targetRunId: seeded.runId,
|
||||
turnId: "turn-1",
|
||||
duplicate: false,
|
||||
});
|
||||
});
|
||||
|
||||
it("keeps the identity pending after a steering timeout, then reconciles it on a later acknowledgement", async () => {
|
||||
const seeded = await seedQueue();
|
||||
await seedDispatchIdentity(seeded);
|
||||
steerNativeSessionMock.mockRejectedValueOnce(
|
||||
new NativeSessionSteeringError("steering_timeout", "The provider did not acknowledge steering in time."),
|
||||
);
|
||||
|
||||
const initial = await request(app(seeded.companyId))
|
||||
.get(`/api/issues/${seeded.issueId}/queued-comments`);
|
||||
const steered = await request(app(seeded.companyId))
|
||||
.post(`/api/issues/${seeded.issueId}/queued-comments/${seeded.commentIds[0]}/steer`)
|
||||
.send({ queueId: seeded.wakeId, targetRunId: seeded.runId, revision: initial.body.revision });
|
||||
|
||||
expect(steered.status).toBe(409);
|
||||
expect(steered.body.details).toMatchObject({ code: "steering_timeout", retryable: true });
|
||||
|
||||
const pending = await db
|
||||
.select()
|
||||
.from(runIdentityContexts)
|
||||
.where(eq(runIdentityContexts.messageId, seeded.commentIds[0]))
|
||||
.then((rows) => rows[0]);
|
||||
expect(pending?.status).toBe("pending");
|
||||
|
||||
await reconcileSteeredIdentity(db, pending!);
|
||||
|
||||
const reconciled = await db
|
||||
.select({ status: runIdentityContexts.status })
|
||||
.from(runIdentityContexts)
|
||||
.where(eq(runIdentityContexts.id, pending!.id))
|
||||
.then((rows) => rows[0]);
|
||||
expect(reconciled?.status).toBe("accepted");
|
||||
const run = await db
|
||||
.select({ activeIdentityContextId: heartbeatRuns.activeIdentityContextId })
|
||||
.from(heartbeatRuns)
|
||||
.where(eq(heartbeatRuns.id, seeded.runId))
|
||||
.then((rows) => rows[0]);
|
||||
expect(run?.activeIdentityContextId).toBe(pending!.id);
|
||||
});
|
||||
|
||||
it("sets the identity to rejected on a definite steering rejection", async () => {
|
||||
const seeded = await seedQueue();
|
||||
await seedDispatchIdentity(seeded);
|
||||
steerNativeSessionMock.mockRejectedValueOnce(
|
||||
new NativeSessionSteeringError("steering_rejected", "The provider rejected the steering message."),
|
||||
);
|
||||
|
||||
const initial = await request(app(seeded.companyId))
|
||||
.get(`/api/issues/${seeded.issueId}/queued-comments`);
|
||||
const steered = await request(app(seeded.companyId))
|
||||
.post(`/api/issues/${seeded.issueId}/queued-comments/${seeded.commentIds[0]}/steer`)
|
||||
.send({ queueId: seeded.wakeId, targetRunId: seeded.runId, revision: initial.body.revision });
|
||||
|
||||
expect(steered.status).toBe(409);
|
||||
expect(steered.body.details).toMatchObject({ code: "steering_rejected", retryable: true });
|
||||
|
||||
const rejected = await db
|
||||
.select({ status: runIdentityContexts.status })
|
||||
.from(runIdentityContexts)
|
||||
.where(eq(runIdentityContexts.messageId, seeded.commentIds[0]))
|
||||
.then((rows) => rows[0]);
|
||||
expect(rejected?.status).toBe("rejected");
|
||||
});
|
||||
|
||||
it("throws queued_comment_order_mismatch and changes no row for an invalid reorder set", async () => {
|
||||
const seeded = await seedQueue();
|
||||
const initial = await request(app(seeded.companyId))
|
||||
.get(`/api/issues/${seeded.issueId}/queued-comments`);
|
||||
|
||||
const invalid = await request(app(seeded.companyId))
|
||||
.put(`/api/issues/${seeded.issueId}/queued-comments/order`)
|
||||
.send({
|
||||
queueId: seeded.wakeId,
|
||||
revision: initial.body.revision,
|
||||
orderedCommentIds: [seeded.commentIds[0], randomUUID()],
|
||||
});
|
||||
|
||||
expect(invalid.status, JSON.stringify(invalid.body)).toBe(409);
|
||||
expect(invalid.body.details?.code).toBe("queued_comment_order_mismatch");
|
||||
|
||||
const wake = await db
|
||||
.select({ payload: agentWakeupRequests.payload })
|
||||
.from(agentWakeupRequests)
|
||||
.where(eq(agentWakeupRequests.id, seeded.wakeId))
|
||||
.then((rows) => rows[0]);
|
||||
expect((wake?.payload as any)?._paperclipWakeContext?.wakeCommentIds).toEqual(seeded.commentIds);
|
||||
const after = await request(app(seeded.companyId))
|
||||
.get(`/api/issues/${seeded.issueId}/queued-comments`);
|
||||
expect(after.body.entries.map((entry: any) => entry.comment.id)).toEqual(seeded.commentIds);
|
||||
});
|
||||
});
|
||||
Reference in new issue
Block a user