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
35 changes: 35 additions & 0 deletions apps/server/src/orchestration-v2/ShellStream.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,8 @@ import type {
OrchestrationV2ShellSnapshot,
OrchestrationV2StoredEvent,
OrchestrationV2ThreadShell,
OrchestrationProjectShell,
RepositoryIdentity,
} from "@t3tools/contracts";
import { ProjectId, ThreadId } from "@t3tools/contracts";
import * as Effect from "effect/Effect";
Expand All @@ -16,6 +18,7 @@ import {
coalesceStoredThreadEvents,
composeShellStreamWithEnrichment,
dedupeShellEnrichment,
projectsWithResolvedRepositoryIdentities,
shellStreamItemFromEnrichmentRefresh,
shellStreamItemFromThreadShell,
shellStreamItemsFromInitialSnapshot,
Expand Down Expand Up @@ -211,6 +214,38 @@ describe("archivedShellStreamItemFromThreadShell", () => {
});
});

describe("projectsWithResolvedRepositoryIdentities", () => {
const identity = (canonicalKey: string) =>
({ canonicalKey, locator: { source: "git-remote" } }) as unknown as RepositoryIdentity;
const projects = [
{ id: "project-a", workspaceRoot: "/workspace/a", repositoryIdentity: null },
{ id: "project-b", workspaceRoot: "/workspace/b", repositoryIdentity: identity("stale-b") },
{ id: "project-c", workspaceRoot: "/workspace/c", repositoryIdentity: identity("kept-c") },
] as unknown as ReadonlyArray<OrchestrationProjectShell>;

it("refreshes only changed roots with the identity each resolution carried", () => {
const refreshed = projectsWithResolvedRepositoryIdentities(projects, [
{ workspaceRoot: "/workspace/a", enrichment: { repositoryIdentity: identity("first-a") } },
{ workspaceRoot: "/workspace/b", enrichment: { repositoryIdentity: null } },
{ workspaceRoot: "/workspace/a", enrichment: { repositoryIdentity: identity("latest-a") } },
{ workspaceRoot: "/workspace/unknown", enrichment: { repositoryIdentity: null } },
]);

expect(refreshed.map((project) => [project.id, project.repositoryIdentity])).toEqual([
["project-a", identity("latest-a")],
["project-b", null],
]);
});

it("emits no projects when no stored project matches a change", () => {
expect(
projectsWithResolvedRepositoryIdentities(projects, [
{ workspaceRoot: "/workspace/other", enrichment: { repositoryIdentity: null } },
]),
).toEqual([]);
});
});

describe("shellStreamItemFromEnrichmentRefresh", () => {
it("batches nearby completion roots onto one snapshot item", () => {
expect(
Expand Down
25 changes: 25 additions & 0 deletions apps/server/src/orchestration-v2/ShellStream.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ import type {
OrchestrationV2ThreadShellSnapshot,
OrchestrationV2ShellStreamItem,
OrchestrationV2StoredEvent,
RepositoryIdentity,
} from "@t3tools/contracts";
import { OrchestrationProjectShell as ProjectShellSchema } from "@t3tools/contracts";
import * as Schema from "effect/Schema";
Expand Down Expand Up @@ -124,6 +125,30 @@ export function dedupeShellEnrichment<E, R>(
});
}

/**
* Projects whose repository identity just resolved, carrying the identity from
* the resolution itself. Re-enriching the projects here would re-request every
* expired root, whose resolution publishes again, so one expiry would keep
* every shell subscriber reloading every project's metadata once a minute.
* The latest change for a root wins.
*/
export function projectsWithResolvedRepositoryIdentities(
projects: ReadonlyArray<OrchestrationProjectShell>,
changes: ReadonlyArray<{
readonly workspaceRoot: string;
readonly enrichment: { readonly repositoryIdentity: RepositoryIdentity | null };
}>,
): ReadonlyArray<OrchestrationProjectShell> {
const identities = new Map(
changes.map((change) => [change.workspaceRoot, change.enrichment.repositoryIdentity]),
);
return projects.flatMap((project) =>
identities.has(project.workspaceRoot)
? [{ ...project, repositoryIdentity: identities.get(project.workspaceRoot) ?? null }]
: [],
);
}

/** Build a shell snapshot stream item for a batched enrichment completion. */
export function shellStreamItemFromEnrichmentRefresh(input: {
readonly snapshot: OrchestrationV2ShellSnapshot;
Expand Down
41 changes: 21 additions & 20 deletions apps/server/src/process/externalLauncher.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ import * as NodeServices from "@effect/platform-node/NodeServices";
import { assert, it } from "@effect/vitest";
import * as Clock from "effect/Clock";
import * as ConfigProvider from "effect/ConfigProvider";
import * as Deferred from "effect/Deferred";
import * as Effect from "effect/Effect";
import * as Fiber from "effect/Fiber";
import * as FileSystem from "effect/FileSystem";
Expand Down Expand Up @@ -1268,26 +1269,24 @@ it.effect("memoizes editor discovery and refreshes after the cache window", () =
);
});

// A client that disconnects mid-scan interrupts the shared discovery effect on
// the connection fiber. The cache must not retain that interrupt: doing so
// replayed it to every later connect for the whole TTL, so `server.getConfig`
// failed and no client could reconnect until the server restarted.
it.effect("rescans after an interrupted discovery instead of caching the interrupt", () => {
// Connects run discovery under a timeout and may disconnect mid-scan. Neither
// may cancel the scan: on a busy host every connect would time out partway
// through, cache nothing, and leave every client without editors.
it.effect("keeps scanning after the caller is interrupted and shares that scan", () => {
const fileInfo = { type: "File" } as FileSystem.File.Info;
let blockFirstScan = true;
let scans = 0;
const release = Deferred.makeUnsafe<void>();
let parkedStats = 0;
const launcherLayer = ExternalLauncher.layer.pipe(
Layer.provide(
Layer.mergeAll(
FileSystem.layerNoop({
// The first scan parks inside `stat` so the interrupt lands while
// discovery is in flight, which is what a client disconnecting
// mid-connect does to the shared effect.
// Scans park inside `stat` until released, so the interrupt lands
// while discovery is in flight.
stat: () =>
Effect.gen(function* () {
scans += 1;
if (blockFirstScan) {
return yield* Effect.never;
if (!Deferred.isDoneUnsafe(release)) {
parkedStats += 1;
yield* Deferred.await(release);
}
return fileInfo;
}),
Expand All @@ -1304,16 +1303,18 @@ it.effect("rescans after an interrupted discovery instead of caching the interru
return Effect.gen(function* () {
const launcher = yield* ExternalLauncher.ExternalLauncher;

const fiber = yield* Effect.forkChild(launcher.resolveAvailableEditors());
const interrupted = yield* Effect.forkChild(launcher.resolveAvailableEditors());
yield* Effect.yieldNow;
yield* Fiber.interrupt(fiber);
yield* Fiber.interrupt(interrupted);

// The next connect must still get a real answer well inside the TTL.
blockFirstScan = false;
scans = 0;
const editors = yield* launcher.resolveAvailableEditors();
// The next connect joins the running scan instead of starting its own.
const next = yield* Effect.forkChild(launcher.resolveAvailableEditors());
yield* Effect.yieldNow;
assert.equal(parkedStats, 1);

yield* Deferred.succeed(release, undefined);
const editors = yield* Fiber.join(next);
assert.equal(editors.includes("vscode"), true);
assert.isAbove(scans, 0);
}).pipe(
Effect.provide(
Layer.mergeAll(
Expand Down
79 changes: 56 additions & 23 deletions apps/server/src/process/externalLauncher.ts
Original file line number Diff line number Diff line change
Expand Up @@ -28,13 +28,16 @@ import {
import * as Clock from "effect/Clock";
import * as Config from "effect/Config";
import * as Context from "effect/Context";
import * as Deferred from "effect/Deferred";
import * as Effect from "effect/Effect";
import * as Encoding from "effect/Encoding";
import * as Exit from "effect/Exit";
import * as FileSystem from "effect/FileSystem";
import * as Layer from "effect/Layer";
import * as Option from "effect/Option";
import * as Path from "effect/Path";
import * as Ref from "effect/Ref";
import * as Scope from "effect/Scope";
import * as ChildProcess from "effect/unstable/process/ChildProcess";
import * as ChildProcessSpawner from "effect/unstable/process/ChildProcessSpawner";

Expand Down Expand Up @@ -481,21 +484,21 @@ const resolveFileManagerRevealKind = Effect.fn("externalLauncher.resolveFileMana
// the discovered set for a bounded window so repeat connects skip even the
// per-command cache lookups in @t3tools/shared/shell.
//
// This deliberately does not use `Effect.cachedWithTTL`: that memoizes the
// first caller's Exit whatever it is, including an interrupt. Callers run this
// on the connection fiber under a timeout (`resolveAvailableEditorsForConfig`),
// so one client disconnecting mid-scan would cache the interrupt and replay it
// to every later connect for the whole TTL, breaking `server.getConfig`
// permanently. Storing only on success means an interrupted scan leaves the
// cache untouched and the next connect simply rescans.
// The scan runs on its own fiber in the service scope, and every caller awaits
// that one scan. Callers apply a timeout (`resolveAvailableEditorsForConfig`)
// and disconnect mid-connect; neither may cancel a scan other connects are
// waiting on, or throw away work a slow host (a busy server at startup, a
// long PATH) needs more than one connect to finish. A failed scan clears the
// entry so the next caller starts over rather than replaying the failure.
// Expiry uses the monotonic clock (Clock.currentTimeNanos), matching the
// command-resolution cache in @t3tools/shared/shell, so a backward wall-clock
// adjustment cannot keep an expired entry alive.
const EDITOR_DISCOVERY_CACHE_TTL_NANOS = 60_000_000_000n;

interface EditorDiscoveryCacheEntry {
readonly editors: ReadonlyArray<EditorId>;
readonly expiresAtNanos: bigint;
readonly scan: Deferred.Deferred<ReadonlyArray<EditorId>>;
/** Undefined while the scan is still running. */
readonly expiresAtNanos: bigint | undefined;
}

/**
Expand Down Expand Up @@ -779,27 +782,57 @@ export const make = Effect.gen(function* () {
Effect.provideService(Path.Path, path),
);

const scope = yield* Scope.Scope;
const editorDiscoveryCache = yield* Ref.make<Option.Option<EditorDiscoveryCacheEntry>>(
Option.none(),
);
const cachedAvailableEditors = Effect.gen(function* () {
const nowNanos = yield* Clock.currentTimeNanos;
const entry = yield* Ref.get(editorDiscoveryCache);
if (Option.isSome(entry) && entry.value.expiresAtNanos > nowNanos) {
return entry.value.editors;
}
const editors = yield* provideCommandResolutionServices(resolveAvailableEditors()).pipe(
const runEditorDiscovery = (scan: Deferred.Deferred<ReadonlyArray<EditorId>>) =>
provideCommandResolutionServices(resolveAvailableEditors()).pipe(
Effect.provideService(ChildProcessSpawner.ChildProcessSpawner, spawner),
Effect.onExit((exit) =>
Effect.gen(function* () {
const expiresAtNanos = (yield* Clock.currentTimeNanos) + EDITOR_DISCOVERY_CACHE_TTL_NANOS;
yield* Ref.update(editorDiscoveryCache, (current) =>
Option.isNone(current) || current.value.scan !== scan
? current
: Exit.isSuccess(exit)
? Option.some({ scan, expiresAtNanos })
: Option.none(),
);
yield* Deferred.done(scan, exit);
}),
),
Effect.interruptible,
Effect.forkIn(scope),
);
yield* Ref.set(
// Claiming the cache entry and starting its scan must not be split by an
// interrupt, or the entry would wait on a scan that never runs.
const acquireEditorDiscovery = Effect.gen(function* () {
const nowNanos = yield* Clock.currentTimeNanos;
const [scan, isNewScan] = yield* Ref.modify(
editorDiscoveryCache,
Option.some({
editors,
expiresAtNanos: nowNanos + EDITOR_DISCOVERY_CACHE_TTL_NANOS,
}),
(
current,
): [
[EditorDiscoveryCacheEntry["scan"], boolean],
Option.Option<EditorDiscoveryCacheEntry>,
] => {
if (
Option.isSome(current) &&
(current.value.expiresAtNanos === undefined || current.value.expiresAtNanos > nowNanos)
) {
return [[current.value.scan, false], current];
}
const scan = Deferred.makeUnsafe<ReadonlyArray<EditorId>>();
return [[scan, true], Option.some({ scan, expiresAtNanos: undefined })];
},
);
return editors;
});
if (isNewScan) {
yield* runEditorDiscovery(scan);
}
return scan;
}).pipe(Effect.uninterruptible);
const cachedAvailableEditors = Effect.flatMap(acquireEditorDiscovery, Deferred.await);

return ExternalLauncher.of({
resolveAvailableEditors: () => cachedAvailableEditors,
Expand Down
1 change: 0 additions & 1 deletion apps/server/src/provider/Layers/PiProvider.ts
Original file line number Diff line number Diff line change
Expand Up @@ -59,7 +59,6 @@ import {

const PI_PRESENTATION = {
displayName: "Pi",
badgeLabel: "Early Access",
showInteractionModeToggle: false,
supportedRuntimeModes: ["approval-required", "auto-accept-edits", "full-access"],
// The adapter reports context usage from Pi's streaming usage while a
Expand Down
24 changes: 24 additions & 0 deletions apps/server/src/relay/AgentAwarenessRelay.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -715,6 +715,30 @@ describe("AgentAwarenessRelay", () => {
}),
);

it.effect.each([
{ label: "live", archived: false },
{ label: "archived", archived: true },
])("never publishes tombstones for $label subagent threads", ({ archived }) =>
Effect.gen(function* () {
const { relay, currentShell, publications } = yield* makeTestRelay();
yield* Ref.set(
currentShell,
shell({
lineage: {
rootThreadId: THREAD_ID,
parentThreadId: THREAD_ID,
relationshipToParent: "subagent",
},
...(archived ? { archivedAt: yield* DateTime.now } : {}),
}),
);
yield* relay.publishThread(THREAD_ID);
yield* TestClock.adjust("5 seconds");
yield* relay.drain;
assert.equal(publications.length, 0);
}),
);

it.effect("confirms a first completed state and respects disabling during confirmation", () =>
Effect.gen(function* () {
const { relay, secrets, currentShell, publications } = yield* makeTestRelay();
Expand Down
10 changes: 10 additions & 0 deletions apps/server/src/relay/AgentAwarenessRelay.ts
Original file line number Diff line number Diff line change
Expand Up @@ -551,6 +551,16 @@ export const make = Effect.gen(function* () {
// domain event, so materializing the full shell here would make the cost
// of one thread's activity proportional to how many threads exist.
const threadShell = yield* threads.getThreadShell(threadId);
if (
threadShell?.lineage.relationshipToParent === "subagent" &&
!(yield* Ref.get(publishedStateByThreadRef)).has(threadId)
) {
// Subagents never project activity, so the relay holds no row to clear.
// Their events would otherwise publish a tombstone each, and every
// publish re-delivers the user's aggregate. Checked before the archive
// filter so archiving one stays quiet too.
return;
}
const thread =
threadShell === null || threadShell.archivedAt !== null
? Option.none<OrchestrationV2ThreadShell>()
Expand Down
Loading
Loading