package api import ( "net/http/httptest" "os" "path/filepath" "strings" "sync" "testing" "time" "codespace/internal/process" "codespace/internal/service" "codespace/internal/shell" "codespace/internal/workspace" "github.com/gin-gonic/gin" "github.com/gorilla/websocket" ) func TestProcessWS(t *testing.T) { tmpDir := t.TempDir() opencodePath := filepath.Join(tmpDir, "opencode") script := []byte("#!/bin/sh\nread line\necho \"got: $line\"\nsleep 0.1\nexit 0\n") if err := os.WriteFile(opencodePath, script, 0o755); err != nil { t.Fatalf("write fake opencode: %v", err) } wsRoot := filepath.Join(tmpDir, "workspaces") wsMgr := workspace.NewLocalManager(wsRoot) procMgr := process.NewManager(opencodePath, nil) shellMgr := shell.NewManager("bash", []string{"-i"}) wsSvc := service.NewWorkspaceService(wsMgr, procMgr, shellMgr, nil) fileSvc := service.NewFileService(wsMgr, 1<<20) procSvc := service.NewProcessService(wsMgr, procMgr, nil) shellSvc := service.NewShellService(wsMgr, shellMgr, nil) acpSvc := service.NewAcpService(procMgr, wsMgr, nil) r := NewRouter(wsSvc, fileSvc, procSvc, shellSvc, acpSvc, nil, gin.TestMode) srv := httptest.NewServer(r) defer srv.Close() ws, err := wsSvc.Create("wstest") if err != nil { t.Fatalf("create workspace: %v", err) } t.Cleanup(func() { if err := wsSvc.Delete(ws.ID); err != nil { t.Logf("cleanup delete workspace: %v", err) } }) if err := procSvc.Start(ws.ID); err != nil { t.Fatalf("start process: %v", err) } wsURL := "ws" + strings.TrimPrefix(srv.URL, "http") + "/api/workspaces/" + ws.ID + "/process/ws" conn, _, err := websocket.DefaultDialer.Dial(wsURL, nil) if err != nil { t.Fatalf("dial websocket: %v", err) } defer conn.Close() if err := conn.WriteMessage(websocket.TextMessage, []byte("hello\n")); err != nil { t.Fatalf("write message: %v", err) } if err := conn.SetReadDeadline(time.Now().Add(3 * time.Second)); err != nil { t.Fatalf("set read deadline: %v", err) } var gotEcho, gotExit bool for { _, data, err := conn.ReadMessage() if err != nil { break } text := string(data) if strings.Contains(text, "got: hello") { gotEcho = true } if strings.Contains(text, "[process exited") { gotExit = true break } } if !gotEcho { t.Fatal("did not receive echo frame containing \"got: hello\"") } if !gotExit { t.Fatal("did not receive exit banner frame") } } func TestProcessWSMultiSubscriber(t *testing.T) { tmpDir := t.TempDir() opencodePath := filepath.Join(tmpDir, "opencode") script := []byte("#!/bin/sh\nread line\nfor i in 1 2 3; do\n echo \"got: $line ($i)\"\n sleep 0.05\ndone\nexit 0\n") if err := os.WriteFile(opencodePath, script, 0o755); err != nil { t.Fatalf("write fake opencode: %v", err) } wsRoot := filepath.Join(tmpDir, "workspaces") wsMgr := workspace.NewLocalManager(wsRoot) procMgr := process.NewManager(opencodePath, nil) shellMgr := shell.NewManager("bash", []string{"-i"}) wsSvc := service.NewWorkspaceService(wsMgr, procMgr, shellMgr, nil) fileSvc := service.NewFileService(wsMgr, 1<<20) procSvc := service.NewProcessService(wsMgr, procMgr, nil) shellSvc := service.NewShellService(wsMgr, shellMgr, nil) acpSvc := service.NewAcpService(procMgr, wsMgr, nil) r := NewRouter(wsSvc, fileSvc, procSvc, shellSvc, acpSvc, nil, gin.TestMode) srv := httptest.NewServer(r) defer srv.Close() ws, err := wsSvc.Create("wsmulti") if err != nil { t.Fatalf("create workspace: %v", err) } t.Cleanup(func() { if err := wsSvc.Delete(ws.ID); err != nil { t.Logf("cleanup delete workspace: %v", err) } }) if err := procSvc.Start(ws.ID); err != nil { t.Fatalf("start process: %v", err) } wsURL := "ws" + strings.TrimPrefix(srv.URL, "http") + "/api/workspaces/" + ws.ID + "/process/ws" ready := make(chan struct{}, 3) start := make(chan struct{}) frames := make([][]string, 3) errs := make([]error, 3) var wg sync.WaitGroup for i := 0; i < 3; i++ { wg.Add(1) go func(idx int) { defer wg.Done() conn, _, err := websocket.DefaultDialer.Dial(wsURL, nil) if err != nil { errs[idx] = err return } defer conn.Close() ready <- struct{}{} <-start if err := conn.WriteMessage(websocket.TextMessage, []byte("hello\n")); err != nil { errs[idx] = err return } if err := conn.SetReadDeadline(time.Now().Add(5 * time.Second)); err != nil { errs[idx] = err return } for { _, data, err := conn.ReadMessage() if err != nil { errs[idx] = err break } text := string(data) frames[idx] = append(frames[idx], text) if strings.Contains(text, "[process exited") { break } } }(i) } for i := 0; i < 3; i++ { <-ready } close(start) wg.Wait() for i, ff := range frames { var got1, got2, got3, gotExit bool var got strings.Builder for _, f := range ff { got.WriteString(f) if strings.Contains(f, "got: hello (1)") { got1 = true } if strings.Contains(f, "got: hello (2)") { got2 = true } if strings.Contains(f, "got: hello (3)") { got3 = true } if strings.Contains(f, "[process exited") { gotExit = true } } if !got1 || !got2 || !got3 || !gotExit || errs[i] != nil { t.Fatalf("subscriber #%d failed: err=%v, frames=%q", i+1, errs[i], got.String()) } } }