Skip to content
Closed
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
97 changes: 97 additions & 0 deletions src/core-mutation-tools.ts
Original file line number Diff line number Diff line change
Expand Up @@ -98,6 +98,92 @@ export function assertCoreMutationRecoveryOwnerClient(
};
}

const rebindOwnerRecoveryEnvelopeSchema = z.strictObject({
schema: z.literal("devspace.core_mutation_rebind_owner_recovery.v1"),
auditEvidence: z.string().min(1).max(1024),
orphanProcessRecovery: z.strictObject({
expectedOriginalActorKey: z.string().regex(/^(?:mcp|openai):[0-9a-f]{64}$/),
expectedSourceHead: z.string().regex(/^[0-9a-f]{40}$/),
expectedSourceTree: z.string().regex(/^[0-9a-f]{40}$/),
expectedCurrentHead: z.string().regex(/^[0-9a-f]{40}$/),
expectedTargetTree: z.string().regex(/^git-tree:[0-9a-f]{40}$/),
expectedDiffHash: z.string().regex(/^sha256:[0-9a-f]{64}$/),
expectedChangedPaths: z.array(z.string().min(1)).max(100),
expectedDeletedPaths: z.array(z.string().min(1)).max(100),
}),
});

type RebindOwnerRecoveryEnvelope = z.infer<typeof rebindOwnerRecoveryEnvelopeSchema>;

export function parseCoreMutationRebindOwnerRecoveryEvidence(
evidence: string,
): RebindOwnerRecoveryEnvelope | undefined {
let parsed: unknown;
try {
parsed = JSON.parse(evidence);
} catch {
return undefined;
}
if (
!parsed ||
typeof parsed !== "object" ||
Array.isArray(parsed) ||
(parsed as Record<string, unknown>).schema !== "devspace.core_mutation_rebind_owner_recovery.v1"
) {
return undefined;
}
const result = rebindOwnerRecoveryEnvelopeSchema.safeParse(parsed);
if (!result.success) {
throw new Error(
"[CORE_MUTATION_REBIND_RECOVERY_EVIDENCE_INVALID] Owner recovery envelope is malformed or contains unexpected fields.",
);
}
return result.data;
}

export async function recoverOrphanedProcessBeforeRebind(input: {
store: CoreMutationSessionStore;
workspaceSessionId: string;
workspaceRoot: string;
sessionId: string;
bindingHash: string;
evidence: string;
extra: CoreMutationToolExtra;
recoveryOwnerClientId?: string;
inspectWriterDomain?: (
session: NonNullable<ReturnType<CoreMutationSessionStore["getById"]>>,
domain: CoreMutationManagedWriterDomain,
) => Promise<"CLEAR" | "ACTIVE" | "UNKNOWN"> | "CLEAR" | "ACTIVE" | "UNKNOWN";
}): Promise<boolean> {
const session = input.store.getById(input.sessionId);
const needsRecovery =
session?.workspaceSessionId === input.workspaceSessionId &&
session.status === "ACTIVE" &&
session.writerReconciliationState === "OUTCOME_UNKNOWN" &&
session.writerDomains.length === 1 &&
session.writerDomains[0] === "PROCESS";
if (!needsRecovery) return false;

const envelope = parseCoreMutationRebindOwnerRecoveryEvidence(input.evidence);
if (!envelope) return false;
if (!input.inspectWriterDomain) {
throw new Error(
"[CORE_MUTATION_RECOVERY_DISABLED] PROCESS writer inspection is unavailable on this DevSpace runtime.",
);
}
const owner = assertCoreMutationRecoveryOwnerClient(input.extra, input.recoveryOwnerClientId);
await input.store.recoverOrphanedProcessEffect({
sessionId: input.sessionId,
workspaceSessionId: input.workspaceSessionId,
workspaceRoot: input.workspaceRoot,
recoveryActorKey: owner.recoveryActorKey,
bindingHash: input.bindingHash,
evidence: envelope.orphanProcessRecovery as CoreMutationOrphanProcessRecoveryEvidence,
inspectProcessWriter: (record) => input.inspectWriterDomain!(record, "PROCESS"),
});
return true;
}

function controllerCallerFingerprints(extra: CoreMutationToolExtra): {
callerIdentityFingerprint?: string;
conversationIdentityFingerprint?: string;
Expand Down Expand Up @@ -496,6 +582,17 @@ export function registerCoreMutationSessionTools(
const workspace = workspaces.getWorkspace(workspaceId);
const newActorKey = actorKeyRequired(extra);
try {
await recoverOrphanedProcessBeforeRebind({
store,
workspaceSessionId: workspaceId,
workspaceRoot: workspace.root,
sessionId,
bindingHash,
evidence,
extra,
recoveryOwnerClientId: options.recoveryOwnerClientId,
inspectWriterDomain,
});
const rebound = await store.rebindActor({
sessionId,
workspaceSessionId: workspaceId,
Expand Down
160 changes: 159 additions & 1 deletion src/server.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -56,7 +56,12 @@ import {
DIRECT_CANDIDATE_EXECUTION_SCHEMA,
validateDirectCandidateExecutionEvidence,
} from "./execution-protocol.js";
import { assertCoreMutationRecoveryOwnerClient, CORE_MUTATION_TEST_ONLY_UNTRUSTED_BYPASS } from "./core-mutation-tools.js";
import {
assertCoreMutationRecoveryOwnerClient,
CORE_MUTATION_TEST_ONLY_UNTRUSTED_BYPASS,
parseCoreMutationRebindOwnerRecoveryEvidence,
recoverOrphanedProcessBeforeRebind,
} from "./core-mutation-tools.js";

const execFileAsync = promisify(execFile);

Expand Down Expand Up @@ -2005,6 +2010,159 @@ test("Core orphan PROCESS recovery tool keeps a stable catalog and exact-owner-c
assert.equal(owner.recoveryActorKey, `mcp:${createHash("sha256").update(ownerClientId).digest("hex")}`);
});

test("existing rebind surface can owner-recover orphan PROCESS before caller handoff without a new tool schema", async (t) => {
const originalConversationScopeId = "core-rebind-owner-recovery-original";
const nextConversationScopeId = "core-rebind-owner-recovery-next";
const ownerClientId = "devspace-core-recovery-owner";
const context = await fixture(t, { git: true, coreMutation: true, coreMutationRecoveryOwnerClientId: ownerClientId });
const opened = await callOpen(context.client, context.project, originalConversationScopeId);
const workspaceId = structuredContent(opened).workspaceId as string;
const bound = await bindTestCoreSession({
fixture: context,
workspaceId,
workspaceRoot: context.project,
conversationScopeId: originalConversationScopeId,
allowedPaths: ["AGENTS.md"],
});
const originalActorKey = `openai:${createHash("sha256").update(originalConversationScopeId).digest("hex")}`;
const nextActorKey = `openai:${createHash("sha256").update(nextConversationScopeId).digest("hex")}`;

await context.coreMutationSessions!.admitEffect({
workspaceSessionId: workspaceId,
workspaceRoot: context.project,
workspaceMode: "checkout",
managed: false,
actorKey: originalActorKey,
pointer: { required: true, sessionId: bound.session.id, bindingHash: bound.session.bindingHash },
paths: ["AGENTS.md"],
pathContainment: "NOT_PROVEN",
writerDomain: "PROCESS",
});
writeFileSync(join(context.project, "AGENTS.md"), "# recovered then rebound\n");
const snapshot = await context.coreMutationSessions!.snapshot({
sessionId: bound.session.id,
workspaceSessionId: workspaceId,
workspaceRoot: context.project,
actorKey: originalActorKey,
});
const evidence = JSON.stringify({
schema: "devspace.core_mutation_rebind_owner_recovery.v1",
auditEvidence: "issue-240 existing rebind surface owner recovery",
orphanProcessRecovery: {
expectedOriginalActorKey: originalActorKey,
expectedSourceHead: bound.head,
expectedSourceTree: bound.tree,
expectedCurrentHead: snapshot.currentHead,
expectedTargetTree: snapshot.targetTree,
expectedDiffHash: snapshot.diffHash,
expectedChangedPaths: snapshot.changedPaths,
expectedDeletedPaths: snapshot.deletedPaths,
},
});
assert.ok(parseCoreMutationRebindOwnerRecoveryEvidence(evidence));

const recovered = await recoverOrphanedProcessBeforeRebind({
store: context.coreMutationSessions!,
workspaceSessionId: workspaceId,
workspaceRoot: context.project,
sessionId: bound.session.id,
bindingHash: bound.session.bindingHash,
evidence,
extra: { authInfo: { clientId: ownerClientId }, _meta: { "openai/session": nextConversationScopeId } },
recoveryOwnerClientId: ownerClientId,
inspectWriterDomain: () => "CLEAR",
});
assert.equal(recovered, true);
assert.deepEqual(context.coreMutationSessions!.getById(bound.session.id)?.writerDomains, []);
assert.equal(context.coreMutationSessions!.getById(bound.session.id)?.writerReconciliationState, "CLEAR");

const rebound = await context.coreMutationSessions!.rebindActor({
sessionId: bound.session.id,
workspaceSessionId: workspaceId,
workspaceRoot: context.project,
expectedActorKey: originalActorKey,
newActorKey: nextActorKey,
evidence,
pointer: { sessionId: bound.session.id, bindingHash: bound.session.bindingHash },
});
assert.equal(rebound.fromActorKey, originalActorKey);
assert.equal(rebound.toActorKey, nextActorKey);
assert.equal(context.coreMutationSessions!.getById(bound.session.id)?.actorKey, nextActorKey);
});

test("rebind owner recovery envelope stays fail-closed for foreign owner client and malformed envelopes", async (t) => {
const originalConversationScopeId = "core-rebind-owner-recovery-foreign";
const ownerClientId = "devspace-core-recovery-owner";
const context = await fixture(t, { git: true, coreMutation: true, coreMutationRecoveryOwnerClientId: ownerClientId });
const opened = await callOpen(context.client, context.project, originalConversationScopeId);
const workspaceId = structuredContent(opened).workspaceId as string;
const bound = await bindTestCoreSession({
fixture: context,
workspaceId,
workspaceRoot: context.project,
conversationScopeId: originalConversationScopeId,
allowedPaths: ["AGENTS.md"],
});
const originalActorKey = `openai:${createHash("sha256").update(originalConversationScopeId).digest("hex")}`;
await context.coreMutationSessions!.admitEffect({
workspaceSessionId: workspaceId,
workspaceRoot: context.project,
workspaceMode: "checkout",
managed: false,
actorKey: originalActorKey,
pointer: { required: true, sessionId: bound.session.id, bindingHash: bound.session.bindingHash },
paths: ["AGENTS.md"],
pathContainment: "NOT_PROVEN",
writerDomain: "PROCESS",
});
const snapshot = await context.coreMutationSessions!.snapshot({
sessionId: bound.session.id,
workspaceSessionId: workspaceId,
workspaceRoot: context.project,
actorKey: originalActorKey,
});
const evidence = JSON.stringify({
schema: "devspace.core_mutation_rebind_owner_recovery.v1",
auditEvidence: "foreign owner must fail closed",
orphanProcessRecovery: {
expectedOriginalActorKey: originalActorKey,
expectedSourceHead: bound.head,
expectedSourceTree: bound.tree,
expectedCurrentHead: snapshot.currentHead,
expectedTargetTree: snapshot.targetTree,
expectedDiffHash: snapshot.diffHash,
expectedChangedPaths: snapshot.changedPaths,
expectedDeletedPaths: snapshot.deletedPaths,
},
});

await assert.rejects(
() => recoverOrphanedProcessBeforeRebind({
store: context.coreMutationSessions!,
workspaceSessionId: workspaceId,
workspaceRoot: context.project,
sessionId: bound.session.id,
bindingHash: bound.session.bindingHash,
evidence,
extra: { authInfo: { clientId: "foreign-client" } },
recoveryOwnerClientId: ownerClientId,
inspectWriterDomain: () => "CLEAR",
}),
/CORE_MUTATION_RECOVERY_OWNER_REQUIRED/,
);
assert.deepEqual(context.coreMutationSessions!.getById(bound.session.id)?.writerDomains, ["PROCESS"]);

assert.throws(
() => parseCoreMutationRebindOwnerRecoveryEvidence(JSON.stringify({
schema: "devspace.core_mutation_rebind_owner_recovery.v1",
auditEvidence: "malformed",
orphanProcessRecovery: {},
unexpected: true,
})),
/CORE_MUTATION_REBIND_RECOVERY_EVIDENCE_INVALID/,
);
});

test("Core mutation snapshot tool returns the physical snapshot over MCP", async (t) => {
const conversationScopeId = "core-snapshot-mcp";
const conversation = { "openai/session": conversationScopeId };
Expand Down
Loading