返回 DeepSeek-Reasonix
remote_serve_fake_test.go
根目录 / desktop / remote_serve_fake_test.go
1 package main
2
3 import (
4 "encoding/json"
5 "fmt"
6 "io"
7 "net/http"
8 "net/http/httptest"
9 "sync"
10 "testing"
11 "time"
12
13 "reasonix/internal/config"
14 )
15
16 // fakeServe is a minimal Serve stand-in for bridge tests: token handshake,
17 // session enter, an SSE feed that emits two frames then holds, and recorded
18 // command endpoints for the proxy bindings.
19 type fakeServe struct {
20 t *testing.T
21 token string
22 server *httptest.Server
23
24 mu sync.Mutex
25 newCalled int
26 newSessionPath string
27 resumePath string
28 resumeSessionID string
29 cookieOnNew bool
30 sessions []serveSessionEntry
31 calls []string // "METHOD /path body" per command request
32 expectedPaths []string // foreground command fence headers
33 failNext string // non-empty ⇒ next command endpoint replies 409 with this text
34 failEnter string // non-empty ⇒ next /new or /resume replies 409
35 enterDelay time.Duration
36 newStarted, newRelease chan struct{}
37 resumeStarted chan string
38 resumeRelease chan struct{}
39 failHistory bool // /history replies 500 when set
40 historyBody string
41 historyStarted, historyRelease chan struct{}
42 failSessions bool // /sessions replies 500 when set
43 sessionsFailCount int // /sessions replies 500 this many times, then recovers
44 resumeDropCount int // /resume commits the switch but drops the connection unanswered this many times
45 sessionsStarted chan struct{}
46 sessionsRelease chan struct{}
47 eventsConns int // /events connections opened
48 eventsQuery string
49 eventFrames []string
50 eventFeed <-chan string
51 eventsStatus int // non-zero makes /events fail before opening
52 eventsFailCount int // refuse this many /events opens with 503, then serve normally
53 eventsCloseEarly bool // return immediately after the initial 200 frames
54 statusPayload string
55 statusAfterCancel string
56 contextExpected string
57 contextStarted, contextRelease chan struct{}
58 }
59
60 func (fs *fakeServe) eventsCount() int { fs.mu.Lock(); defer fs.mu.Unlock(); return fs.eventsConns }
61
62 func (fs *fakeServe) recorded() []string {
63 fs.mu.Lock()
64 defer fs.mu.Unlock()
65 out := make([]string, len(fs.calls))
66 copy(out, fs.calls)
67 return out
68 }
69
70 func (fs *fakeServe) recordedExpectedPaths() []string {
71 fs.mu.Lock()
72 defer fs.mu.Unlock()
73 return append([]string(nil), fs.expectedPaths...)
74 }
75
76 func (fs *fakeServe) record(method, path, body string) {
77 fs.mu.Lock()
78 fs.calls = append(fs.calls, method+" "+path+" "+body)
79 fs.mu.Unlock()
80 }
81
82 // newFakeServe builds a stand-in for the workspace Serve's HTTP surface. The
83 // mux is wrapped in a token-mode gate that mirrors the real authGate's
84 // contract: POST /auth/token is matched on the EXACT path (a "//auth/token"
85 // double slash — what naive base+path joins produce from EnsureServer's
86 // trailing-slash LocalURL — is denied with 401 before routing), and every
87 // other path requires the session cookie the bootstrap installs. A bare mux
88 // cannot catch this: it 301-redirects unclean paths and Go's client follows
89 // preserving POST, so a double-slash request would silently succeed here
90 // while the real Serve rejects it.
91 func newFakeServe(t *testing.T, token string, sessions []serveSessionEntry) *fakeServe {
92 t.Helper()
93 fs := &fakeServe{t: t, token: token, sessions: sessions}
94 mux := http.NewServeMux()
95 mux.HandleFunc("POST /auth/token", func(w http.ResponseWriter, r *http.Request) {
96 var body struct {
97 Token string `json:"token"`
98 }
99 if err := json.NewDecoder(r.Body).Decode(&body); err != nil || body.Token != fs.token {
100 http.Error(w, "denied", http.StatusUnauthorized)
101 return
102 }
103 http.SetCookie(w, &http.Cookie{Name: "reasonix_token", Value: fs.token, Path: "/", HttpOnly: true})
104 w.Header().Set(serveCapabilitiesHeader, "permission-presets-v1,present-files-v1,execution-v2,session-history-v1,session-identity-v1,session-ownership-v1")
105 w.WriteHeader(http.StatusNoContent)
106 })
107 mux.HandleFunc("POST /new", func(w http.ResponseWriter, r *http.Request) {
108 fs.mu.Lock()
109 fs.newCalled++
110 _, cookieErr := r.Cookie("reasonix_token")
111 fs.cookieOnNew = cookieErr == nil
112 fail := fs.failEnter
113 fs.failEnter = ""
114 enterDelay := fs.enterDelay
115 newSessionPath := fs.newSessionPath
116 newStarted, newRelease := fs.newStarted, fs.newRelease
117 if fail != "" {
118 fs.mu.Unlock()
119 http.Error(w, fail, http.StatusConflict)
120 return
121 }
122 // The serve abandons the current session on /new: no file, not listed.
123 for i := range fs.sessions {
124 fs.sessions[i].Current = false
125 }
126 fs.mu.Unlock()
127 if newStarted != nil {
128 select {
129 case newStarted <- struct{}{}:
130 default:
131 }
132 }
133 if newRelease != nil {
134 select {
135 case <-newRelease:
136 case <-r.Context().Done():
137 return
138 }
139 }
140 if enterDelay > 0 {
141 time.Sleep(enterDelay)
142 }
143 if newSessionPath != "" {
144 w.Header().Set("X-Reasonix-Session-Path", newSessionPath)
145 }
146 w.WriteHeader(http.StatusNoContent)
147 })
148 mux.HandleFunc("POST /resume", func(w http.ResponseWriter, r *http.Request) {
149 var body struct {
150 Path string `json:"path"`
151 SessionID string `json:"sessionId"`
152 }
153 // Mirrors the real handler: a canonical row carries only sessionId, so
154 // requiring a path here would hide every identity-route regression.
155 if err := json.NewDecoder(r.Body).Decode(&body); err != nil || body.Path == "" && body.SessionID == "" {
156 http.Error(w, "missing path or sessionId", http.StatusBadRequest)
157 return
158 }
159 fs.mu.Lock()
160 fail := fs.failEnter
161 fs.failEnter = ""
162 enterDelay := fs.enterDelay
163 resumeStarted, resumeRelease := fs.resumeStarted, fs.resumeRelease
164 drop := fs.takeFault(&fs.resumeDropCount)
165 if fail != "" {
166 fs.mu.Unlock()
167 http.Error(w, fail, http.StatusConflict)
168 return
169 }
170 fs.resumePath, fs.resumeSessionID = body.Path, body.SessionID
171 for i := range fs.sessions {
172 fs.sessions[i].Current = body.SessionID != "" && fs.sessions[i].SessionID == body.SessionID ||
173 body.SessionID == "" && fs.sessions[i].Path == body.Path
174 }
175 fs.mu.Unlock()
176 if body.SessionID != "" {
177 w.Header().Set("X-Reasonix-Session-ID", body.SessionID)
178 }
179 if drop {
180 // Serve committed the switch; only the response is lost.
181 dropHTTPConnection(w)
182 return
183 }
184 if resumeStarted != nil {
185 select {
186 case resumeStarted <- body.Path:
187 default:
188 }
189 }
190 if resumeRelease != nil {
191 select {
192 case <-resumeRelease:
193 case <-r.Context().Done():
194 return
195 }
196 }
197 if enterDelay > 0 {
198 time.Sleep(enterDelay)
199 }
200 w.WriteHeader(http.StatusNoContent)
201 })
202 mux.HandleFunc("GET /sessions", func(w http.ResponseWriter, r *http.Request) {
203 fs.record(r.Method, "/sessions", "")
204 fs.mu.Lock()
205 fail := fs.failSessions || fs.takeFault(&fs.sessionsFailCount)
206 started, release := fs.sessionsStarted, fs.sessionsRelease
207 sessions := append([]serveSessionEntry(nil), fs.sessions...)
208 fs.mu.Unlock()
209 if started != nil {
210 select {
211 case started <- struct{}{}:
212 default:
213 }
214 }
215 if release != nil {
216 select {
217 case <-release:
218 case <-r.Context().Done():
219 return
220 }
221 }
222 if fail {
223 http.Error(w, "sessions unavailable", http.StatusInternalServerError)
224 return
225 }
226 writeTestJSON(w, sessions)
227 })
228 mux.HandleFunc("GET /events", func(w http.ResponseWriter, r *http.Request) {
229 fs.mu.Lock()
230 fs.eventsConns++
231 fs.eventsQuery = r.URL.RawQuery
232 eventsStatus := fs.eventsStatus
233 if fs.takeFault(&fs.eventsFailCount) && eventsStatus == 0 {
234 eventsStatus = http.StatusServiceUnavailable
235 }
236 closeEarly := fs.eventsCloseEarly
237 frames := append([]string(nil), fs.eventFrames...)
238 feed := fs.eventFeed
239 fs.mu.Unlock()
240 if eventsStatus != 0 {
241 http.Error(w, "event stream unavailable", eventsStatus)
242 return
243 }
244 w.Header().Set("Content-Type", "text/event-stream")
245 flusher, ok := w.(http.Flusher)
246 if !ok {
247 http.Error(w, "no flusher", http.StatusInternalServerError)
248 return
249 }
250 if len(frames) == 0 {
251 frames = []string{`{"kind":"session_start"}`, `{"kind":"ready"}`}
252 }
253 for _, frame := range frames {
254 fmt.Fprintf(w, "data: %s\n\n", frame)
255 }
256 flusher.Flush()
257 if closeEarly {
258 return
259 }
260 if feed == nil {
261 <-r.Context().Done()
262 return
263 }
264 for {
265 select {
266 case frame := <-feed:
267 fmt.Fprintf(w, "data: %s\n\n", frame)
268 flusher.Flush()
269 case <-r.Context().Done():
270 return
271 }
272 }
273 })
274 command := func(path string) http.HandlerFunc {
275 return func(w http.ResponseWriter, r *http.Request) {
276 data, _ := io.ReadAll(io.LimitReader(r.Body, 4<<10))
277 fs.record(r.Method, path, string(data))
278 fs.mu.Lock()
279 fs.expectedPaths = append(fs.expectedPaths, r.Header.Get(expectedSessionPathHeader))
280 fail := fs.failNext
281 fs.failNext = ""
282 if path == "/cancel" && fs.statusAfterCancel != "" {
283 fs.statusPayload = fs.statusAfterCancel
284 }
285 fs.mu.Unlock()
286 if fail != "" {
287 http.Error(w, fail, http.StatusConflict)
288 return
289 }
290 if path == "/composer-profile" {
291 w.Header().Set("Content-Type", "application/json")
292 _, _ = io.WriteString(w, `{"drainedApprovalIDs":["approval-1"]}`)
293 return
294 }
295 if path == "/permission/preset" {
296 w.Header().Set("Content-Type", "application/json")
297 _, _ = io.WriteString(w, `{"snapshot":{"sessionId":"remote-session","generation":1,"revision":8,"preset":"workspace-write","workspaceRoot":"/workspace","grants":[],"capabilities":{"backend":"seatbelt","enforcement":"full","supportedPresets":["read-only","workspace-write","danger-full-access"]}}}`)
298 return
299 }
300 w.WriteHeader(http.StatusNoContent)
301 }
302 }
303 for _, path := range []string{"/submit", "/cancel", "/approve", "/plan-decision", "/answer", "/extension-form", "/rewind", "/goal", "/goal/edit", "/goal/pause", "/goal/resume", "/jobs/cancel", "/inbox/items", "/permission/preset", "/composer-profile", "/delete-session", "/model", "/effort", "/quality-floor", "/plan", "/compact", "/fork", "/summarize", "/forget", "/clear"} {
304 mux.HandleFunc("POST "+path, command(path))
305 }
306 snapshot := func(path, payload string) {
307 mux.HandleFunc("GET "+path, func(w http.ResponseWriter, r *http.Request) {
308 fs.record(r.Method, path, "")
309 responsePayload := payload
310 if path == "/context" {
311 fs.mu.Lock()
312 fs.contextExpected = r.Header.Get(expectedSessionPathHeader)
313 if id := r.Header.Get(expectedSessionIDHeader); id != "" {
314 fs.contextExpected = remoteSessionIDRoutePrefix + id
315 }
316 started, release := fs.contextStarted, fs.contextRelease
317 fs.mu.Unlock()
318 if started != nil {
319 started <- struct{}{}
320 }
321 if release != nil {
322 select {
323 case <-release:
324 case <-r.Context().Done():
325 return
326 }
327 }
328 }
329 if path == "/history" {
330 fs.mu.Lock()
331 fail := fs.failHistory
332 if fs.historyBody != "" {
333 responsePayload = fs.historyBody
334 }
335 started, release := fs.historyStarted, fs.historyRelease
336 fs.mu.Unlock()
337 if started != nil {
338 started <- struct{}{}
339 }
340 if release != nil {
341 select {
342 case <-release:
343 case <-r.Context().Done():
344 return
345 }
346 }
347 if fail {
348 http.Error(w, "gone", http.StatusInternalServerError)
349 return
350 }
351 }
352 w.Header().Set("Content-Type", "application/json")
353 _, _ = w.Write([]byte(responsePayload))
354 })
355 }
356 snapshot("/history", `[{"role":"user","content":"hi"}]`)
357 snapshot("/context", `{"used":10,"window":128}`)
358 snapshot("/todos", `[]`)
359 snapshot("/checkpoints", `[{"turn":1}]`)
360 snapshot("/models", `{"current":"remote/chat","label":"chat","models":[{"ref":"remote/chat","provider":"remote","model":"chat","active":true}]}`)
361 snapshot("/commands", `[{"name":"remote-review","description":"Review remotely","kind":"custom","group":"skills"}]`)
362 snapshot("/pending-prompts", `[{"kind":"approval_request","approval":{"id":"approval-1","tool":"bash"}}]`)
363 snapshot("/permission", `{"sessionId":"remote-session","generation":1,"revision":7,"preset":"workspace-write","workspaceRoot":"/workspace","grants":[],"capabilities":{"backend":"seatbelt","enforcement":"full","supportedPresets":["read-only","workspace-write","danger-full-access"]}}`)
364 mux.HandleFunc("GET /status", func(w http.ResponseWriter, r *http.Request) {
365 fs.record(r.Method, "/status", "")
366 fs.mu.Lock()
367 payload := fs.statusPayload
368 fs.mu.Unlock()
369 if payload == "" {
370 payload = `{"state":"ready"}`
371 }
372 w.Header().Set("Content-Type", "application/json")
373 _, _ = w.Write([]byte(payload))
374 })
375 snapshot("/branches", `{"branches":[]}`)
376 snapshot("/skills", `[]`)
377 gate := http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
378 if r.URL.Path == "/auth/token" {
379 mux.ServeHTTP(w, r)
380 return
381 }
382 if c, err := r.Cookie("reasonix_token"); err == nil && c.Value == fs.token {
383 mux.ServeHTTP(w, r)
384 return
385 }
386 http.Error(w, "Unauthorized", http.StatusUnauthorized)
387 })
388 fs.server = httptest.NewServer(gate)
389 t.Cleanup(fs.server.Close)
390 return fs
391 }
392
393 func (fs *fakeServe) snapshot() (newCalled int, resumePath string, cookieOnNew bool) {
394 fs.mu.Lock()
395 defer fs.mu.Unlock()
396 return fs.newCalled, fs.resumePath, fs.cookieOnNew
397 }
398
399 func writeTestJSON(w http.ResponseWriter, v any) {
400 w.Header().Set("Content-Type", "application/json")
401 _ = json.NewEncoder(w).Encode(v)
402 }
403
404 func TestRemoteHostFixtureIsolatesCanonicalWorkspaceRegistry(t *testing.T) {
405 isolateDesktopUserDirs(t)
406 prior := NewApp()
407 if _, err := prior.ensureDesktopWorkspace(t.Context(), "global", ""); err != nil {
408 t.Fatal(err)
409 }
410 seedBridgeTestHost(t, "box")
411 fresh := NewApp()
412 if _, err := fresh.ensureDesktopWorkspace(t.Context(), "global", ""); err != nil {
413 t.Fatalf("remote fixture reused another home's canonical registry: %v", err)
414 }
415 }
416
417 func seedBridgeTestHost(t *testing.T, hostID string) {
418 t.Helper()
419 // A new config home also needs its own canonical registry: global workspace
420 // identity contains that home, while REASONIX_STATE_HOME otherwise survives.
421 home := isolateDesktopUserDirs(t)
422 t.Setenv("REASONIX_HOME", home)
423 if err := editUserConfig(func(c *config.Config) error {
424 return c.UpsertRemoteHost(config.RemoteHostEntry{Name: hostID, Host: "127.0.0.1", Port: 22, User: "dev"})
425 }); err != nil {
426 t.Fatal(err)
427 }
428 }
429
430 func waitForRemoteEventCount(t *testing.T, log *eventLog, prefix string, want int) {
431 t.Helper()
432 deadline := time.Now().Add(3 * time.Second)
433 for {
434 if got := log.count(prefix); got >= want {
435 return
436 }
437 if time.Now().After(deadline) {
438 t.Fatalf("event count for %q = %d, want >= %d (events: %v)", prefix, log.count(prefix), want, log.recorded())
439 }
440 time.Sleep(10 * time.Millisecond)
441 }
442 }
443
444 // cleanupRemoteTabPumps cancels every open tab's SSE pump and waits for all
445 // bridge tasks to return. Waiting matters on Windows: an async resume can
446 // publish its final tab snapshot after the assertion succeeds, racing the
447 // temporary user directory's cleanup.
448 func cleanupRemoteTabPumps(t *testing.T, a *App) {
449 t.Helper()
450 t.Cleanup(func() {
451 a.remoteTabMu.Lock()
452 for _, tab := range a.remoteTabs {
453 if tab.cancel != nil {
454 tab.cancel()
455 }
456 }
457 a.remoteTabMu.Unlock()
458 done := make(chan struct{})
459 go func() { a.remoteTabTasks.Wait(); close(done) }()
460 select {
461 case <-done:
462 case <-time.After(5 * time.Second):
463 t.Error("remote tab tasks did not stop after pump cancellation")
464 }
465 })
466 }
467
468 // openReadyRemoteTab opens a tab against the fake serve and waits for ready.
469 func openReadyRemoteTab(t *testing.T, a *App, opts RemoteTabOpenOptions) TabMeta {
470 t.Helper()
471 meta, err := a.OpenRemoteProjectTab("box", "~/app", opts)
472 if err != nil {
473 t.Fatal(err)
474 }
475 waitForTabState(t, a, meta.ID, "ready")
476 return meta
477 }
478
478 lines GO