diff --git a/apps/mobile/src/features/threads/SecretRequestCard.tsx b/apps/mobile/src/features/threads/SecretRequestCard.tsx new file mode 100644 index 000000000000..f1cf25fde12f --- /dev/null +++ b/apps/mobile/src/features/threads/SecretRequestCard.tsx @@ -0,0 +1,170 @@ +import { + SECRET_REQUEST_DEFAULT_PLACEHOLDER, + SECRET_REQUEST_PRIVACY_NOTE, + secretRequestAnswerInput, + secretRequestDisplay, + secretRequestFailureMessage, + type SecretRequestItem, +} from "@t3tools/client-runtime/secret-request"; +import { + isAtomCommandInterrupted, + squashAtomCommandFailure, +} from "@t3tools/client-runtime/state/runtime"; +import type { EnvironmentId, OrchestrationV2ProjectedTurnItem } from "@t3tools/contracts"; +import { useRef, useState } from "react"; +import { Pressable, View, type ColorValue } from "react-native"; + +import { SymbolView, type AppSymbolName } from "../../components/AppSymbol"; +import { AppText as Text, AppTextInput as TextInput } from "../../components/AppText"; +import { serverEnvironment } from "../../state/server"; +import { useAtomCommand } from "../../state/use-atom-command"; +import { RequestActionButton } from "./RequestActionButton"; + +/** + * Feed card for a secret an agent asked the user for. The typed value lives + * only in this component's state and the RPC payload: it is never logged, + * alerted, or persisted, and the field clears once the answer is sent. + */ +const LOCK_SYMBOL: AppSymbolName = { ios: "lock", android: "lock" }; +const PRIVATE_SYMBOL: AppSymbolName = { ios: "checkmark.shield", android: "lock" }; + +export function SecretRequestCard(props: { + readonly environmentId: EnvironmentId; + readonly projectedItem: OrchestrationV2ProjectedTurnItem; + readonly iconColor: ColorValue; +}) { + const { item, visibility } = props.projectedItem; + if (item.type !== "secret_request") return null; + const display = secretRequestDisplay(item, visibility); + if (display.kind === "pending") { + return ( + + ); + } + const icon: AppSymbolName = + display.kind === "pending-elsewhere" + ? LOCK_SYMBOL + : display.outcome === "saved" + ? "checkmark" + : "minus"; + return ( + + + + {item.label} · {display.label} + + + ); +} + +function PendingSecretRequestForm(props: { + readonly environmentId: EnvironmentId; + readonly item: SecretRequestItem; + readonly iconColor: ColorValue; +}) { + const { item } = props; + const answer = useAtomCommand(serverEnvironment.answerSecretRequest, { + label: "answer secret request", + // The failure cause holds the request; keep it out of the console. + reportFailure: false, + reportDefect: false, + }); + const [secret, setSecret] = useState(""); + const [submitting, setSubmitting] = useState(false); + const [error, setError] = useState(null); + // Submit then a tap can both run before a re-render; this guard is synchronous. + const inFlight = useRef(false); + + const send = async ( + reply: { readonly type: "save"; readonly secret: string } | { readonly type: "decline" }, + ) => { + const input = secretRequestAnswerInput(item, reply); + if (input === null || inFlight.current) return; + inFlight.current = true; + setSubmitting(true); + setError(null); + const result = await answer({ environmentId: props.environmentId, input }).finally(() => { + inFlight.current = false; + setSubmitting(false); + }); + if (result._tag === "Success") { + // The card switches to its answered row once the item updates. + setSecret(""); + return; + } + if (!isAtomCommandInterrupted(result)) { + setError(secretRequestFailureMessage(squashAtomCommandFailure(result))); + } + }; + + // Same hierarchy as web: what is asked, why, the field, then the promise + // about where the value goes. + return ( + + + {item.label} + {item.reason.trim() ? ( + {item.reason} + ) : null} + + void send({ type: "save", secret })} + /> + {error !== null ? ( + + {error} + + ) : null} + void send({ type: "save", secret })} + /> + + + + + {SECRET_REQUEST_PRIVACY_NOTE} + + + {/* Quiet like the web card's: the field and Save are the action. */} + void send({ type: "decline" })} + > + Decline + + + + ); +} diff --git a/apps/mobile/src/features/threads/ThreadFeed.tsx b/apps/mobile/src/features/threads/ThreadFeed.tsx index 7c71b4dd6c79..58191639331a 100644 --- a/apps/mobile/src/features/threads/ThreadFeed.tsx +++ b/apps/mobile/src/features/threads/ThreadFeed.tsx @@ -1,5 +1,6 @@ import { ThreadContextDivider } from "./thread-context-divider"; import { ThreadHandoffRow } from "./thread-handoff-row"; +import { SecretRequestCard } from "./SecretRequestCard"; import { WorktreeWorkingHeader, WorktreeSetupCard, @@ -162,6 +163,7 @@ import { threadFeedRunIsUnsettled, isContextCompactionActivityGroup, isContextHandoffActivityGroup, + isSecretRequestActivityGroup, type ThreadFeedEntry, type ThreadFeedLatestRun, } from "../../lib/threadActivity"; @@ -1621,6 +1623,16 @@ function renderFeedEntry( ); } + if (entry.type === "activity-group" && isSecretRequestActivityGroup(entry)) { + return ( + + ); + } + if (entry.type === "activity-group" && isContextCompactionActivityGroup(entry)) { const label = entry.activities[0]!.summary; const active = diff --git a/apps/mobile/src/lib/threadActivity.ts b/apps/mobile/src/lib/threadActivity.ts index 86c4184d1c1f..927fd59b7dec 100644 --- a/apps/mobile/src/lib/threadActivity.ts +++ b/apps/mobile/src/lib/threadActivity.ts @@ -321,6 +321,13 @@ export function isContextHandoffActivityGroup(entry: ThreadFeedActivityGroup): b ); } +export function isSecretRequestActivityGroup(entry: ThreadFeedActivityGroup): boolean { + return ( + entry.activities.length === 1 && + entry.activities[0]?.projectedItem.item.type === "secret_request" + ); +} + function isUserInputActivityGroup(entry: ThreadFeedActivityGroup): boolean { return entry.activities.some((activity) => activity.workEntry.questionAnswer !== undefined); } @@ -420,7 +427,13 @@ function itemIsToolLike(item: OrchestrationV2TurnItem): boolean { } function itemIsProminent(item: OrchestrationV2TurnItem): boolean { - return item.type === "fork" || item.type === "thread_created" || item.type === "system_notice"; + return ( + item.type === "fork" || + item.type === "thread_created" || + item.type === "system_notice" || + // An answerable card: it must stand alone and never fold away with the run. + item.type === "secret_request" + ); } function itemStatus(item: OrchestrationV2TurnItem): ThreadFeedActivity["status"] { @@ -537,6 +550,8 @@ function itemIcon(item: OrchestrationV2TurnItem): ThreadFeedActivity["icon"] { case "fork": case "thread_created": return "zap"; + case "secret_request": + return "lock"; } } @@ -590,6 +605,8 @@ function itemSummary( return "Thread forked"; case "thread_created": return "Thread created"; + case "secret_request": + return item.label; case "dynamic_tool": { const classified = classifyToolActivity({ itemType: "dynamic_tool_call", @@ -647,6 +664,8 @@ function itemPreview(item: OrchestrationV2TurnItem): string | null { case "fork": case "thread_created": return item.targetThreadId; + case "secret_request": + return item.reason || null; case "subagent": return item.result ?? item.progress ?? item.prompt; case "dynamic_tool": diff --git a/apps/server/src/auth/RpcAuthorization.ts b/apps/server/src/auth/RpcAuthorization.ts index 61d7ea279bf4..a2988000ca78 100644 --- a/apps/server/src/auth/RpcAuthorization.ts +++ b/apps/server/src/auth/RpcAuthorization.ts @@ -96,6 +96,7 @@ export const RPC_REQUIRED_SCOPES = { [WS_METHODS.scheduledTasksDelete]: AuthOrchestrationOperateScope, [WS_METHODS.scheduledTasksRunNow]: AuthOrchestrationOperateScope, [WS_METHODS.scheduledTasksRotateWebhookToken]: AuthOrchestrationOperateScope, + [WS_METHODS.secretsAnswerRequest]: AuthOrchestrationOperateScope, // Delivery logs hold request bodies, so they need the same scope as the URL. [WS_METHODS.scheduledTasksListWebhookDeliveries]: AuthOrchestrationOperateScope, [WS_METHODS.scheduledTasksGetWebhookDelivery]: AuthOrchestrationOperateScope, diff --git a/apps/server/src/mcp/OrchestratorMcpService.activity.test.ts b/apps/server/src/mcp/OrchestratorMcpService.activity.test.ts index 2433eb8a704a..b179a2f45822 100644 --- a/apps/server/src/mcp/OrchestratorMcpService.activity.test.ts +++ b/apps/server/src/mcp/OrchestratorMcpService.activity.test.ts @@ -20,6 +20,7 @@ import * as ProviderAdapterRegistry from "../orchestration-v2/ProviderAdapterReg import * as ProviderRegistry from "../provider/Services/ProviderRegistry.ts"; import * as ProjectService from "../project/ProjectService.ts"; import * as ScheduledTaskService from "../scheduledTasks/ScheduledTaskService.ts"; +import * as SecretRequests from "../secrets/SecretRequests.ts"; import * as ThreadManagementService from "../orchestration-v2/ThreadManagementService.ts"; import type * as McpInvocationContext from "./McpInvocationContext.ts"; import * as OrchestratorMcpService from "./OrchestratorMcpService.ts"; @@ -149,6 +150,7 @@ it("readThread prefers activity-run status over a newer cancelled queued run", a getProviders: Effect.succeed([]), } satisfies Partial), Layer.mock(ProjectService.ProjectService)({}), + Layer.mock(SecretRequests.SecretRequests)({}), Layer.mock(ScheduledTaskService.ScheduledTaskService)({ list: () => Effect.succeed({ tasks: [] }), } satisfies Partial), @@ -213,6 +215,7 @@ it("readThread prefers waiting activity status over a newer cancelled queued run getProviders: Effect.succeed([]), } satisfies Partial), Layer.mock(ProjectService.ProjectService)({}), + Layer.mock(SecretRequests.SecretRequests)({}), Layer.mock(ScheduledTaskService.ScheduledTaskService)({ list: () => Effect.succeed({ tasks: [] }), } satisfies Partial), @@ -325,6 +328,7 @@ it("taskStatus returns task.providerInstanceId rather than the driver kind", asy getProviders: Effect.succeed([]), } satisfies Partial), Layer.mock(ProjectService.ProjectService)({}), + Layer.mock(SecretRequests.SecretRequests)({}), Layer.mock(ScheduledTaskService.ScheduledTaskService)({ list: () => Effect.succeed({ tasks: [] }), } satisfies Partial), @@ -447,6 +451,7 @@ it("readThread and sendToThread reach threads in other projects", async () => { getProviders: Effect.succeed([]), } satisfies Partial), Layer.mock(ProjectService.ProjectService)({}), + Layer.mock(SecretRequests.SecretRequests)({}), Layer.mock(ScheduledTaskService.ScheduledTaskService)({ list: () => Effect.succeed({ diff --git a/apps/server/src/mcp/OrchestratorMcpService.test.ts b/apps/server/src/mcp/OrchestratorMcpService.test.ts index a10730200da3..2d255fda123d 100644 --- a/apps/server/src/mcp/OrchestratorMcpService.test.ts +++ b/apps/server/src/mcp/OrchestratorMcpService.test.ts @@ -24,6 +24,7 @@ import * as ProviderRegistry from "../provider/Services/ProviderRegistry.ts"; import { buildUnavailableProviderSnapshot } from "../provider/unavailableProviderSnapshot.ts"; import * as ProjectService from "../project/ProjectService.ts"; import * as ScheduledTaskService from "../scheduledTasks/ScheduledTaskService.ts"; +import * as SecretRequests from "../secrets/SecretRequests.ts"; import type { McpInvocationScope } from "./McpInvocationContext.ts"; import * as OrchestratorMcpService from "./OrchestratorMcpService.ts"; @@ -120,6 +121,7 @@ describe("OrchestratorMcpService", () => { list: () => Effect.succeed([]), }), Layer.mock(ProjectService.ProjectService)({}), + Layer.mock(SecretRequests.SecretRequests)({}), Layer.mock(ScheduledTaskService.ScheduledTaskService)({}), ); const scope: McpInvocationScope = { @@ -209,6 +211,7 @@ describe("OrchestratorMcpService", () => { list: () => Effect.succeed([]), }), Layer.mock(ProjectService.ProjectService)({}), + Layer.mock(SecretRequests.SecretRequests)({}), Layer.mock(ScheduledTaskService.ScheduledTaskService)({}), ); const scope: McpInvocationScope = { @@ -290,6 +293,7 @@ describe("OrchestratorMcpService", () => { list: () => Effect.succeed([]), }), Layer.mock(ProjectService.ProjectService)({}), + Layer.mock(SecretRequests.SecretRequests)({}), Layer.mock(ScheduledTaskService.ScheduledTaskService)({}), ); const scope: McpInvocationScope = { @@ -363,6 +367,7 @@ describe("OrchestratorMcpService", () => { list: () => Effect.succeed([]), }), Layer.mock(ProjectService.ProjectService)({}), + Layer.mock(SecretRequests.SecretRequests)({}), Layer.mock(ScheduledTaskService.ScheduledTaskService)({}), ); const scope: McpInvocationScope = { @@ -444,6 +449,7 @@ describe("OrchestratorMcpService", () => { list: () => Effect.succeed([]), }), Layer.mock(ProjectService.ProjectService)({}), + Layer.mock(SecretRequests.SecretRequests)({}), Layer.mock(ScheduledTaskService.ScheduledTaskService)({}), ); const scope: McpInvocationScope = { @@ -652,6 +658,7 @@ describe("OrchestratorMcpService provider resolution", () => { disabledAntigravityInstanceId, ]), Layer.mock(ProjectService.ProjectService)({}), + Layer.mock(SecretRequests.SecretRequests)({}), Layer.mock(ScheduledTaskService.ScheduledTaskService)({}), ); @@ -796,6 +803,7 @@ describe("OrchestratorMcpService provider resolution", () => { }), adapterRegistryLayer([codexInstanceId, antigravityInstanceId]), Layer.mock(ProjectService.ProjectService)({}), + Layer.mock(SecretRequests.SecretRequests)({}), Layer.mock(ScheduledTaskService.ScheduledTaskService)({}), ); @@ -890,6 +898,7 @@ describe("OrchestratorMcpService provider resolution", () => { }), adapterRegistryLayer([codexInstanceId, antigravityInstanceId]), Layer.mock(ProjectService.ProjectService)({}), + Layer.mock(SecretRequests.SecretRequests)({}), Layer.mock(ScheduledTaskService.ScheduledTaskService)({}), ); @@ -937,6 +946,7 @@ describe("OrchestratorMcpService provider resolution", () => { ]), adapterRegistryLayer([codexInstanceId]), Layer.mock(ProjectService.ProjectService)({}), + Layer.mock(SecretRequests.SecretRequests)({}), Layer.mock(ScheduledTaskService.ScheduledTaskService)({}), ); @@ -1044,6 +1054,7 @@ describe("OrchestratorMcpService provider resolution", () => { adapterRegistryLayer([codexInstanceId, claudeInstanceId]), Layer.mock(ScheduledTaskService.ScheduledTaskService)({}), Layer.mock(ProjectService.ProjectService)({}), + Layer.mock(SecretRequests.SecretRequests)({}), ); yield* Effect.gen(function* () { @@ -1200,6 +1211,7 @@ describe("OrchestratorMcpService provider resolution", () => { ]), adapterRegistryLayer([codexInstanceId, codexAltInstanceId]), Layer.mock(ProjectService.ProjectService)({}), + Layer.mock(SecretRequests.SecretRequests)({}), Layer.mock(ScheduledTaskService.ScheduledTaskService)({}), ); diff --git a/apps/server/src/mcp/OrchestratorMcpService.ts b/apps/server/src/mcp/OrchestratorMcpService.ts index 7a18455c4d70..685f1be0faea 100644 --- a/apps/server/src/mcp/OrchestratorMcpService.ts +++ b/apps/server/src/mcp/OrchestratorMcpService.ts @@ -5,6 +5,7 @@ import { MessageId, type ModelSelection, NodeId, + TurnItemId, type OrchestrationV2Run, type OrchestrationV2Subagent, type OrchestrationV2ThreadProjection, @@ -19,6 +20,8 @@ import { type OrchestratorMcpDelegateTaskResult, type OrchestratorMcpInteractionMode, type OrchestratorMcpDeleteScheduledTaskInput, + type OrchestratorMcpRequestSecretInput, + type OrchestratorMcpRequestSecretResult, type OrchestratorMcpDeleteScheduledTaskResult, type OrchestratorMcpListScheduledTasksResult, type OrchestratorMcpRuntimeMode, @@ -60,6 +63,7 @@ import * as Crypto from "effect/Crypto"; import * as DateTime from "effect/DateTime"; import * as Duration from "effect/Duration"; import * as Effect from "effect/Effect"; +import * as Exit from "effect/Exit"; import * as Layer from "effect/Layer"; import * as Option from "effect/Option"; import * as Schema from "effect/Schema"; @@ -79,6 +83,8 @@ import { type McpThreadInvocationScope, requireThreadScope, } from "./McpInvocationContext.ts"; +import * as Metrics from "../observability/Metrics.ts"; +import * as SecretRequests from "../secrets/SecretRequests.ts"; const DEFAULT_WAIT_TIMEOUT_MS = 10 * 60 * 1_000; const MAX_WAIT_TIMEOUT_MS = 60 * 60 * 1_000; @@ -90,6 +96,8 @@ const TASK_WAKE_EVENTS = [ { thread: "child", eventType: "subagent.updated" }, { thread: "child", eventType: "provider-thread.updated" }, ] as const; +/** A person answers the card, so a slower poll is plenty. */ +const SECRET_REQUEST_POLL_INTERVAL_MS = 500; const DEFAULT_THREAD_LIST_LIMIT = 50; const DEFAULT_THREAD_READ_LIMIT = 50; const DEFAULT_THREAD_RUN_LIMIT = 10; @@ -140,6 +148,14 @@ export interface OrchestratorMcpServiceShape { scope: McpInvocationScope, input: OrchestratorMcpDeleteScheduledTaskInput, ) => Effect.Effect; + /** + * Asks the user for a secret through a card in the calling thread and waits + * for the answer. The value never reaches the agent: the result is a status. + */ + readonly requestSecret: ( + scope: McpInvocationScope, + input: OrchestratorMcpRequestSecretInput, + ) => Effect.Effect; readonly listThreads: ( scope: McpInvocationScope, input: OrchestratorMcpThreadListInput, @@ -220,7 +236,11 @@ function scheduledTaskSummary(task: ScheduledTask): OrchestratorMcpScheduledTask schedule: task.schedule, nextRunAt: task.nextRunAt, lastRunStatus: task.lastRunStatus, - ...(task.webhook === undefined ? {} : { webhookUrl: task.webhook.url ?? task.webhook.path }), + // A bare path is not a URL anyone can call, so agents never get one to share. + ...(task.webhook?.url == null ? {} : { webhookUrl: task.webhook.url }), + ...(task.webhook === undefined + ? {} + : { webhookSignature: task.webhook.hasSecret ? "set" : "none" }), }; } @@ -725,6 +745,8 @@ function turnItemText(item: OrchestrationV2TurnItem): string | null { return `Forked to thread ${item.targetThreadId}.`; case "thread_created": return `Created thread ${item.targetThreadId} with ${item.targetProviderInstanceId} (${item.targetModel}).`; + case "secret_request": + return `Asked the user for ${item.label}: ${item.secretStatus}.`; case "subagent": return item.result ?? item.progress ?? item.prompt; case "dynamic_tool": @@ -806,6 +828,7 @@ const make = Effect.gen(function* () { ), ) : Effect.succeed(project.defaultModelSelection); + const secretRequests = yield* SecretRequests.SecretRequests; const requireCapability = (scope: McpInvocationScope) => scope.capabilities.has("orchestration") @@ -1514,6 +1537,140 @@ const make = Effect.gen(function* () { ); return { scheduledTaskId: existing.id, deleted: true }; }), + requestSecret: (scope, input) => + Effect.gen(function* () { + // The card is shown in, and answered from, the caller's own thread. + const { scope: threadScope, parent } = yield* loadThreadCaller(scope, "request_secret"); + const threadId = threadScope.thread.threadId; + const run = ThreadManagementService.latestActiveRun(parent); + if ( + run === undefined || + run.rootNodeId === null || + run.providerInstanceId !== threadScope.thread.providerInstanceId + ) { + return yield* failure( + "parent_not_active", + "Asking for a secret requires an active run owned by this MCP provider session.", + ); + } + const runId = run.id; + const nodeId = run.rootNodeId; + const key = yield* requestKey(input.clientRequestId); + // Turn item ids are global; scope the key to this thread. A retry with + // the same clientRequestId finds this card, answered or not. + const turnItemId = TurnItemId.make( + `turn-item:secret-request:${stablePart(threadId)}:${stablePart(key)}`, + ); + const record = (secretStatus: "pending" | "cancelled") => + threadManagement + .dispatch({ + type: "secret_request.record", + commandId: stableCommandId({ + scope, + requestKey: key, + operation: `secret-${secretStatus}`, + }), + threadId: threadId, + runId, + nodeId, + turnItemId, + label: input.label, + reason: input.reason, + ...(input.placeholder === undefined ? {} : { placeholder: input.placeholder }), + secretStatus, + }) + .pipe( + Effect.mapError((error) => + failure( + "orchestration_error", + `Could not record the secret request: ${errorMessage(error)}`, + ), + ), + ); + yield* record("pending"); + // Only this call can hand the agent its ref, so the card must not + // outlive it: a timeout, a failed wait or an aborted call closes it as + // cancelled. If even that fails, the server still refuses an answer + // once the run ends, and an unused value expires. + const closeCard = record("cancelled").pipe( + Effect.catch((error) => + Effect.logWarning("Could not close a secret request card", { error: error.message }), + ), + ); + + // The user answers the card (secrets.answerRequest), or it ends with + // the run; poll it like a delegated task. + const answered = yield* Effect.gen(function* () { + while (true) { + const projection = yield* threadManagement + .getThreadRecords(threadId, ["runs", "turnItems"], { + turnItemTypes: ["secret_request"], + messageRoles: [], + }) + .pipe( + Effect.mapError((error) => + failure( + "orchestration_error", + `Unable to read the secret request: ${errorMessage(error)}`, + ), + ), + ); + const item = projection.turnItems.find((candidate) => candidate.id === turnItemId); + if (item?.type === "secret_request" && item.secretStatus !== "pending") { + return item.secretStatus; + } + const current = projection.runs.find((candidate) => candidate.id === runId); + if ( + current === undefined || + ThreadManagementService.isTerminalRunStatus(current.status) + ) { + yield* record("cancelled"); + return "cancelled" as const; + } + yield* Effect.sleep(Duration.millis(SECRET_REQUEST_POLL_INTERVAL_MS)); + } + }).pipe( + Effect.timeoutOption( + Duration.millis( + Math.min(input.timeoutMs ?? DEFAULT_WAIT_TIMEOUT_MS, MAX_WAIT_TIMEOUT_MS), + ), + ), + Effect.onExit((exit) => (Exit.isSuccess(exit) ? Effect.void : closeCard)), + ); + if (Option.isNone(answered)) yield* closeCard; + // An answer that raced the timeout still wins: the card is answered once. + const status = Option.isSome(answered) + ? answered.value + : yield* threadManagement + .getThreadRecords(threadId, ["turnItems"], { + turnItemTypes: ["secret_request"], + messageRoles: [], + }) + .pipe( + Effect.map((records) => { + const item = records.turnItems.find((candidate) => candidate.id === turnItemId); + return item?.type === "secret_request" && + (item.secretStatus === "saved" || item.secretStatus === "declined") + ? item.secretStatus + : ("timed_out" as const); + }), + Effect.orElseSucceed(() => "timed_out" as const), + ); + yield* Effect.annotateCurrentSpan({ "secret_request.status": status }); + yield* Metrics.increment(Metrics.secretRequestsTotal, { status }); + if (status !== "saved") return { status }; + // Saved means the value was stored before the card said so; a missing + // value is a storage fault, not an answer the agent can act on. + const secretRef = yield* secretRequests.savedRef({ threadId: threadId, turnItemId }); + if (Option.isNone(secretRef)) { + return yield* failure( + "orchestration_error", + "The user saved the secret, but it could not be read. Ask again with a new clientRequestId.", + ); + } + return { status, secretRef: secretRef.value }; + }).pipe(Effect.withSpan("OrchestratorMcpService.requestSecret")), + capabilities: (scope) => Effect.gen(function* () { const { parent, limits } = yield* loadCaller(scope); @@ -2192,4 +2349,5 @@ export const layer: Layer.Layer< | ProviderAdapterRegistry.ProviderAdapterRegistryV2 | ScheduledTaskService.ScheduledTaskService | ProjectService.ProjectService + | SecretRequests.SecretRequests > = Layer.effect(OrchestratorMcpService, make); diff --git a/apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts b/apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts index a847c8d081f5..0a1aeed097ab 100644 --- a/apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts +++ b/apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts @@ -44,12 +44,14 @@ import * as Option from "effect/Option"; import * as PubSub from "effect/PubSub"; import * as Ref from "effect/Ref"; import * as Schema from "effect/Schema"; +import * as Fiber from "effect/Fiber"; import * as Stream from "effect/Stream"; import { McpSchema, McpServer } from "effect/ai"; import { ClaudeProviderCapabilitiesV2 } from "../orchestration-v2/Adapters/ClaudeAdapterV2.ts"; import { CodexProviderCapabilitiesV2 } from "../orchestration-v2/Adapters/CodexAdapterV2.ts"; import { CodexOrchestratorReplayHarness } from "../orchestration-v2/Adapters/CodexAdapterV2.testkit.ts"; +import { threadShellFromProjection } from "../orchestration-v2/ProjectionStore.ts"; import * as EventSink from "../orchestration-v2/EventSink.ts"; import * as Orchestrator from "../orchestration-v2/Orchestrator.ts"; import * as ThreadManagementService from "../orchestration-v2/ThreadManagementService.ts"; @@ -73,6 +75,8 @@ import { import { makeProviderRegistryLayer } from "../provider/testUtils/providerRegistryMock.ts"; import * as ProjectService from "../project/ProjectService.ts"; import * as ScheduledTaskService from "../scheduledTasks/ScheduledTaskService.ts"; +import * as SecretRequests from "../secrets/SecretRequests.ts"; +import * as ServerSecretStore from "../auth/ServerSecretStore.ts"; import * as McpHttpServer from "./McpHttpServer.ts"; import * as McpInvocationContext from "./McpInvocationContext.ts"; import { delegatedTaskRun, hasPendingChildRuns } from "./OrchestratorMcpService.ts"; @@ -444,7 +448,19 @@ function scheduledTaskFromUpsert(input: ScheduledTaskUpsertInput): ScheduledTask prompt: input.prompt, enabled: input.enabled, schedule: - input.schedule.type === "webhook" ? { type: "webhook", signature: null } : input.schedule, + input.schedule.type === "webhook" + ? { + type: "webhook", + signature: + input.schedule.signature == null + ? null + : { + header: input.schedule.signature.header, + encoding: input.schedule.signature.encoding, + prefix: input.schedule.signature.prefix, + }, + } + : input.schedule, projectId: input.projectId, threadId: input.threadId ?? null, workspaceStrategy: input.workspaceStrategy, @@ -463,6 +479,25 @@ function scheduledTaskFromUpsert(input: ScheduledTaskUpsertInput): ScheduledTask }; } +/** In-memory server secret store for tests that exercise secret requests. */ +const memorySecretStoreLayer = Layer.sync(ServerSecretStore.ServerSecretStore, () => { + const stored = new Map(); + return ServerSecretStore.ServerSecretStore.of({ + get: (name) => Effect.succeed(Option.fromNullishOr(stored.get(name))), + set: (name, value) => Effect.sync(() => void stored.set(name, value)), + create: (name, value) => Effect.sync(() => void stored.set(name, value)), + getOrCreateRandom: (name, bytes) => + Effect.sync(() => { + const existing = stored.get(name); + if (existing) return existing; + const value = new Uint8Array(bytes).fill(7); + stored.set(name, value); + return value; + }), + remove: (name) => Effect.sync(() => void stored.delete(name)), + }); +}); + const unusedScheduledTaskStubLayer = Layer.succeed( ScheduledTaskService.ScheduledTaskService, ScheduledTaskService.ScheduledTaskService.of({ @@ -658,6 +693,12 @@ describe("orchestrator MCP toolkit", () => { ), }), ), + Layer.provideMerge( + SecretRequests.layer.pipe( + Layer.provide(memorySecretStoreLayer), + Layer.provide(orchestrationLayer), + ), + ), Layer.provide(NodeServices.layer), ); @@ -1472,6 +1513,114 @@ describe("orchestrator MCP toolkit", () => { }); expect(yield* Ref.get(scheduledStore)).toHaveLength(0); + // The agent asks for a secret; the tool waits for the user and + // returns a one-use ref, never the value, which a signed webhook + // task then consumes. + const secretRequests = yield* SecretRequests.SecretRequests; + const secretFiber = yield* invoke("request_secret", { + label: "GitHub webhook secret", + reason: "Signs release webhooks. Enter the same value in GitHub's webhook settings.", + placeholder: "Paste the webhook secret", + clientRequestId: "release-webhook-secret", + }).pipe(Effect.forkChild); + // Polled without the helper's short budget: under load the tool's + // own reads come first. + const asked = yield* Effect.gen(function* () { + while (true) { + const projection = yield* orchestrator.getThreadProjection(parentThreadId); + if ( + projection.turnItems.some( + (item) => item.type === "secret_request" && item.secretStatus === "pending", + ) + ) { + return projection; + } + yield* Effect.sleep("5 millis"); + } + }); + const card = asked.turnItems.find((item) => item.type === "secret_request"); + if (card?.type !== "secret_request") { + return yield* Effect.die(new Error("Secret request card missing.")); + } + expect(card).toMatchObject({ + label: "GitHub webhook secret", + placeholder: "Paste the webhook secret", + }); + // Asking again with the same id leaves the open card exactly as it was. + yield* orchestrator.dispatch({ + type: "secret_request.record", + commandId: CommandId.make("command:test:secret-request-replay"), + threadId: parentThreadId, + runId: card.runId!, + nodeId: card.nodeId!, + turnItemId: card.id, + label: "Something else", + reason: "A different reason.", + secretStatus: "pending", + }); + expect( + (yield* orchestrator.getThreadProjection(parentThreadId)).turnItems.find( + (item) => item.id === card.id, + ), + ).toMatchObject({ label: "GitHub webhook secret", runId: card.runId }); + // The agent is blocked on the user, so the thread asks for input + // like a question does, in both shell paths. + expect( + (yield* orchestrator.getThreadShell(parentThreadId))?.pendingRuntimeRequest, + ).toMatchObject({ kind: "user_input" }); + expect(threadShellFromProjection(asked).pendingRuntimeRequest).toMatchObject({ + kind: "user_input", + }); + // What the card's Save sends (secrets.answerRequest). + yield* secretRequests.answer({ + threadId: parentThreadId, + turnItemId: card.id, + answer: { type: "save", secret: "github-webhook-secret" }, + }); + const secretCall = yield* Fiber.join(secretFiber); + expect(secretCall.isError).toBe(false); + const secretResult = secretCall.structuredContent as { + status: string; + secretRef?: string; + }; + expect(secretResult.status).toBe("saved"); + expect( + (yield* orchestrator.getThreadShell(parentThreadId))?.pendingRuntimeRequest ?? null, + ).toBeNull(); + expect(secretResult.secretRef).toMatch(/^secret-ref:[0-9a-f]{32}$/); + // The value appears nowhere in what the agent received. + const received = [ + ...secretCall.content.map((part) => ("text" in part ? part.text : "")), + ...Object.values(secretResult), + ]; + expect(received.some((value) => value.includes("github-webhook-secret"))).toBe(false); + + // A retry that lost the first result gets the same answer, with no + // second card for the user. + const retried = yield* invoke("request_secret", { + label: "GitHub webhook secret", + reason: "Signs release webhooks. Enter the same value in GitHub's webhook settings.", + clientRequestId: "release-webhook-secret", + }); + expect(retried.structuredContent).toEqual(secretResult); + expect( + (yield* orchestrator.getThreadProjection(parentThreadId)).turnItems.filter( + (item) => item.type === "secret_request", + ), + ).toHaveLength(1); + + // The ref is the secret for exactly one consumer. + expect( + yield* secretRequests.consume({ + ref: secretResult.secretRef as never, + projectId, + }), + ).toBe("github-webhook-secret"); + const reused = yield* secretRequests + .consume({ ref: secretResult.secretRef as never, projectId }) + .pipe(Effect.flip); + expect(reused.message).toContain("already used"); + const delegatedCall = yield* invoke("delegate_task", { task: delegatedPrompt, target: { @@ -3570,6 +3719,12 @@ describe("orchestrator MCP toolkit", () => { Layer.provide(providerRegistryLayer), Layer.provide(unusedScheduledTaskStubLayer), Layer.provide(Layer.mock(ProjectService.ProjectService)({})), + Layer.provideMerge( + SecretRequests.layer.pipe( + Layer.provide(memorySecretStoreLayer), + Layer.provide(orchestrationLayer), + ), + ), Layer.provide(NodeServices.layer), ); diff --git a/apps/server/src/mcp/toolkits/core.test.ts b/apps/server/src/mcp/toolkits/core.test.ts index cb26b1db25e6..feacf73e7003 100644 --- a/apps/server/src/mcp/toolkits/core.test.ts +++ b/apps/server/src/mcp/toolkits/core.test.ts @@ -23,6 +23,7 @@ import * as ProviderAdapterRegistry from "../../orchestration-v2/ProviderAdapter import * as ThreadManagement from "../../orchestration-v2/ThreadManagementService.ts"; import * as ProjectService from "../../project/ProjectService.ts"; import * as ProviderRegistry from "../../provider/Services/ProviderRegistry.ts"; +import * as SecretRequests from "../../secrets/SecretRequests.ts"; import * as ScheduledTaskService from "../../scheduledTasks/ScheduledTaskService.ts"; import * as McpHttpServer from "../McpHttpServer.ts"; import * as McpInvocationContext from "../McpInvocationContext.ts"; @@ -386,6 +387,7 @@ it.effect("refuses act-as-caller tools to a client caller", () => Layer.provide(Layer.mock(ProviderAdapterRegistry.ProviderAdapterRegistryV2)({})), Layer.provide(Layer.mock(ScheduledTaskService.ScheduledTaskService)({})), Layer.provide(Layer.mock(ProjectService.ProjectService)({})), + Layer.provide(Layer.mock(SecretRequests.SecretRequests)({})), ), ), ), @@ -437,6 +439,7 @@ it.effect("a caller cannot rewrite a scheduled task that runs above its own mode }), ), Layer.provide(Layer.mock(ProjectService.ProjectService)({})), + Layer.provide(Layer.mock(SecretRequests.SecretRequests)({})), ), ), ), @@ -505,6 +508,7 @@ it.effect("a caller cannot interrupt a thread that runs above its own modes", () Layer.provide(Layer.mock(ProviderAdapterRegistry.ProviderAdapterRegistryV2)({})), Layer.provide(Layer.mock(ScheduledTaskService.ScheduledTaskService)({})), Layer.provide(Layer.mock(ProjectService.ProjectService)({})), + Layer.provide(Layer.mock(SecretRequests.SecretRequests)({})), ), ), ), diff --git a/apps/server/src/mcp/toolkits/orchestrator/handlers.ts b/apps/server/src/mcp/toolkits/orchestrator/handlers.ts index d45197b22c9c..4aaff0485226 100644 --- a/apps/server/src/mcp/toolkits/orchestrator/handlers.ts +++ b/apps/server/src/mcp/toolkits/orchestrator/handlers.ts @@ -54,6 +54,12 @@ const handlers = { const service = yield* OrchestratorMcpService.OrchestratorMcpService; return yield* service.deleteScheduledTask(scope, input); }), + request_secret: (input) => + Effect.gen(function* () { + const scope = yield* McpInvocationContext.McpInvocationContext; + const service = yield* OrchestratorMcpService.OrchestratorMcpService; + return yield* service.requestSecret(scope, input); + }), create_threads: (input) => Effect.gen(function* () { const scope = yield* McpInvocationContext.McpInvocationContext; diff --git a/apps/server/src/mcp/toolkits/orchestrator/tools.ts b/apps/server/src/mcp/toolkits/orchestrator/tools.ts index a4a4ba728fec..6d4f0366d816 100644 --- a/apps/server/src/mcp/toolkits/orchestrator/tools.ts +++ b/apps/server/src/mcp/toolkits/orchestrator/tools.ts @@ -6,6 +6,8 @@ import { OrchestratorMcpDelegateTaskResult, OrchestratorMcpDeleteScheduledTaskInput, OrchestratorMcpDeleteScheduledTaskResult, + OrchestratorMcpRequestSecretInput, + OrchestratorMcpRequestSecretResult, OrchestratorMcpFailure, OrchestratorMcpListScheduledTasksInput, OrchestratorMcpListScheduledTasksResult, @@ -97,7 +99,7 @@ const TaskCancelTool = Tool.make("task_cancel", { export const ScheduleTaskTool = Tool.make("schedule_task", { description: - "Create persistent recurring work in the app scheduler, which runs even when no turn is active. Pass schedule as a STRUCTURED OBJECT, never JSON text: {type:'interval', everyMs:3600000} means hourly; {type:'fixed_time', timeOfDay:'09:00', weekdays:[1,2,3,4,5]} means weekday mornings; {type:'webhook'} runs on each request to a generated URL (returned as webhookUrl), and its prompt may use {{body.path}}, {{headers.name}}, {{query.name}}, {{body}} or {{request}} placeholders, which are the only request data the run sees. Omit projectId for this thread's project. In this thread's project, runs post into THIS thread by default (bindToCurrentThread=true); use false only when the user wants a fresh top-level thread per run. Elsewhere each run launches a fresh thread. Provider, model, and runtime settings inherit from this thread, or from the project default when there is no calling thread. Report the returned schedule and nextRunAt after success.", + "Create persistent work in the app scheduler that runs even when no turn is active. Pass schedule as a STRUCTURED OBJECT, never JSON text. Timers: {type:'interval', everyMs:3600000} is hourly; {type:'fixed_time', timeOfDay:'09:00', weekdays:[1,2,3,4,5]} is weekday mornings; report the returned nextRunAt. Webhooks: {type:'webhook'} runs once per request to a generated URL. The run sees the request ONLY through prompt placeholders: {{body.path}} (e.g. {{body.action}}, {{body.release.tag_name}}), {{headers.name}}, {{query.name}}, {{body}}, or {{request}} (method, headers with credentials redacted, and body). For a sender that signs requests, first call request_secret so the user enters the secret privately (never ask for it in chat or invent one), then set signature with the returned secretRef, e.g. GitHub: {type:'webhook', signature:{header:'x-hub-signature-256', encoding:'hex', prefix:'sha256=', secretRef}}. The result's webhookUrl is the public URL to give the user; if it is absent, this environment has no T3 Connect managed tunnel, so tell the user to enable T3 Connect remote access rather than sharing a path. Omit projectId for this thread's project. In this thread's project, runs post into THIS thread by default (bindToCurrentThread=true), which suits an orchestrator that sees every trigger, delegates work, and can dedupe against what is in flight; use false only when the user wants a fresh top-level thread per run. Elsewhere each run launches a fresh thread. Provider, model, and runtime settings inherit from this thread, or from the project default when there is no calling thread.", parameters: OrchestratorMcpScheduleTaskInput, success: OrchestratorMcpScheduleTaskResult, failure: OrchestratorMcpFailure, @@ -146,6 +148,18 @@ const DeleteScheduledTaskTool = Tool.make("delete_scheduled_task", { .annotate(Tool.Title, "Delete a scheduled task") .annotate(Tool.Destructive, true); +const RequestSecretTool = Tool.make("request_secret", { + description: + "Ask the user for a secret (a token, API key, signing secret, password) through a private card in this thread, and wait for them to answer. The value is kept by the app and NEVER returned to you or shown in the transcript. When saved, the result carries a secretRef: pass it to a tool that accepts one (e.g. schedule_task's signature.secretRef). It works once. Never ask for secrets in chat, and never invent one.", + parameters: OrchestratorMcpRequestSecretInput, + success: OrchestratorMcpRequestSecretResult, + failure: OrchestratorMcpFailure, + failureMode: "return", + dependencies, +}) + .annotate(Tool.Title, "Request a secret from the user") + .annotate(Tool.Destructive, false); + export const CreateThreadsTool = Tool.make("create_threads", { description: "Needs an agent running inside a T3 thread. Create one or more ORDINARY TOP-LEVEL T3 conversations. This is not delegation and does not create child agents/subagents. For delegated work, choose models from orchestrator_capabilities. Prefer native subagents only when they support the chosen model; otherwise call delegate_task, including for same-provider work. Use create_threads for a batch of separate top-level threads sharing this checkout. Prefer t3_thread_launch for a single thread. Both require the user to request separate/new/top-level threads or conversations. Each entry may override provider, model, options, runtime mode, and interaction mode; omitted settings inherit. Project, branch, and worktree always inherit and cannot be overridden here. For independent implementation or a PR stack in its own worktree, use t3_thread_launch with workspaceStrategy instead of asking the agent to create a worktree in its prompt.", @@ -248,6 +262,7 @@ export const OrchestratorToolkit = Toolkit.make( ListScheduledTasksTool, UpdateScheduledTaskTool, DeleteScheduledTaskTool, + RequestSecretTool, CreateThreadsTool, ThreadListTool, ThreadReadTool, diff --git a/apps/server/src/mcp/toolkits/worktree/registration.test.ts b/apps/server/src/mcp/toolkits/worktree/registration.test.ts index 452302095eaf..1878fa87cc74 100644 --- a/apps/server/src/mcp/toolkits/worktree/registration.test.ts +++ b/apps/server/src/mcp/toolkits/worktree/registration.test.ts @@ -19,6 +19,7 @@ import * as ProjectService from "../../../project/ProjectService.ts"; import * as ProjectSetupScriptRunner from "../../../project/ProjectSetupScriptRunner.ts"; import * as ProviderRegistry from "../../../provider/Services/ProviderRegistry.ts"; import * as ScheduledTaskService from "../../../scheduledTasks/ScheduledTaskService.ts"; +import * as SecretRequests from "../../../secrets/SecretRequests.ts"; import * as ServerSettings from "../../../serverSettings.ts"; import * as VcsStatusBroadcaster from "../../../vcs/VcsStatusBroadcaster.ts"; import * as McpHttpServer from "../../McpHttpServer.ts"; @@ -33,6 +34,7 @@ const StubServicesLive = Layer.mergeAll( Layer.mock(ProviderRegistry.ProviderRegistry)({}), Layer.mock(ProviderAdapterRegistry.ProviderAdapterRegistryV2)({}), Layer.mock(ScheduledTaskService.ScheduledTaskService)({}), + Layer.mock(SecretRequests.SecretRequests)({}), Layer.mock(ProjectService.ProjectService)({}), ServerSettings.layerTest({}), Layer.mock(GitWorkflowService.GitWorkflowService)({}), diff --git a/apps/server/src/observability/Metrics.ts b/apps/server/src/observability/Metrics.ts index f133eda3f8b0..db36e09b2bff 100644 --- a/apps/server/src/observability/Metrics.ts +++ b/apps/server/src/observability/Metrics.ts @@ -94,6 +94,16 @@ export const webhookRunsTotal = Metric.counter("t3_webhook_runs_total", { description: "Runs started from webhook deliveries, by outcome.", }); +/** Secrets agents asked users for, by how each ended: saved, declined, cancelled, timed_out. */ +export const secretRequestsTotal = Metric.counter("t3_secret_requests_total", { + description: "Secrets agents asked users for, by how each request ended.", +}); + +/** One-use secret refs a tool tried to use, by result: used, rejected. */ +export const secretRefsConsumedTotal = Metric.counter("t3_secret_refs_consumed_total", { + description: "Secret refs tools tried to use, by result.", +}); + export const metricAttributes = ( attributes: Readonly>, ): ReadonlyArray<[string, string]> => Object.entries(compactMetricAttributes(attributes)); diff --git a/apps/server/src/orchestration-v2/ClaudeAutomaticDelivery.integration.test.ts b/apps/server/src/orchestration-v2/ClaudeAutomaticDelivery.integration.test.ts index b13a8b5d90c1..5c8e6e689ef7 100644 --- a/apps/server/src/orchestration-v2/ClaudeAutomaticDelivery.integration.test.ts +++ b/apps/server/src/orchestration-v2/ClaudeAutomaticDelivery.integration.test.ts @@ -24,6 +24,7 @@ import * as Queue from "effect/Queue"; import * as Schema from "effect/Schema"; import * as Stream from "effect/Stream"; import * as ScheduledTaskService from "../scheduledTasks/ScheduledTaskService.ts"; +import * as SecretRequests from "../secrets/SecretRequests.ts"; import * as Scheduler from "../scheduling/Scheduler.ts"; import { SqlitePersistenceMemory } from "../persistence/Layers/Sqlite.ts"; import * as ClaudeAdapterV2 from "./Adapters/ClaudeAdapterV2.ts"; @@ -315,6 +316,7 @@ it.effect.each(["child completion", "scheduled message", "user steering"] as con }), ), Layer.provide(Layer.mock(ThreadLaunchService.ThreadLaunchService)({})), + Layer.provide(Layer.mock(SecretRequests.SecretRequests)({})), Layer.provide( Layer.mergeAll(NodeCrypto.layer, Scheduler.layer, SqlitePersistenceMemory), ), diff --git a/apps/server/src/orchestration-v2/Orchestrator.ts b/apps/server/src/orchestration-v2/Orchestrator.ts index 7116741f8e2a..f88ebe79e799 100644 --- a/apps/server/src/orchestration-v2/Orchestrator.ts +++ b/apps/server/src/orchestration-v2/Orchestrator.ts @@ -439,6 +439,8 @@ function commandThreadId(command: OrchestrationV2ServerCommand): ThreadId { case "delegated_task.completion-delivery.dispose": case "thread.created.record": return command.parentThreadId; + case "secret_request.record": + return command.threadId; case "thread.fork": case "thread.merge_back": return command.targetThreadId; @@ -6889,6 +6891,99 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio }, ); + /** + * Records or updates the card for a secret an agent asked the user for. The + * item carries the request and its status only; the value goes straight to + * the server's secret store and never through orchestration. + */ + const dispatchSecretRequestRecord = Effect.fn("orchestrationV2.dispatch.secretRequestRecord")( + function* ( + command: Extract, + events: Ref.Ref>, + ) { + const projection = yield* projectionStore + .getThreadRecords( + command.threadId, + ["runs", "nodes", "turnItems", "attempts", "providerTurns"], + { turnItemTypes: ["secret_request"], messageRoles: [] }, + ) + .pipe( + Effect.mapError( + (cause) => new OrchestratorProjectionError({ threadId: command.threadId, cause }), + ), + ); + const run = projection.runs.find((candidate) => candidate.id === command.runId); + if (run === undefined || run.rootNodeId !== command.nodeId) { + return yield* new OrchestratorDispatchError({ + commandId: command.commandId, + commandType: command.type, + cause: `Node ${command.nodeId} is not the root of run ${command.runId}.`, + }); + } + const existing = projection.turnItems.find((item) => item.id === command.turnItemId); + if (existing !== undefined && existing.type !== "secret_request") { + return yield* new OrchestratorDispatchError({ + commandId: command.commandId, + commandType: command.type, + cause: `Turn item ${command.turnItemId} is not a secret request.`, + }); + } + if (existing !== undefined && existing.runId !== command.runId) { + return yield* new OrchestratorDispatchError({ + commandId: command.commandId, + commandType: command.type, + cause: `Secret request ${command.turnItemId} belongs to another run.`, + }); + } + const now = yield* DateTime.now; + // A request is answered once, and a retry finds its card as it was + // asked: either one records the card unchanged. + const unchanged = + existing !== undefined && + (existing.secretStatus !== "pending" || command.secretStatus === "pending"); + const pending = command.secretStatus === "pending"; + const turnItem: OrchestrationV2TurnItem = unchanged + ? existing + : { + id: command.turnItemId, + threadId: command.threadId, + runId: command.runId, + nodeId: command.nodeId, + providerThreadId: run.providerThreadId, + providerTurnId: providerTurnForRun(projection, run)?.id ?? null, + nativeItemRef: null, + parentItemId: null, + ordinal: existing?.ordinal ?? (yield* nextTurnItemOrdinal(projection)), + status: pending + ? "waiting" + : command.secretStatus === "saved" + ? "completed" + : "cancelled", + title: command.label, + startedAt: existing?.startedAt ?? now, + completedAt: pending ? null : now, + updatedAt: now, + type: "secret_request", + label: command.label, + reason: command.reason, + ...(command.placeholder === undefined ? {} : { placeholder: command.placeholder }), + secretStatus: command.secretStatus, + }; + yield* emit( + events, + command, + )({ + type: "turn-item.updated", + threadId: command.threadId, + runId: command.runId, + nodeId: command.nodeId, + providerInstanceId: run.providerInstanceId, + occurredAt: now, + payload: turnItem, + }); + }, + ); + const dispatchRuntimeRequestRespond = ( command: Extract, events: Ref.Ref>, @@ -10008,6 +10103,9 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio case "thread.created.record": yield* dispatchCreatedThreadRecord(command, events); break; + case "secret_request.record": + yield* dispatchSecretRequestRecord(command, events); + break; default: return yield* dispatchUnsupported(command); } diff --git a/apps/server/src/orchestration-v2/ProjectionSettlement.test.ts b/apps/server/src/orchestration-v2/ProjectionSettlement.test.ts index 2a76341b42e1..e82e1fc4099a 100644 --- a/apps/server/src/orchestration-v2/ProjectionSettlement.test.ts +++ b/apps/server/src/orchestration-v2/ProjectionSettlement.test.ts @@ -521,5 +521,9 @@ it.effect("shell failure lookups stay on the thread's own turn items", () => const itemLookups = plan.filter((row) => row.detail.startsWith("SEARCH item ")); assert.lengthOf(itemLookups, 2); assert.isTrue(itemLookups.every((row) => row.detail.includes("turn_items_thread_run_idx"))); + // The pending secret request lookup is bounded the same way. + const secretLookups = plan.filter((row) => row.detail.startsWith("SEARCH secret ")); + assert.lengthOf(secretLookups, 1); + assert.include(secretLookups[0]!.detail, "turn_items_thread_run_idx"); }).pipe(Effect.provide(SqlLayer)), ); diff --git a/apps/server/src/orchestration-v2/ProjectionStore.ts b/apps/server/src/orchestration-v2/ProjectionStore.ts index 93288a512ac7..a700f557cfbf 100644 --- a/apps/server/src/orchestration-v2/ProjectionStore.ts +++ b/apps/server/src/orchestration-v2/ProjectionStore.ts @@ -28,7 +28,6 @@ import type { ProviderThreadId, ProviderTurnId, RunAttemptId, - RuntimeRequestId, MessageId, } from "@t3tools/contracts"; import { @@ -50,6 +49,7 @@ import { OrchestrationV2TurnItemJson as OrchestrationV2TurnItemJsonSchema, orchestrationV2RunWorkStartedAt, RunId, + RuntimeRequestId, CheckpointScopeId, ThreadId, TurnItemId, @@ -919,6 +919,7 @@ type ShellThreadRow = { readonly blocking_run_completed_at: string | null; readonly blocking_failure_payload_json: string | null; readonly pending_request_payload_json: string | null; + readonly pending_secret_request_payload_json: string | null; readonly latest_user_message_at: string | null; readonly latest_user_authored_message_at: string | null; readonly has_actionable_proposed_plan: number; @@ -1303,6 +1304,29 @@ function buildVisibleTurnItems(input: { ]); } +/** + * An agent waiting on a secret is waiting on the user just like a question, + * so the shell reports it as pending user input. Secret requests have no + * runtime request of their own; this stands one in for the shell summary + * only, keyed by the card's turn item. + */ +function secretRequestAsPendingInput( + item: OrchestrationV2TurnItem | null, +): OrchestrationV2ThreadProjection["runtimeRequests"][number] | null { + if (item?.type !== "secret_request" || item.nodeId === null) return null; + return { + id: RuntimeRequestId.make(item.id), + nodeId: item.nodeId, + providerTurnId: item.providerTurnId, + nativeRequestRef: null, + kind: "user_input", + status: "pending", + responseCapability: { type: "message" }, + createdAt: item.startedAt ?? item.updatedAt, + resolvedAt: null, + }; +} + export function threadShellFromProjection( projection: OrchestrationV2ThreadProjection, ): OrchestrationV2ThreadShell { @@ -1327,13 +1351,28 @@ export function threadShellFromProjection( projection.runs .filter(isActivityRunForShell) .toSorted((left, right) => right.ordinal - left.ordinal)[0] ?? null; + const liveRunIds = new Set(projection.runs.filter(isActivityRunForShell).map((run) => run.id)); const pendingRuntimeRequest = projection.runtimeRequests .filter((request) => request.status === "pending") .toSorted( (left, right) => DateTime.toEpochMillis(right.createdAt) - DateTime.toEpochMillis(left.createdAt), - )[0] ?? null; + )[0] ?? + secretRequestAsPendingInput( + projection.turnItems + .filter( + (item) => + item.type === "secret_request" && + item.status === "waiting" && + item.runId !== null && + liveRunIds.has(item.runId), + ) + .toSorted( + (left, right) => + DateTime.toEpochMillis(right.updatedAt) - DateTime.toEpochMillis(left.updatedAt), + )[0] ?? null, + ); const userMessages = projection.messages .filter((message) => message.role === "user") .toSorted( @@ -4958,6 +4997,17 @@ export const layer: Layer.Layer = ORDER BY request.created_at DESC, request.runtime_request_id DESC LIMIT 1 ) AS pending_request_payload_json, + ( + SELECT secret.payload_json + FROM orchestration_v2_projection_turn_items secret + INDEXED BY orchestration_v2_projection_turn_items_thread_run_idx + INNER JOIN orchestration_v2_projection_runs r ON r.run_id = secret.run_id + WHERE secret.thread_id = t.thread_id + AND secret.type = 'secret_request' AND secret.status = 'waiting' + AND r.status IN ('preparing', 'starting', 'running', 'waiting') + ORDER BY secret.updated_at DESC, secret.turn_item_id DESC + LIMIT 1 + ) AS pending_secret_request_payload_json, ( SELECT message.updated_at FROM orchestration_v2_projection_messages message @@ -5329,9 +5379,13 @@ export const layer: Layer.Layer = } = input; const thread = yield* decodeThreadPayload(row.payload_json); const pendingRuntimeRequest = - row.pending_request_payload_json === null - ? null - : yield* decodeRuntimeRequestPayload(row.pending_request_payload_json); + row.pending_request_payload_json !== null + ? yield* decodeRuntimeRequestPayload(row.pending_request_payload_json) + : row.pending_secret_request_payload_json !== null + ? secretRequestAsPendingInput( + yield* decodeTurnItemPayload(row.pending_secret_request_payload_json), + ) + : null; let terminalFailureItem = row.terminal_failure_payload_json === null ? null diff --git a/apps/server/src/orchestration-v2/ProviderRuntimeRecoveryService.test.ts b/apps/server/src/orchestration-v2/ProviderRuntimeRecoveryService.test.ts index b36239de290d..cb06238395a6 100644 --- a/apps/server/src/orchestration-v2/ProviderRuntimeRecoveryService.test.ts +++ b/apps/server/src/orchestration-v2/ProviderRuntimeRecoveryService.test.ts @@ -578,6 +578,71 @@ it.effect("cancels a stale waiting run when no checkpoint capture can finish it" }).pipe(Effect.provide(layer)); }); +it.effect("closes a secret request card's form when its run is recovered", () => { + const threadId = ThreadId.make("thread_secret_recovery"); + const runId = RunId.make("run_secret_recovery"); + let committedInput: Parameters[0] | null = + null; + const projection = { + thread: { id: threadId }, + runtimeRequests: [], + providerSessions: [], + providerThreads: [], + providerTurns: [], + runs: [{ id: runId, status: "running", providerInstanceId: ProviderInstanceId.make("codex") }], + attempts: [], + nodes: [], + subagents: [], + messages: [], + turnItems: [ + { + id: "turn-item:secret-request:recovery", + threadId, + runId, + nodeId: null, + type: "secret_request", + status: "waiting", + secretStatus: "pending", + label: "GitHub token", + reason: "Used as GH_TOKEN.", + }, + ], + } as unknown as OrchestrationV2ThreadProjection; + const layer = ProviderRuntimeRecovery.layer.pipe( + Layer.provide(ServerSettings.layerTest()), + Layer.provide( + Layer.mergeAll( + Layer.mock(ProjectionStore.ProjectionStoreV2)({ + getRecoveryThreadIds: () => Effect.succeed([threadId]), + getRuntimeRecoveryProjection: () => Effect.succeed(projection), + }), + Layer.mock(EventSink.EventSinkV2)({ + commitCommand: (input) => { + committedInput = input; + return Effect.succeed({ committed: true, cancelledEffectCount: 0 } as never); + }, + }), + IdAllocator.layer, + Layer.mock(EffectWorker.OrchestrationEffectWorkerV2)({ + runRecoveryOnce: Effect.succeed(false), + }), + Layer.mock(EffectOutbox.EffectOutboxV2)({ + listByCommandId: () => Effect.succeed([]), + reconcileAfterProcessLoss: Effect.succeed({ requeued: 0, cancelled: 0 }), + }), + ), + ), + ); + + return Effect.gen(function* () { + yield* (yield* ProviderRuntimeRecovery.ProviderRuntimeRecoveryService).reconcile("startup"); + const itemEvent = committedInput?.events.find((event) => event.type === "turn-item.updated"); + const card = itemEvent?.type === "turn-item.updated" ? itemEvent.payload : null; + assert.equal(card?.status, "cancelled"); + assert.equal(card?.type === "secret_request" ? card.secretStatus : null, "cancelled"); + }).pipe(Effect.provide(layer)); +}); + it.effect("holds accepted queued work without cancelling its execution state after restart", () => { const threadId = ThreadId.make("thread_queued_restart"); const runId = RunId.make("run_queued_restart"); diff --git a/apps/server/src/orchestration-v2/ProviderRuntimeRecoveryService.ts b/apps/server/src/orchestration-v2/ProviderRuntimeRecoveryService.ts index 6eb0ef12e605..27a79b9fa757 100644 --- a/apps/server/src/orchestration-v2/ProviderRuntimeRecoveryService.ts +++ b/apps/server/src/orchestration-v2/ProviderRuntimeRecoveryService.ts @@ -149,6 +149,15 @@ function resolveStaleBackgroundItemProviderInstanceId( return projection.providerThreads[0]?.providerInstanceId ?? projection.thread.providerInstanceId; } +/** An open item closed by recovery. A secret card also closes its form, so it takes no answer. */ +const cancelledItem = ( + item: OrchestrationV2ThreadProjection["turnItems"][number], + now: DateTime.Utc, +): OrchestrationV2ThreadProjection["turnItems"][number] => + item.type === "secret_request" && item.secretStatus === "pending" + ? { ...item, status: "cancelled", secretStatus: "cancelled", completedAt: now, updatedAt: now } + : { ...item, status: "cancelled", completedAt: now, updatedAt: now }; + /** * A provider thread's latest started run: the last turn that provider saw. * Restart recovery records the thread's cancelled background work on it, and @@ -426,7 +435,7 @@ export const make = Effect.gen(function* () { ...(item.nodeId === null ? {} : { nodeId: item.nodeId }), providerInstanceId: run.providerInstanceId, occurredAt: now, - payload: { ...item, status: "cancelled", completedAt: now, updatedAt: now }, + payload: cancelledItem(item, now), }); } } @@ -456,7 +465,7 @@ export const make = Effect.gen(function* () { ...(item.nodeId === null || item.nodeId === undefined ? {} : { nodeId: item.nodeId }), providerInstanceId, occurredAt: now, - payload: { ...item, status: "cancelled", completedAt: now, updatedAt: now }, + payload: cancelledItem(item, now), }); if (item.nodeId !== null && item.nodeId !== undefined) { const staleItemNode = projection.nodes.find( diff --git a/apps/server/src/orchestration-v2/ThreadLaunchService.test.ts b/apps/server/src/orchestration-v2/ThreadLaunchService.test.ts index 9b986f9fe860..82d66edb132f 100644 --- a/apps/server/src/orchestration-v2/ThreadLaunchService.test.ts +++ b/apps/server/src/orchestration-v2/ThreadLaunchService.test.ts @@ -48,6 +48,7 @@ import * as ManagedProjectFolders from "../project/ManagedProjectFolders.ts"; import { makeProviderRegistryLayer } from "../provider/testUtils/providerRegistryMock.ts"; import * as ServerSettings from "../serverSettings.ts"; import * as ScheduledTasks from "../scheduledTasks/ScheduledTaskService.ts"; +import * as SecretRequests from "../secrets/SecretRequests.ts"; import * as TextGeneration from "../textGeneration/TextGeneration.ts"; import { CodexProviderCapabilitiesV2 } from "./Adapters/CodexAdapterV2.ts"; import * as CommandReceiptStore from "./CommandReceiptStore.ts"; @@ -288,7 +289,14 @@ it.effect.each( ({ target, createdBy }) => { const harness = makeHarness(); const scheduledTasks = ScheduledTasks.layer.pipe( - Layer.provide(Layer.mergeAll(harness.layer, NodeCrypto.layer, Scheduler.layer)), + Layer.provide( + Layer.mergeAll( + harness.layer, + NodeCrypto.layer, + Scheduler.layer, + Layer.mock(SecretRequests.SecretRequests)({}), + ), + ), ); return Effect.gen(function* () { const tasks = yield* ScheduledTasks.ScheduledTaskService; diff --git a/apps/server/src/orchestration-v2/runtimeLayer.ts b/apps/server/src/orchestration-v2/runtimeLayer.ts index 8573c6570795..b4e262d08a70 100644 --- a/apps/server/src/orchestration-v2/runtimeLayer.ts +++ b/apps/server/src/orchestration-v2/runtimeLayer.ts @@ -51,6 +51,7 @@ import { layer as threadLifecycleServiceLayer } from "./ThreadLifecycleService.t import { layer as threadForkServiceLayer } from "./ThreadForkService.ts"; import { layer as turnItemPositionStoreLayer } from "./TurnItemPositionStore.ts"; import { layer as scheduledTaskServiceLayer } from "../scheduledTasks/ScheduledTaskService.ts"; +import * as SecretRequests from "../secrets/SecretRequests.ts"; /** The shared application event log and its command receipts. */ export const OrchestrationEventInfrastructureLayerLive = Layer.mergeAll( @@ -256,8 +257,11 @@ const threadLaunchProvided = threadLaunchServiceLayer.pipe( const threadLifecycleProvided = threadLifecycleServiceLayer.pipe( Layer.provide(threadManagementProvided), ); +const secretRequestsProvided = SecretRequests.layer.pipe(Layer.provide(threadManagementProvided)); const scheduledTaskProvided = scheduledTaskServiceLayer.pipe( - Layer.provide(Layer.mergeAll(threadLaunchProvided, threadManagementProvided)), + Layer.provide( + Layer.mergeAll(threadLaunchProvided, threadManagementProvided, secretRequestsProvided), + ), ); const providerContinuationWorkerProvided = providerContinuationWorkerLive.pipe( Layer.provide( @@ -314,6 +318,7 @@ export const OrchestrationV2ProductionLayerLive = Layer.mergeAll( threadLaunchProvided, threadLifecycleProvided, scheduledTaskProvided, + secretRequestsProvided, UsageLimitRecoveryWorker.workerLive.pipe( Layer.provide(Layer.mergeAll(projectionStoreLayer, threadManagementProvided)), ), diff --git a/apps/server/src/provider/T3OrchestrationInstructions.ts b/apps/server/src/provider/T3OrchestrationInstructions.ts index fbdb48f3ae75..edf6093fa12c 100644 --- a/apps/server/src/provider/T3OrchestrationInstructions.ts +++ b/apps/server/src/provider/T3OrchestrationInstructions.ts @@ -9,7 +9,8 @@ The \`t3-code\` MCP server provides app-owned orchestration. Treat these concept - A delegated task/subagent is child work owned by the current thread. Use \`orchestrator_capabilities\` to discover the current provider/model IDs from the same live catalog as the composer, including configured custom models. Do not treat a native tool's model list as the full list of available subagent models. Prefer native subagent tools for same-provider work only when they support the chosen model. Use \`delegate_task\` with that provider instance and model when native tools cannot, including for same-provider work. Also use \`delegate_task\` for cross-provider or explicitly T3-owned child tasks. Retain each returned \`taskId\`, and use \`task_status\` or \`task_cancel\` to manage it. The returned \`childThreadId\` is backing storage for the subagent, not the target for starting another delegated review round. - \`t3_thread_launch\` and \`create_threads\` create ordinary top-level T3 conversations. Use them only when the user explicitly asks for separate/new/top-level threads or conversations. Never use them merely because the user said "subagent" or requested parallel delegated work. - For every T3 delegated review round, call \`delegate_task\` again. Include the original brief, prior findings, responses, and unresolved objections in each new task prompt. Track each round by its own \`taskId\`. Use a distinct \`clientRequestId\` per round, stable across retries of that round. Do not use \`t3_thread_send\` on \`childThreadId\` to continue a delegated review. -- \`schedule_task\` creates persistent recurring work in the app scheduler. Pass \`schedule\` as a structured object, never as JSON text: \`{"type":"interval","everyMs":3600000}\` for an interval, or \`{"type":"fixed_time","timeOfDay":"09:00","weekdays":[1,2,3,4,5]}\` for a wall-clock schedule, or \`{"type":"webhook"}\` to run on each request to the returned \`webhookUrl\` (the prompt may use \`{{body.path}}\`-style placeholders). By default runs return to the current thread; set \`bindToCurrentThread=false\` only when the user wants a fresh thread for every run. After scheduling, report the returned cadence and next run time. +- \`schedule_task\` creates persistent recurring work in the app scheduler. Pass \`schedule\` as a structured object, never as JSON text: \`{"type":"interval","everyMs":3600000}\` for an interval, or \`{"type":"fixed_time","timeOfDay":"09:00","weekdays":[1,2,3,4,5]}\` for a wall-clock schedule, or \`{"type":"webhook"}\` to run on each request to the returned \`webhookUrl\` (the run sees the request only through \`{{body.path}}\`-style placeholders in the prompt). By default runs return to the current thread, which suits orchestrating: each trigger arrives here and you delegate or dedupe; set \`bindToCurrentThread=false\` only when the user wants a fresh thread for every run. After scheduling a timer, report the returned cadence and next run time; for a webhook, report its \`webhookUrl\`, or say T3 Connect remote access is needed if it is missing. +- When you need a secret from the user (a token, API key, or webhook signing secret), call \`request_secret\` so they enter it privately, then pass the returned \`secretRef\` to the tool that needs it, e.g. \`signature.secretRef\` on a webhook task for a sender that signs requests such as GitHub. A \`secretRef\` works once. Never ask for a secret in chat, never invent one, and never repeat one. ### Choose the workspace before starting a new thread diff --git a/apps/server/src/scheduledTasks/ScheduledTaskService.schedule.test.ts b/apps/server/src/scheduledTasks/ScheduledTaskService.schedule.test.ts index c4cd4a816433..b80135c3132a 100644 --- a/apps/server/src/scheduledTasks/ScheduledTaskService.schedule.test.ts +++ b/apps/server/src/scheduledTasks/ScheduledTaskService.schedule.test.ts @@ -10,6 +10,7 @@ import * as TestClock from "effect/testing/TestClock"; import * as ThreadLaunchService from "../orchestration-v2/ThreadLaunchService.ts"; import * as ThreadManagementService from "../orchestration-v2/ThreadManagementService.ts"; +import * as SecretRequests from "../secrets/SecretRequests.ts"; import { SqlitePersistenceMemory } from "../persistence/Layers/Sqlite.ts"; import * as ScheduledTaskService from "./ScheduledTaskService.ts"; @@ -22,6 +23,7 @@ it.effect("rejects a stale form save after deletion while preserving explicit-id Scheduler.layer, Layer.mock(ThreadLaunchService.ThreadLaunchService)({}), Layer.mock(ThreadManagementService.ThreadManagementService)({}), + Layer.mock(SecretRequests.SecretRequests)({}), ); yield* Effect.gen(function* () { const service = yield* ScheduledTaskService.ScheduledTaskService; @@ -64,6 +66,7 @@ it.effect("preserves a due run when a save only pads the scheduled hour", () => Scheduler.layer, Layer.mock(ThreadLaunchService.ThreadLaunchService)({}), Layer.mock(ThreadManagementService.ThreadManagementService)({}), + Layer.mock(SecretRequests.SecretRequests)({}), ); yield* Effect.gen(function* () { const service = yield* ScheduledTaskService.ScheduledTaskService; diff --git a/apps/server/src/scheduledTasks/ScheduledTaskService.test.ts b/apps/server/src/scheduledTasks/ScheduledTaskService.test.ts index 0095d1436c17..e6e27b676141 100644 --- a/apps/server/src/scheduledTasks/ScheduledTaskService.test.ts +++ b/apps/server/src/scheduledTasks/ScheduledTaskService.test.ts @@ -17,6 +17,7 @@ import * as TestClock from "effect/testing/TestClock"; import * as ThreadLaunchService from "../orchestration-v2/ThreadLaunchService.ts"; import * as ThreadManagementService from "../orchestration-v2/ThreadManagementService.ts"; +import * as SecretRequests from "../secrets/SecretRequests.ts"; import { SqlitePersistenceMemory } from "../persistence/Layers/Sqlite.ts"; import * as ScheduledTaskService from "./ScheduledTaskService.ts"; @@ -211,6 +212,7 @@ it.effect( ), }), Layer.mock(ThreadManagementService.ThreadManagementService)({}), + Layer.mock(SecretRequests.SecretRequests)({}), NodeCrypto.layer, Scheduler.layer, ), @@ -316,6 +318,7 @@ it.effect( ), }), Layer.mock(ThreadManagementService.ThreadManagementService)({}), + Layer.mock(SecretRequests.SecretRequests)({}), NodeCrypto.layer, Scheduler.layer, ), diff --git a/apps/server/src/scheduledTasks/ScheduledTaskService.ts b/apps/server/src/scheduledTasks/ScheduledTaskService.ts index f784484ad19b..51403c124a98 100644 --- a/apps/server/src/scheduledTasks/ScheduledTaskService.ts +++ b/apps/server/src/scheduledTasks/ScheduledTaskService.ts @@ -43,6 +43,7 @@ import * as SqlClient from "effect/sql/SqlClient"; import * as ThreadLaunchService from "../orchestration-v2/ThreadLaunchService.ts"; import * as Metrics from "../observability/Metrics.ts"; import * as ThreadManagementService from "../orchestration-v2/ThreadManagementService.ts"; +import * as SecretRequests from "../secrets/SecretRequests.ts"; import * as Scheduler from "../scheduling/Scheduler.ts"; import { isMissedFixedTimeRun, isSameSchedule, nextScheduledRunAt } from "./Schedule.ts"; import { @@ -393,6 +394,7 @@ export const layer = Layer.effect( const crypto = yield* Crypto.Crypto; const threadLaunch = yield* ThreadLaunchService.ThreadLaunchService; const threadManagement = yield* ThreadManagementService.ThreadManagementService; + const secretRequests = yield* SecretRequests.SecretRequests; const scheduler = yield* Scheduler.Scheduler; const readWebhookOrigin = yield* ScheduledTaskWebhookOrigin; // Webhook deliveries for one task dispatch in arrival order rather than @@ -1042,10 +1044,28 @@ export const layer = Layer.effect( const token = existing?.token ?? (yield* newWebhookToken); const signature = input.schedule.type === "webhook" ? input.schedule.signature : null; - const secret = - signature == null ? null : (signature.secret ?? existing?.secret ?? null); + // A secretRef is a value the user entered for an agent; this + // save consumes it, so it cannot be used again. A replay of a + // save that already used it keeps the secret that save stored. + const replay = input.commandId !== undefined && existing?.secret != null; + const fromRef = + signature?.secretRef === undefined + ? undefined + : yield* secretRequests + .consume({ ref: signature.secretRef, projectId: input.projectId }) + .pipe( + Effect.catch((error) => + replay + ? Effect.succeed(undefined) + : Effect.fail(taskError(error.message, { taskId: id })), + ), + ); + // A ref, when given, is the only source: a plain secret sent + // alongside it must not replace what a replayed save stored. + const provided = signature?.secretRef === undefined ? signature?.secret : fromRef; + const secret = signature == null ? null : (provided ?? existing?.secret ?? null); const secretChanged = - signature == null || signature.secret !== undefined || existing === null; + signature == null || provided !== undefined || existing === null; if (signature != null && secret === null) { return yield* taskError("A webhook signature check needs a signing secret.", { taskId: id, @@ -1608,11 +1628,12 @@ export const layer = Layer.effect( ) : runOutcome("started"), ), - Effect.catchTag("WebhookDeliverySkipped", (skipped) => - runOutcome("skipped").pipe( - Effect.andThen(markDeliveryFailed(deliveryId, skipped.reason)), - ), - ), + Effect.catchTags({ + WebhookDeliverySkipped: (skipped) => + runOutcome("skipped").pipe( + Effect.andThen(markDeliveryFailed(deliveryId, skipped.reason)), + ), + }), // The log is readable over RPC, so it gets a fixed reason; the // cause, which can carry request data, stays in the server log. Effect.catchCause((cause) => diff --git a/apps/server/src/scheduledTasks/ScheduledTaskService.webhook.test.ts b/apps/server/src/scheduledTasks/ScheduledTaskService.webhook.test.ts index 5fdaf45dcf3c..86c96b05901b 100644 --- a/apps/server/src/scheduledTasks/ScheduledTaskService.webhook.test.ts +++ b/apps/server/src/scheduledTasks/ScheduledTaskService.webhook.test.ts @@ -2,7 +2,7 @@ import * as NodeCrypto from "node:crypto"; import * as NodePlatformCrypto from "@effect/platform-node/NodeCrypto"; import { assert, it } from "@effect/vitest"; -import { ScheduledTaskUpsertInput } from "@t3tools/contracts"; +import { ScheduledTaskUpsertInput, SecretRequestError } from "@t3tools/contracts"; import * as DateTime from "effect/DateTime"; import * as Deferred from "effect/Deferred"; import * as Effect from "effect/Effect"; @@ -16,6 +16,7 @@ import * as TestClock from "effect/testing/TestClock"; import * as ThreadLaunchService from "../orchestration-v2/ThreadLaunchService.ts"; import * as ThreadManagementService from "../orchestration-v2/ThreadManagementService.ts"; +import * as SecretRequests from "../secrets/SecretRequests.ts"; import { SqlitePersistenceMemory } from "../persistence/Layers/Sqlite.ts"; import * as Scheduler from "../scheduling/Scheduler.ts"; import * as ScheduledTaskService from "./ScheduledTaskService.ts"; @@ -66,11 +67,14 @@ const withService = ( body: (input: { readonly service: ScheduledTaskService.ScheduledTaskService["Service"]; readonly launches: Queue.Queue; + /** Secrets the user entered for an agent, by ref; consuming one removes it. */ + readonly secretsByRef: Map; }) => Effect.Effect, options: { readonly gate?: Deferred.Deferred; readonly relayHookBaseUrl?: string } = {}, ) => Effect.gen(function* () { const launches = yield* Queue.unbounded(); + const secretsByRef = new Map(); const dependencies = Layer.mergeAll( NodePlatformCrypto.layer, Scheduler.layer, @@ -82,6 +86,15 @@ const withService = ( ), }), Layer.mock(ThreadManagementService.ThreadManagementService)({}), + Layer.mock(SecretRequests.SecretRequests)({ + consume: ({ ref }) => { + const value = secretsByRef.get(ref); + secretsByRef.delete(ref); + return value === undefined + ? Effect.fail(new SecretRequestError({ reason: "ref_unavailable" })) + : Effect.succeed(value); + }, + }), Layer.succeed( ScheduledTaskService.ScheduledTaskWebhookOrigin, Effect.succeed({ relayHookBaseUrl: options.relayHookBaseUrl ?? null }), @@ -89,7 +102,7 @@ const withService = ( ); return yield* Effect.gen(function* () { const service = yield* ScheduledTaskService.ScheduledTaskService; - return yield* body({ service, launches }); + return yield* body({ service, launches, secretsByRef }); }).pipe(Effect.provide(ScheduledTaskService.layer.pipe(Layer.provide(dependencies)))); }).pipe(Effect.provide(SqlitePersistenceMemory)); @@ -708,6 +721,117 @@ it.effect("deleting a task removes its delivery log", () => ), ); +const githubSignature = (secret: string) => + `sha256=${NodeCrypto.createHmac("sha256", secret).update(pullRequestBody).digest("hex")}`; + +it.effect("a signature can take the user's secret by ref, which works only once", () => + withService(({ service, launches, secretsByRef }) => + Effect.gen(function* () { + secretsByRef.set("secret-ref:00000000000000000000000000000001", "github-secret"); + const githubSchedule = (secretRef: string) => ({ + type: "webhook", + signature: { header: "x-hub-signature-256", encoding: "hex", prefix: "sha256=", secretRef }, + }); + const { task } = yield* service.upsert( + yield* webhookTaskInput({ + schedule: githubSchedule("secret-ref:00000000000000000000000000000001"), + }), + ); + assert.isTrue(task.webhook!.hasSecret); + const signed = yield* service.triggerWebhook( + requestFor(task, { + headers: { + "content-type": "application/json", + "x-hub-signature-256": githubSignature("github-secret"), + }, + }), + ); + assert.equal(signed._tag, "accepted"); + yield* Queue.take(launches); + + // The ref was consumed by that save. + const reused = yield* service + .upsert( + yield* webhookTaskInput({ + id: "scheduled-task:other", + schedule: githubSchedule("secret-ref:00000000000000000000000000000001"), + }), + ) + .pipe(Effect.flip); + assert.include(reused.message, "already used"); + }), + ), +); + +it.effect("a retried save with an already used secretRef keeps the stored secret", () => + withService(({ service, launches, secretsByRef }) => + Effect.gen(function* () { + secretsByRef.set("secret-ref:00000000000000000000000000000002", "github-secret"); + const save = webhookTaskInput({ + id: undefined, + commandId: "command:mcp:schedule-task:release-hook", + schedule: { + type: "webhook", + signature: { + header: "x-hub-signature-256", + encoding: "hex", + prefix: "sha256=", + secretRef: "secret-ref:00000000000000000000000000000002", + }, + }, + }); + const first = yield* service.upsert(yield* save); + // The agent never saw the first result, so it sends the same call again, + // this time with a plain secret alongside the used ref. + const retried = yield* service.upsert( + yield* webhookTaskInput({ + id: undefined, + commandId: "command:mcp:schedule-task:release-hook", + schedule: { + type: "webhook", + signature: { + header: "x-hub-signature-256", + encoding: "hex", + prefix: "sha256=", + secretRef: "secret-ref:00000000000000000000000000000002", + secret: "made-up-secret", + }, + }, + }), + ); + assert.equal(retried.task.id, first.task.id); + const signed = yield* service.triggerWebhook( + requestFor(retried.task, { + headers: { + "content-type": "application/json", + "x-hub-signature-256": githubSignature("github-secret"), + }, + }), + ); + assert.equal(signed._tag, "accepted"); + yield* Queue.take(launches); + }), + ), +); + +it.effect("a signature without any secret is refused", () => + withService(({ service }) => + Effect.gen(function* () { + const failure = yield* service + .upsert( + yield* webhookTaskInput({ + schedule: { + type: "webhook", + signature: { header: "x-hub-signature-256", encoding: "hex", prefix: "sha256=" }, + }, + }), + ) + .pipe(Effect.flip); + assert.include(failure.message, "needs a signing secret"); + }), + ), +); + const signatureFor = (secret: string) => `sha256=${NodeCrypto.createHmac("sha256", secret).update(pullRequestBody).digest("hex")}`; diff --git a/apps/server/src/scheduling/Scheduler.integration.test.ts b/apps/server/src/scheduling/Scheduler.integration.test.ts index 0984cee10874..0a9d192ad53b 100644 --- a/apps/server/src/scheduling/Scheduler.integration.test.ts +++ b/apps/server/src/scheduling/Scheduler.integration.test.ts @@ -24,6 +24,7 @@ import * as ThreadManagementService from "../orchestration-v2/ThreadManagementSe import * as UsageLimitRecoveryWorker from "../orchestration-v2/UsageLimitRecoveryWorker.ts"; import { SqlitePersistenceMemory } from "../persistence/Layers/Sqlite.ts"; import * as ScheduledTasks from "../scheduledTasks/ScheduledTaskService.ts"; +import * as SecretRequests from "../secrets/SecretRequests.ts"; import * as ServerSettings from "../serverSettings.ts"; import * as Scheduler from "./Scheduler.ts"; @@ -112,6 +113,7 @@ it.effect.each(["on time", "after restart"])( Layer.mock(ServerSettings.ServerSettingsService)({ getSettings: Effect.succeed(DEFAULT_SERVER_SETTINGS), }), + Layer.mock(SecretRequests.SecretRequests)({}), ); const workers = Layer.mergeAll( ScheduledTasks.layer, diff --git a/apps/server/src/secrets/SecretRequests.test.ts b/apps/server/src/secrets/SecretRequests.test.ts new file mode 100644 index 000000000000..7089c826b374 --- /dev/null +++ b/apps/server/src/secrets/SecretRequests.test.ts @@ -0,0 +1,413 @@ +import * as NodeCrypto from "@effect/platform-node/NodeCrypto"; +import * as NodeServices from "@effect/platform-node/NodeServices"; +import { assert, it } from "@effect/vitest"; +import { + ProjectId, + type OrchestrationV2ServerCommand, + ThreadId, + TurnItemId, +} from "@t3tools/contracts"; +import * as Clock from "effect/Clock"; +import * as Effect from "effect/Effect"; +import * as Layer from "effect/Layer"; +import * as Metric from "effect/Metric"; +import * as TestClock from "effect/testing/TestClock"; +import * as Tracer from "effect/Tracer"; +import * as Option from "effect/Option"; +import * as PlatformError from "effect/PlatformError"; + +import * as ServerSecretStore from "../auth/ServerSecretStore.ts"; +import * as ServerConfig from "../config.ts"; +import * as ThreadManagementService from "../orchestration-v2/ThreadManagementService.ts"; +import * as SecretRequests from "./SecretRequests.ts"; + +const threadId = ThreadId.make("thread-orchestrator"); +const turnItemId = TurnItemId.make("turn-item:secret-request:1"); +const projectId = ProjectId.make("project-1"); + +/** Runs `body` against the service with an in-memory store and a thread holding one request. */ +const withService = ( + body: (input: { + readonly service: SecretRequests.SecretRequests["Service"]; + readonly stored: Map; + readonly dispatched: Array; + }) => Effect.Effect, + options: { + readonly threadId?: ThreadId; + readonly runStatus?: string; + readonly removeFails?: boolean; + /** The agent's wait closes the card just before the answer's record lands. */ + readonly closedFirst?: boolean; + /** How many record dispatches fail before one succeeds. */ + readonly failedRecords?: number; + } = {}, +) => + Effect.gen(function* () { + const stored = new Map(); + const dispatched: Array = []; + let secretStatus = "pending"; + let failedRecords = options.failedRecords ?? 0; + const requestThreadId = options.threadId ?? threadId; + const dependencies = Layer.mergeAll( + NodeCrypto.layer, + NodeServices.layer, + Layer.succeed( + ServerSecretStore.ServerSecretStore, + ServerSecretStore.ServerSecretStore.of({ + // Yields like a real file read, so concurrent callers can interleave. + get: (name) => Effect.yieldNow.pipe(Effect.as(Option.fromNullishOr(stored.get(name)))), + set: (name, value) => Effect.sync(() => void stored.set(name, value)), + create: (name, value) => + stored.has(name) + ? Effect.fail( + new ServerSecretStore.SecretStorePersistError({ + resource: name, + cause: new PlatformError.PlatformError( + new PlatformError.SystemError({ + _tag: "AlreadyExists", + module: "FileSystem", + method: "open", + }), + ), + } as never), + ) + : Effect.sync(() => void stored.set(name, value)), + getOrCreateRandom: (name, bytes) => + Effect.sync(() => { + const existing = stored.get(name); + if (existing) return existing; + const value = new Uint8Array(bytes).fill(7); + stored.set(name, value); + return value; + }), + remove: (name) => + options.removeFails + ? Effect.fail( + new ServerSecretStore.SecretStorePersistError({ + resource: name, + cause: new Error("read-only"), + }), + ) + : Effect.sync(() => void stored.delete(name)), + }), + ), + Layer.mock(ThreadManagementService.ThreadManagementService)({ + getThreadRecords: () => + Effect.succeed({ + thread: { projectId }, + runs: [{ id: "run-1", status: options.runStatus ?? "running" }], + turnItems: [ + { + id: turnItemId, + threadId: requestThreadId, + runId: "run-1", + nodeId: "node-root", + type: "secret_request", + label: "GitHub token", + reason: "Used as GH_TOKEN.", + secretStatus, + }, + ], + } as never), + dispatch: (command) => { + if (command.type === "secret_request.record" && failedRecords > 0) { + failedRecords -= 1; + return Effect.fail(new Error("orchestrator unavailable") as never); + } + return Effect.sync(() => { + dispatched.push(command); + // Like the orchestrator, a card that is no longer pending keeps its answer. + if (options.closedFirst && command.type === "secret_request.record") { + secretStatus = "cancelled"; + } else if (command.type === "secret_request.record" && secretStatus === "pending") { + secretStatus = command.secretStatus; + } + return {} as never; + }); + }, + }), + ); + return yield* Effect.gen(function* () { + const service = yield* SecretRequests.SecretRequests; + return yield* body({ service, stored, dispatched }); + }).pipe(Effect.provide(SecretRequests.layer.pipe(Layer.provide(dependencies)))); + }); + +/** The secret values in the store, leaving out the server's own salt. */ +const valuesOf = (stored: Map) => + Array.from(stored.entries()) + .filter(([name]) => name !== "secret-request-salt") + .map(([, bytes]) => new TextDecoder().decode(bytes)); + +it.effect("a saved answer becomes a one-use ref, and the thread only learns it was saved", () => + withService(({ service, stored, dispatched }) => + Effect.gen(function* () { + yield* service.answer({ + threadId, + turnItemId, + answer: { type: "save", secret: "ghp_secret" }, + }); + assert.equal(dispatched.length, 1); + assert.include(dispatched[0] as object, { + type: "secret_request.record", + secretStatus: "saved", + }); + assert.notInclude(Object.values(dispatched[0] as object).map(String), "ghp_secret"); + + const ref = Option.getOrThrow(yield* service.savedRef({ threadId, turnItemId })); + assert.equal(yield* service.consume({ ref, projectId }), "ghp_secret"); + // Used once: the value is gone from the store and the ref fails. + assert.isFalse(valuesOf(stored).some((value) => value.includes("ghp_secret"))); + const again = yield* service.consume({ ref, projectId }).pipe(Effect.flip); + assert.include(again.message, "already used"); + }), + ), +); + +it.effect("two concurrent uses of one ref hand the value out once", () => + withService(({ service }) => + Effect.gen(function* () { + yield* service.answer({ + threadId, + turnItemId, + answer: { type: "save", secret: "ghp_secret" }, + }); + const ref = Option.getOrThrow(yield* service.savedRef({ threadId, turnItemId })); + const results = yield* Effect.all( + [service.consume({ ref, projectId }), service.consume({ ref, projectId })].map( + Effect.result, + ), + { concurrency: "unbounded" }, + ); + assert.deepEqual(results.map((result) => result._tag).toSorted(), ["Failure", "Success"]); + }), + ), +); + +it.effect("a save that loses to the card closing deletes the value and says so", () => + withService( + ({ service, stored }) => + Effect.gen(function* () { + const error = yield* service + .answer({ threadId, turnItemId, answer: { type: "save", secret: "ghp_secret" } }) + .pipe(Effect.flip); + assert.equal(error.reason, "agent_stopped"); + assert.isFalse(valuesOf(stored).some((value) => value.includes("ghp_secret"))); + }), + { closedFirst: true }, + ), +); + +it.effect("a save whose record failed can be saved again", () => + withService( + ({ service }) => + Effect.gen(function* () { + const failed = yield* service + .answer({ threadId, turnItemId, answer: { type: "save", secret: "ghp_secret" } }) + .pipe(Effect.flip); + assert.equal(failed.reason, "record_failed"); + yield* service.answer({ + threadId, + turnItemId, + answer: { type: "save", secret: "ghp_secret" }, + }); + const ref = Option.getOrThrow(yield* service.savedRef({ threadId, turnItemId })); + assert.equal(yield* service.consume({ ref, projectId }), "ghp_secret"); + }), + { failedRecords: 1 }, + ), +); + +it.effect("a save whose record and cleanup both failed is finished by saving again", () => + withService( + ({ service }) => + Effect.gen(function* () { + yield* service + .answer({ threadId, turnItemId, answer: { type: "save", secret: "ghp_secret" } }) + .pipe(Effect.flip); + yield* service.answer({ + threadId, + turnItemId, + answer: { type: "save", secret: "ghp_secret" }, + }); + assert.isTrue(Option.isSome(yield* service.savedRef({ threadId, turnItemId }))); + }), + { failedRecords: 1, removeFails: true }, + ), +); + +it.effect("a ref only works in the project it was entered for", () => + withService(({ service }) => + Effect.gen(function* () { + yield* service.answer({ + threadId, + turnItemId, + answer: { type: "save", secret: "ghp_secret" }, + }); + const ref = Option.getOrThrow(yield* service.savedRef({ threadId, turnItemId })); + const elsewhere = yield* service + .consume({ ref, projectId: ProjectId.make("project-other") }) + .pipe(Effect.flip); + assert.include(elsewhere.message, "does not exist"); + // A failed attempt from another project does not burn the ref. + assert.equal(yield* service.consume({ ref, projectId }), "ghp_secret"); + }), + ), +); + +it.effect("declining stores nothing, and a request is answered once", () => + withService(({ service, stored, dispatched }) => + Effect.gen(function* () { + yield* service.answer({ threadId, turnItemId, answer: { type: "decline" } }); + assert.deepEqual(valuesOf(stored), []); + assert.isTrue(Option.isNone(yield* service.savedRef({ threadId, turnItemId }))); + const late = yield* service + .answer({ threadId, turnItemId, answer: { type: "save", secret: "ghp_secret" } }) + .pipe(Effect.flip); + assert.include(late.message, "already answered"); + assert.equal(dispatched.length, 1); + assert.deepEqual(valuesOf(stored), []); + }), + ), +); + +it.effect("traces and counts a saved answer without ever recording the value", () => + Effect.gen(function* () { + const spans: Array = []; + const tracer = Tracer.make({ + span: (options) => { + const span = new Tracer.NativeSpan(options); + spans.push(span); + return span; + }, + }); + const counted = Metric.snapshot.pipe( + Effect.map((snapshots) => { + const found = snapshots.find( + (snapshot) => + snapshot.id === "t3_secret_refs_consumed_total" && + snapshot.attributes?.result === "used", + ); + return found?.type === "Counter" ? Number(found.state.count) : 0; + }), + ); + const before = yield* counted; + yield* withService(({ service }) => + Effect.gen(function* () { + yield* service.answer({ + threadId, + turnItemId, + answer: { type: "save", secret: "ghp_secret" }, + }); + const ref = Option.getOrThrow(yield* service.savedRef({ threadId, turnItemId })); + yield* service.consume({ ref, projectId }); + }), + ).pipe(Effect.withTracer(tracer)); + assert.equal((yield* counted) - before, 1); + const recorded = spans.flatMap((span) => [ + span.name, + ...Array.from(span.attributes.values(), String), + ]); + assert.include(recorded, "SecretRequests.answer"); + assert.include(recorded, "SecretRequests.consume"); + assert.isFalse(recorded.some((value) => value.includes("ghp_secret"))); + }), +); + +it.effect("works in threads with long ids, such as a delegated subagent's", () => { + const delegatedThreadId = ThreadId.make( + `thread:delegated-task:command%3Amcp%3A${"a".repeat(36)}%3Adelegate-task%3Arelease-notes-v0.2.0-${"b".repeat(40)}`, + ); + return withService( + ({ service, stored }) => + Effect.gen(function* () { + yield* service.answer({ + threadId: delegatedThreadId, + turnItemId, + answer: { type: "save", secret: "ghp_secret" }, + }); + // Store names never grow with the thread id, so they fit any filesystem. + assert.isTrue(Array.from(stored.keys()).every((name) => name.length < 100)); + const ref = Option.getOrThrow( + yield* service.savedRef({ threadId: delegatedThreadId, turnItemId }), + ); + assert.equal(yield* service.consume({ ref, projectId }), "ghp_secret"); + }), + { threadId: delegatedThreadId }, + ); +}); + +it.effect("refuses an answer once the agent that asked has stopped", () => + withService( + ({ service, stored, dispatched }) => + Effect.gen(function* () { + const late = yield* service + .answer({ threadId, turnItemId, answer: { type: "save", secret: "ghp_secret" } }) + .pipe(Effect.flip); + assert.include(late.message, "has stopped"); + assert.deepEqual(valuesOf(stored), []); + assert.equal(dispatched.length, 0); + }), + { runStatus: "completed" }, + ), +); + +it.effect("drops values nobody used once they expire, and keeps the rest", () => + Effect.gen(function* () { + const store = yield* ServerSecretStore.ServerSecretStore; + const encode = (savedAt: number) => + new TextEncoder().encode( + `{"projectId":"project-1","value":"ghp_secret","savedAt":${savedAt}}`, + ); + yield* TestClock.adjust("2 days"); + const now = yield* Clock.currentTimeMillis; + const stale = `secret-request-${"a".repeat(32)}`; + const fresh = `secret-request-${"b".repeat(32)}`; + yield* store.set(stale, encode(now - 25 * 60 * 60 * 1000)); + yield* store.set(fresh, encode(now)); + yield* store.set("unrelated", new Uint8Array([1])); + + // Building the service runs the first sweep. + yield* Effect.gen(function* () { + yield* SecretRequests.SecretRequests; + }).pipe( + Effect.provide( + SecretRequests.layer.pipe( + Layer.provide(Layer.mock(ThreadManagementService.ThreadManagementService)({})), + ), + ), + Effect.scoped, + ); + + assert.isTrue(Option.isNone(yield* store.get(stale))); + assert.isTrue(Option.isSome(yield* store.get(fresh))); + assert.isTrue(Option.isSome(yield* store.get("unrelated"))); + assert.isTrue(Option.isSome(yield* store.get("secret-request-salt"))); + }).pipe( + Effect.provide( + ServerSecretStore.layer.pipe( + Layer.provideMerge( + ServerConfig.layerTest(process.cwd(), { prefix: "t3-secret-requests-" }), + ), + Layer.provideMerge(NodeServices.layer), + ), + ), + ), +); + +it.effect("a used value that cannot be deleted is not handed out", () => + withService( + ({ service }) => + Effect.gen(function* () { + yield* service.answer({ + threadId, + turnItemId, + answer: { type: "save", secret: "ghp_secret" }, + }); + const ref = Option.getOrThrow(yield* service.savedRef({ threadId, turnItemId })); + const failed = yield* service.consume({ ref, projectId }).pipe(Effect.flip); + assert.include(failed.message, "Could not use"); + }), + { removeFails: true }, + ), +); diff --git a/apps/server/src/secrets/SecretRequests.ts b/apps/server/src/secrets/SecretRequests.ts new file mode 100644 index 000000000000..8c2dacd4dcad --- /dev/null +++ b/apps/server/src/secrets/SecretRequests.ts @@ -0,0 +1,279 @@ +/** + * SecretRequests - secrets an agent asks the user for. + * + * The user's answer goes straight to the server's secret store under a + * one-use SecretRef; orchestration only records the request and its status. + * A tool that needs the value takes the ref and consumes it, so the value + * never reaches the transcript, projections, clients, or model context. + * + * @module SecretRequests + */ +import { + CommandId, + SecretRef, + SecretRequestError, + type SecretRequestFailureReason, + type ProjectId, + type SecretRequestAnswerInput, + type ThreadId, +} from "@t3tools/contracts"; +import * as NodeCrypto from "node:crypto"; + +import * as Clock from "effect/Clock"; +import * as Context from "effect/Context"; +import * as Effect from "effect/Effect"; +import * as FileSystem from "effect/FileSystem"; +import * as Layer from "effect/Layer"; +import * as Option from "effect/Option"; +import * as Schedule from "effect/Schedule"; +import * as Schema from "effect/Schema"; +import * as Semaphore from "effect/Semaphore"; + +import * as ServerSecretStore from "../auth/ServerSecretStore.ts"; +import * as Metrics from "../observability/Metrics.ts"; +import * as ThreadManagementService from "../orchestration-v2/ThreadManagementService.ts"; + +const SECRET_REF_PREFIX = "secret-ref:"; +/** Store name for a ref's value; refs are fixed-length hex, so names stay short and safe. */ +const storeName = (ref: SecretRef) => `secret-request-${ref.slice(SECRET_REF_PREFIX.length)}`; +const REF_PATTERN = /^secret-ref:[0-9a-f]{32}$/; +/** A value nobody used within this long is dropped; the agent can ask again. */ +const SECRET_REF_TTL_MS = 24 * 60 * 60 * 1000; +const STORE_NAME_PATTERN = /^(secret-request-[0-9a-f]{32})\.bin$/; + +/** + * Each request has exactly one ref, derived from where it was asked. Its + * length never depends on the thread id, so the store's file names stay + * within filesystem limits, and the requesting tool finds it without a + * second record. Unguessable without the server's own salt. + */ +const refFor = (salt: string, threadId: ThreadId, turnItemId: string) => + SecretRef.make( + `${SECRET_REF_PREFIX}${NodeCrypto.createHmac("sha256", salt) + .update(`${threadId}\u0000${turnItemId}`) + .digest("hex") + .slice(0, 32)}`, + ); + +/** A ref's value, the project it was entered for, and when it was saved. */ +const StoredSecret = Schema.fromJsonString( + Schema.Struct({ projectId: Schema.String, value: Schema.String, savedAt: Schema.Number }), +); +const encodeStored = Schema.encodeEffect(StoredSecret); +const decodeStored = Schema.decodeUnknownOption(StoredSecret); + +const fail = (reason: SecretRequestFailureReason, cause?: unknown) => + new SecretRequestError({ reason, ...(cause === undefined ? {} : { cause }) }); + +export class SecretRequests extends Context.Service< + SecretRequests, + { + /** + * Answers a pending request in a thread. Saving stores the value under a + * new ref that the requesting tool reads from the request's status. + */ + readonly answer: (input: SecretRequestAnswerInput) => Effect.Effect; + /** The ref minted when this request was saved, for the tool that asked. */ + readonly savedRef: (input: { + readonly threadId: ThreadId; + readonly turnItemId: string; + }) => Effect.Effect>; + /** + * Reads and deletes a ref's value. Fails for unknown or used refs, and + * for refs entered in another project. + */ + readonly consume: (input: { + readonly ref: SecretRef; + readonly projectId: ProjectId; + }) => Effect.Effect; + } +>()("t3/secrets/SecretRequests") {} + +const make = Effect.gen(function* () { + const store = yield* ServerSecretStore.ServerSecretStore; + const fileSystem = yield* FileSystem.FileSystem; + const threadManagement = yield* ThreadManagementService.ThreadManagementService; + + const salt = Buffer.from( + yield* store.getOrCreateRandom("secret-request-salt", 32).pipe(Effect.orDie), + ).toString("hex"); + + const answer: SecretRequests["Service"]["answer"] = (input) => + Effect.gen(function* () { + yield* Effect.annotateCurrentSpan({ + "orchestration_v2.thread_id": input.threadId, + "secret_request.answer": input.answer.type, + }); + const records = yield* threadManagement + .getThreadRecords(input.threadId, ["runs", "turnItems"], { + turnItemTypes: ["secret_request"], + messageRoles: [], + }) + .pipe(Effect.mapError((cause) => fail("load_failed", cause))); + const item = records.turnItems.find((candidate) => candidate.id === input.turnItemId); + if (item?.type !== "secret_request" || item.runId === null || item.nodeId === null) { + return yield* fail("not_found"); + } + if (item.secretStatus !== "pending") { + return yield* fail("already_answered"); + } + // The agent is waiting inside the run that asked; once it has ended, + // nobody will ever receive the ref, so a value saved now would be lost. + const run = records.runs.find((candidate) => candidate.id === item.runId); + if (run === undefined || ThreadManagementService.isTerminalRunStatus(run.status)) { + return yield* fail("agent_stopped"); + } + // Store first: the card only says saved once the value is kept. Create, + // not set: a second answer racing this one must not replace the value. + // A value already there on a card still pending is an earlier save whose + // record failed and could not be cleaned up; recording it finishes that save. + if (input.answer.type === "save") { + const encoded = yield* encodeStored({ + projectId: records.thread.projectId, + value: input.answer.secret, + savedAt: yield* Clock.currentTimeMillis, + }).pipe(Effect.orDie); + yield* store + .create( + storeName(refFor(salt, input.threadId, item.id)), + new TextEncoder().encode(encoded), + ) + .pipe( + Effect.catchIf(ServerSecretStore.isSecretAlreadyExistsError, () => Effect.void), + Effect.mapError((error) => fail("store_failed", error)), + ); + } + const secretStatus = input.answer.type === "save" ? "saved" : "declined"; + yield* threadManagement + .dispatch({ + type: "secret_request.record", + commandId: CommandId.make(`secret-request:${item.id}:${secretStatus}`), + threadId: input.threadId, + runId: item.runId, + nodeId: item.nodeId, + turnItemId: item.id, + label: item.label, + reason: item.reason, + ...(item.placeholder === undefined ? {} : { placeholder: item.placeholder }), + secretStatus, + }) + .pipe( + Effect.mapError((cause) => fail("record_failed", cause)), + // The card still says pending, so the user can save again; the value + // stored above would make that retry look already answered. + Effect.tapError(() => + input.answer.type === "save" + ? removeLogged(storeName(refFor(salt, input.threadId, item.id))) + : Effect.void, + ), + ); + if (input.answer.type !== "save") return; + // A request is answered once: if the agent's wait closed the card between + // the checks above and this record, the record changed nothing. Nobody + // will receive the ref, so the value is deleted rather than left to expire. + const recorded = yield* threadManagement + .getThreadRecords(input.threadId, ["turnItems"], { + turnItemTypes: ["secret_request"], + messageRoles: [], + }) + .pipe(Effect.mapError((cause) => fail("record_failed", cause))); + const card = recorded.turnItems.find((candidate) => candidate.id === item.id); + if (card?.type === "secret_request" && card.secretStatus !== "saved") { + yield* removeLogged(storeName(refFor(salt, input.threadId, item.id))); + return yield* fail("agent_stopped"); + } + }).pipe(Effect.withSpan("SecretRequests.answer")); + + const savedRef: SecretRequests["Service"]["savedRef"] = (input) => + Effect.gen(function* () { + const ref = refFor(salt, input.threadId, input.turnItemId); + const stored = yield* store.get(storeName(ref)).pipe(Effect.orElseSucceed(Option.none)); + return Option.map(stored, () => ref); + }); + + /** Removes a stored value; a failure is logged, since the value is still on disk. */ + const removeLogged = (name: string) => + store.remove(name).pipe( + Effect.as(true), + Effect.catch((error) => + Effect.logWarning("Could not delete a secret request value", { + errorTag: error._tag, + }).pipe(Effect.as(false)), + ), + ); + + // get and remove are separate store calls; one consumer at a time keeps two + // concurrent calls from both reading a ref before either deletes it. + const consumeLock = yield* Semaphore.make(1); + const consume: SecretRequests["Service"]["consume"] = (input) => + consumeRef(input).pipe( + consumeLock.withPermits(1), + Effect.tap(() => Metrics.increment(Metrics.secretRefsConsumedTotal, { result: "used" })), + Effect.tapError(() => + Metrics.increment(Metrics.secretRefsConsumedTotal, { result: "rejected" }), + ), + Effect.withSpan("SecretRequests.consume"), + ); + + const consumeRef = (input: { readonly ref: SecretRef; readonly projectId: ProjectId }) => + Effect.gen(function* () { + if (!REF_PATTERN.test(input.ref)) return yield* fail("invalid_ref"); + const stored = yield* store + .get(storeName(input.ref)) + .pipe(Effect.mapError((cause) => fail("read_failed", cause))); + const decoded = Option.flatMap(stored, (bytes) => + decodeStored(new TextDecoder().decode(bytes)), + ); + if (Option.isNone(decoded) || decoded.value.projectId !== input.projectId) { + return yield* fail("ref_unavailable"); + } + if ((yield* Clock.currentTimeMillis) - decoded.value.savedAt > SECRET_REF_TTL_MS) { + // The hourly sweep retries a removal that fails here. + yield* removeLogged(storeName(input.ref)); + return yield* fail("ref_expired"); + } + // One use: the value moves into whatever consumed it. If it cannot be + // deleted, it is not handed out, so a ref is never used twice. + if (!(yield* removeLogged(storeName(input.ref)))) { + return yield* fail("consume_failed"); + } + return decoded.value.value; + }); + + /** + * Drops values nobody used before they expired, so an agent that never + * consumed its ref does not leave the user's secret on disk. + */ + const sweepExpired = Effect.gen(function* () { + if (store.directory === undefined) return; + const now = yield* Clock.currentTimeMillis; + const names = (yield* fileSystem.readDirectory(store.directory)).flatMap((file) => { + const match = STORE_NAME_PATTERN.exec(file); + return match?.[1] === undefined ? [] : [match[1]]; + }); + let removed = 0; + for (const name of names) { + const stored = yield* store.get(name).pipe(Effect.orElseSucceed(Option.none)); + const decoded = Option.flatMap(stored, (bytes) => + decodeStored(new TextDecoder().decode(bytes)), + ); + if (Option.isSome(decoded) && now - decoded.value.savedAt <= SECRET_REF_TTL_MS) continue; + if (yield* removeLogged(name)) removed += 1; + } + yield* Effect.annotateCurrentSpan({ "secret_request.expired_removed": removed }); + }).pipe( + Effect.catchCause((cause) => Effect.logWarning("Could not sweep expired secret refs", cause)), + Effect.withSpan("SecretRequests.sweepExpired"), + ); + // Once at startup, then hourly; a value lingers at most an hour past expiry. + yield* sweepExpired; + yield* sweepExpired.pipe( + Effect.delay("1 hour"), + Effect.repeat(Schedule.spaced("1 hour")), + Effect.forkScoped, + ); + + return SecretRequests.of({ answer, savedRef, consume }); +}); + +export const layer = Layer.effect(SecretRequests, make); diff --git a/apps/server/src/ws.ts b/apps/server/src/ws.ts index a9c2e73933a7..c049ab24d32a 100644 --- a/apps/server/src/ws.ts +++ b/apps/server/src/ws.ts @@ -121,6 +121,7 @@ import * as ThreadLaunchService from "./orchestration-v2/ThreadLaunchService.ts" import * as ThreadMessageIntake from "./orchestration-v2/ThreadMessageIntake.ts"; import * as IdAllocator from "./orchestration-v2/IdAllocator.ts"; import * as ScheduledTasks from "./scheduledTasks/ScheduledTaskService.ts"; +import * as SecretRequests from "./secrets/SecretRequests.ts"; import { archivedShellStreamItemFromThreadShell, buildActiveShellSnapshot, @@ -1220,6 +1221,7 @@ const makeWsRpcLayer = ( const threadLaunch = yield* ThreadLaunchService.ThreadLaunchService; const providerSessionManager = yield* ProviderSessionManager.ProviderSessionManagerV2; const scheduledTasks = yield* ScheduledTasks.ScheduledTaskService; + const secretRequests = yield* SecretRequests.SecretRequests; const pullRequests = yield* PullRequestService.PullRequestService; const pullRequestSync = yield* PullRequestSyncReactor.PullRequestSyncReactor; const deviceService = yield* DeviceService.DeviceService; @@ -2102,6 +2104,11 @@ const makeWsRpcLayer = ( scheduledTasks.rotateWebhookToken(input), { "rpc.aggregate": "scheduledTasks", "scheduled_task.id": input.id }, ), + [WS_METHODS.secretsAnswerRequest]: (input) => + observeRpcEffect(WS_METHODS.secretsAnswerRequest, secretRequests.answer(input), { + "rpc.aggregate": "secrets", + "orchestration_v2.thread_id": input.threadId, + }), [WS_METHODS.scheduledTasksListWebhookDeliveries]: (input) => observeRpcEffect( WS_METHODS.scheduledTasksListWebhookDeliveries, diff --git a/apps/web/src/components/ChatView.tsx b/apps/web/src/components/ChatView.tsx index 58d7618b76b2..656b989e8595 100644 --- a/apps/web/src/components/ChatView.tsx +++ b/apps/web/src/components/ChatView.tsx @@ -740,6 +740,8 @@ function eventPathContainsSelector(event: Event, selector: string): boolean { return path.some((target) => target instanceof Element && target.closest(selector)); } +const SECRET_REQUEST_SELECTOR = '[data-v2-item-type="secret_request"]'; + /** * Whether input that landed outside any editable or interactive element * should be redirected into the composer. Shared by type-to-focus and @@ -747,6 +749,9 @@ function eventPathContainsSelector(event: Event, selector: string): boolean { */ function shouldRedirectInputToComposer(event: Event): boolean { if (event.defaultPrevented) return false; + // Near a pending secret request, input is meant for its private field: it + // must never land in the composer draft, which is persisted and sent. + if (eventPathContainsSelector(event, SECRET_REQUEST_SELECTOR)) return false; if (eventPathContainsSelector(event, TYPE_TO_FOCUS_EDITABLE_SELECTOR)) return false; if (eventPathContainsSelector(event, TYPE_TO_FOCUS_INTERACTIVE_SELECTOR)) return false; if (document.querySelector(TYPE_TO_FOCUS_FLOATING_LAYER_SELECTOR)) return false; diff --git a/apps/web/src/components/chat/MessagesTimeline.tsx b/apps/web/src/components/chat/MessagesTimeline.tsx index 1417c3000f03..71cbce5a755e 100644 --- a/apps/web/src/components/chat/MessagesTimeline.tsx +++ b/apps/web/src/components/chat/MessagesTimeline.tsx @@ -276,6 +276,7 @@ import { V2LifecycleRow, type HandoffTimelineRun, } from "./V2LifecycleRow"; +import { SecretRequestCard } from "./SecretRequestCard"; import { TimelineSystemDivider } from "./TimelineSystemDivider"; import { SkillChipIcon, SkillInlineText } from "./SkillInlineText"; @@ -2808,6 +2809,15 @@ function V2EventTimelineRow({ row }: { row: Extract 1) { return ; } + if (item.type === "secret_request") { + return ( + + ); + } if (isV2LifecycleItem(item)) { return ( } + label={ + <> + {item.label} · {display.label} + > + } + /> + ); + } + return ; +} + +function PendingSecretRequestForm(props: { + readonly environmentId: EnvironmentId; + readonly item: SecretRequestItem; +}) { + const { item } = props; + const inputId = useId(); + const errorId = useId(); + const privacyId = useId(); + const answer = useAtomCommand(serverEnvironment.answerSecretRequest, { + label: "answer secret request", + // The failure cause holds the request; keep it out of the console. + reportFailure: false, + reportDefect: false, + }); + const [secret, setSecret] = useState(""); + const [submitting, setSubmitting] = useState(false); + const [error, setError] = useState(null); + // Enter then a click can both run before a re-render; this guard is synchronous. + const inFlight = useRef(false); + + const send = async ( + reply: { readonly type: "save"; readonly secret: string } | { readonly type: "decline" }, + ) => { + const input = secretRequestAnswerInput(item, reply); + if (input === null || inFlight.current) return; + inFlight.current = true; + setSubmitting(true); + setError(null); + const result = await answer({ environmentId: props.environmentId, input }).finally(() => { + inFlight.current = false; + setSubmitting(false); + }); + if (result._tag === "Success") { + // The card switches to its answered row once the item updates. + setSecret(""); + return; + } + if (!isAtomCommandInterrupted(result)) { + setError(secretRequestFailureMessage(squashAtomCommandFailure(result))); + } + }; + + const onSubmit = (event: FormEvent) => { + event.preventDefault(); + void send({ type: "save", secret }); + }; + + // Same hierarchy as a chat card: what is asked, why, the field, then the + // promise about where the value goes. + return ( + { + if (!(event.target instanceof HTMLElement)) return; + if (event.target.closest("input, button, a")) return; + document.getElementById(inputId)?.focus(); + }} + > + + + {item.label} + + {item.reason.trim() ? {item.reason} : null} + + + + setSecret(event.currentTarget.value)} + /> + + + Save securely + + + {error !== null ? ( + + {error} + + ) : null} + + + + {SECRET_REQUEST_PRIVACY_NOTE} + + void send({ type: "decline" })} + > + Decline + + + + ); +} diff --git a/apps/web/src/session-logic.ts b/apps/web/src/session-logic.ts index 4f3e8513b5b5..984dc8743a8a 100644 --- a/apps/web/src/session-logic.ts +++ b/apps/web/src/session-logic.ts @@ -315,12 +315,15 @@ const STANDALONE_V2_ITEM_TYPES = new Set([ "fork", "thread_created", + // Still answerable after a steer supersedes the attempt that asked. + "secret_request", ]); export function timelineEntryIsPersistentResourceCard(entry: TimelineEntry): boolean { diff --git a/docs/operations/observability.md b/docs/operations/observability.md index 7a39ac3dd601..3e083c269a7d 100644 --- a/docs/operations/observability.md +++ b/docs/operations/observability.md @@ -396,6 +396,10 @@ Webhooks have their own families: deliveries start, which happen after the sender has its answer. - `t3_webhook_held_delay` for how long requests the relay held waited before arriving. +- `t3_secret_requests_total` by `status` (`saved`, `declined`, `cancelled`, `timed_out`) for secrets + agents asked users for, and `t3_secret_refs_consumed_total` by `result` (`used`, `rejected`) for + tools redeeming them. Neither ever carries a value. + `ScheduledTaskService.triggerWebhook` spans carry the same outcome per request, and each run started from a delivery is its own `ScheduledTaskService.runWebhookDelivery` trace. For a request the relay forwarded, the span also goes to the T3 Connect trace export as a child of the relay's diff --git a/packages/client-runtime/package.json b/packages/client-runtime/package.json index 6f95926cd1d1..13cfabf182b8 100644 --- a/packages/client-runtime/package.json +++ b/packages/client-runtime/package.json @@ -35,6 +35,10 @@ "types": "./src/userMessage.ts", "default": "./src/userMessage.ts" }, + "./secret-request": { + "types": "./src/secretRequest.ts", + "default": "./src/secretRequest.ts" + }, "./scheduled-task-webhook": { "types": "./src/scheduledTaskWebhook.ts", "default": "./src/scheduledTaskWebhook.ts" diff --git a/packages/client-runtime/src/secretRequest.test.ts b/packages/client-runtime/src/secretRequest.test.ts new file mode 100644 index 000000000000..3f8fe48885ca --- /dev/null +++ b/packages/client-runtime/src/secretRequest.test.ts @@ -0,0 +1,75 @@ +import { SecretRequestError, ThreadId, TurnItemId } from "@t3tools/contracts"; +import { describe, expect, it } from "vite-plus/test"; + +import { + secretRequestAnswerInput, + secretRequestDisplay, + secretRequestFailureMessage, + type SecretRequestItem, +} from "./secretRequest.ts"; + +const item = { + id: TurnItemId.make("turn-item:secret-request:1"), + threadId: ThreadId.make("thread-1"), +}; + +describe("secretRequestDisplay", () => { + it("shows the form only while pending and maps every answer to its outcome copy", () => { + const display = (secretStatus: SecretRequestItem["secretStatus"]) => + secretRequestDisplay({ secretStatus }, "local"); + expect(display("pending")).toEqual({ kind: "pending" }); + expect(display("saved")).toEqual({ + kind: "answered", + outcome: "saved", + label: "Saved securely and kept private", + }); + expect(display("declined")).toEqual({ + kind: "answered", + outcome: "declined", + label: "Declined", + }); + expect(display("cancelled")).toEqual({ + kind: "answered", + outcome: "ended", + label: "Request ended", + }); + }); + + it("never offers the form for a request inherited from another thread", () => { + expect(secretRequestDisplay({ secretStatus: "pending" }, "inherited")).toEqual({ + kind: "pending-elsewhere", + label: "Waiting for an answer in the original thread", + }); + expect(secretRequestDisplay({ secretStatus: "saved" }, "inherited")).toMatchObject({ + outcome: "saved", + }); + }); +}); + +describe("secretRequestAnswerInput", () => { + it("refuses a blank save and trims the value the server will store", () => { + expect(secretRequestAnswerInput(item, { type: "save", secret: " " })).toBeNull(); + expect(secretRequestAnswerInput(item, { type: "save", secret: " whsec_1 " })).toEqual({ + threadId: item.threadId, + turnItemId: item.id, + answer: { type: "save", secret: "whsec_1" }, + }); + expect(secretRequestAnswerInput(item, { type: "decline" })).toEqual({ + threadId: item.threadId, + turnItemId: item.id, + answer: { type: "decline" }, + }); + }); +}); + +describe("secretRequestFailureMessage", () => { + it("passes through server errors but hides anything that could echo the payload", () => { + expect( + secretRequestFailureMessage(new SecretRequestError({ reason: "already_answered" })), + ).toBe("This secret request was already answered."); + expect(secretRequestFailureMessage(new Error('Expected string, got "whsec_1"'))).toBe( + "Could not answer the request. Try again.", + ); + expect(secretRequestFailureMessage(undefined)).toBe("Could not answer the request. Try again."); + }); +}); diff --git a/packages/client-runtime/src/secretRequest.ts b/packages/client-runtime/src/secretRequest.ts new file mode 100644 index 000000000000..2bee8bcaf2eb --- /dev/null +++ b/packages/client-runtime/src/secretRequest.ts @@ -0,0 +1,101 @@ +import type { OrchestrationV2TurnItem, SecretRequestAnswerInput } from "@t3tools/contracts"; + +export type SecretRequestItem = Extract< + OrchestrationV2TurnItem, + { readonly type: "secret_request" } +>; + +/** Shown under the field: the one promise the card makes about the value. */ +export const SECRET_REQUEST_PRIVACY_NOTE = "Stored securely, never shown to the agent"; +export const SECRET_REQUEST_DEFAULT_PLACEHOLDER = "Paste the secret"; + +/** What a secret request card shows: the form while pending, otherwise a one-line outcome. */ +export type SecretRequestDisplay = + | { readonly kind: "pending" } + | { readonly kind: "pending-elsewhere"; readonly label: string } + | { + readonly kind: "answered"; + readonly outcome: "saved" | "declined" | "ended"; + readonly label: string; + }; + +const SAVED_DISPLAY: SecretRequestDisplay = { + kind: "answered", + outcome: "saved", + label: "Saved securely and kept private", +}; +const DECLINED_DISPLAY: SecretRequestDisplay = { + kind: "answered", + outcome: "declined", + label: "Declined", +}; +const ENDED_DISPLAY: SecretRequestDisplay = { + kind: "answered", + outcome: "ended", + label: "Request ended", +}; +const PENDING_DISPLAY: SecretRequestDisplay = { kind: "pending" }; +const PENDING_ELSEWHERE_DISPLAY: SecretRequestDisplay = { + kind: "pending-elsewhere", + label: "Waiting for an answer in the original thread", +}; + +/** + * `visibility` is the projected row's: a request inherited from another + * thread (a fork) can only be answered where it was asked. + */ +export function secretRequestDisplay( + item: Pick, + visibility: "local" | "inherited" | "synthetic", +): SecretRequestDisplay { + switch (item.secretStatus) { + case "pending": + return visibility === "local" ? PENDING_DISPLAY : PENDING_ELSEWHERE_DISPLAY; + case "saved": + return SAVED_DISPLAY; + case "declined": + return DECLINED_DISPLAY; + case "cancelled": + return ENDED_DISPLAY; + } +} + +/** + * Builds the RPC payload for an answer. A save with a blank value returns null, + * since the server rejects it; callers keep Save disabled instead. + */ +export function secretRequestAnswerInput( + item: Pick, + answer: { readonly type: "save"; readonly secret: string } | { readonly type: "decline" }, +): SecretRequestAnswerInput | null { + if (answer.type === "decline") { + return { threadId: item.threadId, turnItemId: item.id, answer: { type: "decline" } }; + } + const secret = answer.secret.trim(); + if (secret.length === 0) return null; + return { threadId: item.threadId, turnItemId: item.id, answer: { type: "save", secret } }; +} + +/** Failures whose message is written for the user and never echoes the request payload. */ +const USER_FACING_FAILURE_TAGS = new Set(["SecretRequestError", "EnvironmentAuthorizationError"]); + +/** + * Inline error copy for a failed answer. Only known server errors pass their + * message through: anything else (transport or encoding failures) gets the + * generic copy, so the typed value can never surface in the UI. + */ +export function secretRequestFailureMessage(failure: unknown): string { + if ( + typeof failure === "object" && + failure !== null && + "_tag" in failure && + typeof failure._tag === "string" && + USER_FACING_FAILURE_TAGS.has(failure._tag) && + "message" in failure && + typeof failure.message === "string" && + failure.message.trim().length > 0 + ) { + return failure.message; + } + return "Could not answer the request. Try again."; +} diff --git a/packages/client-runtime/src/state/server.ts b/packages/client-runtime/src/state/server.ts index 1e42b6826b3f..17037af3c986 100644 --- a/packages/client-runtime/src/state/server.ts +++ b/packages/client-runtime/src/state/server.ts @@ -1299,6 +1299,17 @@ export function createServerEnvironmentAtoms( scheduler: configScheduler, concurrency: configConcurrency, }), + // Off the config lane: answering a card must not queue behind settings + // edits. One answer per card at a time. + answerSecretRequest: createEnvironmentRpcCommand(runtime, { + label: "environment-data:server:secrets:answer-request", + tag: WS_METHODS.secretsAnswerRequest, + concurrency: { + mode: "singleFlight", + key: ({ environmentId, input }) => + JSON.stringify([environmentId, input.threadId, input.turnItemId]), + }, + }), refreshUsageRates: createEnvironmentRpcCommand(runtime, { label: "environment-data:server:refresh-usage-rates", tag: WS_METHODS.serverRefreshUsageRates, diff --git a/packages/client-runtime/src/t3ToolSummary.ts b/packages/client-runtime/src/t3ToolSummary.ts index 8d3a2fa54f06..b9a46fa65287 100644 --- a/packages/client-runtime/src/t3ToolSummary.ts +++ b/packages/client-runtime/src/t3ToolSummary.ts @@ -285,6 +285,9 @@ export function summarizeT3ToolCalls( quantity(countEntities(entityIds("requestId")), "pending question request"), ); break; + case "secret-request": + label = phrase("Asked for", "ask for", quantity(selected.length, "secret")); + break; case "worktree-handoff": label = phrase( "Handed off to", diff --git a/packages/contracts/src/baseSchemas.ts b/packages/contracts/src/baseSchemas.ts index 00cb38a572a8..464e4ab5c912 100644 --- a/packages/contracts/src/baseSchemas.ts +++ b/packages/contracts/src/baseSchemas.ts @@ -349,6 +349,9 @@ export type RuntimeRequestId = typeof RuntimeRequestId.Type; export const RuntimeTaskId = makeEntityId("RuntimeTaskId"); export type RuntimeTaskId = typeof RuntimeTaskId.Type; export const ScheduledTaskId = makeEntityId("ScheduledTaskId"); +/** A one-use handle to a secret the user entered for an agent; the agent never sees the value. */ +export const SecretRef = makeEntityId("SecretRef"); +export type SecretRef = typeof SecretRef.Type; export type ScheduledTaskId = typeof ScheduledTaskId.Type; export const ApprovalRequestId = makeEntityId("ApprovalRequestId"); export type ApprovalRequestId = typeof ApprovalRequestId.Type; diff --git a/packages/contracts/src/index.ts b/packages/contracts/src/index.ts index 8690bb1b2390..bdc695a4f6ba 100644 --- a/packages/contracts/src/index.ts +++ b/packages/contracts/src/index.ts @@ -60,3 +60,4 @@ export * from "./worktreeMcp.ts"; export * from "./resourceTelemetry.ts"; export * from "./rpc.ts"; export * from "./worktreeSetup.ts"; +export * from "./secretRequest.ts"; diff --git a/packages/contracts/src/orchestrationV2.test.ts b/packages/contracts/src/orchestrationV2.test.ts index 9eca08fe03ae..e8fd216de980 100644 --- a/packages/contracts/src/orchestrationV2.test.ts +++ b/packages/contracts/src/orchestrationV2.test.ts @@ -205,7 +205,7 @@ describe("orchestration V2 contracts", () => { }); const known = item("item-known", "system_notice", { message: "Hello" }); // A type no build of this client knows, standing in for a newer server's item. - const future = item("item-future", "secret_request", { secretRef: "ref-1" }); + const future = item("item-future", "hologram", { beam: "ref-1" }); const projected = (position: number, turnItem: { readonly id: string }) => ({ position, visibility: "local", diff --git a/packages/contracts/src/orchestrationV2.ts b/packages/contracts/src/orchestrationV2.ts index cb2ba540e50e..c26d5a9327e4 100644 --- a/packages/contracts/src/orchestrationV2.ts +++ b/packages/contracts/src/orchestrationV2.ts @@ -1346,6 +1346,26 @@ export const OrchestrationV2WebSearchResult = Schema.Struct({ }); export type OrchestrationV2WebSearchResult = typeof OrchestrationV2WebSearchResult.Type; +export const OrchestrationV2SecretRequestStatus = Schema.Literals([ + "pending", + "saved", + "declined", + "cancelled", +]); +export type OrchestrationV2SecretRequestStatus = typeof OrchestrationV2SecretRequestStatus.Type; + +/** + * A secret an agent asked the user for. The value never passes through + * orchestration: the item carries only what was asked and how it was answered. + */ +const OrchestrationV2SecretRequestFields = { + type: Schema.Literal("secret_request"), + label: TrimmedNonEmptyString, + reason: Schema.String, + placeholder: Schema.optional(Schema.String), + secretStatus: OrchestrationV2SecretRequestStatus, +} as const; + export const OrchestrationV2TurnItem = Schema.Union([ Schema.Struct({ ...OrchestrationV2TurnItemBaseFields, @@ -1525,6 +1545,10 @@ export const OrchestrationV2TurnItem = Schema.Union([ targetProviderInstanceId: ProviderInstanceId, targetModel: TrimmedNonEmptyString, }), + Schema.Struct({ + ...OrchestrationV2TurnItemBaseFields, + ...OrchestrationV2SecretRequestFields, + }), Schema.Struct({ ...OrchestrationV2TurnItemBaseFields, type: Schema.Literal("subagent"), @@ -2296,6 +2320,10 @@ export const OrchestrationV2TurnItemJson = Schema.Union([ targetProviderInstanceId: ProviderInstanceId, targetModel: TrimmedNonEmptyString, }), + Schema.Struct({ + ...OrchestrationV2TurnItemJsonBaseFields, + ...OrchestrationV2SecretRequestFields, + }), Schema.Struct({ ...OrchestrationV2TurnItemJsonBaseFields, type: Schema.Literal("subagent"), @@ -3013,6 +3041,7 @@ export const OrchestrationV2Command = Schema.Union([ targetThreadId: ThreadId, targetRunId: Schema.NullOr(RunId), }), + Schema.Struct({ type: Schema.Literal("provider.switch"), commandId: CommandId, @@ -3081,6 +3110,22 @@ const OrchestrationV2InternalCommand = Schema.Union([ threadId: ThreadId, reason: Schema.optional(Schema.String), }), + /** + * Records or updates a secret an agent asked the user for. Internal so no + * client can mark a request saved without the value being stored. + */ + Schema.Struct({ + type: Schema.Literal("secret_request.record"), + commandId: CommandId, + threadId: ThreadId, + runId: RunId, + nodeId: NodeId, + turnItemId: TurnItemId, + label: TrimmedNonEmptyString, + reason: Schema.String, + placeholder: Schema.optional(Schema.String), + secretStatus: OrchestrationV2SecretRequestStatus, + }), ]); export type OrchestrationV2InternalCommand = typeof OrchestrationV2InternalCommand.Type; diff --git a/packages/contracts/src/orchestratorMcp.ts b/packages/contracts/src/orchestratorMcp.ts index 776cf4be52d5..0a9179c2a8f4 100644 --- a/packages/contracts/src/orchestratorMcp.ts +++ b/packages/contracts/src/orchestratorMcp.ts @@ -12,6 +12,7 @@ import { ProjectId, RunId, ScheduledTaskId, + SecretRef, ThreadId, TrimmedNonEmptyString, TurnItemId, @@ -541,8 +542,14 @@ export const OrchestratorMcpScheduledTask = Schema.Struct({ schedule: ScheduledTaskSchedule, nextRunAt: Schema.NullOr(IsoDateTime), lastRunStatus: ScheduledTaskRunStatus, - /** For webhook tasks: the public T3 Connect URL, or the environment-relative path when the environment is not linked. */ - webhookUrl: Schema.optional(Schema.String), + /** For webhook tasks: the public T3 Connect URL. Absent when this environment has no managed tunnel. */ + webhookUrl: Schema.optional(Schema.String).annotate({ + description: + "Public URL to give the sender. Absent when this environment has no T3 Connect managed tunnel; the user must enable T3 Connect remote access first.", + }), + webhookSignature: Schema.optional(Schema.Literals(["none", "set"])).annotate({ + description: "Whether requests must carry a valid signature.", + }), }); export type OrchestratorMcpScheduledTask = typeof OrchestratorMcpScheduledTask.Type; @@ -577,6 +584,44 @@ export const OrchestratorMcpUpdateScheduledTaskInput = Schema.Struct({ export type OrchestratorMcpUpdateScheduledTaskInput = typeof OrchestratorMcpUpdateScheduledTaskInput.Type; +export const OrchestratorMcpRequestSecretInput = Schema.Struct({ + label: TrimmedNonEmptyString.annotate({ + description: "What you need, shown as the card's title, e.g. 'GitHub webhook secret'.", + }), + reason: TrimmedNonEmptyString.annotate({ + description: + "One or two sentences on what it is for and where the user gets or also enters it.", + }), + placeholder: Schema.optional(TrimmedNonEmptyString).annotate({ + description: "Hint inside the input, e.g. 'Paste your GitHub token'.", + }), + timeoutMs: Schema.optional( + Schema.Int.check(Schema.isBetween({ minimum: 1_000, maximum: 60 * 60 * 1_000 })), + ).annotate({ description: "How long to wait for the user. Default 10 minutes." }), + clientRequestId: Schema.optional(OrchestratorMcpClientRequestId).annotate({ + description: + "Reuse when retrying a call that lost its result, so the user sees one card and its answer is returned again. Use a new id to ask again after timed_out or cancelled.", + }), +}); +export type OrchestratorMcpRequestSecretInput = typeof OrchestratorMcpRequestSecretInput.Type; + +export const OrchestratorMcpRequestSecretResult = Schema.Union([ + Schema.Struct({ + status: Schema.Literal("saved").annotate({ description: "secretRef holds the value." }), + secretRef: SecretRef.annotate({ + description: + "Pass it to a tool that accepts a secretRef; it works once, and you never see the value.", + }), + }), + Schema.Struct({ + status: Schema.Literals(["declined", "cancelled", "timed_out"]).annotate({ + description: + "declined: the user chose not to. cancelled: the request ended with the run. timed_out: the user did not answer in time; the card is closed, so ask again with a new clientRequestId if still needed.", + }), + }), +]); +export type OrchestratorMcpRequestSecretResult = typeof OrchestratorMcpRequestSecretResult.Type; + export const OrchestratorMcpDeleteScheduledTaskInput = Schema.Struct({ scheduledTaskId: ScheduledTaskId, }); diff --git a/packages/contracts/src/rpc.ts b/packages/contracts/src/rpc.ts index 575ad6e32517..74cfd8e9abc7 100644 --- a/packages/contracts/src/rpc.ts +++ b/packages/contracts/src/rpc.ts @@ -320,6 +320,7 @@ import { ScheduledTaskUpsertInput, ScheduledTaskMutationResult, } from "./scheduledTask.ts"; +import { SecretRequestAnswerInput, SecretRequestError } from "./secretRequest.ts"; import { ProjectCloneActionInput, ProjectCloneActionResult, @@ -480,6 +481,7 @@ export const WS_METHODS = { scheduledTasksDelete: "scheduledTasks.delete", scheduledTasksRunNow: "scheduledTasks.runNow", scheduledTasksRotateWebhookToken: "scheduledTasks.rotateWebhookToken", + secretsAnswerRequest: "secrets.answerRequest", scheduledTasksListWebhookDeliveries: "scheduledTasks.listWebhookDeliveries", scheduledTasksGetWebhookDelivery: "scheduledTasks.getWebhookDelivery", @@ -1685,6 +1687,11 @@ const WsScheduledTasksRotateWebhookTokenRpc = Rpc.make( }, ); +const WsSecretsAnswerRequestRpc = Rpc.make(WS_METHODS.secretsAnswerRequest, { + payload: SecretRequestAnswerInput, + error: Schema.Union([SecretRequestError, EnvironmentAuthorizationError]), +}); + const WsScheduledTasksListWebhookDeliveriesRpc = Rpc.make( WS_METHODS.scheduledTasksListWebhookDeliveries, { @@ -1789,6 +1796,7 @@ export const WsRpcGroup = RpcGroup.make( WsScheduledTasksDeleteRpc, WsScheduledTasksRunNowRpc, WsScheduledTasksRotateWebhookTokenRpc, + WsSecretsAnswerRequestRpc, WsScheduledTasksListWebhookDeliveriesRpc, WsScheduledTasksGetWebhookDeliveryRpc, WsServerReportClientActivityRpc, diff --git a/packages/contracts/src/scheduledTask.ts b/packages/contracts/src/scheduledTask.ts index 3140ff8a79cf..126551071278 100644 --- a/packages/contracts/src/scheduledTask.ts +++ b/packages/contracts/src/scheduledTask.ts @@ -6,6 +6,7 @@ import { IsoDateTime, ProjectId, ScheduledTaskId, + SecretRef, ThreadId, TrimmedNonEmptyString, } from "./baseSchemas.ts"; @@ -106,6 +107,10 @@ const ScheduledTaskUpsertWebhookSchedule = Schema.Struct({ secret: Schema.optional(TrimmedNonEmptyString).annotate({ description: "Shared signing secret. Omit to keep the stored secret.", }), + secretRef: Schema.optional(SecretRef).annotate({ + description: + "A secret the user entered through request_secret, used instead of secret. It is consumed by this save.", + }), }), ), ).annotate({ diff --git a/packages/contracts/src/secretRequest.ts b/packages/contracts/src/secretRequest.ts new file mode 100644 index 000000000000..f76a84f6c628 --- /dev/null +++ b/packages/contracts/src/secretRequest.ts @@ -0,0 +1,52 @@ +import * as Schema from "effect/Schema"; + +import { ThreadId, TrimmedNonEmptyString, TurnItemId } from "./baseSchemas.ts"; + +/** + * The user's answer to an agent's request for a secret. A saved value is kept + * by the server under a one-use SecretRef; the thread learns only the status. + */ +export const SecretRequestAnswerInput = Schema.Struct({ + threadId: ThreadId, + turnItemId: TurnItemId, + answer: Schema.Union([ + Schema.Struct({ type: Schema.Literal("save"), secret: TrimmedNonEmptyString }), + Schema.Struct({ type: Schema.Literal("decline") }), + ]), +}); +export type SecretRequestAnswerInput = typeof SecretRequestAnswerInput.Type; + +const SECRET_REQUEST_FAILURE_MESSAGES = { + load_failed: "Could not load the secret request.", + not_found: "This secret request no longer exists.", + already_answered: "This secret request was already answered.", + agent_stopped: "The agent that asked has stopped, so this secret can't be used.", + store_failed: "Could not store the secret.", + record_failed: "Saved the secret, but could not update the request.", + invalid_ref: "That secretRef is not valid.", + read_failed: "Could not read the secret.", + ref_unavailable: + "That secretRef was already used or does not exist. Ask the user again with request_secret.", + ref_expired: "That secretRef expired. Ask the user again with request_secret.", + consume_failed: "Could not use that secretRef. Try again.", +} as const; + +export const SecretRequestFailureReason = Schema.Literals( + Object.keys(SECRET_REQUEST_FAILURE_MESSAGES) as Array< + keyof typeof SECRET_REQUEST_FAILURE_MESSAGES + >, +); +export type SecretRequestFailureReason = typeof SecretRequestFailureReason.Type; + +/** Answering a request or using its ref failed; the message is shown to users and agents. */ +export class SecretRequestError extends Schema.TaggedError()( + "SecretRequestError", + { + reason: SecretRequestFailureReason, + cause: Schema.optionalKey(Schema.Defect()), + }, +) { + override get message(): string { + return SECRET_REQUEST_FAILURE_MESSAGES[this.reason]; + } +} diff --git a/packages/shared/src/t3McpToolPresentation.ts b/packages/shared/src/t3McpToolPresentation.ts index ca3928112c02..2b7366392f23 100644 --- a/packages/shared/src/t3McpToolPresentation.ts +++ b/packages/shared/src/t3McpToolPresentation.ts @@ -38,6 +38,7 @@ export type T3McpToolSummaryAction = | "question-list" | "question-read" | "question-respond" + | "secret-request" | "worktree-handoff" | "worktree-list" | "worktree-status" @@ -130,6 +131,7 @@ const T3_MCP_TOOLS: Readonly> = { ["Delete", "Deleting", "Requested deletion of", "a scheduled task"], "schedule-delete", ), + request_secret: tool(["Ask for", "Asking for", "Asked for", "a secret"], "secret-request"), create_threads: tool(["Create", "Creating", "Created", "T3 threads"], "thread-create"), t3_thread_start: tool(["Start", "Starting", "Started", "a T3 thread"], "thread-create"), t3_thread_list: tool(["List", "Listing", "Listed", "T3 threads"], "thread-list"),
{item.reason}
+ {error} +
+ + {SECRET_REQUEST_PRIVACY_NOTE} +