返回 DeepSeek-Reasonix
remote_tab.go
根目录 / desktop / remote_tab.go
1 package main
2
3 import (
4 "bufio"
5 "bytes"
6 "context"
7 "encoding/json"
8 "fmt"
9 "io"
10 "log"
11 "net/http"
12 "strconv"
13 "strings"
14 "time"
15
16 "reasonix/internal/agent"
17 "reasonix/internal/control"
18 "reasonix/internal/event"
19 )
20
21 // The remote-tab bridge exchanges its pre-shared token for an HttpOnly
22 // session cookie over the loopback tunnel. Subsequent API and SSE requests
23 // use that cookie, keeping the token out of request lines and access logs.
24
25 const remoteTabStreamOpenStability = 50 * time.Millisecond
26
27 func remoteSessionTransitionBusy(err error) bool {
28 if err == nil {
29 return false
30 }
31 message := err.Error()
32 return strings.Contains(message, "while a turn is running") ||
33 strings.Contains(message, "while another session change is in progress") ||
34 strings.Contains(message, "session is finishing background teardown")
35 }
36
37 // attachRemoteTabServe starts the event pump before entering the session so
38 // /new or /resume frames are not missed. The caller's context owns the pump;
39 // handshake and session entry use a bounded child context.
40 func (a *App) attachRemoteTabServe(ctx context.Context, tabID, base, token, instanceID string, opts RemoteTabOpenOptions) (bool, error) {
41 callCtx, cancel := context.WithTimeout(ctx, 30*time.Second)
42 defer cancel()
43
44 client, err := newServeHTTPClient(base)
45 if err != nil {
46 return false, err
47 }
48 capabilities, err := serveHandshakeCapabilities(callCtx, client, base, token)
49 if err != nil {
50 log.Printf("[remote] attachRemoteTabServe: handshake FAILED tab=%s base=%q err=%v", tabID, base, err)
51 return false, err
52 }
53 a.remoteTabMu.Lock()
54 tab := a.remoteTabs[tabID]
55 a.remoteTabMu.Unlock()
56 if tab == nil {
57 return false, fmt.Errorf("remote tab %q closed during bootstrap", tabID)
58 }
59 a.remoteTabMu.Lock()
60 if current := a.remoteTabs[tabID]; current == tab {
61 tab.capabilities = make(map[string]bool, len(capabilities))
62 for _, capability := range capabilities {
63 tab.capabilities[capability] = true
64 }
65 }
66 a.remoteTabMu.Unlock()
67 tab.sessionMu.Lock()
68 defer tab.sessionMu.Unlock()
69
70 // Resolve every non-new target before opening the all-session pump. A
71 // detached controller may replay pending prompts as soon as /resume starts;
72 // publishing its route first keeps those frames on the foreground surface.
73 focusOnly := !opts.NewSession && strings.TrimSpace(opts.SessionName) == "" && strings.TrimSpace(opts.SessionPath) == "" && strings.TrimSpace(opts.SessionID) == ""
74 var target serveSessionEntry
75 if !opts.NewSession {
76 target, err = preflightRemoteSessionTarget(callCtx, client, base, opts)
77 if err != nil {
78 return false, err
79 }
80 }
81
82 pumpCtx, gen, attachPathRevision, err := a.installRemoteTabAttachPump(ctx, tabID, tab, client, base, token, remoteSessionRoute(target), !opts.NewSession)
83 if err != nil {
84 return false, err
85 }
86
87 opened := make(chan error, 1)
88 a.goRemoteTabSafe("remoteTabPump", func() { a.remoteTabPump(pumpCtx, tabID, gen, opened) })
89 select {
90 case err = <-opened:
91 if err != nil {
92 a.retireRemoteTabGeneration(tabID, gen)
93 return false, err
94 }
95 case <-callCtx.Done():
96 a.retireRemoteTabGeneration(tabID, gen)
97 return false, callCtx.Err()
98 }
99 entered := true
100 if !focusOnly {
101 enterOpts := opts
102 if !opts.NewSession {
103 enterOpts.SessionName, enterOpts.SessionPath, enterOpts.SessionID, enterOpts.SessionTitle = target.Name, target.Path, target.SessionID, target.Title
104 }
105 target, err = enterRemoteSessionTarget(callCtx, client, base, enterOpts)
106 entered = err == nil && !target.TakenOver
107 }
108 if err == nil && target.TakenOver {
109 // The serve mounted this caller as a read-only spectator (another
110 // runtime owns the session writer). The tab stays attached to render
111 // the file/mirrored view and the take-back banner drives /reclaim.
112 log.Printf("[remote] attachRemoteTabServe: enterRemoteSession SPECTATOR (writer owned elsewhere) tab=%s session=%q", tabID, remoteSessionRoute(target))
113 }
114 if err != nil {
115 // A busy serve refuses session transitions with 409 but retains its
116 // usable current session. Keep the attach so pending work remains visible.
117 if remoteSessionTransitionBusy(err) {
118 log.Printf("[remote] attachRemoteTabServe: enterRemoteSession BUSY (attached to current session) tab=%s err=%v", tabID, err)
119 entered = false
120 target, _ = serveCurrentSession(callCtx, client, base)
121 } else if remoteSessionTakenOver(err) {
122 // A local runtime on the serve host owns the session. Pin the tab to
123 // the requested session as a read-only spectator: the mirror's frames
124 // route here, /history and /status serve the file-backed view, and
125 // the take-back banner drives /reclaim. No serve-frontend transition
126 // ran, but the tab must stay attached to render the mirror.
127 log.Printf("[remote] attachRemoteTabServe: enterRemoteSession TAKEN OVER (read-only spectator) tab=%s session=%q err=%v", tabID, target.Path, err)
128 entered = false
129 if remoteSessionRoute(target) == "" {
130 current, _ := serveCurrentSession(callCtx, client, base)
131 target = current
132 }
133 target.TakenOver = true
134 } else {
135 log.Printf("[remote] attachRemoteTabServe: enterRemoteSession FAILED tab=%s err=%v", tabID, err)
136 a.retireRemoteTabGeneration(tabID, gen)
137 return false, err
138 }
139 }
140 if !a.commitRemoteTabAttachResponse(tabID, tab, gen, attachPathRevision, target, opts.NewSession) {
141 entered = false
142 }
143 if !a.waitRemoteTabStreamStable(callCtx, tabID, gen) {
144 return false, fmt.Errorf("remote tab %q event stream closed during session attach", tabID)
145 }
146 a.remoteTabMu.Lock()
147 if current := a.remoteTabs[tabID]; current == tab && current.gen == gen {
148 current.session.instanceID = instanceID
149 }
150 a.remoteTabMu.Unlock()
151 // A 200 response is only the stream-open barrier. The stream can still die
152 // while /new or /resume is in flight; publish readiness only if its pump has
153 // not already moved this same generation into reconnecting/error.
154 if !a.markRemoteTabAttached(tabID, gen) {
155 return false, fmt.Errorf("remote tab %q event stream closed during session attach", tabID)
156 }
157 return entered, nil
158 }
159
160 // commitRemoteTabAttachResponse applies an attach response only while it still
161 // owns the foreground route. A session_changed frame for a newer adoption is
162 // authoritative even when the older /new or /resume response arrives later.
163 func (a *App) commitRemoteTabAttachResponse(tabID string, tab *remoteTab, gen, requestPathRevision uint64, target serveSessionEntry, reset bool) bool {
164 tab.routeEventMu.Lock()
165 defer tab.routeEventMu.Unlock()
166 a.remoteTabMu.Lock()
167 defer a.remoteTabMu.Unlock()
168 current := a.remoteTabs[tabID]
169 if current != tab || current.gen != gen {
170 return false
171 }
172 target.Path = strings.TrimSpace(target.Path)
173 route := remoteSessionRoute(target)
174 if current.routing.pathRevision != requestPathRevision && current.routing.currentPath != route {
175 return false
176 }
177 alreadyAdopted := current.routing.pathRevision != requestPathRevision
178 if !alreadyAdopted {
179 commitRemoteTabAttachRoute(current, route, reset)
180 }
181 current.session.takenOver = target.TakenOver
182 current.session.path = target.Path
183 current.session.sessionID = target.SessionID
184 if name := strings.TrimSpace(target.Name); name != "" {
185 current.session.name = name
186 }
187 if title := strings.TrimSpace(target.Title); title != "" {
188 current.topicTitle = title
189 }
190 if current.routing.running == nil {
191 current.routing.running = map[string]bool{}
192 }
193 return true
194 }
195
196 func (a *App) waitRemoteTabStreamStable(ctx context.Context, tabID string, gen uint64) bool {
197 timer := time.NewTimer(remoteTabStreamOpenStability)
198 defer timer.Stop()
199 select {
200 case <-ctx.Done():
201 return false
202 case <-timer.C:
203 return a.remoteTabGenerationCurrent(tabID, gen)
204 }
205 }
206
207 func (a *App) markRemoteTabAttached(tabID string, gen uint64) bool {
208 a.remoteTabMu.Lock()
209 defer a.remoteTabMu.Unlock()
210 tab := a.remoteTabs[tabID]
211 if tab == nil || tab.gen != gen || tab.state != "connecting" {
212 return false
213 }
214 tab.attachedGen = gen
215 return true
216 }
217
218 func (a *App) publishRemoteTabAttachedReady(tabID string, gen uint64) bool {
219 tab := a.lockRemoteTabPublication(tabID)
220 if tab == nil {
221 return false
222 }
223 a.remoteTabMu.Lock()
224 if a.remoteTabs[tabID] != tab || tab.gen != gen || tab.attachedGen != gen || tab.state != "connecting" {
225 a.remoteTabMu.Unlock()
226 tab.routeEventMu.Unlock()
227 return false
228 }
229 tab.attachedGen = 0
230 tab.state = "ready"
231 tab.err = ""
232 a.remoteTabMu.Unlock()
233 a.emitRemoteEvent(fmt.Sprintf("remote-tab:%s:state", tabID), RemoteTabStateView{State: "ready"})
234 tab.routeEventMu.Unlock()
235 a.applyPendingRemoteTabOpenSelection(tabID)
236 return true
237 }
238
239 func (a *App) remoteTabGenerationCurrent(tabID string, gen uint64) bool {
240 a.remoteTabMu.Lock()
241 defer a.remoteTabMu.Unlock()
242 tab := a.remoteTabs[tabID]
243 return tab != nil && tab.gen == gen
244 }
245
246 // remoteTabPump forwards Serve events for one tab generation. Cancellation,
247 // stream death, or a generation mismatch retires the pump.
248 func (a *App) remoteTabPump(ctx context.Context, tabID string, gen uint64, opened chan<- error) {
249 signalOpened := func(err error) {
250 if opened == nil {
251 return
252 }
253 select {
254 case opened <- err:
255 default:
256 }
257 opened = nil
258 }
259 a.remoteTabMu.Lock()
260 tab := a.remoteTabs[tabID]
261 var client *http.Client
262 var base string
263 if tab != nil && tab.gen == gen {
264 client, base = tab.client, tab.base
265 }
266 a.remoteTabMu.Unlock()
267 if client == nil || base == "" {
268 if opened != nil {
269 opened <- fmt.Errorf("remote tab %q event stream was retired before opening", tabID)
270 }
271 return
272 }
273
274 req, err := http.NewRequestWithContext(ctx, http.MethodGet, serveURL(base, "/events?all=1"), nil)
275 if err != nil {
276 signalOpened(err)
277 a.emitRemoteTabStateForGeneration(tabID, gen, "error", err.Error())
278 return
279 }
280 req.Header.Set("Accept", "text/event-stream")
281 resp, err := client.Do(req)
282 if err != nil {
283 // Schedule recovery before signalling the opener: the reattach
284 // retirement bumps the generation first, so the opener's own retire
285 // for this error becomes a no-op instead of racing the recovery.
286 if ctx.Err() == nil {
287 log.Printf("[remote] remoteTabPump: /events DO-FAILED tab=%s err=%v", tabID, err)
288 // A tunnel that just dropped the old stream often refuses the
289 // replacement too; parking in error would strand a healthy tab.
290 // Route through the reattach loop, which re-ensures the server
291 // and retries while the transport heals.
292 a.startRemoteTabReattach(tabID, gen)
293 }
294 signalOpened(err)
295 return
296 }
297 defer resp.Body.Close()
298 if resp.StatusCode != http.StatusOK {
299 err = fmt.Errorf("serve /events: status %d", resp.StatusCode)
300 if ctx.Err() == nil {
301 log.Printf("[remote] remoteTabPump: /events BAD-STATUS tab=%s status=%d", tabID, resp.StatusCode)
302 a.startRemoteTabReattach(tabID, gen)
303 }
304 signalOpened(err)
305 return
306 }
307 signalOpened(nil)
308 scanner := bufio.NewScanner(resp.Body)
309 scanner.Buffer(make([]byte, 64<<10), serveEventMaxBytes)
310 for scanner.Scan() {
311 line := scanner.Text()
312 line = strings.TrimRight(line, "\r\n")
313 if !strings.HasPrefix(line, "data:") {
314 continue // ": ping" keepalives and other SSE fields
315 }
316 frame := strings.TrimSpace(strings.TrimPrefix(line, "data:"))
317 if frame == "" {
318 continue
319 }
320 if !a.remoteTabGenerationCurrent(tabID, gen) {
321 return
322 }
323 kind, framePath, current, reset := probeRemoteTabFrame(frame)
324 if kind == "runtime_state" {
325 a.acceptRemoteRuntimeFrame(tabID, gen, framePath, json.RawMessage(frame))
326 continue
327 }
328 // A takeover notice for the session this tab is viewing flips the
329 // spectator pin live: the entry-time probe only runs when the tab
330 // enters a session, so a mid-view takeover (or its reversal) would
331 // otherwise leave the banner and the composer locked to stale state.
332 if kind == "notice" && framePath != "" &&
333 (strings.Contains(frame, event.NoticeCodeSessionTakenOver) ||
334 strings.Contains(frame, event.NoticeCodeSessionReclaimed)) {
335 a.goRemoteTabSafe("remoteTabTakeoverNoticeProbe", func() {
336 a.probeSpectatorAfterNotice(tabID, gen, client, base, framePath)
337 })
338 }
339 if !a.routeRemoteTabWireFrame(tabID, gen, framePath, kind, current, reset) {
340 continue
341 }
342 if a.bufferRemoteTabResumeFrame(tabID, gen, framePath, kind, json.RawMessage(frame)) {
343 continue
344 }
345 a.publishRemoteTabFrame(tabID, gen, framePath, kind, json.RawMessage(frame))
346 }
347 if err := scanner.Err(); err != nil {
348 log.Printf("[remote] remoteTabPump: READ-EXIT tab=%s gen=%d err=%v ctxErr=%v", tabID, gen, err, ctx.Err())
349 }
350 // Only the current generation reacts to an unexpected stream death.
351 // Reattach now; the host status hook also retries on connection recovery.
352 if ctx.Err() == nil {
353 a.startRemoteTabReattach(tabID, gen)
354 }
355 }
356
357 func (a *App) completeRemoteTabTurn(tabID string, gen uint64) {
358 a.remoteTabMu.Lock()
359 tab := a.remoteTabs[tabID]
360 if tab == nil || tab.gen != gen {
361 a.remoteTabMu.Unlock()
362 return
363 }
364 tab.runtime.revision++
365 tab.pendingEvents = nil
366 tab.runtime.running = false
367 tab.runtime.turnStartedAt = 0
368 tab.runtime.pendingPrompt = false
369 tab.runtime.cancelRequested = false
370 tab.runtime.cancellable = tab.runtime.backgroundJobs > 0
371 // A completed turn makes the fresh session non-blank even when the
372 // best-effort /sessions title lookup fails. New Topic must never reuse
373 // a conversation that already has a completed turn.
374 tab.session.reset = false
375 meta := remoteTabMetaLocked(tab)
376 a.remoteTabMu.Unlock()
377 a.emitRemoteEvent("remote-tab:updated", meta)
378 }
379
380 func (a *App) recordRemoteTabTurnStarted(tabID string, gen uint64, frame json.RawMessage) {
381 var payload struct {
382 TurnStartedAt int64 `json:"turnStartedAt"`
383 }
384 _ = json.Unmarshal(frame, &payload)
385 if payload.TurnStartedAt <= 0 {
386 payload.TurnStartedAt = time.Now().UnixMilli()
387 }
388 a.remoteTabMu.Lock()
389 tab := a.remoteTabs[tabID]
390 if tab == nil || tab.gen != gen {
391 a.remoteTabMu.Unlock()
392 return
393 }
394 tab.runtime.revision++
395 tab.runtime.running = true
396 tab.runtime.turnStartedAt = payload.TurnStartedAt
397 tab.runtime.pendingPrompt = false
398 tab.runtime.cancelRequested = false
399 tab.runtime.cancellable = true
400 meta := remoteTabMetaLocked(tab)
401 a.remoteTabMu.Unlock()
402 a.emitRemoteEvent("remote-tab:updated", meta)
403 }
404
405 func (a *App) cacheRemotePendingEvent(tabID string, gen uint64, kind string, frame json.RawMessage) {
406 key := remotePendingEventKey(kind, frame)
407 a.remoteTabMu.Lock()
408 tab := a.remoteTabs[tabID]
409 if tab == nil || tab.gen != gen {
410 a.remoteTabMu.Unlock()
411 return
412 }
413 tab.runtime.revision++
414 if tab.pendingEvents == nil {
415 tab.pendingEvents = make(map[string]json.RawMessage)
416 }
417 tab.pendingEvents[key] = append(json.RawMessage(nil), frame...)
418 tab.runtime.pendingPrompt = true
419 tab.runtime.cancellable = true
420 meta := remoteTabMetaLocked(tab)
421 a.remoteTabMu.Unlock()
422 a.emitRemoteEvent("remote-tab:updated", meta)
423 }
424
425 func (a *App) clearRemotePendingEvent(tabID, kind, callID string) {
426 a.remoteTabMu.Lock()
427 var meta TabMeta
428 changed := false
429 if tab := a.remoteTabs[tabID]; tab != nil {
430 tab.runtime.revision++
431 delete(tab.pendingEvents, kind+":"+strings.TrimSpace(callID))
432 pending := len(tab.pendingEvents) > 0
433 changed = tab.runtime.pendingPrompt != pending
434 tab.runtime.pendingPrompt = pending
435 meta = remoteTabMetaLocked(tab)
436 }
437 a.remoteTabMu.Unlock()
438 if changed {
439 a.emitRemoteEvent("remote-tab:updated", meta)
440 }
441 }
442
443 // serveGet fetches a JSON member of the tab snapshot, returning the raw
444 // payload for verbatim passthrough.
445 func serveGet(ctx context.Context, client *http.Client, url string, expectedPath ...string) (json.RawMessage, error) {
446 req, err := http.NewRequestWithContext(ctx, http.MethodGet, url, nil)
447 if err != nil {
448 return nil, err
449 }
450 if len(expectedPath) != 0 && strings.TrimSpace(expectedPath[0]) != "" {
451 if sessionID, ok := strings.CutPrefix(expectedPath[0], remoteSessionIDRoutePrefix); ok {
452 req.Header.Set(expectedSessionIDHeader, sessionID)
453 } else {
454 req.Header.Set(expectedSessionPathHeader, expectedPath[0])
455 }
456 }
457 resp, err := client.Do(req)
458 if err != nil {
459 return nil, err
460 }
461 defer resp.Body.Close()
462 data, err := io.ReadAll(io.LimitReader(resp.Body, serveSnapshotMaxBytes+1))
463 if err != nil {
464 return nil, err
465 }
466 if len(data) > serveSnapshotMaxBytes {
467 return nil, fmt.Errorf("%s: response exceeds %d bytes", url, serveSnapshotMaxBytes)
468 }
469 if resp.StatusCode != http.StatusOK {
470 return nil, fmt.Errorf("%s: status %d", url, resp.StatusCode)
471 }
472 return json.RawMessage(data), nil
473 }
474
475 // commandContext bounds one proxied command. Boot context when available;
476 // the timeout keeps a wedged tunnel from hanging the binding call.
477 func commandContext(a *App) (context.Context, context.CancelFunc) {
478 ctx := a.bootContext()
479 if ctx == nil {
480 ctx = context.Background()
481 }
482 return context.WithTimeout(ctx, 15*time.Second)
483 }
484
485 // remoteTabCommandClient resolves a tabID to its live serve client. A tab
486 // that has not finished bootstrap, is reconnecting, or has failed is an
487 // error, not a silent no-op.
488 func (a *App) remoteTabCommandClient(tabID string) (*http.Client, string, error) {
489 client, base, _, err := a.remoteTabCommandTarget(tabID)
490 return client, base, err
491 }
492
493 func (a *App) remoteTabCommandTarget(tabID string) (*http.Client, string, string, error) {
494 a.remoteTabMu.Lock()
495 tab := a.remoteTabs[tabID]
496 var client *http.Client
497 var base, expectedPath string
498 switching := tab != nil && tab.routing.rehydratingPath != ""
499 usable := tab != nil && tab.client != nil && tab.state == "ready" && !switching
500 if usable {
501 client, base = tab.client, tab.base
502 expectedPath = tab.routing.currentPath
503 }
504 a.remoteTabMu.Unlock()
505 if !usable {
506 if switching {
507 return nil, "", "", fmt.Errorf("remote tab %q is switching sessions; wait for it to become ready", tabID)
508 }
509 return nil, "", "", fmt.Errorf("remote tab %q is not connected", tabID)
510 }
511 return client, base, expectedPath, nil
512 }
513
514 // remoteTabAdmissionCurrent reports whether the tab still runs the generation
515 // a run-admission decision (model-settings revision or legacy skip) was made
516 // for. Generation 0 marks an ungated decision; any other replaced generation
517 // must be re-admitted so a reconnect's newer Serve never receives an unfenced
518 // request that was approved against the retired connection.
519 func (a *App) remoteTabAdmissionCurrent(tabID string, generation uint64) bool {
520 if generation == 0 {
521 return true
522 }
523 a.remoteTabMu.Lock()
524 defer a.remoteTabMu.Unlock()
525 tab := a.remoteTabs[tabID]
526 return tab != nil && tab.gen == generation
527 }
528
529 func (a *App) isRemoteTab(tabID string) bool {
530 if strings.TrimSpace(tabID) == "" {
531 return false
532 }
533 a.remoteTabMu.Lock()
534 _, ok := a.remoteTabs[tabID]
535 a.remoteTabMu.Unlock()
536 return ok
537 }
538
539 // remoteTabRefFor returns the host+workspace ref when tabID belongs to a
540 // remote tab; view builders use it to mark remote-shaped metas.
541 func (a *App) remoteTabRefFor(tabID string) (RemoteTabRef, bool) {
542 a.remoteTabMu.Lock()
543 defer a.remoteTabMu.Unlock()
544 if tab := a.remoteTabs[tabID]; tab != nil {
545 return tab.ref, true
546 }
547 return RemoteTabRef{}, false
548 }
549
550 func (a *App) remoteTabCurrentModel(tabID string) (string, bool) {
551 if !a.isRemoteTab(tabID) {
552 return "", false
553 }
554 a.remoteTabMu.Lock()
555 tab := a.remoteTabs[tabID]
556 cur := ""
557 if tab != nil {
558 cur = tab.model
559 }
560 a.remoteTabMu.Unlock()
561 return cur, true
562 }
563
564 // ReclaimRemoteTabSession takes a mirrored session back from the local
565 // runtime that took it over. Serve long-polls until the local writer yields,
566 // so this call can outlast a normal command timeout.
567 func (a *App) ReclaimRemoteTabSession(tabID string) error {
568 if err := a.requireRemoteExecutionProtocol(tabID); err != nil {
569 return err
570 }
571 client, base, expectedPath, err := a.remoteTabCommandTarget(tabID)
572 if err != nil {
573 return err
574 }
575 if strings.TrimSpace(expectedPath) == "" {
576 return fmt.Errorf("remote tab %q has no active session", tabID)
577 }
578 observed, err := a.observeRemoteTabForReclaim(tabID, client)
579 if err != nil {
580 return err
581 }
582 observedTab, observedGen := observed.tab, observed.gen
583 stillCurrent := func(tab *remoteTab) bool {
584 return tab != nil && tab == observedTab && tab.client == client && tab.gen == observed.gen &&
585 tab.runtime.revision == observed.runtimeRevision && tab.selectionRevision == observed.selectionRevision &&
586 agent.CanonicalSessionPath(tab.routing.currentPath) == agent.CanonicalSessionPath(expectedPath)
587 }
588 reconcileOwnership := func() { a.reconcileRemoteTabReclaimOwnership(tabID, client, base, expectedPath, stillCurrent) }
589 // Short timeout: the serve caps un-mirrored reclaims at 10s and mirrored
590 // ones use the writer's cooperative heartbeat (seconds, not minutes). A
591 // long client-side timeout only hangs the UI button.
592 ctx, cancel := context.WithTimeout(context.Background(), 20*time.Second)
593 defer cancel()
594 body, _ := json.Marshal(map[string]any{
595 "sessionPath": expectedPath,
596 "mode": "wait",
597 "timeoutMs": 15000,
598 })
599 resp, err := serveDo(ctx, client, http.MethodPost, serveURL(base, "/reclaim"), body)
600 if err != nil {
601 reconcileOwnership()
602 return fmt.Errorf("reclaim session: %w", err)
603 }
604 defer resp.Body.Close()
605 respBody, _ := io.ReadAll(io.LimitReader(resp.Body, 1<<16))
606 if resp.StatusCode != http.StatusNoContent {
607 errMsg := strings.TrimSpace(string(respBody))
608 // A failed reclaim is not proof that ownership changed — generation
609 // conflicts and transient 5xx included. Keep the spectator pin until a
610 // fenced probe proves this exact binding is no longer locally owned.
611 reconcileOwnership()
612 return fmt.Errorf("reclaim session: %s", errMsg)
613 }
614 // Reclaim succeeded: Serve now owns the session again. Clear the spectator
615 // pin immediately so the composer un-locks without waiting for the next
616 // status poll to observe takenOver=false.
617 observedTab.routeEventMu.Lock()
618 defer observedTab.routeEventMu.Unlock()
619 a.remoteTabMu.Lock()
620 if tab := a.remoteTabs[tabID]; stillCurrent(tab) {
621 tab.session.takenOver = false
622 // Fence status payloads reserved before this reclaim: they may still
623 // be in flight and carry the pre-reclaim takenOver=true, which would
624 // re-pin the spectator banner the moment ownership returned.
625 tab.ownership.reclaimRevision = tab.runtime.revision + 1
626 deferBarrier := tab.runtime.running || tab.runtime.pendingPrompt
627 tab.ownership.readyBarrierPending = deferBarrier
628 meta := remoteTabMetaLocked(tab)
629 a.remoteTabMu.Unlock()
630 a.emitRemoteEvent("remote-tab:updated", meta)
631 // The spectator era froze the projection, so publish the ready barrier
632 // to re-hydrate the view and accept the re-owned writer's frames. Defer
633 // it mid-turn: the barrier bumps the frontend connection generation.
634 if !deferBarrier {
635 a.transitionRemoteTabStateLocked(tab, observedGen, "ready", "ready", "")
636 }
637 } else {
638 a.remoteTabMu.Unlock()
639 }
640 a.goRemoteTabSafe("reclaimStatusRefresh", func() { _, _ = a.RemoteTabStatus(tabID) })
641 return nil
642 }
643
644 func (a *App) CancelRemoteTab(tabID string) error {
645 if err := a.requireRemoteExecutionProtocol(tabID); err != nil {
646 return err
647 }
648 client, base, expectedPath, err := a.remoteTabCommandTarget(tabID)
649 if err != nil {
650 return err
651 }
652 ctx, cancel := commandContext(a)
653 defer cancel()
654 return servePostForSession(ctx, client, serveURL(base, "/cancel"), nil, expectedPath)
655 }
656
657 // ApproveRemoteTab answers a tool-approval request. Only one-shot and scoped
658 // session grants are supported; durable approval rules were intentionally
659 // removed from the permission model.
660 func (a *App) ApproveRemoteTab(tabID, callID, decision string) error {
661 if err := a.requireRemoteExecutionProtocol(tabID); err != nil {
662 return err
663 }
664 client, base, expectedPath, err := a.remoteTabCommandTarget(tabID)
665 if err != nil {
666 return err
667 }
668 ctx, cancel := commandContext(a)
669 defer cancel()
670 decision = strings.ToLower(strings.TrimSpace(decision))
671 allow, session := false, false
672 switch decision {
673 case "allow", "once":
674 allow = true
675 case "session":
676 allow, session = true, true
677 case "persist", "persistent", "project":
678 return fmt.Errorf("permanent remote approval is no longer supported")
679 case "deny":
680 default:
681 return fmt.Errorf("invalid remote approval decision %q", decision)
682 }
683 body, _ := json.Marshal(map[string]any{"id": callID, "allow": allow, "session": session, "persist": false})
684 if err := servePostForSession(ctx, client, serveURL(base, "/approve"), body, expectedPath); err != nil {
685 return err
686 }
687 a.clearRemotePendingEvent(tabID, "approval_request", callID)
688 return nil
689 }
690
691 // ResolveRemoteTabPlanDecision preserves the three distinct exit_plan_mode
692 // outcomes that the generic approval boolean cannot represent. Revision text
693 // travels in the same Serve request so the controller can durably stage it
694 // before resolving the approval; a tunnel failure can no longer split the
695 // decision from the requested revision.
696 func (a *App) ResolveRemoteTabPlanDecision(tabID, callID, action, feedback string) error {
697 if err := a.requireRemoteExecutionProtocol(tabID); err != nil {
698 return err
699 }
700 client, base, expectedPath, err := a.remoteTabCommandTarget(tabID)
701 if err != nil {
702 return err
703 }
704 action = strings.ToLower(strings.TrimSpace(action))
705 switch action {
706 case "start_execution", "revise_plan", "exit_plan":
707 default:
708 return fmt.Errorf("invalid remote plan decision %q", action)
709 }
710 ctx, cancel := commandContext(a)
711 defer cancel()
712 body, _ := json.Marshal(map[string]string{"id": callID, "action": action, "feedback": strings.TrimSpace(feedback)})
713 if err := servePostForSession(ctx, client, serveURL(base, "/plan-decision"), body, expectedPath); err != nil {
714 return err
715 }
716 a.clearRemotePendingEvent(tabID, "approval_request", callID)
717 return nil
718 }
719
720 type RemoteAskAnswer struct {
721 QuestionID string `json:"QuestionID"`
722 Selected []string `json:"Selected"`
723 }
724
725 // AnswerRemoteTab preserves the batch ask id at the top level and sends every
726 // question's own id/selections in the Serve AskAnswer wire shape.
727 func (a *App) AnswerRemoteTab(tabID, callID string, answers []RemoteAskAnswer) error {
728 if err := a.requireRemoteExecutionProtocol(tabID); err != nil {
729 return err
730 }
731 client, base, expectedPath, err := a.remoteTabCommandTarget(tabID)
732 if err != nil {
733 return err
734 }
735 ctx, cancel := commandContext(a)
736 defer cancel()
737 body, _ := json.Marshal(map[string]any{
738 "id": callID,
739 "answers": answers,
740 })
741 if err := servePostForSession(ctx, client, serveURL(base, "/answer"), body, expectedPath); err != nil {
742 return err
743 }
744 a.clearRemotePendingEvent(tabID, "ask_request", callID)
745 return nil
746 }
747
748 func (a *App) SubmitRemoteTabExtensionForm(tabID, pluginID, surfaceID string, values map[string]any) error {
749 if err := a.remoteTabPost(tabID, "/extension-form", map[string]any{
750 "pluginId": pluginID, "surfaceId": surfaceID, "values": values,
751 }); err != nil {
752 return err
753 }
754 a.clearRemotePendingExtensionForm(tabID, pluginID, surfaceID)
755 return nil
756 }
757
758 // RewindRemoteTab rewinds to a checkpoint. Serve identifies checkpoints by
759 // TURN index and takes {turn, scope}; the checkpointID string is that turn.
760 func (a *App) RewindRemoteTab(tabID, checkpointID, scope string) error {
761 if err := a.requireRemoteExecutionProtocol(tabID); err != nil {
762 return err
763 }
764 client, base, expectedPath, err := a.remoteTabCommandTarget(tabID)
765 if err != nil {
766 return err
767 }
768 ctx, cancel := commandContext(a)
769 defer cancel()
770 turn, convErr := strconv.Atoi(strings.TrimSpace(checkpointID))
771 if convErr != nil {
772 return fmt.Errorf("invalid checkpoint id %q: want the turn index", checkpointID)
773 }
774 scope = strings.TrimSpace(scope)
775 switch scope {
776 case "code", "conversation", "both":
777 default:
778 return fmt.Errorf("invalid rewind scope %q", scope)
779 }
780 body, _ := json.Marshal(map[string]any{"turn": turn, "scope": scope})
781 return servePostForSession(ctx, client, serveURL(base, "/rewind"), body, expectedPath)
782 }
783
784 func (a *App) SetRemoteTabToolApprovalMode(tabID, mode string) error {
785 if err := a.requireRemoteExecutionProtocol(tabID); err != nil {
786 return err
787 }
788 if err := a.requireRemotePermissionPresets(tabID); err != nil {
789 return err
790 }
791 ctx, cancel := commandContext(a)
792 defer cancel()
793 client, base, expectedPath, err := a.remoteTabCommandTarget(tabID)
794 if err != nil {
795 return err
796 }
797 snapshot, err := remotePermissionSnapshot(ctx, client, base, expectedPath)
798 if err != nil {
799 return err
800 }
801 _, err = setRemotePermissionPresetAt(ctx, client, base, expectedPath, mode, snapshot.Revision)
802 return err
803 }
804
805 func setRemotePermissionPresetAt(ctx context.Context, client *http.Client, base, expectedPath, mode string, revision uint64) (control.PermissionSnapshot, error) {
806 var snapshot control.PermissionSnapshot
807 body, _ := json.Marshal(map[string]any{"preset": mode, "expectedRevision": revision})
808 resp, err := serveDoForSession(ctx, client, http.MethodPost, serveURL(base, "/permission/preset"), body, expectedPath)
809 if err != nil {
810 return snapshot, err
811 }
812 defer resp.Body.Close()
813 if resp.StatusCode < 200 || resp.StatusCode >= 300 {
814 data, _ := io.ReadAll(io.LimitReader(resp.Body, 4<<10))
815 return snapshot, fmt.Errorf("set remote permission preset: status %d: %s", resp.StatusCode, strings.TrimSpace(string(data)))
816 }
817 var envelope struct {
818 Snapshot control.PermissionSnapshot `json:"snapshot"`
819 }
820 if err := json.NewDecoder(io.LimitReader(resp.Body, 1<<20)).Decode(&envelope); err != nil {
821 return snapshot, fmt.Errorf("decode remote permission update: %w", err)
822 }
823 return envelope.Snapshot, nil
824 }
825
826 func (a *App) SetRemoteTabComposerProfile(tabID, collaborationMode, toolApprovalMode, goal string) ([]string, error) {
827 if err := a.requireRemoteExecutionProtocol(tabID); err != nil {
828 return nil, err
829 }
830 if err := a.requireRemotePermissionPresets(tabID); err != nil {
831 return nil, err
832 }
833 if strings.EqualFold(strings.TrimSpace(collaborationMode), "goal") || strings.TrimSpace(goal) != "" {
834 if err := a.requireRemoteGoalLifecycle(tabID); err != nil {
835 return nil, err
836 }
837 }
838 client, base, expectedPath, err := a.remoteTabCommandTarget(tabID)
839 if err != nil {
840 return nil, err
841 }
842 ctx, cancel := commandContext(a)
843 defer cancel()
844 snapshot, err := remotePermissionSnapshot(ctx, client, base, expectedPath)
845 if err != nil {
846 return nil, err
847 }
848 body, _ := json.Marshal(map[string]any{
849 "collaborationMode": collaborationMode,
850 "toolApprovalMode": toolApprovalMode,
851 "goal": goal,
852 "expectedPermissionRevision": snapshot.Revision,
853 })
854 resp, err := serveDoForSession(ctx, client, http.MethodPost, serveURL(base, "/composer-profile"), body, expectedPath)
855 if err != nil {
856 return nil, err
857 }
858 defer resp.Body.Close()
859 data, _ := io.ReadAll(io.LimitReader(resp.Body, 4<<10))
860 if resp.StatusCode < 200 || resp.StatusCode >= 300 {
861 if message := strings.TrimSpace(string(data)); message != "" {
862 return nil, fmt.Errorf("%s: status %d: %s", serveURL(base, "/composer-profile"), resp.StatusCode, message)
863 }
864 return nil, fmt.Errorf("%s: status %d", serveURL(base, "/composer-profile"), resp.StatusCode)
865 }
866 if len(bytes.TrimSpace(data)) == 0 {
867 return []string{}, nil
868 }
869 var result struct {
870 DrainedApprovalIDs []string `json:"drainedApprovalIDs"`
871 }
872 if err := json.Unmarshal(data, &result); err != nil {
873 return nil, fmt.Errorf("decode remote composer profile response: %w", err)
874 }
875 return result.DrainedApprovalIDs, nil
876 }
877
878 func remotePermissionSnapshot(ctx context.Context, client *http.Client, base, expectedPath string) (control.PermissionSnapshot, error) {
879 var snapshot control.PermissionSnapshot
880 resp, err := serveDoForSession(ctx, client, http.MethodGet, serveURL(base, "/permission"), nil, expectedPath)
881 if err != nil {
882 return snapshot, err
883 }
884 defer resp.Body.Close()
885 if resp.StatusCode < 200 || resp.StatusCode >= 300 {
886 data, _ := io.ReadAll(io.LimitReader(resp.Body, 4<<10))
887 return snapshot, fmt.Errorf("query remote permission snapshot: status %d: %s", resp.StatusCode, strings.TrimSpace(string(data)))
888 }
889 if err := json.NewDecoder(io.LimitReader(resp.Body, 1<<20)).Decode(&snapshot); err != nil {
890 return snapshot, fmt.Errorf("decode remote permission snapshot: %w", err)
891 }
892 return snapshot, nil
893 }
894
895 func revokeRemotePermissionGrantAt(ctx context.Context, client *http.Client, base, expectedPath, scope, target string, revision uint64) (control.PermissionSnapshot, error) {
896 var snapshot control.PermissionSnapshot
897 body, _ := json.Marshal(map[string]any{"scope": scope, "target": target, "expectedRevision": revision})
898 resp, err := serveDoForSession(ctx, client, http.MethodPost, serveURL(base, "/permission/grants/revoke"), body, expectedPath)
899 if err != nil {
900 return snapshot, err
901 }
902 defer resp.Body.Close()
903 if resp.StatusCode < 200 || resp.StatusCode >= 300 {
904 data, _ := io.ReadAll(io.LimitReader(resp.Body, 4<<10))
905 return snapshot, fmt.Errorf("revoke remote permission grant: status %d: %s", resp.StatusCode, strings.TrimSpace(string(data)))
906 }
907 if err := json.NewDecoder(io.LimitReader(resp.Body, 1<<20)).Decode(&snapshot); err != nil {
908 return snapshot, fmt.Errorf("decode remote permission revocation: %w", err)
909 }
910 return snapshot, nil
911 }
912
913 func (a *App) requireRemotePermissionPresets(tabID string) error {
914 a.remoteTabMu.Lock()
915 tab := a.remoteTabs[tabID]
916 supported := tab != nil && tab.capabilities["permission-presets-v1"]
917 a.remoteTabMu.Unlock()
918 if tab == nil {
919 return fmt.Errorf("remote tab %q is not open", tabID)
920 }
921 if !supported {
922 return fmt.Errorf("this remote Reasonix Serve is read-only because it does not support permission-presets-v1; upgrade the remote service to run tools or change permissions")
923 }
924 return nil
925 }
926
927 // requireRemoteExecutionProtocol fences every state-changing command at the
928 // authenticated Serve capability boundary. A legacy Serve remains usable for
929 // history reads, but Desktop never emulates the v3 runtime over older RPCs.
930 func (a *App) requireRemoteExecutionProtocol(tabID string) error {
931 a.remoteTabMu.Lock()
932 tab := a.remoteTabs[tabID]
933 supported := tab != nil && tab.capabilities[serveCapabilityExecutionV2] && tab.capabilities[serveCapabilitySessions] && tab.capabilities[serveCapabilitySessionIdentityV1] && tab.capabilities[serveCapabilitySessionOwnershipV1]
934 a.remoteTabMu.Unlock()
935 if tab == nil {
936 return fmt.Errorf("remote tab %q is not open", tabID)
937 }
938 if !supported {
939 return fmt.Errorf("this remote Reasonix Serve is read-only because it does not support %s, %s, %s, and %s; upgrade the remote service to execute or control a session", serveCapabilityExecutionV2, serveCapabilitySessions, serveCapabilitySessionIdentityV1, serveCapabilitySessionOwnershipV1)
940 }
941 return nil
942 }
943
944 func (a *App) SetRemoteTabGoal(tabID, goal string) error {
945 if err := a.requireRemoteExecutionProtocol(tabID); err != nil {
946 return err
947 }
948 if err := a.requireRemoteGoalLifecycle(tabID); err != nil {
949 return err
950 }
951 client, base, expectedPath, err := a.remoteTabCommandTarget(tabID)
952 if err != nil {
953 return err
954 }
955 ctx, cancel := commandContext(a)
956 defer cancel()
957 body, _ := json.Marshal(map[string]string{"goal": goal})
958 return servePostForSession(ctx, client, serveURL(base, "/goal"), body, expectedPath)
959 }
960
961 func (a *App) requireRemoteGoalLifecycle(tabID string) error {
962 a.remoteTabMu.Lock()
963 tab := a.remoteTabs[tabID]
964 supported := tab != nil && tab.capabilities[serveCapabilityGoalLifecycleV2]
965 a.remoteTabMu.Unlock()
966 if tab == nil {
967 return fmt.Errorf("remote tab %q is not open", tabID)
968 }
969 if !supported {
970 return fmt.Errorf("this remote Reasonix Serve does not support %s; upgrade it before creating or controlling goals", serveCapabilityGoalLifecycleV2)
971 }
972 return nil
973 }
974
975 func (a *App) SetRemoteTabQualityFloor(tabID, floor string) error {
976 return a.validateRemoteQualityFloor(tabID, floor)
977 }
978
978 lines GO