From 3e0e266ee72201a35a5104a4204d35ea10c184e2 Mon Sep 17 00:00:00 2001 From: Luiz Ferraz Date: Sun, 13 Sep 2026 16:24:29 +0000 Subject: [PATCH 1/3] feat(server): add t3 wake command to send a message to an existing threa - Sends a message to a running server or desktop app and requests a turn, using the thread's existing provider, model, and modes - Extracts shared running-server discovery from t3 pair into cli/runningServer.ts - Never starts a server; fails cleanly when none is running --- apps/server/src/bin.test.ts | 173 +++++++++++++++++++++- apps/server/src/bin.ts | 2 + apps/server/src/cli/pair.ts | 214 ++------------------------- apps/server/src/cli/runningServer.ts | 190 ++++++++++++++++++++++++ apps/server/src/cli/wake.ts | 136 +++++++++++++++++ docs/user/composer.md | 17 +++ 6 files changed, 530 insertions(+), 202 deletions(-) create mode 100644 apps/server/src/cli/runningServer.ts create mode 100644 apps/server/src/cli/wake.ts diff --git a/apps/server/src/bin.test.ts b/apps/server/src/bin.test.ts index e1a13d4ce8e7..31c6bee80346 100644 --- a/apps/server/src/bin.test.ts +++ b/apps/server/src/bin.test.ts @@ -19,8 +19,10 @@ import { assert, it } from "@effect/vitest"; import * as Effect from "effect/Effect"; import * as DateTime from "effect/DateTime"; import * as Layer from "effect/Layer"; +import * as Stream from "effect/Stream"; import * as HttpRouter from "effect/unstable/http/HttpRouter"; import * as HttpServer from "effect/unstable/http/HttpServer"; +import * as HttpServerResponse from "effect/unstable/http/HttpServerResponse"; import * as HttpApi from "effect/unstable/httpapi/HttpApi"; import * as HttpApiBuilder from "effect/unstable/httpapi/HttpApiBuilder"; import * as CliError from "effect/unstable/cli/CliError"; @@ -357,14 +359,29 @@ it.layer(NodeServices.layer)("project lookup with unavailable workspaces", (it) ); }); -const withLiveProjectCliServer = (baseDir: string, run: () => Effect.Effect) => +const withLiveProjectCliServer = ( + baseDir: string, + run: () => Effect.Effect, + mode: "web" | "desktop" = "web", +) => Effect.gen(function* () { - const config = yield* makeCliTestServerConfig(baseDir); + const config = { ...(yield* makeCliTestServerConfig(baseDir)), mode }; const routesLayer = HttpApiBuilder.layer(ProjectCliHttpApi).pipe( Layer.provide(orchestrationHttpApiLayer), Layer.provide(environmentAuthenticatedAuthLayer), ); - const appLayer = HttpRouter.serve(routesLayer, { + const descriptorLayer = HttpRouter.add( + "GET", + "/.well-known/t3/environment", + HttpServerResponse.json({ + environmentId: "cli-test-environment", + label: "CLI test", + platform: { os: "linux", arch: "x64" }, + serverVersion: "0.0.1", + capabilities: { repositoryIdentity: true }, + }), + ); + const appLayer = HttpRouter.serve(Layer.merge(routesLayer, descriptorLayer), { disableListenLog: true, disableLogger: true, }).pipe( @@ -405,6 +422,156 @@ const withLiveProjectCliServer = (baseDir: string, run: () => Effect.Ef ); }); +it.layer(NodeServices.layer)("t3 wake", (it) => { + it.effect.each(["web", "desktop"] as const)( + "sends one message and requests a turn through a running %s server", + (mode) => + Effect.gen(function* () { + const { baseDir } = yield* makeProjectLookupFixture(true, false); + const threadId = ThreadId.make("thread-project-lookup"); + yield* withLiveProjectCliServer( + baseDir, + () => + Effect.gen(function* () { + const engine = yield* OrchestrationEngine.OrchestrationEngineService; + const createdAt = DateTime.formatIso(yield* DateTime.now); + yield* engine.dispatch({ + type: "thread.runtime-mode.set", + commandId: CommandId.make("wake-runtime"), + threadId, + runtimeMode: "full-access", + createdAt, + }); + yield* engine.dispatch({ + type: "thread.interaction-mode.set", + commandId: CommandId.make("wake-interaction"), + threadId, + interactionMode: "plan", + createdAt, + }); + yield* engine.dispatch({ + type: "thread.settle", + commandId: CommandId.make("wake-settle"), + threadId, + }); + yield* engine.dispatch({ + type: "thread.snooze", + commandId: CommandId.make("wake-snooze"), + threadId, + snoozedUntil: "2099-01-01T00:00:00.000Z", + }); + const beforeSequence = yield* engine.latestSequence; + const message = " Continue the task.\nKeep the existing settings. "; + const { output } = yield* captureStdout( + runCli(["wake", threadId, message, "--base-dir", baseDir]), + ); + assert.include(output, `Sent message to thread ${threadId}`); + const events = yield* engine.readEvents(beforeSequence).pipe(Stream.runCollect); + const messages = events.filter((event) => event.type === "thread.message-sent"); + assert.equal(messages.length, 1); + assert.equal(messages[0]?.payload.text, message); + const starts = events.filter((event) => event.type === "thread.turn-start-requested"); + assert.equal(starts.length, 1); + assert.equal(starts[0]?.payload.messageId, messages[0]?.payload.messageId); + assert.equal(starts[0]?.payload.runtimeMode, "full-access"); + assert.equal(starts[0]?.payload.interactionMode, "plan"); + const query = yield* ProjectionSnapshotQuery.ProjectionSnapshotQuery; + const snapshot = yield* query.getSnapshot(); + const thread = snapshot.threads.find((entry) => entry.id === threadId)!; + assert.deepEqual(thread.modelSelection, { + instanceId: ProviderInstanceId.make("codex"), + model: "gpt-5-codex", + }); + assert.isNull(thread.settledOverride); + assert.isNull(thread.snoozedUntil); + const auth = yield* EnvironmentAuth.EnvironmentAuth; + const sessions = yield* auth.listSessions(); + assert.equal(sessions.length, 0); + }), + mode, + ); + }), + ); + + it.effect("fails without creating state when no server is running", () => + Effect.gen(function* () { + const baseDir = NodeFS.mkdtempSync(NodePath.join(NodeOS.tmpdir(), "t3-wake-none-")); + const error = yield* runCliWithRuntime([ + "wake", + "thread-missing", + "Continue", + "--base-dir", + baseDir, + ]).pipe(Effect.flip); + assert.include(error.message, "No running T3 Code server found."); + assert.deepEqual(NodeFS.readdirSync(baseDir), []); + }), + ); + + it.effect.each(["", " ", "x".repeat(120_001)])( + "rejects an invalid message before accessing server state", + (message) => + Effect.gen(function* () { + const baseDir = NodeFS.mkdtempSync(NodePath.join(NodeOS.tmpdir(), "t3-wake-invalid-")); + const error = yield* runCliWithRuntime([ + "wake", + "thread-missing", + message, + "--base-dir", + baseDir, + ]).pipe(Effect.flip); + assert.include(error.message, "Message cannot"); + assert.deepEqual(NodeFS.readdirSync(baseDir), []); + }), + ); + + it.effect("fails on stale runtime state without changing the persisted thread", () => + Effect.gen(function* () { + const { baseDir } = yield* makeProjectLookupFixture(true, false); + yield* withLiveProjectCliServer(baseDir, () => Effect.void); + const before = yield* readPersistedSnapshot(baseDir); + const error = yield* runCliWithRuntime([ + "wake", + "thread-project-lookup", + "Continue", + "--base-dir", + baseDir, + ]).pipe(Effect.flip); + assert.include(error.message, "No running T3 Code server found."); + assert.deepEqual(yield* readPersistedSnapshot(baseDir), before); + }), + ); + + it.effect.each([false, true])("rejects missing or deleted threads (deleted=%s)", (deleted) => + Effect.gen(function* () { + const { baseDir } = yield* makeProjectLookupFixture(true, false); + yield* withLiveProjectCliServer(baseDir, () => + Effect.gen(function* () { + const engine = yield* OrchestrationEngine.OrchestrationEngineService; + const threadId = ThreadId.make(deleted ? "thread-project-lookup" : "missing-thread"); + if (deleted) { + yield* engine.dispatch({ + type: "thread.delete", + commandId: CommandId.make("wake-delete"), + threadId, + }); + } + const beforeSequence = yield* engine.latestSequence; + const error = yield* runCliWithRuntime([ + "wake", + threadId, + "Continue", + "--base-dir", + baseDir, + ]).pipe(Effect.flip); + assert.include(error.message, threadId); + assert.equal(yield* engine.latestSequence, beforeSequence); + }), + ); + }), + ); +}); + it.layer(NodeServices.layer)("bin cli parsing", (it) => { it.effect("accepts the built-in lowercase log-level flag values", () => Effect.gen(function* () { diff --git a/apps/server/src/bin.ts b/apps/server/src/bin.ts index 52cc363ed04f..d488cf70c1f4 100644 --- a/apps/server/src/bin.ts +++ b/apps/server/src/bin.ts @@ -20,6 +20,7 @@ import { serviceCommand } from "./cli/service.ts"; import { servicePreflightCommand } from "./cli/servicePreflight.ts"; import { themeCommand } from "./cli/theme.ts"; import { triageCommand } from "./cli/triage.ts"; +import { wakeCommand } from "./cli/wake.ts"; const CliRuntimeLayer = Layer.mergeAll(NodeServices.layer, NetService.layer); @@ -62,6 +63,7 @@ export const makeCli = ({ cloudEnabled = hasCloudPublicConfig } = {}) => servicePreflightCommand, themeCommand, triageCommand, + wakeCommand, cloudEnabled ? connectCommand : connectUnavailableCommand, ]), ); diff --git a/apps/server/src/cli/pair.ts b/apps/server/src/cli/pair.ts index 7fd376c6f881..67bd4bbe5a51 100644 --- a/apps/server/src/cli/pair.ts +++ b/apps/server/src/cli/pair.ts @@ -9,19 +9,13 @@ * shared T3 home. `--tailscale` publishes the server over Tailscale Serve * HTTPS and pairs through the tailnet URL instead. */ -import { - AuthStandardClientScopes, - ExecutionEnvironmentDescriptor, - PortSchema, -} from "@t3tools/contracts"; -import { resolveWorktreeT3Home } from "@t3tools/shared/devHome"; +import { AuthStandardClientScopes, PortSchema } from "@t3tools/contracts"; import { buildTailscaleHttpsBaseUrl, DEFAULT_TAILSCALE_SERVE_PORT, ensureTailscaleServe, readTailscaleStatus, } from "@t3tools/tailscale"; -import * as Config from "effect/Config"; import * as Console from "effect/Console"; import * as DateTime from "effect/DateTime"; import * as Duration from "effect/Duration"; @@ -31,21 +25,11 @@ import * as Option from "effect/Option"; import * as References from "effect/References"; import * as Schema from "effect/Schema"; import { Command, Flag, GlobalFlag } from "effect/unstable/cli"; -import { - FetchHttpClient, - HttpClient, - HttpClientRequest, - HttpClientResponse, -} from "effect/unstable/http"; +import { FetchHttpClient } from "effect/unstable/http"; import * as EnvironmentAuth from "../auth/EnvironmentAuth.ts"; import * as ServerConfig from "../config.ts"; -import { resolveBaseDir } from "../os-jank.ts"; -import { - type PersistedServerRuntimeState, - isProcessAlive, - readPersistedServerRuntimeState, -} from "../serverRuntimeState.ts"; +import type { PersistedServerRuntimeState } from "../serverRuntimeState.ts"; import { buildPairingUrl, formatHostForUrl, @@ -55,35 +39,18 @@ import { resolveHeadlessConnectionString, } from "../startupAccess.ts"; import { baseDirFlag, DurationFromString } from "./config.ts"; - -const WELL_KNOWN_ENVIRONMENT_PATH = "/.well-known/t3/environment"; -const PAIR_PROBE_TIMEOUT = Duration.millis(2_500); -// Tailscale provisions an HTTPS certificate on the first request to a fresh -// serve mapping, which can take a few seconds. +import { + type DiscoveredServer, + type EnvironmentProbeResult, + discoverRunningServer, + makeDiscoveredServerConfig, + probeEnvironmentDescriptor, +} from "./runningServer.ts"; + +// Tailscale can take a few seconds to provision its first HTTPS certificate. const TAILSCALE_PROBE_ATTEMPTS = 5; const TAILSCALE_PROBE_RETRY_DELAY = Duration.seconds(1); -export type PairStateVariant = "userdata" | "dev"; - -// deriveServerPaths only checks devUrl for undefined-ness when picking the -// dev-vs-userdata state directory; the value itself is not used. -const DEV_VARIANT_PLACEHOLDER_URL = new URL("http://localhost"); - -export class NoRunningServerError extends Schema.TaggedError()( - "NoRunningServerError", - { - checkedStatePaths: Schema.Array(Schema.String), - }, -) { - override get message(): string { - return [ - "No running T3 Code server found.", - ...this.checkedStatePaths.map((statePath) => ` checked ${statePath}`), - "Start one with `npx t3 serve`, or connect this machine with T3 Connect: `npx t3 connect`.", - ].join("\n"); - } -} - // Each tailscale failure gets its own class (same reasoning as // scripts/lib/dev-share.ts): distinct caller-visible message, distinct remedy. export class TailscaleUnavailableError extends Schema.TaggedError()( @@ -193,157 +160,6 @@ const formatPairOutput = (input: { "", ].join("\n"); -/** - * Three outcomes, because they drive different decisions: a T3 descriptor - * (pair with it), nothing answering (safe to configure Tailscale Serve), or - * something answering that is not a T3 server (do NOT overwrite its mapping). - */ -type EnvironmentProbeResult = - | { readonly _tag: "descriptor"; readonly descriptor: ExecutionEnvironmentDescriptor } - | { readonly _tag: "unreachable" } - | { readonly _tag: "not-a-t3-server" }; - -const probeEnvironmentDescriptor = ( - baseUrl: string, -): Effect.Effect => - Effect.gen(function* () { - const client = yield* HttpClient.HttpClient; - const request = HttpClientRequest.get(new URL(WELL_KNOWN_ENVIRONMENT_PATH, baseUrl).toString()); - const response = yield* client.execute(request).pipe( - Effect.timeout(PAIR_PROBE_TIMEOUT), - // Transport failure or timeout: nothing (reachable) is listening there. - Effect.mapError(() => ({ _tag: "unreachable" }) as const), - ); - // Bad-gateway family means a proxy (Tailscale Serve) answered for a - // backend that is gone — a stale mapping, not a live occupant. Treating - // it as unreachable lets `t3 pair --tailscale` repair its own mapping - // after the server's port changed. - if (response.status === 502 || response.status === 503 || response.status === 504) { - return { _tag: "unreachable" } as const; - } - // Anything else that answered HTTP but not with a valid descriptor is - // some other service. - const descriptor = yield* HttpClientResponse.filterStatusOk(response).pipe( - Effect.flatMap(HttpClientResponse.schemaBodyJson(ExecutionEnvironmentDescriptor)), - Effect.mapError(() => ({ _tag: "not-a-t3-server" }) as const), - ); - return { _tag: "descriptor", descriptor } as const; - }).pipe(Effect.catch((outcome) => Effect.succeed(outcome))); - -interface DiscoveredPairTarget { - readonly baseDir: string; - readonly variant: PairStateVariant; - readonly state: PersistedServerRuntimeState; - readonly descriptor: ExecutionEnvironmentDescriptor; -} - -const discoverPairTarget = Effect.fn("pair.discoverPairTarget")(function* ( - explicitBaseDir: string | undefined, -) { - const bases: Array = []; - if (explicitBaseDir !== undefined && explicitBaseDir.trim().length > 0) { - bases.push(yield* resolveBaseDir(explicitBaseDir)); - } else { - // Same precedence as dev-runner: inside a linked worktree its own `.t3` - // outranks the shared home, so `t3 pair` in a worktree pairs with the dev - // server under test rather than the daily-driver install. - const worktreeHome = yield* resolveWorktreeT3Home(process.cwd()); - if (worktreeHome !== undefined) { - bases.push(worktreeHome); - } - const envHome = yield* Config.string("T3CODE_HOME").pipe(Config.option); - bases.push(yield* resolveBaseDir(Option.getOrUndefined(envHome))); - } - - const checkedStatePaths: Array = []; - for (const baseDir of new Set(bases)) { - for (const variant of ["userdata", "dev"] as const) { - const derivedPaths = yield* ServerConfig.deriveServerPaths( - baseDir, - variant === "dev" ? DEV_VARIANT_PLACEHOLDER_URL : undefined, - {}, - ); - const statePath = derivedPaths.serverRuntimeStatePath; - checkedStatePaths.push(statePath); - const state = yield* readPersistedServerRuntimeState(statePath); - if (Option.isNone(state)) { - continue; - } - // The pid check guards against a dead server's state file whose port - // was since reused by a different server: pairing would then mint a - // token in the old database while the QR code points at the new server. - if (!isProcessAlive(state.value.pid)) { - continue; - } - const probed = yield* probeEnvironmentDescriptor(state.value.origin); - if (probed._tag !== "descriptor") { - continue; - } - return { - baseDir, - variant, - state: state.value, - descriptor: probed.descriptor, - } satisfies DiscoveredPairTarget; - } - } - return yield* new NoRunningServerError({ checkedStatePaths }); -}); - -/** - * Server config pointed at the discovered server's state directory, so the - * minted token lands in the database the running server reads from. Built by - * hand rather than through `resolveServerConfig` to keep the dev-vs-userdata - * choice pinned to where the runtime state was actually found, independent of - * ambient environment variables. - */ -const makePairServerConfig = Effect.fn(function* (input: { - readonly target: DiscoveredPairTarget; - readonly logLevel: ServerConfig.ServerConfig["Service"]["logLevel"]; -}) { - const { baseDir, variant, state } = input.target; - // The state-dir variant does not imply dev-ness: a worktree dev server uses - // an explicit home and therefore lands in `userdata`. The recorded devUrl is - // what actually marks a dev server. - const devUrl = state.devUrl !== undefined ? new URL(state.devUrl) : undefined; - const derivedPaths = yield* ServerConfig.deriveServerPaths( - baseDir, - variant === "dev" ? DEV_VARIANT_PLACEHOLDER_URL : undefined, - {}, - ); - return ServerConfig.make({ - logLevel: input.logLevel, - traceMinLevel: "Info", - traceTimingEnabled: false, - traceBatchWindowMs: 1_000, - traceMaxBytes: 10 * 1024 * 1024, - traceMaxFiles: 10, - otlpTracesUrl: undefined, - otlpMetricsUrl: undefined, - otlpExportIntervalMs: 10_000, - otlpServiceName: "t3-server", - mode: "web", - port: state.port, - host: state.host, - cwd: process.cwd(), - baseDir, - ...derivedPaths, - staticDir: undefined, - devUrl, - devAllowedOrigins: [], - noBrowser: true, - startupPresentation: "headless", - desktopBootstrapToken: undefined, - desktopTelemetryFd: undefined, - desktopTelemetryControlFd: undefined, - resourceMonitorPath: undefined, - autoBootstrapProjectFromCwd: false, - logWebSocketEvents: false, - tailscaleServeEnabled: false, - tailscaleServePort: DEFAULT_TAILSCALE_SERVE_PORT, - }); -}); - const awaitEnvironmentDescriptor = Effect.fn(function* (baseUrl: string) { let last: EnvironmentProbeResult = { _tag: "unreachable" }; for (let attempt = 0; attempt < TAILSCALE_PROBE_ATTEMPTS; attempt += 1) { @@ -357,7 +173,7 @@ const awaitEnvironmentDescriptor = Effect.fn(function* (baseUrl: string) { }); const resolveTailscalePairingBase = Effect.fn("pair.resolveTailscalePairingBase")( - function* (input: { readonly target: DiscoveredPairTarget; readonly servePort: number }) { + function* (input: { readonly target: DiscoveredServer; readonly servePort: number }) { const notes: Array = []; const status = yield* readTailscaleStatus.pipe( Effect.mapError((cause) => new TailscaleUnavailableError({ cause })), @@ -488,7 +304,7 @@ export const pairCommand = Command.make("pair", { // an explicit --log-level still wins. const logLevel = Option.getOrElse(cliLogLevel, () => "Warn" as const); - const target = yield* discoverPairTarget(Option.getOrUndefined(flags.baseDir)); + const target = yield* discoverRunningServer(Option.getOrUndefined(flags.baseDir)); const notes: Array = []; let pairingBaseUrl: string; @@ -513,7 +329,7 @@ export const pairCommand = Command.make("pair", { } } - const config = yield* makePairServerConfig({ target, logLevel }); + const config = yield* makeDiscoveredServerConfig({ target, logLevel }); const issued = yield* mintPairingLink({ config, ttl: flags.ttl, label: flags.label }); const pairingUrl = buildPairingUrl(pairingBaseUrl, issued.credential); diff --git a/apps/server/src/cli/runningServer.ts b/apps/server/src/cli/runningServer.ts new file mode 100644 index 000000000000..3a2743c3bba7 --- /dev/null +++ b/apps/server/src/cli/runningServer.ts @@ -0,0 +1,190 @@ +import { ExecutionEnvironmentDescriptor } from "@t3tools/contracts"; +import { resolveWorktreeT3Home } from "@t3tools/shared/devHome"; +import { DEFAULT_TAILSCALE_SERVE_PORT } from "@t3tools/tailscale"; +import * as Config from "effect/Config"; +import * as Duration from "effect/Duration"; +import * as Effect from "effect/Effect"; +import * as Option from "effect/Option"; +import * as Schema from "effect/Schema"; +import { HttpClient, HttpClientRequest, HttpClientResponse } from "effect/unstable/http"; + +import * as ServerConfig from "../config.ts"; +import { resolveBaseDir } from "../os-jank.ts"; +import { + type PersistedServerRuntimeState, + isProcessAlive, + readPersistedServerRuntimeState, +} from "../serverRuntimeState.ts"; + +const WELL_KNOWN_ENVIRONMENT_PATH = "/.well-known/t3/environment"; +const SERVER_PROBE_TIMEOUT = Duration.millis(2_500); +type ServerStateVariant = "userdata" | "dev"; + +// deriveServerPaths only checks devUrl for undefined-ness when picking the +// dev-vs-userdata state directory; the value itself is not used. +const DEV_VARIANT_PLACEHOLDER_URL = new URL("http://localhost"); + +export class NoRunningServerError extends Schema.TaggedError()( + "NoRunningServerError", + { + checkedStatePaths: Schema.Array(Schema.String), + }, +) { + override get message(): string { + return [ + "No running T3 Code server found.", + ...this.checkedStatePaths.map((statePath) => ` checked ${statePath}`), + "Open the T3 Code desktop app, start a server with `npx t3 serve`, or connect this machine with T3 Connect: `npx t3 connect`.", + ].join("\n"); + } +} + +/** + * Three outcomes, because they drive different decisions: a T3 descriptor + * (pair with it), nothing answering (safe to configure Tailscale Serve), or + * something answering that is not a T3 server (do NOT overwrite its mapping). + */ +export type EnvironmentProbeResult = + | { readonly _tag: "descriptor"; readonly descriptor: ExecutionEnvironmentDescriptor } + | { readonly _tag: "unreachable" } + | { readonly _tag: "not-a-t3-server" }; + +export const probeEnvironmentDescriptor = ( + baseUrl: string, +): Effect.Effect => + Effect.gen(function* () { + const client = yield* HttpClient.HttpClient; + const request = HttpClientRequest.get(new URL(WELL_KNOWN_ENVIRONMENT_PATH, baseUrl).toString()); + const response = yield* client.execute(request).pipe( + Effect.timeout(SERVER_PROBE_TIMEOUT), + // Transport failure or timeout: nothing (reachable) is listening there. + Effect.mapError(() => ({ _tag: "unreachable" }) as const), + ); + // Bad-gateway family means a proxy (Tailscale Serve) answered for a + // backend that is gone — a stale mapping, not a live occupant. Treating + // it as unreachable lets `t3 pair --tailscale` repair its own mapping + // after the server's port changed. + if (response.status === 502 || response.status === 503 || response.status === 504) { + return { _tag: "unreachable" } as const; + } + // Anything else that answered HTTP but not with a valid descriptor is + // some other service. + const descriptor = yield* HttpClientResponse.filterStatusOk(response).pipe( + Effect.flatMap(HttpClientResponse.schemaBodyJson(ExecutionEnvironmentDescriptor)), + Effect.mapError(() => ({ _tag: "not-a-t3-server" }) as const), + ); + return { _tag: "descriptor", descriptor } as const; + }).pipe(Effect.catch((outcome) => Effect.succeed(outcome))); + +export interface DiscoveredServer { + readonly baseDir: string; + readonly variant: ServerStateVariant; + readonly state: PersistedServerRuntimeState; + readonly descriptor: ExecutionEnvironmentDescriptor; +} + +export const discoverRunningServer = Effect.fn("discoverRunningServer")(function* ( + explicitBaseDir: string | undefined, +) { + const bases: Array = []; + if (explicitBaseDir !== undefined && explicitBaseDir.trim().length > 0) { + bases.push(yield* resolveBaseDir(explicitBaseDir)); + } else { + // Same precedence as dev-runner: inside a linked worktree its own `.t3` + // outranks the shared home, so commands target the worktree's dev server. + const worktreeHome = yield* resolveWorktreeT3Home(process.cwd()); + if (worktreeHome !== undefined) { + bases.push(worktreeHome); + } + const envHome = yield* Config.string("T3CODE_HOME").pipe(Config.option); + bases.push(yield* resolveBaseDir(Option.getOrUndefined(envHome))); + } + + const checkedStatePaths: Array = []; + for (const baseDir of new Set(bases)) { + for (const variant of ["userdata", "dev"] as const) { + const derivedPaths = yield* ServerConfig.deriveServerPaths( + baseDir, + variant === "dev" ? DEV_VARIANT_PLACEHOLDER_URL : undefined, + {}, + ); + const statePath = derivedPaths.serverRuntimeStatePath; + checkedStatePaths.push(statePath); + const state = yield* readPersistedServerRuntimeState(statePath); + if (Option.isNone(state)) { + continue; + } + // The pid check guards against a dead server's state file whose port + // was since reused by a different server: pairing would then mint a + // token in the old database while the QR code points at the new server. + if (!isProcessAlive(state.value.pid)) { + continue; + } + const probed = yield* probeEnvironmentDescriptor(state.value.origin); + if (probed._tag !== "descriptor") { + continue; + } + return { + baseDir, + variant, + state: state.value, + descriptor: probed.descriptor, + } satisfies DiscoveredServer; + } + } + return yield* new NoRunningServerError({ checkedStatePaths }); +}); + +/** + * Server config pointed at the discovered server's state directory, so the + * minted token lands in the database the running server reads from. Built by + * hand rather than through `resolveServerConfig` to keep the dev-vs-userdata + * choice pinned to where the runtime state was actually found, independent of + * ambient environment variables. + */ +export const makeDiscoveredServerConfig = Effect.fn(function* (input: { + readonly target: DiscoveredServer; + readonly logLevel: ServerConfig.ServerConfig["Service"]["logLevel"]; +}) { + const { baseDir, variant, state } = input.target; + // The state-dir variant does not imply dev-ness: a worktree dev server uses + // an explicit home and therefore lands in `userdata`. The recorded devUrl is + // what actually marks a dev server. + const devUrl = state.devUrl !== undefined ? new URL(state.devUrl) : undefined; + const derivedPaths = yield* ServerConfig.deriveServerPaths( + baseDir, + variant === "dev" ? DEV_VARIANT_PLACEHOLDER_URL : undefined, + {}, + ); + return ServerConfig.make({ + logLevel: input.logLevel, + traceMinLevel: "Info", + traceTimingEnabled: false, + traceBatchWindowMs: 1_000, + traceMaxBytes: 10 * 1024 * 1024, + traceMaxFiles: 10, + otlpTracesUrl: undefined, + otlpMetricsUrl: undefined, + otlpExportIntervalMs: 10_000, + otlpServiceName: "t3-server", + mode: "web", + port: state.port, + host: state.host, + cwd: process.cwd(), + baseDir, + ...derivedPaths, + staticDir: undefined, + devUrl, + devAllowedOrigins: [], + noBrowser: true, + startupPresentation: "headless", + desktopBootstrapToken: undefined, + desktopTelemetryFd: undefined, + desktopTelemetryControlFd: undefined, + resourceMonitorPath: undefined, + autoBootstrapProjectFromCwd: false, + logWebSocketEvents: false, + tailscaleServeEnabled: false, + tailscaleServePort: DEFAULT_TAILSCALE_SERVE_PORT, + }); +}); diff --git a/apps/server/src/cli/wake.ts b/apps/server/src/cli/wake.ts new file mode 100644 index 000000000000..94c2ac85dd21 --- /dev/null +++ b/apps/server/src/cli/wake.ts @@ -0,0 +1,136 @@ +import { + AuthOrchestrationOperateScope, + AuthOrchestrationReadScope, + CommandId, + EnvironmentHttpApi, + MessageId, + PROVIDER_SEND_TURN_MAX_INPUT_CHARS, + ThreadId, +} from "@t3tools/contracts"; +import * as Console from "effect/Console"; +import * as Crypto from "effect/Crypto"; +import * as DateTime from "effect/DateTime"; +import * as Effect from "effect/Effect"; +import * as Layer from "effect/Layer"; +import * as Option from "effect/Option"; +import * as References from "effect/References"; +import * as Schema from "effect/Schema"; +import { Argument, Command, GlobalFlag } from "effect/unstable/cli"; +import { FetchHttpClient } from "effect/unstable/http"; +import * as HttpApiClient from "effect/unstable/httpapi/HttpApiClient"; + +import * as EnvironmentAuth from "../auth/EnvironmentAuth.ts"; +import * as ServerConfig from "../config.ts"; +import { baseDirFlag } from "./config.ts"; +import { discoverRunningServer, makeDiscoveredServerConfig } from "./runningServer.ts"; + +class WakeCommandError extends Schema.TaggedError()("WakeCommandError", { + message: Schema.String, + cause: Schema.optional(Schema.Defect()), +}) {} + +export const wakeCommand = Command.make("wake", { + baseDir: baseDirFlag, + threadId: Argument.string("thread-id").pipe( + Argument.withSchema(ThreadId), + Argument.withDescription("ID of the existing thread to wake."), + ), + message: Argument.string("message").pipe( + Argument.withDescription("Message to send to the agent (quote messages containing spaces)."), + ), +}).pipe( + Command.withDescription( + "Send a message and start a turn in an existing thread. Requires a running server or desktop app.", + ), + Command.withHandler( + Effect.fn("wakeCommand")(function* (flags) { + if (flags.message.trim().length === 0) { + return yield* new WakeCommandError({ message: "Message cannot be empty." }); + } + if (flags.message.length > PROVIDER_SEND_TURN_MAX_INPUT_CHARS) { + return yield* new WakeCommandError({ + message: `Message cannot exceed ${PROVIDER_SEND_TURN_MAX_INPUT_CHARS} characters.`, + }); + } + const logLevel = Option.getOrElse(yield* GlobalFlag.LogLevel, () => "Info" as const); + const target = yield* discoverRunningServer(Option.getOrUndefined(flags.baseDir)); + const config = yield* makeDiscoveredServerConfig({ target, logLevel }); + + yield* Effect.gen(function* () { + const auth = yield* EnvironmentAuth.EnvironmentAuth; + const client = yield* HttpApiClient.make(EnvironmentHttpApi, { + baseUrl: target.state.origin, + }); + const crypto = yield* Crypto.Crypto; + yield* Effect.acquireUseRelease( + auth.issueSession({ + scopes: [AuthOrchestrationReadScope, AuthOrchestrationOperateScope], + label: "t3 wake cli", + }), + (session) => + Effect.gen(function* () { + const headers = { authorization: `Bearer ${session.token}` }; + const { thread } = yield* client.orchestration + .threadSnapshot({ + headers, + params: { threadId: flags.threadId }, + payload: { turnLimit: 1 }, + }) + .pipe( + Effect.timeout("5 seconds"), + Effect.mapError( + (cause) => + new WakeCommandError({ + message: `Could not read thread '${flags.threadId}' from the running server. Check the thread ID and that the server is still running.`, + cause, + }), + ), + ); + if (thread.deletedAt !== null) { + return yield* new WakeCommandError({ + message: `Thread '${flags.threadId}' has been deleted.`, + }); + } + yield* client.orchestration + .dispatch({ + headers, + payload: { + type: "thread.turn.start", + commandId: CommandId.make(yield* crypto.randomUUIDv4), + threadId: thread.id, + message: { + messageId: MessageId.make(yield* crypto.randomUUIDv4), + role: "user", + text: flags.message, + attachments: [], + }, + runtimeMode: thread.runtimeMode, + interactionMode: thread.interactionMode, + createdAt: DateTime.formatIso(yield* DateTime.now), + }, + }) + .pipe( + Effect.timeout("10 seconds"), + Effect.mapError( + (cause) => + new WakeCommandError({ + message: `Could not confirm delivery to thread '${flags.threadId}'. Check the thread before retrying.`, + cause, + }), + ), + ); + }), + (session) => auth.revokeSession(session.sessionId).pipe(Effect.ignore({ log: true })), + ); + }).pipe( + Effect.provide( + EnvironmentAuth.runtimeLayer.pipe( + Layer.provide(ServerConfig.layer(config)), + Layer.provide(Layer.succeed(References.MinimumLogLevel, logLevel)), + ), + ), + ); + yield* Console.log(`Sent message to thread ${flags.threadId}; turn requested.`); + }, Effect.provide(FetchHttpClient.layer)), + ), +); diff --git a/docs/user/composer.md b/docs/user/composer.md index 97ed8025f3c7..da90d8dbfc8e 100644 --- a/docs/user/composer.md +++ b/docs/user/composer.md @@ -12,6 +12,23 @@ becomes an attachment when inserting it would exceed the message limit. On a hardware keyboard, use `Cmd+Shift+V` on Apple devices or `Ctrl+Shift+V` elsewhere to keep a large paste editable in the composer instead. +## Wake a thread from the command line + +Send a message to an existing thread and have its agent start working: + +```sh +t3 wake "Continue the task and run the tests" +``` + +Run this on the environment's machine while its T3 Code server or desktop app is +running. The command fails if no server is running; it never starts one. It uses +the thread's existing provider, model, permission mode, and interaction mode, and +returns after the server accepts the turn request. + +Use `--base-dir ` to select a server using a different T3 home. Inside a +linked worktree, discovery checks the worktree's `.t3` first, then `T3CODE_HOME` +or the default T3 home. + ## Attach files Attach up to eight files per message. Images can be up to 10 MB; other files can From 7a05cb11d32a3cd4cce80f7477b317553ff4b0d4 Mon Sep 17 00:00:00 2001 From: Luiz Ferraz Date: Sun, 13 Sep 2026 17:21:20 +0000 Subject: [PATCH 2/3] feat(server): expose T3CODE_THREAD_ID to agent processes - Set the thread ID env var across Claude, Codex, Cursor, Grok, OpenCode, OhMyPi, and Antigravity sessions, including Codex shell tools and ACP child commands - Document the variable and its limits in docs/user/composer.md --- .../src/provider/Drivers/AntigravityDriver.ts | 5 +- .../src/provider/Layers/AntigravityAdapter.ts | 1 + .../src/provider/Layers/ClaudeAdapter.test.ts | 25 ++++++++ .../src/provider/Layers/ClaudeAdapter.ts | 5 +- .../CodexCollabRuntime.integration.test.ts | 59 +++++++++++++++++++ .../provider/Layers/CodexSessionRuntime.ts | 9 ++- .../src/provider/Layers/CursorAdapter.test.ts | 59 +++++++++++++++++++ .../src/provider/Layers/CursorAdapter.ts | 15 +++-- .../server/src/provider/Layers/GrokAdapter.ts | 15 +++-- .../src/provider/Layers/OhMyPiAdapter.ts | 15 +++-- .../src/provider/Layers/OpenCodeAdapter.ts | 11 ++-- .../src/provider/acp/AntigravityAcpSupport.ts | 3 + docs/user/composer.md | 8 +++ 13 files changed, 199 insertions(+), 31 deletions(-) diff --git a/apps/server/src/provider/Drivers/AntigravityDriver.ts b/apps/server/src/provider/Drivers/AntigravityDriver.ts index 0082f3cbdc30..43052c1f093f 100644 --- a/apps/server/src/provider/Drivers/AntigravityDriver.ts +++ b/apps/server/src/provider/Drivers/AntigravityDriver.ts @@ -163,7 +163,10 @@ export const AntigravityDriver: ProviderDriver { ); }); + it.effect("isolates the thread environment across sessions and resumes", () => { + const environment = { ...process.env, T3CODE_THREAD_ID: "inherited-thread" }; + const harness = makeHarness({ environment }); + return Effect.gen(function* () { + const adapter = yield* ClaudeAdapter; + yield* adapter.startSession({ threadId: THREAD_ID, runtimeMode: "full-access" }); + const firstEnvironment = harness.getLastCreateQueryInput()?.options.env; + assert.equal(firstEnvironment?.T3CODE_THREAD_ID, THREAD_ID); + yield* adapter.startSession({ threadId: RESUME_THREAD_ID, runtimeMode: "full-access" }); + assert.equal( + harness.getLastCreateQueryInput()?.options.env?.T3CODE_THREAD_ID, + RESUME_THREAD_ID, + ); + assert.equal(firstEnvironment?.T3CODE_THREAD_ID, THREAD_ID); + yield* adapter.stopSession(THREAD_ID); + yield* adapter.startSession({ + threadId: THREAD_ID, + runtimeMode: "full-access", + resumeCursor: { sessionId: "claude-native-session" }, + }); + assert.equal(harness.getLastCreateQueryInput()?.options.env?.T3CODE_THREAD_ID, THREAD_ID); + assert.equal(environment.T3CODE_THREAD_ID, "inherited-thread"); + }).pipe(Effect.provide(harness.layer)); + }); + it.effect("derives bypass permission mode from full-access runtime policy", () => { const harness = makeHarness(); return Effect.gen(function* () { diff --git a/apps/server/src/provider/Layers/ClaudeAdapter.ts b/apps/server/src/provider/Layers/ClaudeAdapter.ts index fc52fe38dc68..957f03e0d976 100644 --- a/apps/server/src/provider/Layers/ClaudeAdapter.ts +++ b/apps/server/src/provider/Layers/ClaudeAdapter.ts @@ -4743,7 +4743,10 @@ export const makeClaudeAdapter = Effect.fn("makeClaudeAdapter")(function* ( canUseTool, onUserDialog, supportedDialogKinds: ["resume_return"], - env: McpProviderSession.withAgentDeviceEnvironment(claudeEnvironment, mcpSession), + env: { + ...McpProviderSession.withAgentDeviceEnvironment(claudeEnvironment, mcpSession), + T3CODE_THREAD_ID: input.threadId, + }, additionalDirectories, ...(Object.keys(extraArgs).length > 0 ? { extraArgs } : {}), ...(mcpSession diff --git a/apps/server/src/provider/Layers/CodexCollabRuntime.integration.test.ts b/apps/server/src/provider/Layers/CodexCollabRuntime.integration.test.ts index 02b7a45f33ad..9e1d61610c69 100644 --- a/apps/server/src/provider/Layers/CodexCollabRuntime.integration.test.ts +++ b/apps/server/src/provider/Layers/CodexCollabRuntime.integration.test.ts @@ -24,8 +24,19 @@ import { assert, describe } from "vite-plus/test"; import wireFixture from "../testFixtures/codexMultiAgentWire.json" with { type: "json" }; import { makeCodexSessionRuntime } from "./CodexSessionRuntime.ts"; +import { writeFakeCli } from "../../testUtils/fakeCli.ts"; import { HostProcessPlatform } from "@t3tools/shared/hostProcess"; +const decodeRecordedEnvironment = Schema.decodeUnknownEffect( + Schema.fromJsonString( + Schema.Struct({ + threadId: Schema.String, + childThreadId: Schema.String, + args: Schema.Array(Schema.String), + }), + ), +); + const ROOT = wireFixture.rootThreadId; const [CHILD_A, CHILD_B] = wireFixture.childThreadIds as [string, string]; const MEMORY = "memory-consolidation-thread"; @@ -166,6 +177,54 @@ const peerPath = NodePath.join( ); describe("CodexSessionRuntime collab integration", () => { + it.effect("passes the T3 thread environment to Codex and its child commands", () => + Effect.gen(function* () { + const directory = NodeFS.mkdtempSync(NodePath.join(NodeOS.tmpdir(), "codex-thread-env-")); + yield* Effect.addFinalizer(() => + Effect.sync(() => NodeFS.rmSync(directory, { recursive: true, force: true })), + ); + const fixturePath = NodePath.join(directory, "script.json"); + const logPath = NodePath.join(directory, "environment.json"); + // @effect-diagnostics-next-line preferSchemaOverJson:off - Scripted peer fixture. + NodeFS.writeFileSync(fixturePath, JSON.stringify({ rootThreadId: ROOT, notifications: [] })); + const binaryPath = writeFakeCli({ + directory, + name: "codex-env-peer", + source: [ + 'import { execFileSync } from "node:child_process";', + 'import { writeFileSync } from "node:fs";', + 'const childThreadId = execFileSync(process.execPath, ["-p", "process.env.T3CODE_THREAD_ID"], { encoding: "utf8" }).trim();', + // @effect-diagnostics-next-line preferSchemaOverJson:off - Quote a filesystem path in generated JavaScript. + `writeFileSync(${JSON.stringify(logPath)}, JSON.stringify({ threadId: process.env.T3CODE_THREAD_ID, childThreadId, args: process.argv.slice(2) }));`, + // @effect-diagnostics-next-line preferSchemaOverJson:off - Quote an import URL in generated JavaScript. + `await import(${JSON.stringify(new URL("../testFixtures/codexCollabMockPeer.mjs", import.meta.url).href)});`, + ].join("\n"), + }); + const threadId = ThreadId.make("t3-thread-environment"); + const environment = { + ...process.env, + T3_CODEX_COLLAB_SCRIPT: fixturePath, + T3CODE_THREAD_ID: "inherited-thread", + }; + const runtime = yield* makeCodexSessionRuntime({ + threadId, + binaryPath, + cwd: directory, + runtimeMode: "full-access", + environment, + }); + yield* runtime.start(); + const recorded = yield* decodeRecordedEnvironment(NodeFS.readFileSync(logPath, "utf8")); + assert.equal(recorded.threadId, threadId); + assert.equal(recorded.childThreadId, threadId); + assert.include( + recorded.args, + 'shell_environment_policy.set.T3CODE_THREAD_ID="t3-thread-environment"', + ); + assert.equal(environment.T3CODE_THREAD_ID, "inherited-thread"); + }).pipe(Effect.scoped, Effect.provide(NodeServices.layer)), + ); + it.effect("looks up child model metadata once after activity registration", () => Effect.gen(function* () { const script = { diff --git a/apps/server/src/provider/Layers/CodexSessionRuntime.ts b/apps/server/src/provider/Layers/CodexSessionRuntime.ts index a04db1405912..05b1c38239ae 100644 --- a/apps/server/src/provider/Layers/CodexSessionRuntime.ts +++ b/apps/server/src/provider/Layers/CodexSessionRuntime.ts @@ -1312,10 +1312,17 @@ export const makeCodexSessionRuntime = ( const resolvedHomePath = options.homePath ? expandHomePath(options.homePath) : undefined; const env = { ...options.environment, + T3CODE_THREAD_ID: options.threadId, ...(resolvedHomePath ? { CODEX_HOME: resolvedHomePath } : {}), }; const extendEnv = options.environment === undefined; - const appServerArgs = codexSessionAppServerArgs(options.appServerArgs, options.launchArgs); + // Keep the thread ID available to shell tools even with inherit = "core" or "none". + const appServerArgs = [ + ...codexSessionAppServerArgs(options.appServerArgs, options.launchArgs), + "-c", + // @effect-diagnostics-next-line preferSchemaOverJson:off - Encode a TOML basic string in a CLI argument, not a JSON payload. + `shell_environment_policy.set.T3CODE_THREAD_ID=${JSON.stringify(options.threadId)}`, + ]; const spawnCommand = yield* resolveSpawnCommand(options.binaryPath, appServerArgs, { env, extendEnv, diff --git a/apps/server/src/provider/Layers/CursorAdapter.test.ts b/apps/server/src/provider/Layers/CursorAdapter.test.ts index bdc818994a9a..bd725aea27a5 100644 --- a/apps/server/src/provider/Layers/CursorAdapter.test.ts +++ b/apps/server/src/provider/Layers/CursorAdapter.test.ts @@ -162,6 +162,65 @@ const cursorAdapterTestLayer = it.layer( ); cursorAdapterTestLayer("CursorAdapterLive", (it) => { + it.effect( + "passes the thread environment to ACP agents and their child commands on start and resume", + () => + Effect.gen(function* () { + const adapter = yield* CursorAdapter; + const settings = yield* ServerSettingsService; + const workspace = yield* Effect.promise(() => + NodeFSP.mkdtemp(NodePath.join(NodeOS.tmpdir(), "cursor-thread-env-")), + ); + const logPath = NodePath.join(workspace, "environment.jsonl"); + const binaryPath = writeFakeCli({ + directory: workspace, + name: "thread-environment-agent", + source: [ + 'import { execFileSync } from "node:child_process";', + 'import { appendFileSync as recordEnvironment } from "node:fs";', + 'const childThreadId = execFileSync(process.execPath, ["-p", "process.env.T3CODE_THREAD_ID"], { encoding: "utf8" }).trim();', + // @effect-diagnostics-next-line preferSchemaOverJson:off - Quote a filesystem path in generated JavaScript. + `recordEnvironment(${JSON.stringify(logPath)}, JSON.stringify([process.env.T3CODE_THREAD_ID, childThreadId]) + "\\n");`, + execScriptSource({ scriptPath: mockAgentPath }), + ].join("\n"), + }); + yield* settings.updateSettings({ providers: { cursor: { binaryPath } } }); + const firstId = ThreadId.make("cursor-thread-env-first"); + const secondId = ThreadId.make("cursor-thread-env-second"); + const first = yield* adapter.startSession({ + threadId: firstId, + cwd: workspace, + runtimeMode: "full-access", + }); + yield* adapter.startSession({ + threadId: secondId, + cwd: workspace, + runtimeMode: "full-access", + }); + yield* adapter.stopSession(firstId); + yield* adapter.startSession({ + threadId: firstId, + cwd: workspace, + runtimeMode: "full-access", + resumeCursor: first.resumeCursor, + }); + const lines = yield* Effect.promise(() => NodeFSP.readFile(logPath, "utf8")); + assert.deepEqual( + lines + .trim() + .split("\n") + .map((line) => JSON.parse(line)), + [ + [firstId, firstId], + [secondId, secondId], + [firstId, firstId], + ], + ); + yield* adapter.stopSession(firstId); + yield* adapter.stopSession(secondId); + }), + ); + it.effect("rejects rollback without discarding the provider conversation", () => Effect.gen(function* () { const adapter = yield* CursorAdapter; diff --git a/apps/server/src/provider/Layers/CursorAdapter.ts b/apps/server/src/provider/Layers/CursorAdapter.ts index 925d585e5838..08bd24d4efc6 100644 --- a/apps/server/src/provider/Layers/CursorAdapter.ts +++ b/apps/server/src/provider/Layers/CursorAdapter.ts @@ -543,14 +543,13 @@ export function makeCursorAdapter( const mcpSession = McpProviderSession.readMcpProviderSession(input.threadId); const acp = yield* makeCursorAcpRuntime({ cursorSettings: effectiveCursorSettings, - ...(options?.environment || mcpSession?.agentDeviceEnvironment - ? { - environment: McpProviderSession.withAgentDeviceEnvironment( - options?.environment ?? process.env, - mcpSession, - ), - } - : {}), + environment: { + ...McpProviderSession.withAgentDeviceEnvironment( + options?.environment ?? process.env, + mcpSession, + ), + T3CODE_THREAD_ID: input.threadId, + }, childProcessSpawner, cwd, runtimeMode: input.runtimeMode, diff --git a/apps/server/src/provider/Layers/GrokAdapter.ts b/apps/server/src/provider/Layers/GrokAdapter.ts index 01ec3d118a3e..637be4e6c621 100644 --- a/apps/server/src/provider/Layers/GrokAdapter.ts +++ b/apps/server/src/provider/Layers/GrokAdapter.ts @@ -995,14 +995,13 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte const mcpSession = McpProviderSession.readMcpProviderSession(input.threadId); const acp = yield* makeGrokAcpRuntime({ grokSettings, - ...(options?.environment || mcpSession?.agentDeviceEnvironment - ? { - environment: McpProviderSession.withAgentDeviceEnvironment( - options?.environment ?? process.env, - mcpSession, - ), - } - : {}), + environment: { + ...McpProviderSession.withAgentDeviceEnvironment( + options?.environment ?? process.env, + mcpSession, + ), + T3CODE_THREAD_ID: input.threadId, + }, childProcessSpawner, cwd, runtimeMode: input.runtimeMode, diff --git a/apps/server/src/provider/Layers/OhMyPiAdapter.ts b/apps/server/src/provider/Layers/OhMyPiAdapter.ts index 7065ed074772..89b8988b0fcb 100644 --- a/apps/server/src/provider/Layers/OhMyPiAdapter.ts +++ b/apps/server/src/provider/Layers/OhMyPiAdapter.ts @@ -410,14 +410,13 @@ export function makeOhMyPiAdapter( const mcpSession = McpProviderSession.readMcpProviderSession(input.threadId); const acp = yield* makeOhMyPiAcpRuntime({ ohMyPiSettings, - ...(options?.environment || mcpSession?.agentDeviceEnvironment - ? { - environment: McpProviderSession.withAgentDeviceEnvironment( - options?.environment ?? process.env, - mcpSession, - ), - } - : {}), + environment: { + ...McpProviderSession.withAgentDeviceEnvironment( + options?.environment ?? process.env, + mcpSession, + ), + T3CODE_THREAD_ID: input.threadId, + }, childProcessSpawner, cwd, runtimeMode: input.runtimeMode, diff --git a/apps/server/src/provider/Layers/OpenCodeAdapter.ts b/apps/server/src/provider/Layers/OpenCodeAdapter.ts index 3d216bb1167b..27043a54837d 100644 --- a/apps/server/src/provider/Layers/OpenCodeAdapter.ts +++ b/apps/server/src/provider/Layers/OpenCodeAdapter.ts @@ -2833,10 +2833,13 @@ export function makeOpenCodeAdapter( directory, serverUrl, ...(serverPassword ? { serverPassword } : {}), - environment: McpProviderSession.withAgentDeviceEnvironment( - options?.environment ?? process.env, - mcpSession, - ), + environment: { + ...McpProviderSession.withAgentDeviceEnvironment( + options?.environment ?? process.env, + mcpSession, + ), + T3CODE_THREAD_ID: input.threadId, + }, }); const client = openCodeRuntime.createOpenCodeSdkClient({ baseUrl: server.url, diff --git a/apps/server/src/provider/acp/AntigravityAcpSupport.ts b/apps/server/src/provider/acp/AntigravityAcpSupport.ts index f2f370068181..d48875f479be 100644 --- a/apps/server/src/provider/acp/AntigravityAcpSupport.ts +++ b/apps/server/src/provider/acp/AntigravityAcpSupport.ts @@ -5,6 +5,7 @@ import { PROVIDER_SEND_TURN_MAX_IMAGE_BYTES, type ProviderSendTurnInput, type RuntimeMode, + type ThreadId, } from "@t3tools/contracts"; import * as Crypto from "effect/Crypto"; import * as Effect from "effect/Effect"; @@ -35,6 +36,8 @@ export interface AntigravityAcpRuntimeInput extends Omit< | "transformSessionUpdate" | "transformStdout" > { + /** Absent for disposable authentication and model discovery sessions. */ + readonly threadId?: ThreadId; /** Device CLI environment supplied for this provider session. */ readonly agentDeviceEnvironment?: Readonly>; readonly childProcessSpawner: ChildProcessSpawner.ChildProcessSpawner["Service"]; diff --git a/docs/user/composer.md b/docs/user/composer.md index da90d8dbfc8e..0f52524dd07f 100644 --- a/docs/user/composer.md +++ b/docs/user/composer.md @@ -29,6 +29,14 @@ Use `--base-dir ` to select a server using a different T3 home. Inside linked worktree, discovery checks the worktree's `.t3` first, then `T3CODE_HOME` or the default T3 home. +Agent processes launched by T3 Code receive `T3CODE_THREAD_ID`, the ID of their +T3 thread. Scripts and child commands can read it from their environment instead +of accepting a thread ID argument. The value is set when the provider session +starts or resumes; already-running processes need a session restart to receive it. +This does not apply to externally managed OpenCode servers, whose process +environment T3 Code does not control. Custom provider environment filters may +also restrict which variables reach tools. + ## Attach files Attach up to eight files per message. Images can be up to 10 MB; other files can From 1fe94d43cf2ccd931eb13cfff7887a487aa83c84 Mon Sep 17 00:00:00 2001 From: Luiz Ferraz Date: Sun, 13 Sep 2026 20:11:36 +0000 Subject: [PATCH 3/3] fix(server): reject wake requests for inactive threads --- .../Layers/OrchestrationEngine.test.ts | 67 +++++++++++++++++++ apps/server/src/orchestration/decider.ts | 8 ++- 2 files changed, 74 insertions(+), 1 deletion(-) diff --git a/apps/server/src/orchestration/Layers/OrchestrationEngine.test.ts b/apps/server/src/orchestration/Layers/OrchestrationEngine.test.ts index 7ea1d588acfb..d573ab0d8358 100644 --- a/apps/server/src/orchestration/Layers/OrchestrationEngine.test.ts +++ b/apps/server/src/orchestration/Layers/OrchestrationEngine.test.ts @@ -758,6 +758,73 @@ describe("OrchestrationEngine", () => { await system.dispose(); }); + it.each(["thread.archive", "thread.delete"] as const)( + "rejects a stale turn request after %s without persisting its message", + async (type) => { + const system = await createOrchestrationSystem(); + try { + const projectId = ProjectId.make("wake-race-project"); + const threadId = ThreadId.make("wake-race-thread"); + await system.run( + system.engine.dispatch({ + type: "project.create", + commandId: CommandId.make("wake-race-project-create"), + projectId, + title: "Wake race project", + workspaceRoot: "/tmp/wake-race-project", + createdAt: now(), + }), + ); + await system.run( + system.engine.dispatch({ + type: "thread.create", + commandId: CommandId.make("wake-race-thread-create"), + threadId, + projectId, + title: "Wake race thread", + modelSelection: { instanceId: ProviderInstanceId.make("codex"), model: "gpt-5" }, + runtimeMode: "full-access", + interactionMode: "default", + branch: null, + worktreePath: null, + createdAt: now(), + }), + ); + const observed = (await system.readModel()).threads[0]!; + const command = { + type: "thread.turn.start", + commandId: CommandId.make("wake-race-turn"), + threadId, + message: { + messageId: MessageId.make("wake-race-message"), + role: "user", + text: "Continue", + attachments: [], + }, + runtimeMode: observed.runtimeMode, + interactionMode: observed.interactionMode, + createdAt: now(), + } satisfies OrchestrationCommand; + await system.run( + system.engine.dispatch({ + type, + commandId: CommandId.make("wake-race-deactivate"), + threadId, + }), + ); + const before = await system.readModel(); + const sequence = await system.run(system.engine.latestSequence); + const error = await system.run(system.engine.dispatch(command).pipe(Effect.flip)); + expect(error._tag).toBe("OrchestrationCommandInvariantError"); + expect(error.message).toContain(type === "thread.archive" ? "archived" : "deleted"); + expect(await system.run(system.engine.latestSequence)).toBe(sequence); + expect(await system.readModel()).toEqual(before); + } finally { + await system.dispose(); + } + }, + ); + it("archives and unarchives threads through orchestration commands", async () => { const system = await createOrchestrationSystem(); const { engine } = system; diff --git a/apps/server/src/orchestration/decider.ts b/apps/server/src/orchestration/decider.ts index 8703a80c1bc1..6e29144baa45 100644 --- a/apps/server/src/orchestration/decider.ts +++ b/apps/server/src/orchestration/decider.ts @@ -1275,11 +1275,17 @@ export const decideOrchestrationCommand = Effect.fn("decideOrchestrationCommand" detail: `Message id '${command.message.messageId}' uses the reserved imported-session namespace.`, }); } - const targetThread = yield* requireThread({ + const targetThread = yield* requireThreadNotArchived({ readModel, command, threadId: command.threadId, }); + if (targetThread.deletedAt !== null) { + return yield* new OrchestrationCommandInvariantError({ + commandType: command.type, + detail: `Thread '${command.threadId}' was deleted and cannot start a turn.`, + }); + } const sourceProposedPlan = command.sourceProposedPlan; const sourceThread = sourceProposedPlan ? yield* requireThread({