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
131 changes: 128 additions & 3 deletions apps/server/src/provider/prime/PrimeAgentDaemonAdapter.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -172,6 +172,9 @@ interface FakeCaptures {
followUpFailure: boolean;
inputRecoveryPending: boolean;
inputAdmissionBusy: boolean;
correlatedPromptLifecycleAdmissionBlocked: boolean;
correlatedPromptLifecycleAdmissionBlockAfterReads: number | undefined;
correlatedPromptLifecycleAdmissionReads: number;
correlatedPromptLifecycleAvailable: boolean;
readonly correlatedPromptSubmissions: Array<{
readonly text: string;
Expand Down Expand Up @@ -327,6 +330,9 @@ interface FakeCaptures {
rlmQuiescenceFailure: boolean;
rlmConnectionGeneration: number;
rlmContinuityValid: boolean;
correlatedRecoveryProofCurrent: boolean;
correlatedRecoveryProofEpoch: number;
reconnectSnapshotResolutionAccepted: boolean;
readonly reconnectResolutions: Array<{
readonly generation: number;
readonly reconciled: boolean;
Expand Down Expand Up @@ -356,6 +362,9 @@ function makeCaptures(): FakeCaptures {
followUpFailure: false,
inputRecoveryPending: false,
inputAdmissionBusy: false,
correlatedPromptLifecycleAdmissionBlocked: false,
correlatedPromptLifecycleAdmissionBlockAfterReads: undefined,
correlatedPromptLifecycleAdmissionReads: 0,
correlatedPromptLifecycleAvailable: false,
correlatedPromptSubmissions: [],
correlatedPromptCancellations: [],
Expand Down Expand Up @@ -477,6 +486,9 @@ function makeCaptures(): FakeCaptures {
rlmQuiescenceFailure: false,
rlmConnectionGeneration: 0,
rlmContinuityValid: true,
correlatedRecoveryProofCurrent: true,
correlatedRecoveryProofEpoch: 1,
reconnectSnapshotResolutionAccepted: true,
reconnectResolutions: [],
retryWorkerRecoverySnapshots: false,
retryWorkerRecoverySnapshotCalls: [],
Expand Down Expand Up @@ -836,7 +848,14 @@ function fakeRuntimeFactory(
const resolution = { generation, reconciled, terminalResponseObserved };
captures.reconnectResolutions.push(resolution);
captures.reconnectSnapshotResolutionObserved?.(resolution);
if (generation !== captures.rlmConnectionGeneration) return false;
if (
generation !== captures.rlmConnectionGeneration ||
!captures.reconnectSnapshotResolutionAccepted ||
(captures.correlatedPromptLifecycleAvailable &&
!captures.correlatedRecoveryProofCurrent)
) {
return false;
}
captures.rlmContinuityValid = reconciled;
return true;
},
Expand All @@ -848,8 +867,12 @@ function fakeRuntimeFactory(
noteWorkerRecoveryTerminalResponse: () => {
captures.workerRecoveryTerminalResponseObserved?.();
},
isConnectionGenerationCurrent: (generation) =>
generation === captures.rlmConnectionGeneration,
isConnectionGenerationCurrent: (generation, proofEpoch) =>
generation === captures.rlmConnectionGeneration &&
(!captures.correlatedPromptLifecycleAvailable ||
(captures.correlatedRecoveryProofCurrent &&
(proofEpoch ?? captures.correlatedRecoveryProofEpoch) ===
captures.correlatedRecoveryProofEpoch)),
get correlatedPromptLifecycleAvailable() {
return captures.correlatedPromptLifecycleAvailable;
},
Expand Down Expand Up @@ -897,6 +920,15 @@ function fakeRuntimeFactory(
}
);
}),
get correlatedPromptLifecycleAdmissionBlocked() {
captures.correlatedPromptLifecycleAdmissionReads += 1;
return (
captures.correlatedPromptLifecycleAdmissionBlocked ||
(captures.correlatedPromptLifecycleAdmissionBlockAfterReads !== undefined &&
captures.correlatedPromptLifecycleAdmissionReads >=
captures.correlatedPromptLifecycleAdmissionBlockAfterReads)
);
},
get inputAdmissionBusy() {
return captures.inputAdmissionBusy;
},
Expand Down Expand Up @@ -1121,6 +1153,32 @@ function offer(captures: FakeCaptures, event: PrimeDaemonEvent) {
}

describe("PrimeAgentDaemonAdapter", () => {
it.effect("rechecks correlated recovery before committing a new strict turn", () =>
Effect.scoped(
Effect.gen(function* () {
const captures = makeCaptures();
captures.correlatedPromptLifecycleAvailable = true;
captures.correlatedPromptLifecycleAdmissionBlockAfterReads = 2;
const adapter = yield* makePrimeAgentDaemonAdapter(decodeSettings({}), manager, {
instanceId,
runtimeFactory: fakeRuntimeFactory(captures),
});
yield* adapter.startSession({ threadId, cwd: process.cwd(), runtimeMode: "full-access" });

const error = yield* adapter
.sendTurn({ threadId, input: "must not cross pending resync" })
.pipe(Effect.flip);

expect(error).toMatchObject({
_tag: "ProviderAdapterValidationError",
reason: "busy",
});
expect(captures.correlatedPromptLifecycleAdmissionReads).toBeGreaterThanOrEqual(2);
expect(captures.correlatedPromptSubmissions).toEqual([]);
}),
).pipe(Effect.provide(testLayer)),
);

it.effect("settles only the delivered correlated owner and uses terminal lifecycle usage", () =>
Effect.scoped(
Effect.gen(function* () {
Expand Down Expand Up @@ -1376,6 +1434,8 @@ describe("PrimeAgentDaemonAdapter", () => {
...initialSnapshot(),
lastEventSequence: 2,
replayContinuity: "complete",
connectionGeneration: 0,
correlatedProofEpoch: 1,
promptLifecycles: {
records: [lifecycleSnapshot(correlationId, "failed", 2, { usage })],
expired: [],
Expand Down Expand Up @@ -1428,6 +1488,8 @@ describe("PrimeAgentDaemonAdapter", () => {
...initialSnapshot(),
lastEventSequence: 2,
replayContinuity: "complete",
connectionGeneration: 0,
correlatedProofEpoch: 1,
children: [{ id: "background-child", label: "background", status: "running" }],
promptLifecycles: {
records: [lifecycleSnapshot(correlationId, "queued", 2)],
Expand Down Expand Up @@ -1931,6 +1993,69 @@ describe("PrimeAgentDaemonAdapter", () => {
).pipe(Effect.provide(testLayer)),
);

it.effect("does not apply a terminal correlated snapshot when proof settlement is rejected", () =>
Effect.scoped(
Effect.gen(function* () {
const captures = makeCaptures();
captures.correlatedPromptLifecycleAvailable = true;
captures.correlatedPromptObserved = yield* Queue.unbounded<string>();
captures.rlmConnectionGeneration = 1;
captures.rlmContinuityValid = false;
captures.reconnectSnapshotResolutionAccepted = false;
const adapter = yield* makePrimeAgentDaemonAdapter(decodeSettings({}), manager, {
instanceId,
runtimeFactory: fakeRuntimeFactory(captures),
});
const subscription = yield* subscribe(adapter);
yield* adapter.startSession({ threadId, cwd: process.cwd(), runtimeMode: "full-access" });
const turnFiber = yield* adapter
.sendTurn({ threadId, input: "retired proof must not settle" })
.pipe(Effect.forkChild);
const correlationId = yield* Queue.take(captures.correlatedPromptObserved);
yield* offer(captures, {
_tag: "PromptLifecycleUpdated",
lifecycle: lifecycleSnapshot(correlationId, "delivered", 2),
});
const answer = assistantMessage("stale terminal answer");
yield* offer(captures, {
_tag: "MessageCompleted",
message: answer,
attribution: { scope: "prompt", correlationId },
});

yield* offer(captures, {
...initialSnapshot(),
state: { ...initialSnapshot().state, messageCount: 1 },
messages: [answer],
replayContinuity: "complete",
connectionGeneration: 1,
correlatedProofEpoch: 1,
promptLifecycles: {
records: [lifecycleSnapshot(correlationId, "completed", 3, { usage })],
expired: [],
},
});
yield* offer(captures, {
_tag: "SessionClosed",
error: "Prime Agent correlated prompt capability proof was lost during recovery.",
});

const result = yield* Fiber.join(turnFiber);
expect(captures.reconnectResolutions).toContainEqual({
generation: 1,
reconciled: true,
terminalResponseObserved: false,
});
const terminal = subscription.events.findLast(
(event) => event.turnId === result.turnId && event.type === "turn.completed",
);
expect(terminal).toMatchObject({ payload: { state: "failed" } });
expect(terminal?.payload).not.toHaveProperty("usage");
expect(terminal?.payload).not.toHaveProperty("totalCostUsd");
}),
).pipe(Effect.provide(testLayer)),
);

it.effect("keeps control mismatch busy and scopes pre-delivery cancellation", () =>
Effect.scoped(
Effect.gen(function* () {
Expand Down
83 changes: 50 additions & 33 deletions apps/server/src/provider/prime/PrimeAgentDaemonAdapter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2004,6 +2004,16 @@ export function makePrimeAgentDaemonAdapter(
context.threadId,
Effect.gen(function* () {
if (sessions.get(context.threadId) === context && !context.stopped) {
const reconnectGeneration = event.connectionGeneration;
if (
reconnectGeneration !== undefined &&
!context.runtime.isConnectionGenerationCurrent(
reconnectGeneration,
event.correlatedProofEpoch,
)
) {
return;
}
if (!managedSourceVerified) {
const activeTurn = context.activeTurn;
if (activeTurn !== undefined) {
Expand All @@ -2019,13 +2029,6 @@ export function makePrimeAgentDaemonAdapter(
reconnectRecoveryFailed = true;
return;
}
const reconnectGeneration = event.connectionGeneration;
if (
reconnectGeneration !== undefined &&
!context.runtime.isConnectionGenerationCurrent(reconnectGeneration)
) {
return;
}
context.managedPlanProjectionEnabled = true;
const activeTurn = context.activeTurn;
if (context.runtime.correlatedPromptLifecycleAvailable) {
Expand All @@ -2047,36 +2050,39 @@ export function makePrimeAgentDaemonAdapter(
reconnectRecoveryFailed = true;
return;
}
if (activeTurn?.correlationId !== undefined) {
const lifecycle = event.promptLifecycles?.records.find(
(candidate) => candidate.correlationId === activeTurn.correlationId,
);
if (lifecycle === undefined) {
if (reconnectGeneration !== undefined) {
context.runtime.resolveReconnectSnapshot(
reconnectGeneration,
false,
false,
const lifecycle =
activeTurn?.correlationId === undefined
? undefined
: event.promptLifecycles?.records.find(
(candidate) => candidate.correlationId === activeTurn.correlationId,
);
}
yield* settleActiveTurnLocked(context, activeTurn, {
state: "failed",
errorMessage:
"Prime Agent could not recover the correlated prompt lifecycle after synchronizing.",
runtimeErrorMessage:
"Prime Agent could not recover the correlated prompt lifecycle after synchronizing.",
});
context.stopRequested = true;
reconnectRecoveryFailed = true;
return;
if (activeTurn?.correlationId !== undefined && lifecycle === undefined) {
if (reconnectGeneration !== undefined) {
context.runtime.resolveReconnectSnapshot(reconnectGeneration, false, false);
}
yield* settleActiveTurnLocked(context, activeTurn, {
state: "failed",
errorMessage:
"Prime Agent could not recover the correlated prompt lifecycle after synchronizing.",
runtimeErrorMessage:
"Prime Agent could not recover the correlated prompt lifecycle after synchronizing.",
});
context.stopRequested = true;
reconnectRecoveryFailed = true;
return;
}
if (
reconnectGeneration === undefined ||
!context.runtime.resolveReconnectSnapshot(reconnectGeneration, true, false)
) {
reconnectRecoveryFailed = true;
return;
}
if (lifecycle !== undefined) {
yield* applyCorrelatedPromptLifecycleLocked(context, lifecycle, {
authoritativeSnapshot: true,
});
}
if (reconnectGeneration !== undefined) {
context.runtime.resolveReconnectSnapshot(reconnectGeneration, true, false);
}
}
} else if (reconnectGeneration !== undefined) {
const pendingRunCompletionBefore = activeTurn?.pendingRunCompletionHandoff;
Expand Down Expand Up @@ -3876,13 +3882,16 @@ export function makePrimeAgentDaemonAdapter(
const knownCompactionBusy =
context.activeCompactionScope !== undefined ||
context.manualCompactionRequestActive;
const correlatedRecoveryBusy =
context.runtime.correlatedPromptLifecycleAdmissionBlocked;
const nativeInputBusy = context.runtime.inputAdmissionBusy;
const admissionBusy =
nativeInputBusy ||
(context.runtime.correlatedPromptLifecycleAvailable && knownCompactionBusy);
if (
admissionBusy &&
(!context.runtime.correlatedPromptLifecycleAvailable || !controlsMatchCurrent)
correlatedRecoveryBusy ||
(admissionBusy &&
(!context.runtime.correlatedPromptLifecycleAvailable || !controlsMatchCurrent))
) {
return yield* new ProviderAdapterValidationError({
provider: PROVIDER,
Expand Down Expand Up @@ -3961,6 +3970,14 @@ export function makePrimeAgentDaemonAdapter(
projectedPlanToolCallIds: new Set(),
...(/^\/compact(?:\s|$)/.test(text) ? { command: "compact" as const } : {}),
};
if (context.runtime.correlatedPromptLifecycleAdmissionBlocked) {
return yield* new ProviderAdapterValidationError({
provider: PROVIDER,
operation: "sendTurn",
reason: "busy",
issue: "Prime Agent is still reconciling its correlated prompt lifecycle.",
});
}
context.activeTurn = turn;
context.session = {
...context.session,
Expand Down
Loading
Loading