208 lines
5.1 KiB
Go
208 lines
5.1 KiB
Go
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)
|
|
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)
|
|
|
|
r := NewRouter(wsSvc, fileSvc, procSvc, shellSvc, 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)
|
|
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)
|
|
|
|
r := NewRouter(wsSvc, fileSvc, procSvc, shellSvc, 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())
|
|
}
|
|
}
|
|
}
|