Repository navigation
Expand file tree
/
Copy paththreadAccess.ts
More file actions
204 lines (190 loc) · 7.33 KB
/
Copy paththreadAccess.ts
File metadata and controls
204 lines (190 loc) · 7.33 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
import type { ProjectionRecordField } from "../orchestration-v2/ProjectionStore.ts";
import {
CommandId,
OrchestratorMcpFailure,
type ProjectId,
type ProviderInteractionMode,
type RuntimeMode,
type ThreadId,
type OrchestrationV2ThreadShell,
} from "@t3tools/contracts";
import * as Crypto from "effect/Crypto";
import * as Effect from "effect/Effect";
import type { OrchestratorV2Error } from "../orchestration-v2/Orchestrator.ts";
import * as ThreadManagement from "../orchestration-v2/ThreadManagementService.ts";
import * as OrchestrationMcp from "./OrchestratorMcpService.ts";
import * as McpInvocationContext from "./McpInvocationContext.ts";
export const unavailable = () =>
new OrchestratorMcpFailure({
code: "orchestration_error",
message: "The operation could not be completed.",
});
/** Decider string rejections are public; wrapped storage and hydration causes are not. */
export const dispatchFailure = (error: OrchestratorV2Error) =>
(error._tag === "OrchestratorDispatchError" ||
error._tag === "OrchestratorCommandRejectedError") &&
typeof error.cause === "string" &&
error.cause.length > 0
? new OrchestratorMcpFailure({
code: "orchestration_error",
message: Array.from(error.cause).slice(0, 1000).join(""),
})
: unavailable();
/**
* The most a caller may hand to the threads it targets. A thread caller is
* capped by its own thread's modes; an OAuth client by the ceiling chosen
* when it was approved.
*/
export interface CallerLimits {
readonly runtimeMode: RuntimeMode;
readonly interactionMode: ProviderInteractionMode;
}
export interface Caller {
readonly scope: McpInvocationContext.McpInvocationScope;
readonly threads: ThreadManagement.ThreadManagementService["Service"];
/** The calling thread, absent for a client signed in from outside T3. */
readonly caller: OrchestrationV2ThreadShell | undefined;
readonly limits: CallerLimits;
}
export const readCaller = Effect.fn("mcp.readCaller")(function* () {
const scope = yield* McpInvocationContext.McpInvocationContext;
if (!scope.capabilities.has("orchestration")) {
return yield* new OrchestratorMcpFailure({
code: "capability_denied",
message: "This credential cannot control threads.",
});
}
return yield* loadCaller();
});
/**
* The caller and its limits, whichever tools its credential grants. Tools
* check their own capability; `McpToolAccess` uses this for every tool.
*/
export const loadCaller = Effect.fn("mcp.loadCaller")(function* () {
const scope = yield* McpInvocationContext.McpInvocationContext;
const threads = yield* ThreadManagement.ThreadManagementService;
if (scope.thread === undefined) {
return {
scope,
threads,
caller: undefined,
limits: {
runtimeMode: McpInvocationContext.clientRuntimeModeCeiling(scope.client),
interactionMode: "default",
},
} satisfies Caller;
}
const caller = yield* threads
.getThreadShell(scope.thread.threadId)
.pipe(Effect.mapError(unavailable));
if (caller === null || caller.deletedAt !== null) {
return yield* new OrchestratorMcpFailure({
code: "thread_not_found",
message: "The calling thread was not found.",
});
}
return {
scope,
threads,
caller,
limits: { runtimeMode: caller.runtimeMode, interactionMode: caller.interactionMode },
} satisfies Caller;
});
/**
* A caller may change another thread only if that thread runs within the
* caller's own modes. Its own thread is always within them.
*/
export const assertTargetWithinLimits = (
limits: CallerLimits,
target: { readonly runtimeMode: RuntimeMode; readonly interactionMode: ProviderInteractionMode },
) =>
OrchestrationMcp.resolveRuntimeMode(limits.runtimeMode, target.runtimeMode).pipe(
Effect.andThen(
OrchestrationMcp.resolveInteractionMode(limits.interactionMode, target.interactionMode),
),
Effect.asVoid,
);
/**
* A thread caller acts only while it owns a live run of a thread that is not
* archived, so a provider token that outlived its session cannot act. An
* OAuth client has no run; its session and ceiling are its authority.
*/
export function assertLiveCaller({ caller, scope }: Caller) {
if (caller === undefined) return Effect.void;
return caller.archivedAt !== null ||
caller.activeRunId === null ||
caller.providerInstanceId !== scope.thread?.providerInstanceId
? Effect.fail(
new OrchestratorMcpFailure({
code: "parent_not_active",
message: "The calling provider no longer owns an active thread run.",
}),
)
: Effect.void;
}
/**
* Actions that change the environment itself (projects, preferences) need full
* access: a thread caller in full-access/default mode, or a client approved
* with a full-access ceiling.
*/
export const assertFullAccess = (context: Caller, message: string) =>
context.limits.runtimeMode === "full-access" && context.limits.interactionMode === "default"
? Effect.void
: Effect.fail(new OrchestratorMcpFailure({ code: "capability_denied", message }));
/** A target project: the one passed, else the calling thread's. */
export const resolveProjectId = (context: Caller, projectId: ProjectId | undefined) =>
projectId !== undefined
? Effect.succeed(projectId)
: context.caller !== undefined
? Effect.succeed(context.caller.projectId)
: Effect.fail(
new OrchestratorMcpFailure({
code: "target_required",
message: "Pass projectId: this MCP client is not running inside a T3 thread.",
}),
);
/** A target thread: the one passed, else the calling thread. */
const resolveThreadId = (context: Caller, threadId: ThreadId | undefined) =>
threadId !== undefined
? Effect.succeed(threadId)
: context.caller !== undefined
? Effect.succeed(context.caller.id)
: Effect.fail(
new OrchestratorMcpFailure({
code: "target_required",
message: "Pass threadId: this MCP client is not running inside a T3 thread.",
}),
);
/** Load a target thread anywhere in the environment; an omitted id means the calling thread. */
export const readThread = Effect.fn("mcp.readThread")(function* <
K extends ProjectionRecordField = never,
>(threadId?: ThreadId, fields: ReadonlyArray<K> = []) {
const context = yield* readCaller();
const targetId = yield* resolveThreadId(context, threadId);
const shell = yield* context.threads.getThreadShell(targetId).pipe(Effect.mapError(unavailable));
if (shell === null || shell.deletedAt !== null) {
return yield* new OrchestratorMcpFailure({
code: "thread_not_found",
message: "The thread was not found.",
});
}
const projection = yield* context.threads
.getProjectThreadRecords({ projectId: shell.projectId, threadId: targetId }, fields, {
turnItemTypes: ["user_input_request"],
})
.pipe(
Effect.mapError((error) =>
error._tag === "ThreadManagementThreadNotFoundError"
? new OrchestratorMcpFailure({
code: "thread_not_found",
message: "The thread was not found.",
})
: unavailable(),
),
);
return { ...context, projection };
});
export const newCommandId = Effect.fn("mcp.newCommandId")(function* () {
const crypto = yield* Crypto.Crypto;
return CommandId.make(`mcp:${yield* crypto.randomUUIDv4.pipe(Effect.orDie)}`);
});