Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
19 commits
Select commit Hold shift + click to select a range
99a4a25
fix(web): make bulk thread deletion reliable
0x00-sys Sep 17, 2026
7fd630b
fix: preserve worktrees when thread deletion fails
0x00-sys Sep 17, 2026
0634558
fix: protect archived deletion and worktree cleanup retries
0x00-sys Sep 17, 2026
3b13860
fix: prevent worktree adoption during cleanup
0x00-sys Sep 17, 2026
23dccb3
fix(web): check deletion support before stopping resources
0x00-sys Sep 17, 2026
0c7fdba
Merge upstream main into bulk deletion fix
0x00-sys Sep 17, 2026
2cbe17f
Merge upstream main into bulk deletion fix
0x00-sys Sep 17, 2026
6b28aca
chore: merge upstream main into bulk deletion fix
0x00-sys Sep 17, 2026
d54ee39
chore: merge upstream main into bulk deletion fix
0x00-sys Sep 17, 2026
6995364
chore: merge upstream main into bulk deletion fix
0x00-sys Sep 18, 2026
c3b225e
chore: merge upstream main into bulk deletion fix
0x00-sys Sep 18, 2026
90f4555
Merge remote-tracking branch 'upstream/main' into t3code/fix-bulk-del…
0x00-sys Sep 19, 2026
3d37672
Merge remote-tracking branch 'upstream/main' into t3code/fix-bulk-del…
0x00-sys Sep 20, 2026
5df7d69
Merge remote-tracking branch 'upstream/main' into t3code/fix-bulk-del…
0x00-sys Sep 22, 2026
7c88573
Merge remote-tracking branch 'upstream/main' into t3code/fix-bulk-del…
0x00-sys Sep 23, 2026
5115f7b
Merge remote-tracking branch 'upstream/main' into t3code/fix-bulk-del…
0x00-sys Sep 24, 2026
0c4224f
Merge remote-tracking branch 'upstream/main' into t3code/fix-bulk-del…
0x00-sys Sep 25, 2026
07beccf
Merge remote-tracking branch 'upstream/main' into t3code/fix-bulk-del…
0x00-sys Sep 27, 2026
6e111f0
Merge remote-tracking branch 'upstream/main' into t3code/fix-bulk-del…
0x00-sys Sep 30, 2026
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
1 change: 1 addition & 0 deletions apps/server/src/environment/ServerEnvironment.ts
Original file line number Diff line number Diff line change
Expand Up @@ -223,6 +223,7 @@ export const make = Effect.gen(function* () {
inlineMessageContext: true,
requiredWorktreeBootstrap: true,
threadSettlement: true,
recoverableThreadDeletion: true,
threadAutoSettlement: true,
storageCleanup: true,
projectWorktreeCleanup: true,
Expand Down
81 changes: 81 additions & 0 deletions apps/server/src/orchestration/Layers/OrchestrationEngine.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,8 @@ import {
import * as NodeServices from "@effect/platform-node/NodeServices";
import { it as effectIt } from "@effect/vitest";
import * as Effect from "effect/Effect";
import * as Deferred from "effect/Deferred";
import * as Fiber from "effect/Fiber";
import * as Layer from "effect/Layer";
import * as ManagedRuntime from "effect/ManagedRuntime";
import * as Metric from "effect/Metric";
Expand Down Expand Up @@ -54,6 +56,7 @@ import {
} from "../Services/ProjectionPipeline.ts";
import { ProjectionSnapshotQuery } from "../Services/ProjectionSnapshotQuery.ts";
import { ServerConfig } from "../../config.ts";
import { removeUnusedWorktree } from "../threadWorktreeDeletion.ts";

const asProjectId = (value: string): ProjectId => ProjectId.make(value);
const asMessageId = (value: string): MessageId => MessageId.make(value);
Expand Down Expand Up @@ -130,6 +133,84 @@ const hasMetricSnapshot = (
);

describe("OrchestrationEngine", () => {
effectIt.effect.each(["thread.create", "thread.meta.update"] as const)(
"prevents concurrent %s adoption during cleanup without blocking unrelated commands",
(type) =>
Effect.gen(function* () {
const engine = yield* OrchestrationEngineService;
const snapshots = yield* ProjectionSnapshotQuery;
const projectId = ProjectId.make("cleanup-project");
const threadId = ThreadId.make("cleanup-adopter");
yield* engine.dispatch({
type: "project.create",
commandId: CommandId.make("cleanup-project"),
projectId,
title: "Test",
workspaceRoot: "/repo",
createdAt: now(),
});
const create = {
type: "thread.create",
commandId: CommandId.make("cleanup-create"),
projectId,
threadId,
title: "Test",
modelSelection: { instanceId: ProviderInstanceId.make("codex"), model: "test" },
runtimeMode: "full-access",
interactionMode: "default",
branch: null,
worktreePath: "/repo/staged",
createdAt: now(),
} as const;
if (type === "thread.meta.update")
yield* engine.dispatch({ ...create, worktreePath: null });
const started = yield* Deferred.make<void>();
const release = yield* Deferred.make<void>();
const cleanup = yield* removeUnusedWorktree(
{ cwd: "/repo", path: create.worktreePath, force: true },
Deferred.succeed(started, undefined).pipe(
Effect.andThen(Deferred.await(release)),
Effect.andThen(
type === "thread.meta.update" ? Effect.die("cleanup failed") : Effect.void,
),
),
).pipe(Effect.exit, Effect.forkChild({ startImmediately: true }));
yield* Deferred.await(started);
const adoption: OrchestrationCommand =
type === "thread.create"
? create
: {
type,
commandId: CommandId.make("cleanup-adopt"),
threadId,
worktreePath: create.worktreePath,
};
expect((yield* engine.dispatch(adoption).pipe(Effect.exit))._tag).toBe("Failure");
expect(
(yield* snapshots.getCommandReadModel()).threads.some(
(thread) => thread.worktreePath === create.worktreePath,
),
).toBe(false);
// This receipt must arrive while filesystem removal is still blocked.
yield* engine.dispatch({
type: "project.create",
commandId: CommandId.make("unrelated-project"),
projectId: ProjectId.make("unrelated-project"),
title: "Other work",
workspaceRoot: "/other",
createdAt: now(),
});
yield* Deferred.succeed(release, undefined);
expect((yield* Fiber.join(cleanup))._tag).toBe(
type === "thread.meta.update" ? "Failure" : "Success",
);
yield* engine.dispatch({ ...adoption, commandId: CommandId.make("retry-adoption") });
expect((yield* snapshots.getCommandReadModel()).threads[0]?.worktreePath).toBe(
create.worktreePath,
);
}).pipe(Effect.provide(makeOrchestrationLayer())),
);

it.each(["running", "stopped"] as const)(
"sends async answers with a %s session and rejects old duplicate replies",
async (status) => {
Expand Down
63 changes: 59 additions & 4 deletions apps/server/src/orchestration/Layers/OrchestrationEngine.ts
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@ import * as Exit from "effect/Exit";
import * as Layer from "effect/Layer";
import * as Metric from "effect/Metric";
import * as Option from "effect/Option";
import * as Path from "effect/Path";
import * as PubSub from "effect/PubSub";
import * as Queue from "effect/Queue";
import * as Schema from "effect/Schema";
Expand Down Expand Up @@ -89,11 +90,13 @@ const makeOrchestrationEngine = Effect.gen(function* () {
const projectionSnapshotQuery = yield* ProjectionSnapshotQuery;
const threadBackgroundLiveness = yield* ThreadBackgroundLivenessService;
const crypto = yield* Crypto.Crypto;
const path = yield* Path.Path;
const worktreeCleanups = new Map<string, number>();

const nowIso = Effect.map(DateTime.now, DateTime.formatIso);
let commandReadModel = createEmptyReadModel(yield* nowIso);

const commandQueue = yield* Queue.unbounded<CommandEnvelope>();
const commandQueue = yield* Queue.unbounded<Effect.Effect<void>>();
const eventPubSub = yield* PubSub.unbounded<OrchestrationEvent>();

const projectEventsOntoReadModel = (
Expand Down Expand Up @@ -242,6 +245,17 @@ const makeOrchestrationEngine = Effect.gen(function* () {
envelope.command.type === "thread.user-input.dismiss"
? yield* projectionSnapshotQuery.getUserInputActivity(envelope.command)
: Option.none();
if (
(envelope.command.type === "thread.create" ||
envelope.command.type === "thread.meta.update") &&
envelope.command.worktreePath &&
worktreeCleanups.has(path.resolve(envelope.command.worktreePath))
) {
return yield* new OrchestrationCommandInvariantError({
commandType: envelope.command.type,
detail: "Worktree cleanup is in progress. Try again after it finishes.",
});
}
const eventBase = yield* decideOrchestrationCommand({
command: envelope.command,
readModel: commandReadModel,
Expand Down Expand Up @@ -413,7 +427,7 @@ const makeOrchestrationEngine = Effect.gen(function* () {
yield* projectionPipeline.bootstrap;
commandReadModel = yield* projectionSnapshotQuery.getCommandReadModel();

const worker = Effect.forever(Queue.take(commandQueue).pipe(Effect.flatMap(processEnvelope)));
const worker = Effect.forever(Queue.take(commandQueue).pipe(Effect.flatten));
yield* Effect.forkScoped(worker);
yield* Effect.logDebug("orchestration engine started").pipe(
Effect.annotateLogs({ sequence: commandReadModel.snapshotSequence }),
Expand All @@ -438,20 +452,61 @@ const makeOrchestrationEngine = Effect.gen(function* () {
const dispatch: OrchestrationEngineShape["dispatch"] = (command, options) =>
Effect.gen(function* () {
const result = yield* Deferred.make<{ sequence: number }, OrchestrationDispatchError>();
yield* Queue.offer(commandQueue, {
const envelope: CommandEnvelope = {
command,
origin: options?.origin,
result,
startedAtMs: yield* Clock.currentTimeMillis,
});
};
yield* Queue.offer(
commandQueue,
Effect.suspend(() => processEnvelope(envelope)),
);
return yield* Deferred.await(result);
});

const runSerialized = <A, E>(effect: Effect.Effect<A, E>) =>
Effect.gen(function* () {
const result = yield* Deferred.make<A, E>();
yield* Queue.offer(
commandQueue,
effect.pipe(
Effect.exit,
Effect.flatMap((exit) => Deferred.done(result, exit)),
Effect.asVoid,
),
);
return yield* Deferred.await(result);
});

const withWorktreeCleanup: OrchestrationEngineShape["withWorktreeCleanup"] = (paths, cleanup) => {
const keys = [...new Set(paths.map((entry) => path.resolve(entry)))];
return Effect.acquireUseRelease(
runSerialized(
Effect.sync(() => {
for (const key of keys) worktreeCleanups.set(key, (worktreeCleanups.get(key) ?? 0) + 1);
}),
),
() => cleanup,
() =>
runSerialized(
Effect.sync(() => {
for (const key of keys) {
const remaining = worktreeCleanups.get(key)! - 1;
if (remaining === 0) worktreeCleanups.delete(key);
else worktreeCleanups.set(key, remaining);
}
}),
),
);
};

return {
readEvents,
readThreadEvents,
getThreadReplayStats,
dispatch,
withWorktreeCleanup,
subscribeDomainEvents: PubSub.subscribe(eventPubSub).pipe(Effect.map(Stream.fromSubscription)),
// Each access creates a fresh PubSub subscription so that multiple
// consumers (wsServer, ProviderRuntimeIngestion, CheckpointReactor, etc.)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -425,6 +425,7 @@ describe("ProviderCommandReactor", () => {
Effect.gen(function* () {
const engine = yield* OrchestrationEngineService;
return {
withWorktreeCleanup: engine.withWorktreeCleanup,
readEvents: engine.readEvents,
readThreadEvents: engine.readThreadEvents,
getThreadReplayStats: engine.getThreadReplayStats,
Expand Down
6 changes: 6 additions & 0 deletions apps/server/src/orchestration/Services/OrchestrationEngine.ts
Original file line number Diff line number Diff line change
Expand Up @@ -75,6 +75,12 @@ export interface OrchestrationEngineShape {
options?: { readonly origin?: OrchestrationClientOrigin },
) => Effect.Effect<{ sequence: number }, OrchestrationDispatchError, never>;

/** Reserve paths on the command queue, preventing adoption while filesystem cleanup runs outside it. */
readonly withWorktreeCleanup: <A, E>(
paths: ReadonlyArray<string>,
cleanup: Effect.Effect<A, E>,
) => Effect.Effect<A, E>;

/**
* Stream persisted domain events in dispatch order.
*
Expand Down
22 changes: 21 additions & 1 deletion apps/server/src/orchestration/decider.ts
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@ import {
threadPullRequestKeysEqual,
} from "@t3tools/shared/threadPullRequests";
import { compareDateTimeStrings } from "@t3tools/shared/dateTime";
import { normalizeProjectPathForComparison } from "@t3tools/shared/path";
import * as DateTime from "effect/DateTime";
import * as Crypto from "effect/Crypto";
import * as Effect from "effect/Effect";
Expand Down Expand Up @@ -411,11 +412,30 @@ export const decideOrchestrationCommand = Effect.fn("decideOrchestrationCommand"
}

case "thread.delete": {
yield* requireThread({
const thread = yield* requireThread({
readModel,
command,
threadId: command.threadId,
});
if (command.deleteWorktreePath) {
const worktreePath = normalizeProjectPathForComparison(command.deleteWorktreePath);
if (
thread.worktreePath === null ||
normalizeProjectPathForComparison(thread.worktreePath) !== worktreePath ||
readModel.threads.some(
(entry) =>
entry.id !== thread.id &&
entry.deletedAt === null &&
entry.worktreePath !== null &&
normalizeProjectPathForComparison(entry.worktreePath) === worktreePath,
)
) {
return yield* new OrchestrationCommandInvariantError({
commandType: command.type,
detail: "Worktree references changed during deletion. Try deleting the thread again.",
});
}
}
const occurredAt = yield* nowIso;
return {
...(yield* withEventBase({
Expand Down
Loading
Loading