diff --git a/server/src/modules/wake-queue/adapters/queued-comment-postgres.test.ts b/server/src/modules/wake-queue/adapters/queued-comment-postgres.test.ts index d8f434adfc..7f3fa19214 100644 --- a/server/src/modules/wake-queue/adapters/queued-comment-postgres.test.ts +++ b/server/src/modules/wake-queue/adapters/queued-comment-postgres.test.ts @@ -14,12 +14,21 @@ import { QueuedCommentMutationError } from "../application/queued-comment-use-ca // The steering mutation delivers through the live native-runtime transport. // This mock stands in for that transport, so the test proves the module's -// own read/write behavior without a live provider connection. +// own read/write behavior without a live provider connection. The steering +// mutation also asks the same transport for the live steering state right +// after a successful steer, so it can answer with the authoritative +// disposition instead of the static "temporarily_unavailable" fallback. const steerNativeSessionMock = vi.hoisted(() => vi.fn()); +const getNativeSessionSteeringStateMock = vi.hoisted(() => vi.fn()); vi.mock("../../../services/native-runtime/native-session-executor.js", async (importOriginal) => { const actual = await importOriginal(); steerNativeSessionMock.mockImplementation(actual.steerNativeSession); - return { ...actual, steerNativeSession: steerNativeSessionMock }; + getNativeSessionSteeringStateMock.mockImplementation(actual.getNativeSessionSteeringState); + return { + ...actual, + steerNativeSession: steerNativeSessionMock, + getNativeSessionSteeringState: getNativeSessionSteeringStateMock, + }; }); // Proves the same two properties the release-half adapter test proves for @@ -400,6 +409,40 @@ describeEmbeddedPostgres("queued-comment postgres adapter", () => { const companyId = await seedCompany(); const agentId = await seedAgent({ companyId }); const issueId = await seedIssue({ companyId, assigneeAgentId: agentId }); + const commentId = await seedComment({ companyId, issueId, authorUserId: "user-1" }); + const wakeId = await seedDeferredWake({ companyId, agentId, issueId, commentIds: [commentId] }); + + const issueLock = createQueuedCommentIssueLockWriter(db, noopDeps); + const queue = await issueLock.withLockedQueue( + { + issue: { id: issueId, companyId, assigneeAgentId: agentId, executionRunId: null }, + actor: { actorType: "user", actorId: "user-1", agentId: null }, + queueId: wakeId, + }, + async (locked, transaction) => { + // An edit never delivers same-turn steering itself, so it must never + // probe the live provider for the answer, unlike the steer mutation + // below. + return transaction.buildQueueSnapshot({ + issue: { id: issueId, companyId, assigneeAgentId: agentId, executionRunId: null }, + actor: { actorType: "user", actorId: "user-1", agentId: null }, + wake: locked.wake, + state: locked.state, + queueRun: locked.queueRun, + activeRun: { id: randomUUID(), status: "running", runtimeMode: "native", contextSnapshot: {} }, + }); + }, + ); + + expect(queue.protocol).toBe("paperclip_runner_v1"); + expect(queue.steeringDisposition).toBe("temporarily_unavailable"); + expect(getNativeSessionSteeringStateMock).not.toHaveBeenCalled(); + }); + + it("reports the live steering disposition after a successful steer leaves more messages queued", async () => { + const companyId = await seedCompany(); + const agentId = await seedAgent({ companyId, adapterType: "paperclip_runner" }); + const issueId = await seedIssue({ companyId, assigneeAgentId: agentId }); const firstCommentId = await seedComment({ companyId, issueId, authorUserId: "user-1" }); const secondCommentId = await seedComment({ companyId, issueId, authorUserId: "user-1" }); const wakeId = await seedDeferredWake({ @@ -411,6 +454,9 @@ describeEmbeddedPostgres("queued-comment postgres adapter", () => { const targetRunId = await seedRunningNativeRun({ companyId, agentId, issueId }); steerNativeSessionMock.mockResolvedValueOnce({ turnId: "turn-1" }); + // The live provider still has an active turn to steer, since the second + // queued message has not been delivered yet. + getNativeSessionSteeringStateMock.mockResolvedValueOnce({ disposition: "available", activeTurnId: "turn-1" }); const issueLock = createQueuedCommentIssueLockWriter(db, noopDeps); const peeked = await issueLock.withLockedQueue( @@ -433,9 +479,30 @@ describeEmbeddedPostgres("queued-comment postgres adapter", () => { expect(result.queue.entries).toHaveLength(1); expect(result.queue.entries[0]!.comment.id).toBe(secondCommentId); - // A queue mutation never probes the live provider for the steering - // answer. Only the read path probes the provider, and it corrects - // this value the next time a caller reads the queue. - expect(result.queue.steeringDisposition).toBe("temporarily_unavailable"); + // The steer just delivered a message, so it must ask the live provider + // for the real answer instead of falling back to + // "temporarily_unavailable" and leaving the next steering action + // disabled until an unrelated read refreshes the state. + expect(result.queue.steeringDisposition).toBe("available"); + expect(getNativeSessionSteeringStateMock).toHaveBeenCalledWith(targetRunId); + + // A second, consecutive steer on the same run must keep reporting the + // live disposition, not the stale value from the first steer. + getNativeSessionSteeringStateMock.mockResolvedValueOnce({ disposition: "temporarily_unavailable", activeTurnId: null }); + steerNativeSessionMock.mockResolvedValueOnce({ turnId: "turn-2" }); + const second = await issueLock.steerQueuedWakeComment({ + issue: { id: issueId, companyId, assigneeAgentId: agentId, executionRunId: null }, + actor: { actorType: "user", actorId: "user-1", agentId: null }, + commentId: secondCommentId, + queueId: wakeId, + targetRunId, + revision: result.queue.revision, + }); + + // The queue is now empty, so the wake is cancelled and no run remains to + // probe -- the shared rule answers "temporarily_unavailable" directly, + // without a live call. + expect(second.queue.entries).toHaveLength(0); + expect(second.queue.steeringDisposition).toBe("temporarily_unavailable"); }); }); diff --git a/server/src/modules/wake-queue/adapters/queued-comment-postgres.ts b/server/src/modules/wake-queue/adapters/queued-comment-postgres.ts index fa53a21b7f..fb9c38c19c 100644 --- a/server/src/modules/wake-queue/adapters/queued-comment-postgres.ts +++ b/server/src/modules/wake-queue/adapters/queued-comment-postgres.ts @@ -11,6 +11,7 @@ import { } from "../../../services/issue-queued-comment-queue.js"; import { logActivity as persistActivityLogRow, type ActivityPublication } from "../../../services/activity-log.js"; import { + getNativeSessionSteeringState, NativeSessionSteeringError, steerNativeSession, } from "../../../services/native-runtime/native-session-executor.js"; @@ -125,7 +126,7 @@ function buildTransaction(tx: Db, companyId: string, deps: QueuedCommentQueuePos .where(and(eq(issues.id, issueId), eq(issues.companyId, companyId), eq(issues.executionRunId, executionRunId))); }, - async buildQueueSnapshot({ issue, actor, wake, state, queueRun, activeRun }): Promise { + async buildQueueSnapshot({ issue, actor, wake, state, queueRun, activeRun, probeLiveSteering }): Promise { const commentIds = queuedCommentIdsFromWakePayload(wake?.payload ?? null); const rows = commentIds.length > 0 @@ -155,10 +156,13 @@ function buildTransaction(tx: Db, companyId: string, deps: QueuedCommentQueuePos .then((agentRows) => agentRows[0] ?? null) : null; - // A queue mutation never delivers same-turn steering itself, so this - // adapter never probes the live runner: it answers + // A queue mutation never delivers same-turn steering itself, so by + // default this adapter never probes the live runner: it answers // "temporarily_unavailable" wherever the shared rule says a caller - // may probe. Only the read path probes the live provider. + // may probe. The steer mutation is the one exception: right after it + // delivers a message, it already knows the live provider is reachable, + // so it asks `probeLiveSteering: true` for the snapshot it returns to + // its own caller, matching what a fresh read would report. const steering = decideQueuedCommentQueueSteering({ state, queueRunRuntimeMode: queueRun?.runtimeMode ?? null, @@ -167,7 +171,13 @@ function buildTransaction(tx: Db, companyId: string, deps: QueuedCommentQueuePos queuedCommentCount: comments.length, }); const steeringDisposition: IssueQueuedCommentQueue["steeringDisposition"] = - steering.kind === "probe" ? "temporarily_unavailable" : steering.kind; + steering.kind !== "probe" + ? steering.kind + : probeLiveSteering + ? await getNativeSessionSteeringState(steering.steeringRunId) + .then((liveState) => liveState.disposition) + .catch(() => "temporarily_unavailable" as const) + : "temporarily_unavailable"; return buildQueuedCommentQueueSnapshot({ issueId: issue.id, @@ -461,6 +471,7 @@ export function createQueuedCommentIssueLockWriter(db: Db, deps: QueuedCommentQu state: current?.state ?? null, queueRun: current?.queueRun ? toRunRow(current.queueRun) : null, activeRun: retryRunRow.status === "running" ? toRunRow(retryRunRow) : null, + probeLiveSteering: true, }); } @@ -547,7 +558,15 @@ export function createQueuedCommentIssueLockWriter(db: Db, deps: QueuedCommentQu if (priorAcknowledgement.status === "acknowledged" && priorAcknowledgement.queueId === queueId) { duplicate = true; turnId = typeof priorAcknowledgement.turnId === "string" ? priorAcknowledgement.turnId : null; - return lockedQueue; + return transaction.buildQueueSnapshot({ + issue, + actor, + wake: toWakeRow(wakeRow), + state, + queueRun: queueRunRow ? toRunRow(queueRunRow) : null, + activeRun: toRunRow(activeRunRow), + probeLiveSteering: true, + }); } requireMutationTarget(lockedQueue, queueId, revision); @@ -614,6 +633,7 @@ export function createQueuedCommentIssueLockWriter(db: Db, deps: QueuedCommentQu state: nextWakeRow ? "deferred" : null, queueRun: null, activeRun: toRunRow(activeRunRow), + probeLiveSteering: true, }); }); return { queue, turnId, duplicate }; diff --git a/server/src/modules/wake-queue/application/queued-comment-ports.ts b/server/src/modules/wake-queue/application/queued-comment-ports.ts index dbf3e34f65..3b046f5074 100644 --- a/server/src/modules/wake-queue/application/queued-comment-ports.ts +++ b/server/src/modules/wake-queue/application/queued-comment-ports.ts @@ -144,6 +144,15 @@ export interface QueuedCommentQueueTransaction { state: "deferred" | "queued" | null; queueRun: QueuedCommentRunRow | null; activeRun: QueuedCommentRunRow | null; + /** + * A queue mutation does not itself deliver same-turn steering, so it must + * leave `false` (the default) and report the static + * "temporarily_unavailable" answer instead of a live probe. Only the + * steer mutation sets this `true`, and only for a snapshot it returns to + * the caller, because a steer that just ran knows the live disposition + * the read path would otherwise have to probe for. + */ + probeLiveSteering?: boolean; }): Promise; syncCommentReferences(commentId: string): Promise; deleteCommentReferenceSource(commentId: string): Promise;