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
4 changes: 0 additions & 4 deletions apps/mobile/src/features/terminal/terminalMenu.ts
Original file line number Diff line number Diff line change
Expand Up @@ -156,10 +156,6 @@ export function resolveProjectScriptTerminalId(input: {
return nextTerminalId(input.existingTerminalIds, input.uniqueSuffix);
}

export function projectScriptMenuLabel(script: ProjectScript): string {
return script.runOnWorktreeCreate ? `${script.name} (setup)` : script.name;
}

export function projectScriptMenuIcon(icon: ProjectScript["icon"]) {
if (icon === "test") return "flask";
if (icon === "lint") return "checklist";
Expand Down
2 changes: 1 addition & 1 deletion apps/mobile/src/features/threads/ThreadGitControls.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -23,9 +23,9 @@ import {
basename,
getTerminalStatusLabel,
projectScriptMenuIcon,
projectScriptMenuLabel,
type TerminalMenuSession,
} from "../terminal/terminalMenu";
import { projectScriptMenuLabel } from "@t3tools/shared/projectScripts";

function truncateMiddle(value: string, maxLength: number): string {
if (value.length <= maxLength) {
Expand Down
134 changes: 132 additions & 2 deletions apps/server/src/orchestration-v2/ThreadSettlementService.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,7 @@ import * as Stream from "effect/Stream";
import { TestClock } from "effect/testing";

import * as GitManager from "../git/GitManager.ts";
import * as ProjectSetupScriptRunner from "../project/ProjectSetupScriptRunner.ts";
import * as PullRequestService from "../pullRequest/PullRequestService.ts";
import * as ServerActivation from "../serverActivation.ts";
import * as ServerSettings from "../serverSettings.ts";
Expand Down Expand Up @@ -486,6 +487,12 @@ interface HarnessOptions {
readonly pullRequestSummary?: PullRequestService.PullRequestService["Service"]["summary"];
readonly existingWorktreePaths?: ReadonlyArray<string>;
readonly onDispatch?: (command: AutoSettleCommand) => Effect.Effect<void>;
/** Runs before each settle script start; a failure fails that start. */
readonly onScriptRun?: (
input: ProjectSetupScriptRunner.ProjectSetupScriptRunnerInput,
) => Effect.Effect<void>;
/** Runs inside each `closeIdle` call. */
readonly onCloseIdle?: () => Effect.Effect<void>;
/** Threads `getThread` returns when a `thread.settled` event is handled. */
readonly currentThreads?: ReadonlyArray<OrchestrationV2AppThread>;
}
Expand Down Expand Up @@ -514,6 +521,8 @@ const makeHarness = Effect.fn("makeThreadSettlementHarness")(function* (options:
const invalidatedCwds = yield* Ref.make<ReadonlyArray<string>>([]);
const domainEvents = yield* PubSub.unbounded<OrchestrationV2DomainEvent>();
const closedIdle = yield* Queue.unbounded<{ readonly threadId: string }>();
const scriptRuns =
yield* Queue.unbounded<ProjectSetupScriptRunner.ProjectSetupScriptRunnerInput>();

const updateSettings = (patch: ServerSettingsPatch) =>
Effect.gen(function* () {
Expand Down Expand Up @@ -599,7 +608,15 @@ const makeHarness = Effect.fn("makeThreadSettlementHarness")(function* (options:
dispatch,
}),
Layer.mock(TerminalManager.TerminalManager)({
closeIdle: (input) => Queue.offer(closedIdle, input).pipe(Effect.asVoid),
closeIdle: (input) =>
Queue.offer(closedIdle, input).pipe(Effect.andThen(options.onCloseIdle?.() ?? Effect.void)),
}),
Layer.mock(ProjectSetupScriptRunner.ProjectSetupScriptRunner)({
runForThread: (input) =>
Queue.offer(scriptRuns, input).pipe(
Effect.andThen(options.onScriptRun?.(input) ?? Effect.void),
Effect.as({ status: "no-script" } as const),
),
}),
Layer.mock(GitManager.GitManager)({
branchPullRequest,
Expand Down Expand Up @@ -631,6 +648,7 @@ const makeHarness = Effect.fn("makeThreadSettlementHarness")(function* (options:
summaryRecovery,
invalidatedCwds,
closedIdle,
scriptRuns,
publishEvent: (event: OrchestrationV2DomainEvent) => PubSub.publish(domainEvents, event),
updateSettings,
publishMerge: PubSub.publish(mergedPullRequests, {
Expand Down Expand Up @@ -1010,6 +1028,7 @@ describe("ThreadSettlementServiceV2 terminals", () => {
const appThread = (
id: string,
settledOverride: OrchestrationV2AppThread["settledOverride"],
worktreePath: string | null = null,
): OrchestrationV2AppThread => {
const threadId = ThreadId.make(id);
return {
Expand All @@ -1023,7 +1042,7 @@ describe("ThreadSettlementServiceV2 terminals", () => {
runtimeMode: "full-access",
interactionMode: "default",
branch: null,
worktreePath: null,
worktreePath,
activeProviderThreadId: null,
lineage: { parentThreadId: null, relationshipToParent: null, rootThreadId: threadId },
forkedFrom: null,
Expand Down Expand Up @@ -1086,6 +1105,117 @@ describe("ThreadSettlementServiceV2 terminals", () => {
}),
),
);

it.effect("runs the settle script in the thread's own worktree", () =>
Effect.scoped(
Effect.gen(function* () {
const shared = appThread("shared-checkout-thread", "settled");
const missing = appThread("removed-worktree-thread", "settled", "/worktrees/removed");
const worktree = appThread("worktree-thread", "settled", "/worktrees/thread");
const fixture = yield* makeHarness({
snapshot: makeSnapshot([]),
currentThreads: [shared, missing, worktree],
existingWorktreePaths: ["/worktrees/thread"],
});

yield* Effect.gen(function* () {
const service = yield* ThreadSettlementService.ThreadSettlementServiceV2;
yield* startHarness(service, fixture.activation, fixture.snapshotReads);
// Events run in order, so the first run belongs to the last event
// only if the shared checkout and the removed worktree were skipped.
yield* fixture.publishEvent(settledEvent(shared));
yield* fixture.publishEvent(settledEvent(missing));
yield* fixture.publishEvent(settledEvent(worktree));
const run = yield* Queue.take(fixture.scriptRuns);
assert.strictEqual(run.threadId, worktree.id);
assert.strictEqual(run.worktreePath, "/worktrees/thread");
assert.strictEqual(run.trigger, "settle");
}).pipe(Effect.provide(fixture.layer));
}),
),
);

it.effect("skips the settle script for a thread re-engaged while its shells close", () =>
Effect.scoped(
Effect.gen(function* () {
const settled = appThread("reengaging-thread", "settled", "/worktrees/reengaging");
const marker = appThread("marker-thread", "settled", "/worktrees/marker");
const currentThreads = [settled, marker];
const fixture = yield* makeHarness({
snapshot: makeSnapshot([]),
currentThreads,
existingWorktreePaths: ["/worktrees/reengaging", "/worktrees/marker"],
// The user sends a message while the first close checks processes.
onCloseIdle: () =>
Effect.sync(() => {
currentThreads[0] = { ...settled, settledOverride: "active" };
}),
});

yield* Effect.gen(function* () {
const service = yield* ThreadSettlementService.ThreadSettlementServiceV2;
yield* startHarness(service, fixture.activation, fixture.snapshotReads);
yield* fixture.publishEvent(settledEvent(settled));
yield* fixture.publishEvent(settledEvent(marker));
assert.strictEqual((yield* Queue.take(fixture.scriptRuns)).threadId, marker.id);
}).pipe(Effect.provide(fixture.layer));
}),
),
);

it.effect("runs the settle script once when a settlement is re-emitted", () =>
Effect.scoped(
Effect.gen(function* () {
const repeated = appThread("repeated-thread", "settled", "/worktrees/repeated");
const marker = appThread("marker-thread", "settled", "/worktrees/marker");
const fixture = yield* makeHarness({
snapshot: makeSnapshot([]),
currentThreads: [repeated, marker],
existingWorktreePaths: ["/worktrees/repeated", "/worktrees/marker"],
});

yield* Effect.gen(function* () {
const service = yield* ThreadSettlementService.ThreadSettlementServiceV2;
yield* startHarness(service, fixture.activation, fixture.snapshotReads);
// Settling an already settled thread emits the same settlement again.
yield* fixture.publishEvent(settledEvent(repeated));
yield* fixture.publishEvent(settledEvent(repeated));
yield* fixture.publishEvent(settledEvent(marker));
assert.strictEqual((yield* Queue.take(fixture.scriptRuns)).threadId, repeated.id);
assert.strictEqual((yield* Queue.take(fixture.scriptRuns)).threadId, marker.id);
}).pipe(Effect.provide(fixture.layer));
}),
),
);

it.effect("retries the settle script on a re-emitted settlement after a failed start", () =>
Effect.scoped(
Effect.gen(function* () {
const thread = appThread("failed-start-thread", "settled", "/worktrees/failed-start");
const starts = yield* Ref.make(0);
const fixture = yield* makeHarness({
snapshot: makeSnapshot([]),
currentThreads: [thread],
existingWorktreePaths: ["/worktrees/failed-start"],
onScriptRun: () =>
Ref.updateAndGet(starts, (count) => count + 1).pipe(
Effect.flatMap((count) =>
count === 1 ? Effect.die(new Error("terminal failed to open")) : Effect.void,
),
),
});

yield* Effect.gen(function* () {
const service = yield* ThreadSettlementService.ThreadSettlementServiceV2;
yield* startHarness(service, fixture.activation, fixture.snapshotReads);
yield* fixture.publishEvent(settledEvent(thread));
yield* fixture.publishEvent(settledEvent(thread));
assert.strictEqual((yield* Queue.take(fixture.scriptRuns)).threadId, thread.id);
assert.strictEqual((yield* Queue.take(fixture.scriptRuns)).threadId, thread.id);
}).pipe(Effect.provide(fixture.layer));
}),
),
);
});

describe("ThreadSettlementServiceV2 single-thread sweeps", () => {
Expand Down
38 changes: 33 additions & 5 deletions apps/server/src/orchestration-v2/ThreadSettlementService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@ import type * as Scope from "effect/Scope";
import * as Stream from "effect/Stream";

import * as GitManager from "../git/GitManager.ts";
import * as ProjectSetupScriptRunner from "../project/ProjectSetupScriptRunner.ts";
import * as PullRequestService from "../pullRequest/PullRequestService.ts";
import * as ServerSettings from "../serverSettings.ts";
import { forkParked } from "../serverActivation.ts";
Expand Down Expand Up @@ -265,6 +266,10 @@ export const make = Effect.gen(function* () {
const crypto = yield* Crypto.Crypto;
const fileSystem = yield* FileSystem.FileSystem;
const terminals = yield* TerminalManager.TerminalManager;
const projectScripts = yield* ProjectSetupScriptRunner.ProjectSetupScriptRunner;
// Settling a settled thread re-emits thread.settled with the same settledAt,
// so this keeps the settle action to one run per settlement.
const settleActionRunAt = new Map<ThreadId, number>();
Comment thread
macroscopeapp[bot] marked this conversation as resolved.

const sweep = Effect.fn("ThreadSettlementServiceV2.sweep")(function* (
mergedPullRequest: PullRequestService.PullRequestMergeEvent | null,
Expand Down Expand Up @@ -514,20 +519,43 @@ export const make = Effect.gen(function* () {

// Settling closes the thread's shells that sit at an idle prompt, so they stop
// holding the worktree. A terminal running a command (a dev server, an
// editor) stays for the user to close.
const closeIdleTerminals = Effect.fn("ThreadSettlementServiceV2.closeIdleTerminals")(
// editor) stays for the user to close. Then the project's settle script runs
// in the thread's own worktree; a thread in the shared checkout skips it,
// because other threads may still be working there.
const cleanUpSettledThread = Effect.fn("ThreadSettlementServiceV2.cleanUpSettledThread")(
function* (threadId: ThreadId) {
// A thread re-engaged before this event ran keeps its shells.
const settled = yield* projections.getThread(threadId);
if (settled.settledOverride !== "settled") return;
yield* terminals.closeIdle({ threadId });
const worktreePath = settled.worktreePath;
if (worktreePath === null || !(yield* fileSystem.exists(worktreePath))) return;
// Closing and the worktree check wait on I/O. A thread re-engaged
// meanwhile is working again, so its worktree is no place for cleanup.
const thread = yield* projections.getThread(threadId);
if (thread.settledOverride !== "settled") return;
yield* terminals.closeIdle({ threadId });
const settledAtMs = toMillis(thread.settledAt);
if (settledAtMs === null || settleActionRunAt.get(threadId) === settledAtMs) return;
Comment thread
macroscopeapp[bot] marked this conversation as resolved.
const run = yield* projectScripts.runForThread({
Comment thread
macroscopeapp[bot] marked this conversation as resolved.
threadId,
projectId: thread.projectId,
worktreePath,
trigger: "settle",
// A clean exit closes the script's shell so it does not hold the worktree.
observeCompletion: {},
});
// Recorded after a successful start, so a failed start retries on the next event.
settleActionRunAt.set(threadId, settledAtMs);
if (run.status === "started" && run.completion) {
yield* run.completion.pipe(Effect.forkDetach);
}
},
(effect, threadId) =>
effect.pipe(
Effect.catchCause((cause) =>
Cause.hasInterruptsOnly(cause)
? Effect.failCause(cause)
: Effect.logWarning("closing idle terminals after settlement failed", {
: Effect.logWarning("cleaning up a settled thread failed", {
threadId,
cause: Cause.pretty(cause),
}),
Expand All @@ -538,7 +566,7 @@ export const make = Effect.gen(function* () {
const processEvent = (event: OrchestrationV2DomainEvent) => {
switch (event.type) {
case "thread.settled":
return closeIdleTerminals(event.threadId);
return cleanUpSettledThread(event.threadId);
Comment thread
macroscopeapp[bot] marked this conversation as resolved.
case "thread.pull-request-synced":
case "provider-session.detached":
return worker.enqueue(event.threadId);
Expand Down
63 changes: 62 additions & 1 deletion apps/server/src/project/ProjectSetupScriptRunner.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ import * as NodeCrypto from "@effect/platform-node/NodeCrypto";
import { assert, it, vi } from "@effect/vitest";
import { ProjectId } from "@t3tools/contracts";
import * as Effect from "effect/Effect";
import * as Fiber from "effect/Fiber";
import * as Layer from "effect/Layer";
import * as Option from "effect/Option";

Expand Down Expand Up @@ -29,6 +30,9 @@ it.effect("resolves setup scripts through the standalone project service", () =>
const write = vi.fn(
(_input: Parameters<TerminalManager.TerminalManager["Service"]["write"]>[0]) => Effect.void,
);
const closeIdle = vi.fn(
(_input: Parameters<TerminalManager.TerminalManager["Service"]["closeIdle"]>[0]) => Effect.void,
);
const listeners: Array<Parameters<TerminalManager.TerminalManager["Service"]["subscribe"]>[0]> =
[];
const subscribe: TerminalManager.TerminalManager["Service"]["subscribe"] = (listener) =>
Expand All @@ -52,6 +56,14 @@ it.effect("resolves setup scripts through the standalone project service", () =>
icon: "configure" as const,
runOnWorktreeCreate: true,
},
{
id: "clean",
name: "Clean",
command: "cargo clean",
icon: "build" as const,
runOnWorktreeCreate: false,
runOnSettle: true,
},
],
createdAt: "2026-06-20T00:00:00.000Z",
updatedAt: "2026-06-20T00:00:00.000Z",
Expand All @@ -63,7 +75,7 @@ it.effect("resolves setup scripts through the standalone project service", () =>
Layer.mock(ProjectService.ProjectService)({
getById: () => Effect.succeed(Option.some(project)),
}),
Layer.mock(TerminalManager.TerminalManager)({ open, write, subscribe }),
Layer.mock(TerminalManager.TerminalManager)({ open, write, subscribe, closeIdle }),
ServerSettings.layerTest(),
NodeCrypto.layer,
),
Expand Down Expand Up @@ -117,5 +129,54 @@ it.effect("resolves setup scripts through the standalone project service", () =>
});
assert.deepEqual(lines, ["Downloading 10%", "Downloading 20%", "Done"]);
yield* listener({ type: "closed", threadId: "thread-1", terminalId: "setup-setup" });

const settle = yield* runner.runForThread({
threadId: "thread-1",
projectId,
worktreePath: "/repo-worktree",
trigger: "settle",
});
const settleTerminalId = settle.status === "started" ? settle.terminalId : "";
assert.match(settleTerminalId, /^settle-clean-/);
assert.equal(write.mock.calls.at(-1)?.[0].data, "cargo clean\r");

// A clean run closes its shell once the prompt is back, not at the
// sentinel, so the prompt redraw is not taken for new activity.
const observedSettle = yield* runner.runForThread({
threadId: "thread-1",
projectId,
worktreePath: "/repo-worktree",
trigger: "settle",
observeCompletion: {},
});
const observedTerminalId = observedSettle.status === "started" ? observedSettle.terminalId : "";
// Each settle gets its own shell, so a busy one is never typed into.
assert.notEqual(observedTerminalId, settleTerminalId);
const token = /__T3_SETUP_DONE___(\w+):/.exec(write.mock.calls.at(-1)?.[0].data ?? "")?.[1];
const settleListener = listeners.at(-1)!;
const completion = yield* Effect.forkChild(
observedSettle.status === "started" && observedSettle.completion
? observedSettle.completion
: Effect.die("no completion"),
);
yield* settleListener({
type: "output",
threadId: "thread-1",
terminalId: observedTerminalId,
data: `\r\n__T3_SETUP_DONE___${token}:0\r\n`,
});
yield* Effect.yieldNow;
assert.equal(closeIdle.mock.calls.length, 0);
yield* settleListener({
type: "output",
threadId: "thread-1",
terminalId: observedTerminalId,
data: "$ ",
});
assert.deepEqual((yield* Fiber.join(completion)).exitCode, 0);
assert.deepEqual(closeIdle.mock.calls[0]?.[0], {
threadId: "thread-1",
terminalId: observedTerminalId,
});
}).pipe(Effect.provide(layer));
});
Loading
Loading