返回 DeepSeek-Reasonix
browser_broker.go
根目录 / desktop / browser_broker.go
1 package main
2
3 import (
4 "context"
5 "crypto/rand"
6 "encoding/hex"
7 "encoding/json"
8 "fmt"
9 "net"
10 "net/http"
11 "os"
12 "strings"
13 "sync"
14 "time"
15
16 "reasonix/internal/browser"
17 "reasonix/internal/remote/forward"
18 )
19
20 func (s *brokerSessionExecutor) BrowserCapability(ctx context.Context, name string, args json.RawMessage) (json.RawMessage, error) {
21 res, err := s.resolve(ctx)
22 if err != nil {
23 return nil, err
24 }
25 exec, ok := res.exec.(browser.CapabilityExecutor)
26 if !ok {
27 return nil, fmt.Errorf("capability_unsupported: %s", name)
28 }
29 out, err := exec.BrowserCapability(ctx, name, args)
30 if err != nil {
31 return out, err
32 }
33 if current, err := s.resolve(ctx); err != nil || current.workspace != res.workspace {
34 return nil, browser.ErrUnknownOutcome
35 }
36 if name != "record" {
37 return out, nil
38 }
39 var result map[string]json.RawMessage
40 if err := json.Unmarshal(out, &result); err != nil {
41 return nil, err
42 }
43 var localPath, state string
44 _ = json.Unmarshal(result["path"], &localPath)
45 _ = json.Unmarshal(result["state"], &state)
46 delete(result, "path")
47 if state == "completed" && localPath != "" {
48 remote, err := s.relay(ctx, res.workspace, localPath)
49 if err != nil {
50 return nil, err
51 }
52 result["path"], _ = json.Marshal(remote)
53 }
54 if current, err := s.resolve(ctx); err != nil || current.workspace != res.workspace {
55 return nil, browser.ErrUnknownOutcome
56 }
57 return json.Marshal(result)
58 }
59
60 // The desktop browser broker is the local end of the remote browser channel:
61 // a loopback listener behind SSH reverse forwards, keyed by per-generation
62 // tokens that a reconnect revokes at once.
63
64 // browserBrokerForwardName prefixes the per-host reverse forward the remote
65 // serve's REASONIX_BROWSER_BROKER endpoint points at.
66 const browserBrokerForwardName = "browser-broker:"
67
68 // browserBrokerRoute binds one token to one host connection generation.
69 type browserBrokerRoute struct {
70 hostID string
71 gen *managedHost
72 ctx context.Context
73 cancel context.CancelFunc
74 }
75
76 // browserSessionResolution is what the broker resolves one request's session
77 // header into: the desktop executor bound to that session's tab and the
78 // workspace whose SFTP scratch area relays captures back.
79 type browserSessionResolution struct {
80 exec browser.Executor
81 workspace string
82 }
83
84 // browserSessionResolver maps (host, remote session path) to the desktop
85 // executor that owns it. An unknown or foreign session must fail with
86 // browser.ErrNoGrant so the wire handler answers 409 no_grant.
87 type browserSessionResolver func(hostID, sessionPath string) (browserSessionResolution, error)
88
89 type browserBroker struct {
90 lifecycleMu sync.Mutex
91 mu sync.Mutex
92 ln net.Listener
93 server *http.Server
94 port int
95 routes map[string]*browserBrokerRoute
96 byHost map[string]string
97 resolve browserSessionResolver
98 // current reports whether gen is still the live connection for hostID;
99 // a replaced generation's token stops authenticating immediately.
100 current func(hostID string, gen *managedHost) bool
101 // connFor returns the generation's SSH client for the capture relay.
102 connFor func(hostID string, gen *managedHost) sftpConn
103 // newRelay builds the capture relay for a connection; nil uses the SFTP
104 // relay. Tests substitute a fake.
105 newRelay func(conn sftpConn) FileRelay
106 onRevoke func(hostID string)
107 }
108
109 func newBrowserBroker(resolve browserSessionResolver, current func(string, *managedHost) bool, connFor func(string, *managedHost) sftpConn) *browserBroker {
110 return &browserBroker{
111 routes: map[string]*browserBrokerRoute{},
112 byHost: map[string]string{},
113 resolve: resolve,
114 current: current,
115 connFor: connFor,
116 }
117 }
118
119 // register mints a fresh token for (hostID, gen), replacing the host's
120 // previous token. Returns the token and the broker's loopback port.
121 func (b *browserBroker) register(hostID string, gen *managedHost) (string, int, error) {
122 b.lifecycleMu.Lock()
123 defer b.lifecycleMu.Unlock()
124 buf := make([]byte, 32)
125 if _, err := rand.Read(buf); err != nil {
126 return "", 0, fmt.Errorf("browser broker: mint token: %w", err)
127 }
128 token := hex.EncodeToString(buf)
129 b.mu.Lock()
130 if b.ln == nil {
131 b.mu.Unlock()
132 return "", 0, fmt.Errorf("browser broker: not running")
133 }
134 replaced := false
135 if old := b.byHost[hostID]; old != "" {
136 if route := b.routes[old]; route != nil && route.cancel != nil {
137 route.cancel()
138 }
139 delete(b.routes, old)
140 replaced = true
141 }
142 ctx, cancel := context.WithCancel(context.Background())
143 b.routes[token] = &browserBrokerRoute{hostID: hostID, gen: gen, ctx: ctx, cancel: cancel}
144 b.byHost[hostID] = token
145 port := b.port
146 b.mu.Unlock()
147 if replaced && b.onRevoke != nil {
148 b.onRevoke(hostID)
149 }
150 return token, port, nil
151 }
152
153 // revokeHost drops every token minted for hostID (serve stop, disconnect).
154 func (b *browserBroker) revokeHost(hostID string) {
155 b.lifecycleMu.Lock()
156 defer b.lifecycleMu.Unlock()
157 b.mu.Lock()
158 if token := b.byHost[hostID]; token != "" {
159 if route := b.routes[token]; route != nil && route.cancel != nil {
160 route.cancel()
161 }
162 delete(b.routes, token)
163 delete(b.byHost, hostID)
164 }
165 b.mu.Unlock()
166 if b.onRevoke != nil {
167 b.onRevoke(hostID)
168 }
169 }
170
171 func (b *browserBroker) close() {
172 b.lifecycleMu.Lock()
173 defer b.lifecycleMu.Unlock()
174 b.mu.Lock()
175 server, listener := b.server, b.ln
176 hosts := make([]string, 0, len(b.byHost))
177 for _, route := range b.routes {
178 if route.cancel != nil {
179 route.cancel()
180 }
181 hosts = append(hosts, route.hostID)
182 }
183 b.server, b.ln = nil, nil
184 b.routes = map[string]*browserBrokerRoute{}
185 b.byHost = map[string]string{}
186 b.mu.Unlock()
187 if b.onRevoke != nil {
188 for _, hostID := range hosts {
189 b.onRevoke(hostID)
190 }
191 }
192 if server != nil {
193 _ = server.Close()
194 }
195 if listener != nil {
196 _ = listener.Close()
197 }
198 }
199
200 func (b *browserBroker) ServeHTTP(w http.ResponseWriter, r *http.Request) {
201 // Unauthenticated liveness for the reverse-tunnel probe, mirroring the
202 // credential proxy: the listener is only reachable through the tunnel.
203 if r.URL.Path == "/healthz" {
204 w.WriteHeader(http.StatusNoContent)
205 return
206 }
207 token := bearerToken(r.Header.Get("Authorization"))
208 b.mu.Lock()
209 route := b.routes[token]
210 b.mu.Unlock()
211 if token == "" || route == nil || (b.current != nil && !b.current(route.hostID, route.gen)) {
212 w.Header().Set("WWW-Authenticate", `Bearer realm="reasonix-browser-broker"`)
213 http.Error(w, "invalid or stale browser broker token", http.StatusUnauthorized)
214 return
215 }
216 exec := &brokerSessionExecutor{broker: b, route: route}
217 if route.ctx != nil {
218 ctx, cancel := context.WithCancel(r.Context())
219 body := r.Body
220 stop := context.AfterFunc(route.ctx, func() { cancel(); _ = body.Close() })
221 defer stop()
222 defer cancel()
223 r = r.WithContext(ctx)
224 }
225 browser.NewHTTPHandler(exec, token).ServeHTTP(w, r)
226 }
227
228 // brokerSessionExecutor is the per-request executor the broker serves: every
229 // method resolves the request's session header to the desktop tab that owns
230 // it, so one host token can never drive another session's browser.
231 type brokerSessionExecutor struct {
232 broker *browserBroker
233 route *browserBrokerRoute
234 }
235
236 func (s *brokerSessionExecutor) resolve(ctx context.Context) (browserSessionResolution, error) {
237 if !s.current(ctx) {
238 return browserSessionResolution{}, browser.ErrNoGrant
239 }
240 res, err := s.broker.resolve(s.route.hostID, browser.SessionFromContext(ctx))
241 if err != nil {
242 return browserSessionResolution{}, err
243 }
244 if res.exec == nil || !s.current(ctx) {
245 return browserSessionResolution{}, browser.ErrNoGrant
246 }
247 return res, nil
248 }
249
250 func (s *brokerSessionExecutor) current(ctx context.Context) bool {
251 return ctx.Err() == nil && (s.route.ctx == nil || s.route.ctx.Err() == nil) && (s.broker.current == nil || s.broker.current(s.route.hostID, s.route.gen))
252 }
253
254 func (s *brokerSessionExecutor) Available(ctx context.Context) bool {
255 res, err := s.resolve(ctx)
256 if err != nil {
257 return false
258 }
259 if a, ok := res.exec.(browser.Availability); ok {
260 return a.Available(ctx)
261 }
262 return true
263 }
264
265 func (s *brokerSessionExecutor) Tabs(ctx context.Context) ([]browser.Tab, error) {
266 res, err := s.resolve(ctx)
267 if err != nil {
268 return nil, err
269 }
270 return res.exec.Tabs(ctx)
271 }
272
273 func (s *brokerSessionExecutor) Open(ctx context.Context, req browser.OpenRequest) (browser.Tab, error) {
274 res, err := s.resolve(ctx)
275 if err != nil {
276 return browser.Tab{}, err
277 }
278 return res.exec.Open(ctx, req)
279 }
280
281 func (s *brokerSessionExecutor) Navigate(ctx context.Context, req browser.NavigateRequest) (browser.Tab, error) {
282 res, err := s.resolve(ctx)
283 if err != nil {
284 return browser.Tab{}, err
285 }
286 return res.exec.Navigate(ctx, req)
287 }
288
289 func (s *brokerSessionExecutor) Snapshot(ctx context.Context, req browser.SnapshotRequest) (browser.Snapshot, error) {
290 res, err := s.resolve(ctx)
291 if err != nil {
292 return browser.Snapshot{}, err
293 }
294 return res.exec.Snapshot(ctx, req)
295 }
296
297 // Screenshot relays the capture file onto the remote host before answering:
298 // the path the serve receives must be local to the serve, never a desktop
299 // path it cannot read.
300 func (s *brokerSessionExecutor) Screenshot(ctx context.Context, req browser.ScreenshotRequest) (browser.Screenshot, error) {
301 res, err := s.resolve(ctx)
302 if err != nil {
303 return browser.Screenshot{}, err
304 }
305 shot, err := res.exec.Screenshot(ctx, req)
306 if err != nil {
307 return browser.Screenshot{}, err
308 }
309 shot.Path, err = s.relay(ctx, res.workspace, shot.Path)
310 if err != nil {
311 return browser.Screenshot{}, err
312 }
313 return shot, nil
314 }
315
316 func (s *brokerSessionExecutor) Downloads(ctx context.Context, req browser.DownloadsRequest) ([]browser.Download, error) {
317 res, err := s.resolve(ctx)
318 if err != nil {
319 return nil, err
320 }
321 downloads, err := res.exec.Downloads(ctx, req)
322 if err != nil {
323 return nil, err
324 }
325 for i, d := range downloads {
326 if strings.TrimSpace(d.Path) == "" {
327 continue
328 }
329 downloads[i].Path, err = s.relay(ctx, res.workspace, d.Path)
330 if err != nil {
331 return nil, err
332 }
333 }
334 return downloads, nil
335 }
336
337 func (s *brokerSessionExecutor) Act(ctx context.Context, req browser.ActRequest) (browser.ActResult, error) {
338 res, err := s.resolve(ctx)
339 if err != nil {
340 return browser.ActResult{}, err
341 }
342 if req.Action == browser.ActionUpload {
343 owner, ok := res.exec.(interface{ captureDir() (string, error) })
344 if !ok || s.broker.connFor == nil {
345 return browser.ActResult{}, fmt.Errorf("browser upload: no staging owner")
346 }
347 conn := s.broker.connFor(s.route.hostID, s.route.gen)
348 if conn == nil {
349 return browser.ActResult{}, browser.ErrNoGrant
350 }
351 newRelay := s.broker.newRelay
352 if newRelay == nil {
353 newRelay = func(c sftpConn) FileRelay { return sftpFileRelay{conn: c} }
354 }
355 relay, ok := newRelay(conn).(browserUploadRelay)
356 if !ok {
357 return browser.ActResult{}, fmt.Errorf("browser upload: relay cannot receive remote files")
358 }
359 scratch, err := owner.captureDir()
360 if err != nil {
361 return browser.ActResult{}, err
362 }
363 dir, err := os.MkdirTemp(scratch, "remote-upload-")
364 if err != nil {
365 return browser.ActResult{}, err
366 }
367 defer os.RemoveAll(dir)
368 files := make([]string, 0, len(req.Files))
369 for _, remote := range req.Files {
370 local, err := relay.Fetch(ctx, res.workspace, remote, dir)
371 if err != nil {
372 return browser.ActResult{}, err
373 }
374 files = append(files, local)
375 }
376 req.Files = files
377 }
378 if !s.current(ctx) {
379 return browser.ActResult{}, browser.ErrNoGrant
380 }
381 return res.exec.Act(ctx, req)
382 }
383
384 func (s *brokerSessionExecutor) Close(ctx context.Context, req browser.CloseRequest) error {
385 res, err := s.resolve(ctx)
386 if err != nil {
387 return err
388 }
389 return res.exec.Close(ctx, req)
390 }
391
392 // relay stages one desktop capture file onto the remote host through the
393 // connection generation's SFTP channel.
394 func (s *brokerSessionExecutor) relay(ctx context.Context, workspace, localPath string) (string, error) {
395 if strings.TrimSpace(localPath) == "" {
396 return "", nil
397 }
398 if s.broker.connFor == nil {
399 return "", fmt.Errorf("browser broker: no file relay for this connection")
400 }
401 conn := s.broker.connFor(s.route.hostID, s.route.gen)
402 if conn == nil {
403 return "", fmt.Errorf("browser broker: host %q connection is gone", s.route.hostID)
404 }
405 newRelay := s.broker.newRelay
406 if newRelay == nil {
407 newRelay = func(c sftpConn) FileRelay { return sftpFileRelay{conn: c} }
408 }
409 return newRelay(conn).Stage(ctx, workspace, localPath)
410 }
411
412 // browserBrokerPort returns the broker's loopback port, starting the listener
413 // on first use. The broker serves every remote host off one port; tokens keep
414 // the hosts apart.
415 func (a *App) browserBrokerPort() (int, error) {
416 a.browserBrokerMu.Lock()
417 defer a.browserBrokerMu.Unlock()
418 if a.browserBroker != nil {
419 return a.browserBroker.port, nil
420 }
421 if !a.hostMode() {
422 return 0, fmt.Errorf("browser broker: the desktop shell is not attached")
423 }
424 ln, err := net.Listen("tcp", "127.0.0.1:0")
425 if err != nil {
426 return 0, fmt.Errorf("browser broker: listen: %w", err)
427 }
428 b := newBrowserBroker(a.resolveRemoteBrowserSession, a.remoteHostGenerationCurrent, a.remoteHostGenerationClient)
429 b.onRevoke = a.revokeRemoteBrowserHost
430 b.ln = ln
431 b.port = ln.Addr().(*net.TCPAddr).Port
432 b.server = &http.Server{
433 Handler: b,
434 ReadHeaderTimeout: 10 * time.Second,
435 IdleTimeout: 2 * time.Minute,
436 MaxHeaderBytes: 1 << 20,
437 }
438 a.browserBroker = b
439 a.goSafe("browserBroker", func() { _ = b.server.Serve(ln) })
440 return b.port, nil
441 }
442
443 // registerBrowserBrokerRoute starts the broker if needed and mints the
444 // (host, generation) token a remote serve bootstrap hands over.
445 func (a *App) registerBrowserBrokerRoute(hostID string, gen *managedHost) (string, int, error) {
446 if _, err := a.browserBrokerPort(); err != nil {
447 return "", 0, err
448 }
449 a.browserBrokerMu.Lock()
450 b := a.browserBroker
451 a.browserBrokerMu.Unlock()
452 if b == nil {
453 return "", 0, fmt.Errorf("browser broker: not running")
454 }
455 return b.register(hostID, gen)
456 }
457
458 func (a *App) revokeBrowserBrokerRoutes(hostID string) {
459 a.browserBrokerMu.Lock()
460 b := a.browserBroker
461 a.browserBrokerMu.Unlock()
462 if b != nil {
463 b.revokeHost(hostID)
464 }
465 }
466
467 func (a *App) closeBrowserBroker() {
468 a.browserBrokerMu.Lock()
469 b := a.browserBroker
470 a.browserBroker = nil
471 a.browserBrokerMu.Unlock()
472 if b != nil {
473 b.close()
474 }
475 }
476
477 // closeRemoteBrokers tears down the loopback brokers that serve remote hosts.
478 func (a *App) closeRemoteBrokers() {
479 a.closeCredentialProxy()
480 a.closeBrowserBroker()
481 }
482
483 // remoteHostGenerationCurrent fences broker routes to their connection
484 // generation: once the manager swaps or drops the host, minted tokens die.
485 func (a *App) remoteHostGenerationCurrent(hostID string, gen *managedHost) bool {
486 a.remoteMu.Lock()
487 rt := a.remoteRuntime
488 a.remoteMu.Unlock()
489 m, ok := rt.(*desktopRemoteManager)
490 if !ok || m == nil {
491 return false
492 }
493 return m.isCurrent(hostID, gen)
494 }
495
496 func (a *App) remoteHostGenerationClient(hostID string, gen *managedHost) sftpConn {
497 a.remoteMu.Lock()
498 rt := a.remoteRuntime
499 a.remoteMu.Unlock()
500 m, ok := rt.(*desktopRemoteManager)
501 if !ok || m == nil {
502 return nil
503 }
504 m.mu.Lock()
505 defer m.mu.Unlock()
506 if m.hosts[hostID] != gen {
507 return nil
508 }
509 return gen.client
510 }
511
512 // resolveRemoteBrowserSession maps a remote serve's session path to the
513 // desktop executor of the remote tab displaying it. Any session this desktop
514 // does not show for that host is refused with browser.ErrNoGrant, so a token
515 // can never reach a foreign session's tabs.
516 func (a *App) resolveRemoteBrowserSession(hostID, sessionPath string) (browserSessionResolution, error) {
517 sessionPath = strings.TrimSpace(sessionPath)
518 if sessionPath == "" || !a.hostMode() {
519 return browserSessionResolution{}, browser.ErrNoGrant
520 }
521 a.remoteTabMu.Lock()
522 defer a.remoteTabMu.Unlock()
523 for _, t := range a.remoteTabs {
524 if t == nil || t.ref.HostID != hostID {
525 continue
526 }
527 path := strings.TrimSpace(t.session.path)
528 if path != "" && path == sessionPath {
529 // Resolve identity and mint the immutable executor in the same lock
530 // epoch. sessionMu serializes transitions, not reads of these fields;
531 // acquiring it here would invert the resume path's lock order.
532 if exec := a.browserExecutorForRemoteTabLocked(t, sessionPath); exec != nil {
533 return browserSessionResolution{exec: exec, workspace: t.ref.Workspace}, nil
534 }
535 }
536 }
537 return browserSessionResolution{}, fmt.Errorf("%w: no stable desktop tab serves session %s", browser.ErrNoGrant, sessionPath)
538 }
539
540 // browserExecutorForRemoteTab returns the cached executor for one remote
541 // tab's browser surface; a session rotation re-scopes the grant.
542 func (a *App) browserExecutorForRemoteTab(tab *remoteTab, sessionPath string) browser.Executor {
543 a.remoteTabMu.Lock()
544 defer a.remoteTabMu.Unlock()
545 return a.browserExecutorForRemoteTabLocked(tab, sessionPath)
546 }
547
548 // Caller holds remoteTabMu. A provisional foreground route is not evidence
549 // that Serve has transferred this session's browser ownership.
550 func (a *App) browserExecutorForRemoteTabLocked(tab *remoteTab, sessionPath string) browser.Executor {
551 if tab == nil || !a.hostMode() || a.browserControl.off() {
552 return nil
553 }
554 if a.remoteTabs[tab.id] != tab || tab.session.path != sessionPath || tab.routing.rehydratingPath != "" {
555 return nil
556 }
557 canonicalID := tab.session.sessionID
558 if canonicalID != "" {
559 if tab.routing.currentPath != remoteSessionIDRoutePrefix+canonicalID {
560 return nil
561 }
562 } else if tab.routing.currentPath != "" && tab.routing.currentPath != sessionPath {
563 // Legacy peers may only have a path; never infer an ID from UI intent.
564 return nil
565 }
566 diagnosticScope := browserDiagnosticScope(tab.ref.HostID, canonicalID)
567 key := "remote/" + tab.id
568 a.browserExecMu.Lock()
569 defer a.browserExecMu.Unlock()
570 if a.browserExecutors == nil {
571 a.browserExecutors = map[string]*hostBrowserExecutor{}
572 }
573 if exec, ok := a.browserExecutors[key]; ok {
574 if exec.sessionKey == sessionPath && exec.diagnosticScope == diagnosticScope {
575 return exec
576 }
577 // A session rotation creates a new immutable owner; mutating the old
578 // executor races in-flight calls and lets them inherit the new grant.
579 a.revokeBrowserExecutor(exec)
580 }
581 exec := &hostBrowserExecutor{
582 app: a, host: a.hostShell.server, tabID: tab.id,
583 grantID: newBrowserGrantID(), sessionKey: sessionPath,
584 diagnosticScope: diagnosticScope,
585 }
586 a.browserExecutors[key] = exec
587 return exec
588 }
589
590 // ensureBrowserBrokerForward opens (idempotently) the reverse tunnel from the
591 // remote loopback to the desktop broker, mirroring the credential proxy's
592 // forward. Returns the actually bound remote port.
593 func ensureBrowserBrokerForward(c desktopSSHClient, hostID string, desktopPort int) (int, error) {
594 name := browserBrokerForwardName + hostID
595 target := fmt.Sprintf("127.0.0.1:%d", desktopPort)
596 for _, f := range c.Forwards().List() {
597 if f.Spec.Name == name && f.Spec.TargetAddr == target && f.Up {
598 if port, ok := portOfAddr(f.BoundAddr); ok {
599 return port, nil
600 }
601 }
602 }
603 bound, err := c.Forwards().Replace(forward.Spec{
604 Name: name,
605 Direction: forward.Remote,
606 BindAddr: "127.0.0.1:0",
607 TargetAddr: target,
608 })
609 if err != nil {
610 return 0, err
611 }
612 port, ok := portOfAddr(bound)
613 if !ok {
614 return 0, fmt.Errorf("browser broker: reverse tunnel bound unexpected address %q", bound)
615 }
616 return port, nil
617 }
618
618 lines GO