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
28 changes: 28 additions & 0 deletions apps/server/src/provider/prime/PrimeAgentDaemonBridge.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -143,6 +143,34 @@ afterEach(() => {
});

describe("PrimeAgentDaemonBridge", () => {
it.effect("loads settlement observation only with its frozen SDK contract", () =>
Effect.gen(function* () {
for (const registry of [undefined, "mutable", "frozen"] as const) {
const pkg = makePackage({
moduleSource:
daemonModuleSource({
...(registry === undefined ? {} : { sdkFeatureRegistry: registry }),
extraSdkFeatures: ["owned_session_settlement_observation_v1"],
}) +
'\nexport async function observeOwnedSessionSettlement() { return { feature: "owned_session_settlement_observation_v1", status: "settled" }; }',
});
const bridge = yield* loadPrimeAgentDaemonBridge(pkg.cliPath);
expect(typeof bridge.observeOwnedSessionSettlement).toBe(
registry === "frozen" ? "function" : "undefined",
);
}
const missing = makePackage({
moduleSource: daemonModuleSource({
sdkFeatureRegistry: "frozen",
extraSdkFeatures: ["owned_session_settlement_observation_v1"],
}),
});
expect(
(yield* loadPrimeAgentDaemonBridge(missing.cliPath)).observeOwnedSessionSettlement,
).toBeUndefined();
}),
);

it.effect("reports a missing configured executable", () =>
Effect.gen(function* () {
const error = yield* Effect.flip(
Expand Down
14 changes: 14 additions & 0 deletions apps/server/src/provider/prime/PrimeAgentDaemonBridge.ts
Original file line number Diff line number Diff line change
Expand Up @@ -368,6 +368,11 @@ export interface PrimeAgentDaemonBridge extends PrimeAgentPublicPackage {
/** Exact frozen client capability evidence used with the post-connect daemon gates. */
readonly sdkFeatures?: ReadonlyArray<string>;
readonly recoverableOwnedSessionAdoptionAvailable?: boolean;
readonly observeOwnedSessionSettlement?: (input: {
readonly agentDir: string;
readonly activeSessionId: string;
readonly contractProof: PrimeAgentOwnedSessionContractProof;
}) => Promise<unknown>;
readonly createRecoverableOwnedSession?: (
client: PrimeAgentDaemonClient,
options: PrimeAgentRecoverableOwnedSessionCreateOptions,
Expand Down Expand Up @@ -742,6 +747,15 @@ function requireDaemonExports(input: {
negotiatedDaemonSessionCapabilitiesAvailable,
sdkFeatures,
recoverableOwnedSessionAdoptionAvailable,
...(sdkFeatures.includes("owned_session_settlement_observation_v1") &&
Predicate.isFunction(input.loadedModule.observeOwnedSessionSettlement)
? {
observeOwnedSessionSettlement: input.loadedModule
.observeOwnedSessionSettlement as NonNullable<
PrimeAgentDaemonBridge["observeOwnedSessionSettlement"]
>,
}
: {}),
...(recoverableOwnedSessionAdoptionAvailable
? {
createRecoverableOwnedSession: createRecoverableOwnedSession as NonNullable<
Expand Down
139 changes: 139 additions & 0 deletions apps/server/src/provider/prime/PrimeAgentLegacySettlement.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,139 @@
// @effect-diagnostics nodeBuiltinImport:off
import * as NodeFSP from "node:fs/promises";
import * as NodeOS from "node:os";
import * as NodePath from "node:path";
import { afterEach, describe, expect, it, vi } from "vite-plus/test";
import { recoverPrimeAgentLegacySettlement } from "./PrimeAgentLegacySettlement.ts";
import { PrimeAgentOwnershipReceiptStore } from "./PrimeAgentOwnershipReceipt.ts";
import type { PrimeAgentDaemonBridge } from "./PrimeAgentDaemonBridge.ts";

const roots: string[] = [];
afterEach(async () => {
await Promise.all(
roots.splice(0).map((root) => NodeFSP.rm(root, { recursive: true, force: true })),
);
});
async function fixture() {
const root = await NodeFSP.realpath(
await NodeFSP.mkdtemp(NodePath.join(NodeOS.tmpdir(), "prime-legacy-")),
);
roots.push(root);
const store = new PrimeAgentOwnershipReceiptStore(root, {
inspectProcessIdentity: async (pid) => `test:${pid}`,
});
const handle = await store.begin({
instanceId: "primeAgent",
configRevision: "revision",
effectiveHome: root,
});
await store.markAcquired(handle, {
activeSessionId: "previous-active",
nativeSessionId: "previous-native",
attachProof: {
feature: "caller_owned_session_environment_cleanup_v1",
status: "attached",
daemon: {
protocolName: "prime-agent.daemon",
protocolVersion: 7,
schemaRevision: 31,
supervisorGeneration: "previous-daemon",
transportGeneration: 1,
},
},
});
const path = NodePath.join(store.directory, `${handle.attemptId}.json`);
const scanned = (await store.scan()).receipts[0];
if (scanned?.state !== "acquired") throw new Error("Missing acquired fixture");
const saved = { ...scanned, ownerProcessId: "previous-process" };
await NodeFSP.writeFile(path, JSON.stringify(saved));
const observe = vi.fn(async () => ({
feature: "owned_session_settlement_observation_v1",
status: "settled",
}));
// Recovery only consumes the frozen feature list and the observer, not an adapter/session.
const bridge = {
sdkFeatures: Object.freeze(["owned_session_settlement_observation_v1"]),
observeOwnedSessionSettlement: observe,
} satisfies Pick<PrimeAgentDaemonBridge, "sdkFeatures" | "observeOwnedSessionSettlement">;
const loadBridge = vi.fn(async () => bridge);
const run = () =>
recoverPrimeAgentLegacySettlement({ instanceId: "primeAgent", store, loadBridge });
return { root, store, path, saved, observe, bridge, loadBridge, run };
}
describe("legacy Prime settlement recovery", () => {
it("clears only the exact legacy receipt after public SDK settlement", async () => {
const f = await fixture();
expect(await f.run()).toBe(true);
expect(f.observe).toHaveBeenCalledWith({
agentDir: f.root,
activeSessionId: "previous-active",
contractProof: f.saved.attachProof,
});
expect((await f.store.scan()).receipts).toEqual([]);
});
it.each(["registered", "unavailable", "completed", "unknown"])(
"retains the receipt for %s",
async (status) => {
const f = await fixture();
f.observe.mockResolvedValue({ feature: "owned_session_settlement_observation_v1", status });
await expect(f.run()).rejects.toThrow("settlement could not be proved");
expect((await f.store.scan()).receipts).toHaveLength(1);
},
);
it("does not treat another cleanup feature as settlement", async () => {
const f = await fixture();
f.observe.mockResolvedValue({ feature: "other", status: "settled" });
await expect(f.run()).rejects.toThrow("settlement could not be proved");
});
it("keeps records when the observer throws", async () => {
const f = await fixture();
f.observe.mockRejectedValue(new Error("offline"));
await expect(f.run()).rejects.toThrow("offline");
expect((await f.store.scan()).receipts).toHaveLength(1);
});
it("never substitutes settlement observation for recoverable adoption", async () => {
const f = await fixture();
f.saved.recovery = {
threadId: "thread",
sessionIncarnationId: "incarnation",
admissionRequestId: "request",
recoveryHandle: "private",
ownershipGeneration: 1,
};
await NodeFSP.writeFile(f.path, JSON.stringify(f.saved));
await expect(f.run()).rejects.toThrow("recovery authority");
expect(f.observe).not.toHaveBeenCalled();
});
it("requires the SDK feature, not method presence", async () => {
const f = await fixture();
const loadBridge = async () => ({ ...f.bridge, sdkFeatures: [] });
await expect(
recoverPrimeAgentLegacySettlement({ instanceId: "primeAgent", store: f.store, loadBridge }),
).rejects.toThrow("settlement observation support");
expect(f.observe).not.toHaveBeenCalled();
});
it("cannot cross-clear another instance", async () => {
const f = await fixture();
expect(
await recoverPrimeAgentLegacySettlement({
instanceId: "other",
store: f.store,
loadBridge: f.loadBridge,
}),
).toBe(false);
expect(f.observe).not.toHaveBeenCalled();
expect((await f.store.scan()).receipts).toHaveLength(1);
});
it("does not clear a receipt replaced while observation was pending", async () => {
const f = await fixture();
f.observe.mockImplementation(async () => {
await NodeFSP.writeFile(
f.path,
JSON.stringify({ ...f.saved, nativeSessionId: "replacement" }),
);
return { feature: "owned_session_settlement_observation_v1", status: "settled" };
});
await expect(f.run()).rejects.toThrow("settlement could not be proved");
expect((await f.store.scan()).receipts).toHaveLength(1);
});
});
71 changes: 71 additions & 0 deletions apps/server/src/provider/prime/PrimeAgentLegacySettlement.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,71 @@
import * as Schema from "effect/Schema";

import type { PrimeAgentDaemonBridge } from "./PrimeAgentDaemonBridge.ts";
import {
PrimeAgentOwnershipReceiptStore,
primeAgentOwnershipReceiptIsSafeLive,
} from "./PrimeAgentOwnershipReceipt.ts";

const isSettled = Schema.is(
Schema.Struct({
feature: Schema.Literal("owned_session_settlement_observation_v1"),
status: Schema.Literal("settled"),
}),
);

/** Runs only from managed maintenance with a freshly receipt-verified target SDK. */
export async function recoverPrimeAgentLegacySettlement(input: {
readonly instanceId: string;
readonly store: PrimeAgentOwnershipReceiptStore;
readonly loadBridge: () => Promise<
Pick<PrimeAgentDaemonBridge, "sdkFeatures" | "observeOwnedSessionSettlement">
>;
}): Promise<boolean> {
const scan = await input.store.scan();
if (scan.corrupt)
throw new Error("Prime ownership recovery is blocked by an unreadable ownership record.");
const receipts = scan.receipts.filter(
(receipt) =>
receipt.instanceId === input.instanceId && !primeAgentOwnershipReceiptIsSafeLive(receipt),
);
if (receipts.length === 0) return false;
if (receipts.some((receipt) => receipt.state !== "acquired" || receipt.recovery !== undefined)) {
throw new Error("Prime ownership recovery requires the original session's recovery authority.");
}
const bridge = await input.loadBridge();
const observe = bridge.observeOwnedSessionSettlement;
if (!bridge.sdkFeatures?.includes("owned_session_settlement_observation_v1") || !observe) {
throw new Error(
"Prime ownership recovery needs a managed build with settlement observation support.",
);
}
let changed = false;
for (const receipt of receipts) {
if (receipt.state !== "acquired") continue;
const { appVersion, buildId, ...daemon } = receipt.attachProof.daemon;
const contractProof = {
...receipt.attachProof,
daemon: {
...daemon,
...(appVersion === undefined ? {} : { appVersion }),
...(buildId === undefined ? {} : { buildId }),
},
};
const cleared = await input.store.clearLegacyAfterSettlement(receipt, async () =>
isSettled(
await observe({
agentDir: receipt.effectiveHome,
activeSessionId: receipt.activeSessionId,
contractProof,
}),
),
);
if (!cleared) {
throw new Error(
"Prime ownership remains quarantined: prior session settlement could not be proved.",
);
}
changed = true;
}
return changed;
}
42 changes: 42 additions & 0 deletions apps/server/src/provider/prime/PrimeAgentManagedToolStore.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -376,6 +376,7 @@ async function makeHarness(
readonly installMode?: "seam" | "production";
readonly crashAfterCommitOnce?: boolean;
readonly installationBarrier?: () => Promise<void>;
readonly recoverLegacyOwnership?: PrimeManagedToolStoreDependencies["recoverLegacyOwnership"];
readonly loadMetadata?: PrimeManagedToolStoreDependencies["loadLatestVerifiedPublicationMetadata"];
} = {},
) {
Expand Down Expand Up @@ -404,6 +405,9 @@ async function makeHarness(
if (loaderError) throw loaderError;
return currentBundle;
}),
...(input.recoverLegacyOwnership === undefined
? {}
: { recoverLegacyOwnership: input.recoverLegacyOwnership }),
readBinding: async () => binding,
listBindings: async () => [{ instanceId: "primeAgent", binding }],
listOwnedRuntimeBuildReferences: async () => ownedRuntimeBuildReferences,
Expand Down Expand Up @@ -498,6 +502,44 @@ function commandId(prefix: string): string {
}

describe("Pylon-managed Prime tool store", () => {
it("observes legacy settlement on verified staged bytes before reserving and switching", async () => {
const recover = vi.fn(async (_instanceId: string, binaryPath: string) => {
expect(harness.binding.binaryPath).toBe(harness.stock);
expect(await NodeFSP.realpath(binaryPath)).toContain(
"node_modules/prime-agent/dist/bundle/cli.js",
);
});
const harness = await makeHarness({ recoverLegacyOwnership: recover });
const result = await harness.store.command({
commandId: commandId("recover"),
instanceId: "primeAgent",
action: "install",
channel: "stable",
scheduleIfBusy: false,
});
expect(result.status).toBe("succeeded");
expect(recover).toHaveBeenCalledOnce();
expect(recover).toHaveBeenCalledWith("primeAgent", harness.binding.binaryPath);
});

it("preserves the configured runtime when legacy settlement remains unproved", async () => {
const harness = await makeHarness({
recoverLegacyOwnership: async () => {
throw new Error("quarantined");
},
});
const result = await harness.store.command({
commandId: commandId("unproved"),
instanceId: "primeAgent",
action: "install",
channel: "stable",
scheduleIfBusy: false,
});
expect(result).toMatchObject({ status: "failed", message: "quarantined" });
expect(harness.binding.binaryPath).toBe(harness.stock);
expect((await harness.store.status("primeAgent")).selectedBuildId).toBeNull();
});

it.each(["absent", "failed"] as const)(
"reports %s publications without loading an install archive",
async (outcome) => {
Expand Down
8 changes: 8 additions & 0 deletions apps/server/src/provider/prime/PrimeAgentManagedToolStore.ts
Original file line number Diff line number Diff line change
Expand Up @@ -99,6 +99,11 @@ export interface PrimeManagedToolStoreDependencies {
instanceId: string,
expected: PrimeManagedBinding,
) => Promise<PrimeManagedReservationResult>;
/** Called only with a receipt-verified managed target, before the normal quiescence fence. */
readonly recoverLegacyOwnership?: (
instanceId: string,
verifiedBinaryPath: string,
) => Promise<void>;
/** Compare-and-set the complete expected binding while the quiescence fence is held. */
readonly commitBinding: (input: {
readonly instanceId: string;
Expand Down Expand Up @@ -1170,6 +1175,9 @@ export class PrimeAgentManagedToolStore {
message: "Waiting for exact provider-instance quiescence before changing its binary.",
});
state = await this.#readState();
if (input.buildId !== null) {
await this.#dependencies.recoverLegacyOwnership?.(receipt.instanceId, input.targetBinaryPath);
}
const reservation = await this.#dependencies.reserveQuiescentBinding(
receipt.instanceId,
input.expected,
Expand Down
26 changes: 26 additions & 0 deletions apps/server/src/provider/prime/PrimeAgentOwnershipReceipt.ts
Original file line number Diff line number Diff line change
Expand Up @@ -939,6 +939,32 @@ export class PrimeAgentOwnershipReceiptStore {
);
}

/** A replacement observes settlement; it never claims the prior live owner's authority. */
async clearLegacyAfterSettlement(
receipt: PrimeAgentAcquiredOwnershipReceipt,
observe: () => Promise<boolean>,
): Promise<boolean> {
if (receipt.recovery !== undefined || receipt.ownerProcessId === PROCESS_OWNER_ID) return false;
const filePath = this.filePath(receipt.attemptId);
return await withReceiptLock(
this.directory,
{ attemptId: receipt.attemptId, effectiveHome: receipt.effectiveHome },
this.lockRuntime,
async () => {
const previous = await readReceiptFile(filePath);
if (previous.state !== "acquired" || !sameAcquiredReceipt(previous, receipt)) return false;
if (!(await observe())) return false;
// Keep the observation under the same cross-process lock as the exact receipt CAS.
const current = await readReceiptFile(filePath);
if (current.state !== "acquired" || !sameAcquiredReceipt(current, receipt)) return false;
await NodeFSP.unlink(filePath);
await syncDirectory(this.directory);
liveSafeAttempts.delete(receipt.attemptId);
return true;
},
);
}

async claimForAdoption(input: {
readonly receipt: PrimeAgentAcquiredOwnershipReceipt;
readonly nextConfigRevision: string;
Expand Down
Loading
Loading