Repository navigation
Expand file tree
/
Copy pathOpenCode2AdapterV2.ts
More file actions
4339 lines (4210 loc) · 176 KB
/
Copy pathOpenCode2AdapterV2.ts
File metadata and controls
4339 lines (4210 loc) · 176 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
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
879
880
881
882
883
884
885
886
887
888
889
890
891
892
893
894
895
896
897
898
899
900
901
902
903
904
905
906
907
908
909
910
911
912
913
914
915
916
917
918
919
920
921
922
923
924
925
926
927
928
929
930
931
932
933
934
935
936
937
938
939
940
941
942
943
944
945
946
947
948
949
950
951
952
953
954
955
956
957
958
959
960
961
962
963
964
965
966
967
968
969
970
971
972
973
974
975
976
977
978
979
980
981
982
983
984
985
986
987
988
989
990
991
992
993
994
995
996
997
998
999
1000
/**
* The OpenCode 2 runtime behind the `opencode` driver. It talks to the
* instance's `opencode serve` process through the HTTP client and reads that
* server's `/api/event` stream, routed here by session id.
*
* A turn is one `session.prompt` (or `session.command`, or `session.compact`
* for `/compact`); the session's next `session.execution.*` terminal ends it.
* A steer is another prompt with `delivery: "steer"`, which OpenCode reads at
* the running execution's next step boundary, so it ends with the turn it
* joined. Fork and rollback cut the native history before the user message a
* turn prompted with (its `nativeTurnRef`). Each runtime mode is a set of
* session permission rules, plan mode is also OpenCode's `plan` agent
* (switched before the prompt like the model), and OpenCode's permission asks
* and question forms become runtime requests on the asking session's thread, a
* subagent's on its parent's.
*
* A `subagent` call runs in a child session shown as a subagent thread. A
* background one outlives its turn; when it ends, OpenCode starts a parent
* execution T3 did not ask for, which waits here for the continuation turn T3
* opens for it.
*
* @module orchestration-v2/Adapters/OpenCode2AdapterV2
*/
import {
AbsolutePath,
Agent,
Form,
Location,
Model,
Permission,
Provider,
Session,
SessionMessage,
Skill,
type OpenCodeEvent,
} from "@opencode/client/effect";
import { Mcp } from "@opencode/schema/mcp";
import {
DESKTOP_MCP_SERVER_NAME,
isOrchestrationV2WorkActive,
type OrchestrationV2AppThread,
type OrchestrationV2ConversationMessage,
type OrchestrationV2ExecutionNode,
type OrchestrationV2ProviderCapabilities,
type OrchestrationV2ProviderSession,
type OrchestrationV2ProviderThread,
type OrchestrationV2ProviderTurn,
type OrchestrationV2RuntimeRequest,
type OrchestrationV2Subagent,
type OrchestrationV2TurnItem,
type OrchestrationV2UserInputQuestion,
type ModelSelection,
type ProviderApprovalDecision,
type ProviderInstanceId,
type RunId,
type RuntimeRequestId,
} from "@t3tools/contracts";
import type * as Cause from "effect/Cause";
import * as DateTime from "effect/DateTime";
import * as Deferred from "effect/Deferred";
import * as Crypto from "effect/Crypto";
import * as Effect from "effect/Effect";
import * as Hex from "effect/encoding/Hex";
import * as Exit from "effect/Exit";
import * as Option from "effect/Option";
import * as Queue from "effect/Queue";
import * as Schedule from "effect/Schedule";
import * as Scope from "effect/Scope";
import * as Schema from "effect/Schema";
import * as Semaphore from "effect/Semaphore";
import * as Stream from "effect/Stream";
import * as ServerConfig from "../../config.ts";
import { makeOptionalResolveEnabledDesktopMcp } from "../../desktopControl/desktopMcpLaunch.ts";
import { paginate, type OpenCode2StreamEvent } from "../../provider/opencode2/OpenCode2Client.ts";
import * as OpenCode2Server from "../../provider/opencode2/OpenCode2Server.ts";
import {
parseOpenCodeModelSlug,
type OpenCodeRuntimeError,
} from "../../provider/opencodeRuntime.ts";
import * as McpProviderSession from "@t3tools/provider-core/server/mcpSession";
import { buildRuntimeInstructions } from "../../provider/RuntimeInstructions.ts";
import { t3OrchestrationSystemPrompt } from "@t3tools/provider-core/server/orchestrationInstructions";
import { SKILL_MENTION_PATTERN } from "@t3tools/shared/composerInlineTokens";
import * as KeyedLock from "@t3tools/shared/KeyedLock";
import { getModelSelectionStringOptionValue, modelSelectionsEqual } from "@t3tools/shared/model";
import { causeErrorTag } from "@t3tools/shared/observability";
import { providerMessageTextWithAttachmentPaths } from "../AttachmentPrompt.ts";
import * as IdAllocator from "@t3tools/provider-core/server/IdAllocator";
import { backgroundWorkNotification, type BackgroundWorkReport } from "../Notification.ts";
import * as ProviderContinuationRequests from "@t3tools/provider-core/server/continuationRequests";
import { makeProviderFailure } from "@t3tools/provider-core/server/failure";
import {
makeSubagentChildThread,
makeSubagentConversationArtifacts,
subagentThreadTitle,
} from "../SubagentProjection.ts";
import * as ProviderAdapter from "@t3tools/provider-core/server/ProviderAdapter";
import { turnScopedSelectionTransition } from "@t3tools/provider-core/server/selectionTransition";
import { OPENCODE_PROVIDER, openCodePermissionRequestKind } from "./OpenCodeAdapterV2.ts";
import { openCodeToolTurnItem } from "./OpenCodeToolItems.ts";
const OpenCode2ProviderCapabilities = {
sessions: {
// One server serves every location, so one session runtime owns them all.
supportsMultipleProviderThreadsPerSession: true,
supportsModelSwitchInSession: true,
supportsProviderSwitchingViaHandoff: true,
// A mode change rewrites the session's rules when its next turn resumes it.
supportsRuntimeModeSwitchInSession: true,
pendingRequestsSurviveRestart: false,
},
threads: {
canCreateEmptyThread: true,
canReadThreadSnapshot: true,
canRollbackThread: true,
canForkThread: true,
canForkFromTurn: true,
canForkFromSubagentThread: false,
exposesNativeThreadId: true,
},
turns: {
exposesNativeTurnId: false,
emitsTurnStarted: true,
emitsTurnCompleted: true,
supportsInterrupt: true,
supportsActiveSteering: true,
supportsSteeringByInterruptRestart: true,
// T3 holds queued messages and starts each as its own turn once the one
// before it ends, so OpenCode's own `queue` delivery is never used.
supportsQueuedMessages: true,
terminalStatusQuality: "strong",
},
streaming: {
streamsAssistantText: true,
streamsReasoning: true,
streamsToolOutput: false,
streamsPlanText: false,
emitsMessageCompleted: true,
},
tools: {
exposesToolItemIds: true,
emitsToolStarted: true,
emitsToolCompleted: true,
emitsToolOutput: true,
supportsMcpTools: false,
supportsDynamicToolCallbacks: false,
},
approvals: {
supportsCommandApproval: true,
supportsFileReadApproval: true,
supportsFileChangeApproval: true,
supportsApplyPatchApproval: true,
approvalsHaveNativeRequestIds: true,
approvalCallbacksAreLiveOnly: true,
approvalsCanOriginateFromSubagents: true,
},
planning: {
emitsPlanUpdated: false,
emitsTodoList: false,
emitsProposedPlan: false,
supportsStructuredQuestions: true,
planDeltasHaveItemIds: false,
},
subagents: {
supportsSubagents: true,
exposesSubagentThreadIds: true,
emitsSubagentLifecycle: true,
canWaitForSubagents: true,
canCloseSubagents: false,
canForkSubagentThread: false,
},
context: {
acceptsSystemContext: false,
acceptsDeveloperContext: false,
acceptsSyntheticUserContext: true,
canGenerateSummaries: false,
canConsumeHandoffSummaries: true,
supportsDeltaHandoff: true,
supportsFullThreadHandoff: true,
maxRecommendedHandoffChars: null,
},
checkpointing: {
appCanCheckpointFilesystem: true,
supportsNestedCheckpointScopes: true,
providerCanRollbackConversation: true,
providerRollbackReturnsSnapshot: true,
providerCanReadConversationSnapshot: true,
},
identity: {
nativeThreadIds: "strong",
nativeTurnIds: "weak",
nativeItemIds: "strong",
nativeRequestIds: "strong",
},
// OpenCode enforces each runtime mode through the session's permission rules.
runtimePolicy: { enforcement: "native" },
} satisfies OrchestrationV2ProviderCapabilities;
type EventOf<T extends OpenCodeEvent["type"]> = Extract<OpenCodeEvent, { readonly type: T }>;
type Tokens = EventOf<"session.step.ended">["data"]["tokens"];
/** What a turn reports against: a run's own turn, or a subagent session's runless one. */
interface TurnOwner {
readonly appThread: OrchestrationV2AppThread;
readonly threadId: OrchestrationV2AppThread["id"];
readonly runId: RunId | null;
readonly runOrdinal: number;
readonly rootNodeId: ProviderAdapter.ProviderAdapterV2TurnInput["rootNodeId"];
readonly modelSelection: ModelSelection;
/** Where the turn runs: a subagent runs in its parent turn's directory. */
readonly runtimePolicy: Pick<ProviderAdapter.ProviderAdapterV2RuntimePolicy, "cwd">;
}
interface ActiveTurn {
readonly input: TurnOwner;
providerTurn: OrchestrationV2ProviderTurn;
/** Open text and reasoning blocks, keyed `<assistantMessageID>:<kind>:<ordinal>`. */
readonly texts: Map<string, OpenBlock>;
readonly tools: Map<string, { readonly name: string; input: Record<string, unknown> }>;
readonly startedAt: Map<string, DateTime.Utc>;
readonly ordinals: Map<string, number>;
nextOrdinal: number;
/** Input includes cache reads and writes, as 1.x reports it; the parts are also kept apart. */
readonly usage: {
input: number;
cached: number;
cacheWrite: number;
output: number;
reasoning: number;
};
steps: number;
lastStep: Tokens | undefined;
/**
* The history item this turn's execution follows, for a turn with no prompt
* id of T3's: the session's newest item before a `/name` command (null when
* the history was empty), since `session.command` takes no id and answers
* without one; the report a continuation turn's execution answers; or, on a
* subagent's session, the first item queued for its turn. A reconnect
* backfills that turn from everything after this. Other turns backfill from
* their own prompt's id (`nativeTurnRef`).
*/
before: string | null | undefined;
/** The compaction running in this turn, `/compact` or OpenCode's own when the context fills. */
compaction: { readonly nativeId: string; readonly startedAt: DateTime.Utc } | undefined;
compactions: number;
interrupted: boolean;
/**
* Set until the turn's prompt, command or compaction is sent. A Stop before
* then has nothing on the server to stop, so it ends the turn here and the
* request is never sent.
*/
unsent: boolean;
/**
* Steers sent into this turn that OpenCode has not delivered yet, by inbox
* id. Each one wakes the session, so an execution that ends before reading
* one is followed by another that does, and the turn spans both.
*/
readonly steers: Set<string>;
/** This turn's steers OpenCode delivered or cancelled: a retried one is not sent again. */
readonly settledInbox: Set<string>;
/** How the last execution ended, while the turn waits for the one that delivers its steers. */
heldEnd: TurnTerminal | undefined;
/**
* Set on a turn started after a timed-out Stop's run left the server: that
* run's tail can still be on the stream, and everything before this turn's
* own `session.execution.started` belongs to it.
*/
awaitingStart: boolean;
/**
* Prefixes a subagent session's tool ids in native item ids: the model
* names tool calls, so a child's could repeat its parent's.
*/
readonly scope: string;
/** Set once the turn calls `subagent`: its usage then leaves out the subagents'. */
usedSubagents: boolean;
}
/** A `subagent` tool call. A background one outlives the turn that made it. */
interface SubagentCall {
readonly toolId: string;
readonly nativeId: string;
/** The calling session and the turn whose item the call is. */
readonly state: ThreadState;
readonly turn: ActiveTurn;
readonly ordinal: number;
readonly startedAt: DateTime.Utc;
prompt: string;
title: string | null;
agent: string | undefined;
model: string | null;
/** The tool returned while the subagent runs on; the report OpenCode gives its parent settles it. */
background: boolean;
child: ThreadState | undefined;
status: OrchestrationV2Subagent["status"];
result: string | null;
completedAt: DateTime.Utc | null;
}
/**
* An execution OpenCode ran on a thread's session without T3 asking: the
* parent's answer once a background subagent ended. Its events are held until
* the continuation turn T3 opens for it takes them.
*/
interface Wake {
readonly events: Array<OpenCode2StreamEvent>;
running: boolean;
/** The background subagents whose end it answers, as its turn's notification names them. */
readonly reports: Array<BackgroundWorkReport>;
/** What OpenCode told the model, the continuation's prompt text. */
readonly detail: string | null;
/** Stopped, or taken by a user turn: the continuation it asked for is not needed. */
dropped: boolean;
/** The last report it delivered (inbox ids are history ids): its execution follows it. */
readonly after: string | undefined;
/** The first report it delivered, where fork and rollback cut before its turn. */
readonly first: string | undefined;
}
interface OpenBlock {
readonly block: { readonly assistantMessageID: string; readonly ordinal: number };
readonly kind: "text" | "reasoning";
readonly startedAt: DateTime.Utc;
text: string;
}
interface ThreadState {
readonly sessionId: string;
providerThread: OrchestrationV2ProviderThread;
readonly providerTurns: Map<string, OrchestrationV2ProviderTurn>;
active: ActiveTurn | undefined;
/** What the native session runs now, so a changed selection is switched before prompting. */
model: ModelRef | undefined;
/**
* Set when a turn ended here while OpenCode may still be running it: a Stop
* that timed out, a prompt whose request failed without a clear answer, or a
* request T3 could not answer. Execution events carry only the session id,
* so the next execution end belongs to that run; it clears this and ends no turn.
*/
unsettled: boolean;
/** The session's location, where its agents' path rules are read. */
directory: string;
/**
* The agent the session runs (`build`, or `plan` in plan mode); its own path
* rules stay in force, and a changed mode switches it before prompting.
*/
agent: string;
/** The native session's rules as T3 last read or wrote them, and the policy they are for. */
rules: ReadonlyArray<Rule> | undefined;
policy: RulesPolicy;
/** "Allow … this session" answers, kept in the session's rules while T3 has it open. */
readonly grants: Array<Rule>;
/** The session's `subagent` calls still running, by tool call id. */
readonly calls: Map<string, SubagentCall>;
/**
* Set on a subagent's session: the call that runs it and the thread it
* shows in. Each of its executions is a runless turn there.
*/
subagent:
| {
call: SubagentCall;
readonly appThread: OrchestrationV2AppThread;
turns: number;
/** The prompt its next execution answers, shown as that turn's user message. */
prompt: string | undefined;
/**
* The first item queued for the session since its last turn began (its
* prompt, or a report it answers): where the next turn's history
* begins, which a reconnect reads back.
*/
queued: string | undefined;
}
| undefined;
/** Executions OpenCode started on its own, oldest first, each waiting for its turn. */
readonly wakes: Array<Wake>;
/**
* Background subagents whose end OpenCode queued for this session and has
* not delivered yet, by inbox id. The wake that delivers them names them.
*/
readonly reports: Map<
string,
{
readonly inboxId: string;
readonly childId: string;
readonly report: BackgroundWorkReport;
readonly text: string;
}
>;
/**
* Background subagent sessions T3 stopped. OpenCode wakes the parent to
* report them; a wake that reports only these is stopped as well.
*/
readonly stoppedChildren: Set<string>;
/**
* Steers a turn ended without (a Stop, a failure). They stay in OpenCode's
* inbox, where the next prompt would deliver them first, so it cancels them.
*/
readonly strandedSteers: Set<string>;
/** T3's MCP server as registered for this thread, and the instructions entry sent with it. */
mcp:
| { readonly name: string; readonly directory: string; readonly credential: string }
| undefined;
instructions: string | undefined;
}
type TurnTerminal =
| { readonly status: "completed" | "interrupted" }
| { readonly status: "failed"; readonly failure: ReturnType<typeof makeProviderFailure> };
type Rule = Permission.Rule;
type RulesPolicy = Pick<
ProviderAdapter.ProviderAdapterV2RuntimePolicy,
"runtimeMode" | "interactionMode"
>;
type NativeForm = EventOf<"form.created">["data"]["form"];
/** A permission ask or question form shown to the user and not answered yet. */
interface PendingRequest {
readonly request: OrchestrationV2RuntimeRequest;
readonly item: OrchestrationV2TurnItem;
readonly node: OrchestrationV2ExecutionNode;
/** The thread the request is shown on: the asking session's, or its parent's. */
readonly state: ThreadState;
readonly turn: ActiveTurn;
/** The session that asked: the thread's own, or one of its subagents'. */
readonly sessionId: string;
/** Set once T3 sends its answer; the orchestrator has already recorded it. */
answering: boolean;
readonly native:
| {
readonly type: "permission";
readonly id: string;
readonly action: string;
readonly resources: ReadonlyArray<string>;
readonly save: ReadonlyArray<string>;
}
| { readonly type: "form"; readonly id: string; readonly form: NativeForm };
}
const rule = (action: string, effect: Rule["effect"]): Rule => ({ action, resource: "*", effect });
/**
* A session's permission rules. OpenCode checks the agent's rules and then
* these, and the last rule that matches decides, so these override the
* agent's. `paths` are the agent's own allows for its directories (saved tool
* output, the plan agent's plan directory), which the blanket rules here would
* otherwise override; `grants` are "Always allow this session" answers.
*/
/**
* T3's MCP server is registered per directory, not per session, so each thread
* gets its own `t3-code-<thread>` entry with its own credential. OpenCode names
* an MCP tool's permission `<server>_<tool>` (non-alphanumerics become `_`).
* OpenCode skips every tool of a server whose name is over 64 characters (its
* tool namespace limit), and its router rejects adding one over 100, so a
* thread id that does not fit is replaced by a digest of it.
*/
export const t3McpServerName = Effect.fn("t3McpServerName")(function* (threadId: string) {
const name = `t3-code-${threadId.replaceAll(/[^a-zA-Z0-9_-]/g, "_")}`;
if (name.length <= 64) return name;
const crypto = yield* Crypto.Crypto;
const digest = yield* crypto
.digest("SHA-256", new TextEncoder().encode(threadId))
.pipe(Effect.orDie);
return `t3-code-${Hex.encode(digest).slice(0, 16)}`;
});
/**
* The rules that keep T3's MCP servers to their own thread, after the mode's:
* the last matching rule wins, so every thread's T3 server is denied and then
* this thread's own is allowed again, in every mode. A subagent's session
* inherits the thread's.
*/
const mcpRules = (mcpServerName: string | null): ReadonlyArray<Rule> =>
mcpServerName === null
? []
: [
{ action: "t3-code-*", resource: "*", effect: "deny" },
{ action: `${mcpServerName}_*`, resource: "*", effect: "allow" },
];
const sessionRules = (
policy: RulesPolicy,
paths: ReadonlyArray<Rule>,
grants: ReadonlyArray<Rule>,
mcpServerName: string | null,
): ReadonlyArray<Rule> => [
...(policy.runtimeMode === "full-access"
? [rule("*", "allow")]
: [
rule("shell", "ask"),
rule("edit", policy.runtimeMode === "auto-accept-edits" ? "allow" : "ask"),
rule("external_directory", "ask"),
]),
...grants,
// Plan mode writes only its plan, which `paths` allows again. Shell and read
// are never denied: the free tier refuses sessions whose rules deny them.
...(policy.interactionMode === "plan" ? [rule("edit", "deny")] : []),
...paths,
...mcpRules(mcpServerName),
];
const sameRules = (left: ReadonlyArray<Rule> | undefined, right: ReadonlyArray<Rule>) =>
left?.length === right.length &&
left.every(
(entry, index) =>
entry.action === right[index]?.action &&
entry.resource === right[index]?.resource &&
entry.effect === right[index]?.effect,
);
/**
* Sent with a declined permission. A reject without a message is OpenCode's
* "stop": it ends the whole execution, which is Cancel.
*/
const DECLINED = "The user declined this request.";
/**
* Steered into the session before a decline: a declined shell call reaches the
* model only as "Unable to execute command", which it retries.
*/
const declinedNote = (action: string, resources: ReadonlyArray<string>) =>
`The user declined the ${action} request${resources.length === 0 ? "" : ` (${resources.join(", ")})`}. Do not retry it; continue without it or ask the user how to proceed.`;
/** The session-wide choice for a request whose `save` patterns OpenCode would remember. */
const sessionGrantLabel = (action: string, save: ReadonlyArray<string>) =>
save.every((pattern) => pattern === "*")
? `Allow every ${action} request this session`
: `Allow ${save.join(", ")} this session`;
const text = (value: string | undefined, fallback: string) => value?.trim() || fallback;
/**
* A form's fields as T3 questions, or why T3 cannot ask them: a link to open,
* a field shown only for another answer, a hidden field, or a number or yes/no
* value. OpenCode's question tool only asks text and multi-select fields.
*/
const formQuestions = (
form: NativeForm,
):
| { readonly questions: ReadonlyArray<OrchestrationV2UserInputQuestion> }
| { readonly unsupported: string } => {
const questions: Array<OrchestrationV2UserInputQuestion> = [];
for (const [index, field] of form.fields.entries()) {
if (field.type === "external") return { unsupported: "a link to open" };
if (field.type !== "string" && field.type !== "multiselect") {
return { unsupported: `a ${field.type} value` };
}
if (field.hidden === true) return { unsupported: "a hidden field" };
if ((field.when?.length ?? 0) > 0) return { unsupported: "a field that depends on another" };
const header = text(field.title, `Question ${index + 1}`);
const options = (field.options ?? []).map((option) => {
const label = text(option.label, text(option.value, "Option"));
return { label, description: text(option.description, label), value: option.value };
});
questions.push({
id: field.key,
header,
question: text(field.description, header),
options,
multiSelect: field.type === "multiselect",
allowCustomAnswer: field.custom === true || options.length === 0,
});
}
return { questions };
};
/** T3's answers in OpenCode's shape: a list for a multi-select, text otherwise. */
const formAnswer = (form: NativeForm, answers: Readonly<Record<string, unknown>>) => {
const answer: Record<string, string | ReadonlyArray<string>> = {};
for (const field of form.fields) {
const raw = answers[field.key];
const values = (Array.isArray(raw) ? raw : [raw]).filter(
(value): value is string => typeof value === "string" && value.trim().length > 0,
);
if (values.length > 0) {
answer[field.key] = field.type === "multiselect" ? values : values.join(", ");
}
}
return answer;
};
const approves = (decision: ProviderApprovalDecision) =>
decision === "accept" || decision === "acceptForSession" || decision === "acceptAlways";
/**
* An answer to a request OpenCode already dropped (its execution ended, or
* its session is gone): nothing waits on it, so it counts as delivered.
*/
const permissionGone = {
PermissionNotFoundError: () => Effect.void,
SessionNotFoundError: () => Effect.void,
};
const formGone = {
FormNotFoundError: () => Effect.void,
FormAlreadySettledError: () => Effect.void,
SessionNotFoundError: () => Effect.void,
};
/** The session rules an agent keeps for its own directories, which T3's blanket rules would override. */
const agentPaths = (rules: ReadonlyArray<Rule>) =>
rules.filter(
(entry) =>
entry.effect === "allow" &&
entry.resource !== "*" &&
(entry.action === "edit" || entry.action === "external_directory"),
);
const ref = (nativeId: string, strength: "strong" | "weak" = "strong") => ({
driver: OPENCODE_PROVIDER,
nativeId,
strength,
});
const sessionIdOf = (providerThread: OrchestrationV2ProviderThread) => {
const nativeId = providerThread.nativeThreadRef?.nativeId;
return nativeId === undefined || nativeId === null
? Effect.fail(
new ProviderAdapter.ProviderAdapterProtocolError({
driver: OPENCODE_PROVIDER,
detail: `Provider thread ${providerThread.id} has no OpenCode session`,
}),
)
: Effect.succeed(nativeId);
};
const textOf = (content: ReadonlyArray<{ readonly type: string; readonly text?: string }>) =>
content.flatMap((part) => (part.type === "text" && part.text ? [part.text] : [])).join("\n");
const stringField = (record: Readonly<Record<string, unknown>> | undefined, key: string) => {
const value = record?.[key];
return typeof value === "string" && value.trim().length > 0 ? value : undefined;
};
/** How a background subagent ended, from the `state` of OpenCode's report to its parent. */
const reportOutcome = (state: string | undefined) =>
state === "completed"
? ("completed" as const)
: state === "error"
? ("failed" as const)
: state === "cancelled"
? ("cancelled" as const)
: ("unknown" as const);
/** A subagent's answer, without the `<subagent …>` wrapper OpenCode gives the model. */
const subagentOutput = (text: string) =>
/^<subagent\b[^>]*>\n?([\s\S]*?)\n?<\/subagent>$/.exec(text.trim())?.[1] ?? text;
const isContinuation = (turnInput: ProviderAdapter.ProviderAdapterV2TurnInput) =>
turnInput.message.createdBy === "agent" && turnInput.message.creationSource === "provider";
/** The native turn id of a continuation turn, which answers a wake and has no prompt. */
const wakeTurnId = (sessionId: string, attemptId: string) => `${sessionId}:wake:${attemptId}`;
const isWakeTurn = (turn: OrchestrationV2ProviderTurn) =>
turn.nativeTurnRef?.nativeId?.includes(":wake:") === true;
const INTERRUPT_TIMEOUT = "10 seconds";
/** The session instructions entry T3 writes its per-turn system prompt to. */
const INSTRUCTIONS_KEY = "t3-code";
/** A lost event stream is resubscribed this many times, this far apart, before the session breaks. */
const RECONNECT_ATTEMPTS = 5;
const RECONNECT_DELAY = "2 seconds";
const RECONCILE_TIMEOUT = "15 seconds";
/** How long a new turn waits for a reconnect in progress. */
const RECONNECT_WAIT = "30 seconds";
/** A background subagent's result when its end was lost with the event stream. */
const LOST_BACKGROUND =
"T3 Code lost its connection to OpenCode while this subagent ran, so its result is not shown.";
/** How long a turn waits on the directory's commands or skills before sending the text as is. */
const INVENTORY_TIMEOUT = "5 seconds";
const ACTIVE_CHECK_TIMEOUT = "5 seconds";
/** Answers that mean the server refused a prompt; any other failure may have been accepted. */
const CLEAR_PROMPT_REJECTIONS: ReadonlySet<string> = new Set([
"InvalidRequestError",
"ConflictError",
"UnauthorizedError",
"CommandNotFoundError",
"SkillNotFoundError",
]);
export const OPENCODE_2_STILL_STOPPING =
"OpenCode is still stopping the previous turn. Send the message again in a moment.";
const REQUEST_REPLY_TIMEOUT = "10 seconds";
/** Whether an answer to a paused request reached the server, trying twice. */
const deliver = <E>(answer: Effect.Effect<void, E>) =>
answer.pipe(
Effect.retry({ times: 1 }),
Effect.timeout(REQUEST_REPLY_TIMEOUT),
Effect.exit,
Effect.map(Exit.isSuccess),
);
type ModelRef = ReturnType<typeof Model.Ref.make>;
/**
* The user message ids T3 prompts under. OpenCode takes a client id (it must
* start with `msg_`) and answers a repeat of one in the same session with the
* item it already has, so a retried request never queues a second message,
* and a turn knows the message fork and rollback cut at before OpenCode
* answers. The id is unique across the whole server, which refuses it in any
* other session (409), so it names the session: another T3 database or
* environment on the same server repeats thread ids and run ordinals, never
* session ids. A turn keeps its id in `nativeTurnRef`.
*/
const turnPromptId = (sessionId: string, attemptId: string) =>
SessionMessage.ID.make(`msg_t3_turn_${sessionId}:${attemptId}`);
const steerPromptId = (sessionId: string, messageId: string) =>
SessionMessage.ID.make(`msg_t3_steer_${sessionId}:${messageId}`);
/**
* The user message a turn prompted with. Turns from before T3 chose prompt ids
* recorded `<session>:attempt:<id>`, which is no message.
*/
const promptOf = (turn: OrchestrationV2ProviderTurn) => {
const nativeId = turn.nativeTurnRef?.nativeId;
return nativeId?.startsWith("msg_") === true ? nativeId : undefined;
};
/**
* Where fork and rollback cut to keep the session's turns up to `kept` (all of
* them before the first turn when undefined): before the message that starts
* the next turn that reached OpenCode, which is in `prompts`. That is a turn's
* prompt, or for a continuation the report it answers. A turn refused before
* it prompted is not in the session and is passed over, and so is a
* continuation that took no report. `null` means nothing follows, so there is
* no cut; a later turn from before T3 chose prompt ids has no known message,
* so no cut is safe.
*/
const boundaryAfter = (
turns: ReadonlyArray<OrchestrationV2ProviderTurn>,
providerThreadId: OrchestrationV2ProviderThread["id"],
kept: OrchestrationV2ProviderTurn | undefined,
prompts: ReadonlySet<string>,
) => {
const later = turns
.filter(
(turn) =>
turn.providerThreadId === providerThreadId &&
(kept === undefined || turn.ordinal > kept.ordinal),
)
.toSorted((left, right) => left.ordinal - right.ordinal);
if (later.some((turn) => promptOf(turn) === undefined && !isWakeTurn(turn))) {
return Effect.fail(
new ProviderAdapter.ProviderAdapterProtocolError({
driver: OPENCODE_PROVIDER,
detail:
"This OpenCode conversation has turns from an earlier T3 Code version, so it can't be cut there.",
}),
);
}
return Effect.succeed(
later.map(promptOf).find((messageId) => messageId !== undefined && prompts.has(messageId)) ??
null,
);
};
// Errors already in the adapter channel keep their tag; only lower-level ones are wrapped.
const isProviderAdapterError = Schema.is(ProviderAdapter.ProviderAdapterV2Error);
/**
* The model OpenCode should run for a `provider/model` slug and its reasoning
* variant, or undefined for any other slug: sending none would run OpenCode's
* default while T3 records the requested model.
*/
const modelRef = (selection: ProviderAdapter.ProviderAdapterV2TurnInput["modelSelection"]) => {
const parsed = parseOpenCodeModelSlug(selection.model);
if (parsed === null) return undefined;
const variant = getModelSelectionStringOptionValue(selection, "variant");
return Model.Ref.make({
providerID: Provider.ID.make(parsed.providerID),
id: Model.ID.make(parsed.modelID),
...(variant === undefined ? {} : { variant: Model.VariantID.make(variant) }),
});
};
const malformedModel = (model: string) =>
`OpenCode model '${model}' must use provider/model format`;
const sameModel = (left: ModelRef, right: ModelRef | undefined) =>
left.providerID === right?.providerID &&
left.id === right?.id &&
(left.variant ?? "default") === (right?.variant ?? "default");
/** OpenCode's own agents for T3's interaction modes; plan mode is its read-only `plan` agent. */
const agentFor = (input: ProviderAdapter.ProviderAdapterV2TurnInput) =>
input.runtimePolicy.interactionMode === "plan" ? "plan" : "build";
/** `/name args` naming one of the workspace's commands, which OpenCode expands itself. */
const commandOf = (text: string) => {
const match = /^\/([^\s/]+)(?:\s+([\s\S]*))?$/.exec(text.trim());
return match === null ? undefined : { name: match[1]!, text: match[2] ?? "" };
};
/** Whether a prompt names any skill at all, before the directory's skills are read. */
const SKILL_MENTION = new RegExp(SKILL_MENTION_PATTERN.source, "u");
/** The workspace skills a prompt names with the composer's `$skill` tokens. */
const skillsNamed = (text: string, known: ReadonlySet<string>) => [
...new Set(
[...text.matchAll(SKILL_MENTION_PATTERN)].flatMap((match) =>
known.has(match[2] ?? "") ? [match[2]!] : [],
),
),
];
/**
* The turn's own tokens: steps add up, and the last step's input is the live
* context size. A subagent's tokens are its own session's, never these.
*/
const turnTokenUsage = (turn: ActiveTurn, status: OrchestrationV2ProviderTurn["status"]) =>
turn.steps === 0
? {
usageScope: "main_agent" as const,
usageStatus: "unavailable" as const,
hasSubagents: turn.usedSubagents,
}
: {
usageScope: "main_agent" as const,
usageStatus: status === "completed" ? ("complete" as const) : ("partial" as const),
inputTokens: turn.usage.input,
cachedInputTokens: turn.usage.cached,
cacheCreationTokens: turn.usage.cacheWrite,
outputTokens: turn.usage.output,
reasoningTokens: turn.usage.reasoning,
hasSubagents: turn.usedSubagents,
};
/**
* The adapter for one provider instance. It talks to the instance's
* {@link OpenCode2Server.OpenCode2Server}, which the driver builds from the instance's settings.
*/
export const make = Effect.fn("OpenCode2Adapter.make")(function* (instanceId: ProviderInstanceId) {
const server = yield* OpenCode2Server.OpenCode2Server;
const idAllocator = yield* IdAllocator.IdAllocatorV2;
const serverConfig = yield* ServerConfig.ServerConfig;
const continuationRequests = yield* ProviderContinuationRequests.ProviderContinuationRequests;
const crypto = yield* Crypto.Crypto;
const driver = OPENCODE_PROVIDER;
const mcpServerNameFor = (threadId: string) =>
t3McpServerName(threadId).pipe(Effect.provideService(Crypto.Crypto, crypto));
// MT Code: Computer Use (`mt-desktop`) needs server settings; narrow layers
// omit it. The instance's server is shared by its sessions, so the
// directories this instance registered the server in are tracked here.
const resolveDesktopMcp = yield* makeOptionalResolveEnabledDesktopMcp();
const desktopMcpDirectories = new Set<string>();
/**
* Lends the instance's server to a session until its scope closes. A spawned
* server that died is started again on the next borrow.
*/
const borrow = Effect.gen(function* () {
const lent = yield* Deferred.make<OpenCode2Server.OpenCode2Connection, OpenCodeRuntimeError>();
yield* server
.withConnection((connection) =>
Deferred.succeed(lent, connection).pipe(Effect.andThen(Effect.never)),
)
.pipe(
Effect.catch((error) => Deferred.fail(lent, error)),
Effect.forkScoped,
);
return yield* Deferred.await(lent);
});
const openSession = Effect.fn("OpenCode2Adapter.openSession")(function* (
input: Parameters<ProviderAdapter.ProviderAdapterV2Shape["openSession"]>[0],
initial: {
readonly connection: OpenCode2Server.OpenCode2Connection;
readonly scope: Scope.Closeable;
},
) {
let connection = initial.connection;
// Replaced when the session reconnects to a restarted server.
let client = connection.client;
const sessionScope = yield* Effect.scope;
// Context windows by directory, then `provider/model`: a project's own
// OpenCode config can change a model's limits, and this one runtime serves
// the instance's threads in every directory.
const contextWindows = new Map<string, Map<string, number>>();
/** A thread without a worktree runs where T3 does, as its session is created. */
const directoryOf = (cwd: string | null | undefined) => cwd ?? serverConfig.cwd;
const windowOf = (cwd: string | null | undefined, model: string) =>
contextWindows.get(directoryOf(cwd))?.get(model);
const now = yield* DateTime.now;
let session: OrchestrationV2ProviderSession = {
id: input.providerSessionId,
driver,
providerInstanceId: instanceId,
status: "ready",
cwd: input.runtimePolicy.cwd ?? serverConfig.cwd,
model: input.modelSelection.model,
capabilities: OpenCode2ProviderCapabilities,
createdAt: now,
updatedAt: now,
lastError: null,
};
const events = yield* Queue.unbounded<ProviderAdapter.ProviderAdapterV2Event, Cause.Done>();
// Thread sessions, and the sessions of their subagents.
const threads = new Map<string, ThreadState>();
// A subagent's session, by its id, to the thread whose session started it.
const childOwners = new Map<string, ThreadState>();
// What OpenCode announced about a subagent's session, until its call names it.
const announced = new Map<string, EventOf<"session.created">["data"]>();
// Sessions of these threads with an execution running, seen on the stream.
const busy = new Set<string>();
const pending = new Map<RuntimeRequestId, PendingRequest>();
// Events and a wake's replay into its turn are handled one at a time, in order.
const lock = yield* Semaphore.make(1);
// Sessions whose `revert.clear` has not had its empty execution yet.
const clearing = new Map<string, Deferred.Deferred<void>>();
// Sessions a failed rollback may have left with a staged revert.
const stagedReverts = new Set<string>();
// Starting a turn and cutting the history take turns on a session: each
// checks that the other is not running before its own requests yield.
const sessionGates = yield* KeyedLock.make<string>();
const exclusive =
(providerThread: OrchestrationV2ProviderThread) =>
<A, E, R>(effect: Effect.Effect<A, E, R>) => {
const sessionId = providerThread.nativeThreadRef?.nativeId;
if (sessionId == null) return effect;
return sessionGates.withLock(sessionId, effect);
};
const emit = (event: ProviderAdapter.ProviderAdapterV2Event) =>
Queue.offer(events, event).pipe(Effect.asVoid);
const ownerOf = (sessionId: string) => childOwners.get(sessionId) ?? threads.get(sessionId);
const newThreadState = (
sessionId: string,
providerThread: OrchestrationV2ProviderThread,
directory: string,
subagent: ThreadState["subagent"],
): ThreadState => ({
sessionId,
providerThread,
providerTurns: new Map(),
active: undefined,
model: undefined,
unsettled: false,
directory,
agent: "build",
rules: undefined,
policy: input.runtimePolicy,
grants: [],
calls: new Map(),
subagent,
wakes: [],
reports: new Map(),
stoppedChildren: new Set(),
strandedSteers: new Set(),
mcp: undefined,
instructions: undefined,
});
/** The thread whose session started this one, through any nesting. */
const rootOf = (state: ThreadState): ThreadState => {
let current = state;
while (current.subagent !== undefined) current = current.subagent.call.state;
return current;
};
/** A thread's session and its subagents' sessions, through any nesting. */
const sessionsOf = (thread: ThreadState): ReadonlyArray<ThreadState> => [
thread,
...[...threads.values()].filter((state) => {
for (let above = state.subagent?.call.state; above; above = above.subagent?.call.state) {
if (above === thread) return true;
}
return false;
}),
];
/**
* A thread's `subagent` calls still running, its subagents' included: a
* subagent's background call runs on after the call that made it ended.
*/
const runningCalls = (thread: ThreadState): ReadonlyArray<SubagentCall> =>
sessionsOf(thread).flatMap((state) => [...state.calls.values()]);
/** The `subagent` calls that lead to a session, from its own up to the thread's. */
const callsAbove = (sessionId: string) => {
const calls: Array<SubagentCall> = [];
let current = threads.get(sessionId);
while (current?.subagent !== undefined) {
calls.push(current.subagent.call);
current = current.subagent.call.state;
}
return calls;
};
/**
* The turn a request from `sessionId` is shown under. A background
* subagent's goes to the turn that started it, whose run waits on the
* subagent; anything else goes to the thread's running turn.
*/
const requestTurn = (sessionId: string) => {
const state = ownerOf(sessionId);
if (state === undefined) return undefined;
const background = callsAbove(sessionId).findLast((call) => call.background);
if (background !== undefined) {
return isOrchestrationV2WorkActive(background.status)
? { state, turn: background.turn }