From 886a31860b8aebbcc6c2066d5b38881f1960e9e4 Mon Sep 17 00:00:00 2001 From: Xon333 <221869142+Xon333@users.noreply.github.com> Date: Fri, 9 Oct 2026 04:15:40 +0200 Subject: [PATCH] fix(server): expose thread executor to native preferences tools --- .../mcp/toolkits/environment/handlers.test.ts | 150 ++++++++++-------- .../src/orchestration-v2/runtimeLayer.test.ts | 144 ++++++++++++++++- .../src/orchestration-v2/runtimeLayer.ts | 1 + 3 files changed, 225 insertions(+), 70 deletions(-) diff --git a/apps/server/src/mcp/toolkits/environment/handlers.test.ts b/apps/server/src/mcp/toolkits/environment/handlers.test.ts index e4ae660cc004..be4b9c651d87 100644 --- a/apps/server/src/mcp/toolkits/environment/handlers.test.ts +++ b/apps/server/src/mcp/toolkits/environment/handlers.test.ts @@ -26,75 +26,89 @@ import { EnvironmentToolkit } from "./tools.ts"; const environmentId = EnvironmentId.make("environment:preferences"); const threadId = ThreadId.make("thread:preferences"); -it.effect("refuses a preferences update when the caller's turn ends while it waits", () => - Effect.gen(function* () { - const caller = yield* Ref.make(liveThreadShell(threadId)); - const updates = yield* Ref.make(0); - // Completes once the declaration's own check has read the caller. - const checked = yield* Deferred.make(); - const layerDependencies = Layer.mergeAll( - ThreadCommandExecutor.layer, - Layer.succeed(McpInvocationContext.McpInvocationContext, { - environmentId, - requestNamespace: "provider:preferences", - thread: { - threadId, - providerSessionId: "provider:preferences", - providerInstanceId: ProviderInstanceId.make("codex"), - }, - client: undefined, - capabilities: new Set(["orchestration" as const]), - issuedAt: 0, - }), - Layer.mock(ThreadManagement.ThreadManagementService)({ - getThreadShell: () => - Ref.get(caller).pipe(Effect.tap(() => Deferred.succeed(checked, undefined))), - }), - Layer.mock(Environment.ServerEnvironment)({ - getDescriptor: Effect.succeed({ +it.effect.each([ + { change: "turn ends", patch: { activeRunId: null }, code: "parent_not_active" }, + { + change: "runtime mode changes", + patch: { runtimeMode: "approval-required" }, + code: "capability_denied", + }, + { + change: "interaction mode changes", + patch: { interactionMode: "plan" }, + code: "capability_denied", + }, +] as const)( + "refuses a preferences update when the caller's $change while it waits", + ({ patch, code }) => + Effect.gen(function* () { + const caller = yield* Ref.make(liveThreadShell(threadId)); + const updates = yield* Ref.make(0); + // Completes once the declaration's own check has read the caller. + const checked = yield* Deferred.make(); + const layerDependencies = Layer.mergeAll( + ThreadCommandExecutor.layer, + Layer.succeed(McpInvocationContext.McpInvocationContext, { environmentId, - label: "Test", - platform: { os: "linux", arch: "x64" }, - serverVersion: "0.0.0", - capabilities: { repositoryIdentity: false }, + requestNamespace: "provider:preferences", + thread: { + threadId, + providerSessionId: "provider:preferences", + providerInstanceId: ProviderInstanceId.make("codex"), + }, + client: undefined, + capabilities: new Set(["orchestration" as const]), + issuedAt: 0, }), - }), - Layer.mock(Settings.ServerSettingsService)({ - getSettings: Effect.succeed(DEFAULT_SERVER_SETTINGS), - updateSettings: () => - Ref.update(updates, (count) => count + 1).pipe(Effect.as(DEFAULT_SERVER_SETTINGS)), - }), - ); - yield* Effect.gen(function* () { - const toolkit = yield* EnvironmentToolkit; - const executor = yield* ThreadCommandExecutor.ThreadCommandExecutor; - const held = yield* Deferred.make(); - const release = yield* Deferred.make(); - // A turn-completion command holds the thread's lock. - const holder = yield* executor - .withLock( - threadId, - Deferred.succeed(held, undefined).pipe(Effect.andThen(Deferred.await(release))), - ) - .pipe(Effect.forkChild); - yield* Deferred.await(held); - const update = yield* toolkit - .handle("t3_environment_preferences_update", { newWorktreesStartFromOrigin: true }) - .pipe(Stream.unwrap, Stream.runCollect, Effect.forkChild); - // The update passed its first check and waits for the lock; the turn then ends. - yield* Deferred.await(checked); - yield* Ref.update(caller, (shell) => ({ ...shell, activeRunId: null })); - yield* Deferred.succeed(release, undefined); - yield* Fiber.join(holder); - const result = yield* Fiber.join(update); - expect(result.at(-1)?.result).toMatchObject({ code: "parent_not_active" }); - expect(yield* Ref.get(updates)).toBe(0); - }).pipe( - Effect.provide( - McpToolAccess.HandlersLayer.layer(EnvironmentHandlers.layer).pipe( - Layer.provideMerge(layerDependencies), + Layer.mock(ThreadManagement.ThreadManagementService)({ + getThreadShell: () => + Ref.get(caller).pipe(Effect.tap(() => Deferred.succeed(checked, undefined))), + }), + Layer.mock(Environment.ServerEnvironment)({ + getDescriptor: Effect.succeed({ + environmentId, + label: "Test", + platform: { os: "linux", arch: "x64" }, + serverVersion: "0.0.0", + capabilities: { repositoryIdentity: false }, + }), + }), + Layer.mock(Settings.ServerSettingsService)({ + getSettings: Effect.succeed(DEFAULT_SERVER_SETTINGS), + updateSettings: () => + Ref.update(updates, (count) => count + 1).pipe(Effect.as(DEFAULT_SERVER_SETTINGS)), + }), + ); + yield* Effect.gen(function* () { + const toolkit = yield* EnvironmentToolkit; + const executor = yield* ThreadCommandExecutor.ThreadCommandExecutor; + const held = yield* Deferred.make(); + const release = yield* Deferred.make(); + // A turn-completion command holds the thread's lock. + const holder = yield* executor + .withLock( + threadId, + Deferred.succeed(held, undefined).pipe(Effect.andThen(Deferred.await(release))), + ) + .pipe(Effect.forkChild); + yield* Deferred.await(held); + const update = yield* toolkit + .handle("t3_environment_preferences_update", { newWorktreesStartFromOrigin: true }) + .pipe(Stream.unwrap, Stream.runCollect, Effect.forkChild); + // The caller changes after the first check, while the update waits for the lock. + yield* Deferred.await(checked); + yield* Ref.update(caller, (shell) => ({ ...shell, ...patch })); + yield* Deferred.succeed(release, undefined); + yield* Fiber.join(holder); + const result = yield* Fiber.join(update); + expect(result.at(-1)?.result).toMatchObject({ code }); + expect(yield* Ref.get(updates)).toBe(0); + }).pipe( + Effect.provide( + McpToolAccess.HandlersLayer.layer(EnvironmentHandlers.layer).pipe( + Layer.provideMerge(layerDependencies), + ), ), - ), - ); - }), + ); + }), ); diff --git a/apps/server/src/orchestration-v2/runtimeLayer.test.ts b/apps/server/src/orchestration-v2/runtimeLayer.test.ts index aed7f73e6c0f..abbf2dbfc080 100644 --- a/apps/server/src/orchestration-v2/runtimeLayer.test.ts +++ b/apps/server/src/orchestration-v2/runtimeLayer.test.ts @@ -10,6 +10,7 @@ import { CommandId, ContextTransferId, EventId, + EnvironmentId, MessageId, NodeId, RuntimeRequestId, @@ -38,6 +39,7 @@ import * as Option from "effect/Option"; import * as Queue from "effect/Queue"; import * as Stream from "effect/Stream"; import * as TestClock from "effect/testing/TestClock"; +import { McpSchema, McpServer } from "effect/ai"; import * as SqlClient from "effect/sql/SqlClient"; import * as CheckpointStore from "../checkpointing/CheckpointStore.ts"; @@ -49,6 +51,12 @@ import * as ProjectEnrichmentService from "../project/ProjectEnrichmentService.t import * as ProjectService from "../project/ProjectService.ts"; import * as ServerSettings from "../serverSettings.ts"; import * as McpSessionRegistryTestkit from "../mcp/McpSessionRegistry.testkit.ts"; +import * as McpHttpServer from "../mcp/McpHttpServer.ts"; +import * as McpInvocationContext from "../mcp/McpInvocationContext.ts"; +import * as McpToolAccessTestkit from "../mcp/McpToolAccess.testkit.ts"; +import * as EnvironmentHandlers from "../mcp/toolkits/environment/handlers.ts"; +import { EnvironmentToolkit } from "../mcp/toolkits/environment/tools.ts"; +import * as ServerEnvironment from "../environment/ServerEnvironment.ts"; import * as ProviderInstanceRegistry from "../provider/ProviderInstanceRegistry.ts"; import type { ProviderInstance } from "@t3tools/provider-core/server/driver"; import * as VcsDriverRegistry from "../vcs/VcsDriverRegistry.ts"; @@ -272,7 +280,6 @@ const layerTest = Layer.mergeAll( ProjectStore.layer, ProjectionStore.layer, EffectOutbox.layer, - ThreadCommandExecutor.layer, ).pipe( Layer.provide(McpSessionRegistryTestkit.layer), Layer.provide(SqlitePersistence.layerMemory), @@ -301,7 +308,6 @@ const layerProjectDeletionTest = Layer.mergeAll( RuntimeLayer.layer.pipe(Layer.provide(RuntimeLayer.layerProjectService)), RuntimeLayer.layerProjectService, RuntimeLayer.layerEventSink, - ThreadCommandExecutor.layer, ).pipe( Layer.provide( Layer.mock(ProjectEnrichmentService.ProjectEnrichmentService)({ @@ -480,6 +486,140 @@ const layerSharedApplicationDataPlaneTest = Layer.mergeAll( ); it.layer(layerTest)("OrchestrationV2LayerLive", (it) => { + it.effect("serves preferences with the runtime's shared thread lock and OAuth access", () => + Effect.gen(function* () { + const executor = yield* ThreadCommandExecutor.ThreadCommandExecutor; + const environmentId = EnvironmentId.make("environment:runtime-preferences"); + const threadId = ThreadId.make("thread:runtime-preferences"); + const invocation: McpInvocationContext.McpInvocationScope = { + environmentId, + requestNamespace: "provider:runtime-preferences", + thread: { + threadId, + providerSessionId: "provider:runtime-preferences", + providerInstanceId: ProviderInstanceId.make("codex"), + }, + client: undefined, + capabilities: new Set(["orchestration"]), + issuedAt: 0, + }; + const client = McpSchema.McpServerClient.of({ + clientId: 1, + clientCapabilities: {}, + clientInfo: { name: "preferences-test", version: "1.0.0" }, + protocolVersion: "2025-06-18", + initializePayload: { + protocolVersion: "2025-06-18", + capabilities: {}, + clientInfo: { name: "preferences-test", version: "1.0.0" }, + }, + getClient: Effect.die("unused"), + }); + const registration = McpHttpServer.toolkitRegistration( + EnvironmentToolkit, + EnvironmentHandlers.layer, + ).pipe( + Layer.provideMerge(McpServer.McpServer.layer), + Layer.provideMerge( + ServerSettings.layerTest({ + sourceControlWritingStyle: { + mode: "custom", + customInstructions: "Preserve these instructions.", + followChangeRequestTemplates: true, + }, + }), + ), + Layer.provide(McpToolAccessTestkit.liveThreadsLayer), + Layer.provide( + Layer.mock(ServerEnvironment.ServerEnvironment)({ + getDescriptor: Effect.succeed({ + environmentId, + label: "Preferences test", + platform: { os: "linux", arch: "x64" }, + serverVersion: "0.0.0", + capabilities: { repositoryIdentity: false }, + }), + }), + ), + ); + yield* Effect.gen(function* () { + const server = yield* McpServer.McpServer; + const settings = yield* ServerSettings.ServerSettingsService; + const call = ( + scope: McpInvocationContext.McpInvocationScope, + name: string, + parameters: Record = {}, + ) => + server + .callTool({ name, arguments: parameters }) + .pipe( + Effect.provideService(McpInvocationContext.McpInvocationContext, scope), + Effect.provideService(McpSchema.McpServerClient, client), + ); + const before = yield* settings.getSettings; + const read = yield* call(invocation, "t3_environment_read"); + assert.equal(read.isError, false); + const held = yield* Deferred.make(); + const release = yield* Deferred.make(); + const requested = yield* Deferred.make(); + const holder = yield* executor + .withLock( + threadId, + Deferred.succeed(held, undefined).pipe(Effect.andThen(Deferred.await(release))), + ) + .pipe(Effect.forkChild); + yield* Deferred.await(held); + const withLock = executor.withLock; + const observeLock: ThreadCommandExecutor.ThreadCommandExecutor["Service"]["withLock"] = ( + key, + effect, + ) => Deferred.succeed(requested, undefined).pipe(Effect.andThen(withLock(key, effect))); + const spy = vi.spyOn(executor, "withLock").mockImplementation(observeLock); + yield* Effect.gen(function* () { + const update = yield* call(invocation, "t3_environment_preferences_update", { + newWorktreesStartFromOrigin: !before.newWorktreesStartFromOrigin, + }).pipe(Effect.forkChild); + yield* Deferred.await(requested); + assert.deepEqual(yield* settings.getSettings, before); + const oauth: McpInvocationContext.McpInvocationScope = { + ...invocation, + requestNamespace: "client:runtime-preferences", + thread: undefined, + client: { sessionId: "runtime-preferences", label: "Test", access: "full-access" }, + }; + // A client has no thread lock, even while a provider's update is waiting. + const reapplied = yield* call(oauth, "t3_environment_preferences_update", { + sourceControlWritingStyle: { + customInstructions: before.sourceControlWritingStyle.customInstructions, + }, + }); + assert.equal(reapplied.isError, false); + assert.equal(spy.mock.calls.length, 1); + assert.deepEqual(yield* settings.getSettings, before); + for (const access of ["read-only", "approval-required", "auto"] as const) { + const denied = yield* call( + { ...oauth, client: { ...oauth.client!, access } }, + "t3_environment_preferences_update", + { newWorktreesStartFromOrigin: true }, + ); + assert.equal(denied.isError, true); + assert.match(JSON.stringify(denied.content), /capability_denied/); + } + yield* Deferred.succeed(release, undefined); + yield* Fiber.join(holder); + assert.equal((yield* Fiber.join(update)).isError, false); + assert.deepEqual(yield* settings.getSettings, { + ...before, + newWorktreesStartFromOrigin: !before.newWorktreesStartFromOrigin, + }); + }).pipe( + Effect.ensuring(Deferred.succeed(release, undefined)), + Effect.ensuring(Effect.sync(() => spy.mockRestore())), + ); + }).pipe(Effect.provide(registration)); + }), + ); + it.effect("emits model updates separately from provider switches", () => Effect.gen(function* () { const orchestrator = yield* Orchestrator.OrchestratorV2; diff --git a/apps/server/src/orchestration-v2/runtimeLayer.ts b/apps/server/src/orchestration-v2/runtimeLayer.ts index 5d736f1564be..0ffddfbcdf32 100644 --- a/apps/server/src/orchestration-v2/runtimeLayer.ts +++ b/apps/server/src/orchestration-v2/runtimeLayer.ts @@ -330,6 +330,7 @@ const layerMcpAppRequestsProvided = McpAppRequests.layer.pipe( ); export const layer = Layer.mergeAll( + ThreadCommandExecutor.layer, layerOrchestratorProvided, layerMcpAppRequestsProvided, layerThreadManagementProvided,