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
8 changes: 5 additions & 3 deletions apps/server/src/checkpointing/Utils.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,10 +3,12 @@ import { CheckpointRef, ProjectId, type ThreadId } from "@t3tools/contracts";

const CHECKPOINT_REFS_PREFIX = "refs/t3/checkpoints";

export function checkpointRefPrefixForThread(threadId: ThreadId): string {
return `${CHECKPOINT_REFS_PREFIX}/${Encoding.encodeBase64Url(threadId)}/`;
}

export function checkpointRefForThreadTurn(threadId: ThreadId, turnCount: number): CheckpointRef {
return CheckpointRef.make(
`${CHECKPOINT_REFS_PREFIX}/${Encoding.encodeBase64Url(threadId)}/turn/${turnCount}`,
);
return CheckpointRef.make(`${checkpointRefPrefixForThread(threadId)}turn/${turnCount}`);
}

export function resolveThreadWorkspaceCwd(input: {
Expand Down
20 changes: 18 additions & 2 deletions apps/server/src/rollback/RollbackSagaRunner.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ import * as Option from "effect/Option";
import * as Stream from "effect/Stream";

import { CheckpointStore } from "../checkpointing/CheckpointStore.ts";
import { checkpointRefForThreadTurn } from "../checkpointing/Utils.ts";
import { OrchestrationEngineService } from "../orchestration/Services/OrchestrationEngine.ts";
import {
RuntimeReceiptBus,
Expand Down Expand Up @@ -53,9 +54,9 @@ const makeState = (operationId: string): RollbackSagaState => ({
targetRevision: 1,
sourceTurnId: null,
targetTurnId: null,
sourceCheckpointRef: CheckpointRef.make("refs/t3/checkpoints/thread-runner/turn/2"),
sourceCheckpointRef: checkpointRefForThreadTurn(threadId, 2),
sourceCheckpointOid: "a".repeat(40),
targetCheckpointRef: CheckpointRef.make("refs/t3/checkpoints/thread-runner/turn/1"),
targetCheckpointRef: checkpointRefForThreadTurn(threadId, 1),
targetCheckpointOid: "b".repeat(40),
targetCheckpointDigest: "target-tree-digest",
providerInstanceId,
Expand Down Expand Up @@ -550,6 +551,21 @@ it.effect("compensates workspace and provider when the provider stays at source"
}),
);

it.effect("does not replay a persisted preimage under a mismatched checkpoint owner", () =>
Effect.gen(function* () {
const environment = makeEnvironment("operation-owner-mismatch", "stayed-source");
environment.setPersistedState({
sourceCheckpointRef: CheckpointRef.make("refs/t3/checkpoints/foreign/turn/2"),
});
const runner = yield* environment.makeRunner();
yield* runner.run("operation-owner-mismatch", false);
const snapshot = environment.snapshot();
assert.equal(snapshot.record.state.phase, "manual-recovery");
assert.isFalse(snapshot.workspaceCalls.includes("restorePreimage"));
assert.equal(snapshot.workspaceDigest, "workspace-target");
}),
);

it.effect(
"fails closed in manual recovery when inspection finds an unrelated provider anchor",
() =>
Expand Down
33 changes: 21 additions & 12 deletions apps/server/src/rollback/RollbackSagaRunner.ts
Original file line number Diff line number Diff line change
Expand Up @@ -189,18 +189,26 @@ export const make = Effect.gen(function* () {
record = workspaceStarted.value;
const preimage = privatePreimage(record.state);
if (preimage !== null) {
workspaceProved = yield* workspace
.restorePreimage({
cwd: record.state.workspaceCwd,
preimage,
})
.pipe(
Effect.tap(() => after("side-effect:workspace-compensated", record.operationId)),
Effect.match({
onFailure: () => false,
onSuccess: (receipt) => receipt.digest === preimage.digest,
}),
);
const validOwner =
record.state.sourceCheckpointRef ===
checkpointRefForThreadTurn(record.state.threadId, record.state.sourceRevision) &&
record.state.targetCheckpointRef ===
checkpointRefForThreadTurn(record.state.threadId, record.state.targetRevision);
workspaceProved = validOwner
? yield* workspace
.restorePreimage({
cwd: record.state.workspaceCwd,
threadId: record.state.threadId,
preimage,
})
.pipe(
Effect.tap(() => after("side-effect:workspace-compensated", record.operationId)),
Effect.match({
onFailure: () => false,
onSuccess: (receipt) => receipt.digest === preimage.digest,
}),
)
: false;
}

const workspaceComplete = yield* update(record, {
Expand Down Expand Up @@ -332,6 +340,7 @@ export const make = Effect.gen(function* () {
.capturePreimage({
operationId: state.operationId,
cwd: state.workspaceCwd,
threadId: state.threadId,
targetCheckpointOid: state.targetCheckpointOid,
})
.pipe(Effect.result);
Expand Down
83 changes: 69 additions & 14 deletions apps/server/src/rollback/RollbackWorkspace.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,8 @@ import * as NodeFSP from "node:fs/promises";
import * as NodePath from "node:path";
import * as NodeOS from "node:os";
import * as NodeUtil from "node:util";
import { ThreadId } from "@t3tools/contracts";
import { checkpointRefForThreadTurn } from "../checkpointing/Utils.ts";

import * as ServerConfig from "../config.ts";
import { RollbackWorkspace, layer as RollbackWorkspaceLive } from "./RollbackWorkspace.ts";
Expand All @@ -33,6 +35,9 @@ layer("RollbackWorkspace", (it) => {
(cwd) =>
Effect.gen(function* () {
const workspace = yield* RollbackWorkspace;
const threadId = ThreadId.make("thread-test");
const targetRef = checkpointRefForThreadTurn(threadId, 1);
const sourceRef = checkpointRefForThreadTurn(threadId, 2);

yield* Effect.promise(async () => {
await run(cwd, ["init", "-b", "main"]);
Expand Down Expand Up @@ -64,7 +69,7 @@ layer("RollbackWorkspace", (it) => {
"-m",
"target checkpoint",
]);
await run(cwd, ["update-ref", "refs/t3/checkpoints/thread-test/turn/1", targetOid]);
await run(cwd, ["update-ref", targetRef, targetOid]);
await run(cwd, ["reset", "--hard", "HEAD"]);

await NodeFSP.writeFile(NodePath.join(cwd, "tracked.txt"), "staged\n");
Expand All @@ -83,7 +88,7 @@ layer("RollbackWorkspace", (it) => {
NodePath.join(cwd, "ignored-target.txt"),
"source ignored path\n",
);
await run(cwd, ["update-ref", "refs/t3/checkpoints/thread-test/turn/2", "HEAD"]);
await run(cwd, ["update-ref", sourceRef, "HEAD"]);
});

const sourceStatus = yield* Effect.promise(() =>
Expand All @@ -95,11 +100,12 @@ layer("RollbackWorkspace", (it) => {
);
const target = yield* workspace.resolveCheckpoint({
cwd,
checkpointRef: "refs/t3/checkpoints/thread-test/turn/1",
checkpointRef: targetRef,
});
const preimage = yield* workspace.capturePreimage({
operationId: "operation-workspace",
cwd,
threadId,
targetCheckpointOid: target.oid,
});
assert.notEqual(preimage.digest.length, 0);
Expand Down Expand Up @@ -149,17 +155,11 @@ layer("RollbackWorkspace", (it) => {
"refs/heads/main",
);
assert.equal(
yield* Effect.promise(() =>
run(cwd, ["rev-parse", "refs/t3/checkpoints/thread-test/turn/2"]),
),
yield* Effect.promise(() => run(cwd, ["rev-parse", sourceRef])),
yield* Effect.promise(() => run(cwd, ["rev-parse", "HEAD"])),
);

yield* Effect.promise(async () => {
await run(cwd, ["update-ref", "-d", "refs/t3/checkpoints/thread-test/turn/2"]);
await run(cwd, ["update-ref", "refs/t3/checkpoints/rogue/turn/99", "HEAD"]);
});
const restored = yield* workspace.restorePreimage({ cwd, preimage });
const restored = yield* workspace.restorePreimage({ cwd, threadId, preimage });
assert.equal(restored.digest, preimage.digest);
assert.equal(
yield* Effect.promise(() => run(cwd, ["status", "--porcelain=v1", "-uall"])),
Expand All @@ -179,6 +179,16 @@ layer("RollbackWorkspace", (it) => {
),
sourceRefs,
);
const driftRef = checkpointRefForThreadTurn(threadId, 3);
yield* Effect.promise(() => run(cwd, ["update-ref", driftRef, "HEAD"]));
const driftRestore = yield* workspace
.restorePreimage({ cwd, threadId, preimage })
.pipe(Effect.exit);
assert.equal(driftRestore._tag, "Failure");
assert.equal(
yield* Effect.promise(() => run(cwd, ["rev-parse", driftRef])),
yield* Effect.promise(() => run(cwd, ["rev-parse", "HEAD"])),
);
assert.equal(
yield* Effect.promise(() =>
NodeFSP.readFile(NodePath.join(cwd, "nested", "untracked.txt"), "utf8"),
Expand Down Expand Up @@ -222,6 +232,11 @@ layer("RollbackWorkspace", (it) => {
(root) =>
Effect.gen(function* () {
const workspace = yield* RollbackWorkspace;
const threadId = ThreadId.make("linked");
const siblingThreadId = ThreadId.make("sibling");
const targetRef = checkpointRefForThreadTurn(threadId, 1);
const siblingFirstRef = checkpointRefForThreadTurn(siblingThreadId, 1);
const siblingSecondRef = checkpointRefForThreadTurn(siblingThreadId, 2);
const main = NodePath.join(root, "main");
const linked = NodePath.join(root, "linked");
yield* Effect.promise(async () => {
Expand All @@ -237,7 +252,7 @@ layer("RollbackWorkspace", (it) => {
await run(linked, ["add", "tracked.txt"]);
const tree = await run(linked, ["write-tree"]);
const checkpoint = await run(linked, ["commit-tree", tree, "-m", "linked checkpoint"]);
await run(linked, ["update-ref", "refs/t3/checkpoints/linked/turn/1", checkpoint]);
await run(linked, ["update-ref", targetRef, checkpoint]);
await run(linked, ["reset", "--hard", "HEAD"]);
await NodeFSP.writeFile(NodePath.join(linked, "tracked.txt"), "linked pre-image\n");
await run(linked, ["add", "tracked.txt"]);
Expand All @@ -249,24 +264,64 @@ layer("RollbackWorkspace", (it) => {
assert.equal(mainIdentity.gitCommonDir, linkedIdentity.gitCommonDir);
assert.notEqual(mainIdentity.workspaceKey, linkedIdentity.workspaceKey);

const siblingFirst = yield* Effect.promise(async () => {
const oid = await run(main, [
"commit-tree",
"HEAD^{tree}",
"-m",
"sibling first checkpoint",
]);
await run(main, ["update-ref", siblingFirstRef, oid]);
return oid;
});

const checkpoint = yield* workspace.resolveCheckpoint({
cwd: linked,
checkpointRef: "refs/t3/checkpoints/linked/turn/1",
checkpointRef: targetRef,
});
const preimage = yield* workspace.capturePreimage({
operationId: "operation-linked",
cwd: linked,
threadId,
targetCheckpointOid: checkpoint.oid,
});
assert.deepEqual(preimage.ownedRefs, [{ ref: targetRef, oid: checkpoint.oid }]);
const siblingAdvanced = yield* Effect.promise(() =>
run(main, ["commit-tree", "HEAD^{tree}", "-m", "sibling advanced checkpoint"]),
);
yield* Effect.promise(async () => {
// Both worktrees share the Git common dir, but hold different
// rollback leases. B may write its refs while A is compensating.
await run(main, ["update-ref", siblingFirstRef, siblingAdvanced]);
await run(main, ["update-ref", siblingSecondRef, siblingFirst]);
});
yield* workspace.applyCheckpoint({ cwd: linked, checkpointOid: checkpoint.oid });
assert.equal(
yield* Effect.promise(() =>
NodeFSP.readFile(NodePath.join(linked, "tracked.txt"), "utf8"),
),
"checkpoint\n",
);
const restored = yield* workspace.restorePreimage({ cwd: linked, preimage });
// Simulate a pre-upgrade persisted preimage containing all repository
// checkpoint refs. B's saved OID must have no recovery authority.
const legacyPreimage = {
...preimage,
ownedRefs: [...preimage.ownedRefs, { ref: siblingFirstRef, oid: siblingFirst }],
};
const restored = yield* workspace.restorePreimage({
cwd: linked,
threadId,
preimage: legacyPreimage,
});
assert.equal(restored.digest, preimage.digest);
assert.equal(
yield* Effect.promise(() => run(main, ["rev-parse", siblingFirstRef])),
siblingAdvanced,
);
assert.equal(
yield* Effect.promise(() => run(main, ["rev-parse", siblingSecondRef])),
siblingFirst,
);
assert.equal(
yield* Effect.promise(() =>
NodeFSP.readFile(NodePath.join(linked, "tracked.txt"), "utf8"),
Expand Down
Loading
Loading