Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
77 changes: 62 additions & 15 deletions apps/server/src/provider/prime/PrimeAgentDaemonAdapter.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2831,14 +2831,17 @@ describe("PrimeAgentDaemonAdapter", () => {
);
}

for (const { withImage, hidden } of [
for (const { withImage, hidden, terminal } of [
{ withImage: false, hidden: "none" },
{ withImage: true, hidden: "none" },
{ withImage: false, hidden: "snapshot" },
{ withImage: false, hidden: "observed" },
] as const) {
].flatMap((variant) => [
{ ...variant, terminal: false },
{ ...variant, terminal: true },
])) {
it.effect(
`reconciles the delivered submitted user from the first complete resync (image: ${withImage}, hidden: ${hidden})`,
`reconciles the delivered submitted user from the first complete resync (image: ${withImage}, hidden: ${hidden}, terminal: ${terminal})`,
() =>
Effect.scoped(
Effect.gen(function* () {
Expand Down Expand Up @@ -2896,37 +2899,48 @@ describe("PrimeAgentDaemonAdapter", () => {
attribution: { scope: "prompt", correlationId },
});
}
const messages = hidden === "none" ? [prompt] : [digest, prompt];
const answer = assistantMessage(
"The directory contains README.md. It describes the fixture.",
);
const messages = [
...(hidden === "none" ? [] : [digest]),
prompt,
...(terminal ? [answer] : []),
];
const resync = {
...initialSnapshot(),
state: {
...initialSnapshot().state,
messageCount: messages.length,
isStreaming: true,
isStreaming: !terminal,
},
messages,
replayContinuity: "complete",
connectionGeneration: 0,
correlatedProofEpoch: 0,
promptLifecycles: { records: [delivered], expired: [] },
promptLifecycles: {
records: [
terminal
? lifecycleSnapshot(correlationId, "completed", 3, { usage })
: delivered,
],
expired: [],
},
} satisfies PrimeDaemonEvent;
yield* offer(captures, { ...resync, connectionGeneration: -1 });
yield* offer(captures, { ...resync, correlatedProofEpoch: -1 });
yield* offer(captures, resync);
expect(yield* Effect.promise(() => resolutionObserved)).toMatchObject({
reconciled: true,
});
expect(turnFiber.pollUnsafe()).toBeUndefined();
if (!terminal) expect(turnFiber.pollUnsafe()).toBeUndefined();
expect(captures.reconnectResolutions).toHaveLength(1);
yield* offer(captures, resync);
yield* offer(captures, {
_tag: "MessageCompleted",
message: prompt,
attribution: { scope: "prompt", correlationId },
});
const answer = assistantMessage(
"The directory contains README.md. It describes the fixture.",
);
yield* offer(captures, {
_tag: "MessageCompleted",
message: answer,
Expand Down Expand Up @@ -2957,7 +2971,7 @@ describe("PrimeAgentDaemonAdapter", () => {
);
}

for (const rejection of [
for (const { rejection, terminal } of [
"hidden unknown replay",
"hidden unavailable replay",
"hidden changed observed digest",
Expand Down Expand Up @@ -2986,8 +3000,17 @@ describe("PrimeAgentDaemonAdapter", () => {
"repeated user boundary",
"changed observed transcript",
"missing observed transcript",
] as const) {
it.effect(`rejects submitted-user resync with ${rejection}`, () =>
"unfinished tool answer",
"terminal still streaming",
].flatMap((rejection) =>
rejection === "unfinished tool answer" || rejection === "terminal still streaming"
? [{ rejection, terminal: true }]
: [
{ rejection, terminal: false },
{ rejection, terminal: true },
],
)) {
it.effect(`rejects submitted-user resync with ${rejection} (terminal: ${terminal})`, () =>
Effect.scoped(
Effect.gen(function* () {
const recovering =
Expand Down Expand Up @@ -3093,11 +3116,17 @@ describe("PrimeAgentDaemonAdapter", () => {
if (rejection === "changed observed transcript")
messages.unshift({ ...prompt, text: "changed prior user" });
if (rejection === "missing observed transcript") messages.length = 0;
if (terminal) {
messages.push({
...assistantMessage("completed answer"),
...(rejection === "unfinished tool answer" ? { stopReason: "toolUse" as const } : {}),
});
}
if (recovering) {
yield* offer(captures, { _tag: "ConnectionStatus", status: "reconnecting" });
}
const lifecycle = {
...delivered,
...(terminal ? lifecycleSnapshot(correlationId, "completed", 3, { usage }) : delivered),
...(rejection === "different owner" ? { correlationId: "another-owner" } : {}),
...(rejection === "undelivered snapshot"
? { phase: "owned" as const, deliveryCrossed: false }
Expand All @@ -3110,7 +3139,7 @@ describe("PrimeAgentDaemonAdapter", () => {
state: {
...initialSnapshot().state,
messageCount: messages.length + (rejection === "hidden count mismatch" ? 1 : 0),
isStreaming: true,
isStreaming: !terminal || rejection === "terminal still streaming",
},
messages,
replayContinuity:
Expand Down Expand Up @@ -10403,6 +10432,24 @@ describe("PrimeAgentDaemonAdapter", () => {
payload: { presentation: { kind: "status", key: "build", text: "Running" } },
});

for (const message of ["Starting Python kernel...", undefined, ""]) {
yield* offer(captures, {
_tag: "ExtensionRequest",
request: { id: "native-working-message", method: "setWorkingMessage", message },
});
const working = yield* awaitObservedType(
subscription.observed,
"session-presentation.updated",
);
expect(working.payload).toEqual({
presentation: {
kind: "status",
key: "prime-working-message",
...(message ? { text: message } : {}),
},
});
}

yield* offer(captures, {
_tag: "ExtensionRequest",
request: {
Expand Down
30 changes: 27 additions & 3 deletions apps/server/src/provider/prime/PrimeAgentDaemonAdapter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -376,6 +376,17 @@ function projectExtensionRequest(
}),
(presentation) => ({ _tag: "Presentation", presentation }),
);
case "setWorkingMessage":
return Option.map(
decodeSessionPresentation({
kind: "status",
key: "prime-working-message",
...(request.message === undefined || request.message.trim().length === 0
? {}
: { text: request.message }),
}),
(presentation) => ({ _tag: "Presentation", presentation }),
);
case "setWidget":
return Option.map(
decodeSessionPresentation({
Expand Down Expand Up @@ -2770,6 +2781,18 @@ export function makePrimeAgentDaemonAdapter(
missingMessages[0]?.role === "harnessDigest"
? 1
: 0;
const submittedAnswer = missingMessages[submittedUserIndex + 1];
// A fast prompt may finish before either message event arrives.
// Only its exact submitted boundary and one terminal answer can
// be recovered under the same completed, delivered lifecycle.
const snapshotHasSubmittedAnswer =
missingMessages.length === submittedUserIndex + 2 &&
submittedAnswer?.role === "assistant" &&
submittedAnswer.stopReason === "stop" &&
submittedAnswer.toolCalls.length === 0 &&
lifecycle?.phase === "completed" &&
snapshotEvent.state.isStreaming === false &&
snapshotEvent.streamingMessage === undefined;
const snapshotRecoversSubmittedUser =
activeTurn !== undefined &&
snapshotEvent.connectionGeneration !== undefined &&
Expand All @@ -2778,7 +2801,8 @@ export function makePrimeAgentDaemonAdapter(
(context.nativeTranscriptMessageCount ===
activeTurn.nativeTranscriptBaselineMessageCount ||
firstHarnessDigestObserved) &&
missingMessages.length === submittedUserIndex + 1 &&
(missingMessages.length === submittedUserIndex + 1 ||
snapshotHasSubmittedAnswer) &&
missingMessages[submittedUserIndex] !== undefined &&
matchesSubmittedUserMessage(
activeTurn,
Expand All @@ -2787,8 +2811,8 @@ export function makePrimeAgentDaemonAdapter(
currentLifecycle?.kind === "model_prompt" &&
currentLifecycle.phase === "delivered" &&
currentLifecycle.deliveryCrossed &&
lifecycle?.phase === "delivered" &&
lifecycle.deliveryCrossed &&
(lifecycle?.phase === "delivered" || snapshotHasSubmittedAnswer) &&
lifecycle?.deliveryCrossed === true &&
(primeAgentPromptLifecycleIsSame(currentLifecycle, lifecycle) ||
primeAgentPromptLifecycleCanAdvance(currentLifecycle, lifecycle));
// Completed tool cycles can reach the durable snapshot before their
Expand Down
1 change: 1 addition & 0 deletions apps/server/src/provider/prime/PrimeAgentDaemonEvents.ts
Original file line number Diff line number Diff line change
Expand Up @@ -607,6 +607,7 @@ const extensionMethod = Schema.Literals([
"editor",
"notify",
"setStatus",
"setWorkingMessage",
"setWidget",
]);

Expand Down
Loading
Loading