fix(server): probe the live steering state after a successful queue steer
A successful same-turn steer that leaves more messages queued answered "temporarily_unavailable" for its own response, since the queue snapshot builder never asked the live provider. The UI stored that answer and disabled the next steering action until an unrelated queue read refreshed the state. Add an opt-in probe to the snapshot builder. A queue mutation that does not itself deliver same-turn steering keeps the static answer. The steer mutation asks the live provider for the snapshot it returns to its own caller, on both a fresh steer and a replayed duplicate, matching the disposition a fresh read would report. Co-authored-by: Paperclip <noreply@paperclip.ing>
This commit is contained in:
parent
d7f53896e1
commit
b324f7592a
|
|
@ -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<typeof import("../../../services/native-runtime/native-session-executor.js")>();
|
||||
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");
|
||||
});
|
||||
});
|
||||
|
|
|
|||
|
|
@ -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<IssueQueuedCommentQueue> {
|
||||
async buildQueueSnapshot({ issue, actor, wake, state, queueRun, activeRun, probeLiveSteering }): Promise<IssueQueuedCommentQueue> {
|
||||
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 };
|
||||
|
|
|
|||
|
|
@ -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<QueuedCommentQueueSnapshot>;
|
||||
syncCommentReferences(commentId: string): Promise<void>;
|
||||
deleteCommentReferenceSource(commentId: string): Promise<void>;
|
||||
|
|
|
|||
Loading…
Reference in New Issue