This repository has been archived on 2026-07-17. You can view files and clone it. You cannot open issues or pull requests or push a commit.
Files
codespace/internal/service/acp_service.go
T
tao.chen 031451fe3e feat(acp): streaming WebSocket route for prompt chunks
Adds GET /api/workspaces/:id/acp/stream. Client opens a WebSocket,
sends {"type":"prompt","content":"..."}, and receives a stream of
{"type":"chunk","messageId","text"} events followed by exactly
one {"type":"complete","stopReason"} or {"type":"error","error"}.
Closing the WS early triggers session/cancel.

- internal/acp/messages.go: StreamEvent wire shape.
- internal/acp/client.go:
  - streamChs []chan StreamEvent set; AddStream / RemoveStream.
  - sendStream non-blocking fanout.
  - Client.Stream(ctx, content, out) registers out, sends prompt,
    emits complete/error after the prompt response, unregisters.
  - handleNotification fans chunk events to all stream consumers.
  - notifyWG ensures chunk ordering vs the terminal event.
- internal/acp/service.go: Service.Stream(workspaceID, content, out)
  mirrors Prompt (per-workspace lock, 5-min timeout, EnsureReady).
- internal/service/acp_service.go: thin AcpService.Stream wrapper
  that maps acp.StreamEvent -> model.AcpStreamEvent.
- internal/model/acp.go: AcpStreamRequest, AcpStreamEvent DTOs.
- internal/api/acp_handler.go: stream WS handler (upgrade, read
  prompt, run Stream in a goroutine, write events, ping/pong, Cancel
  on client close).
- internal/api/router.go: register the new route.
- internal/acp/transport.go: dispatch notifications synchronously
  (vs. goroutine per notification) so chunks preserve order before
  the session/prompt response.
- internal/acp/client_test.go: TestClientStreamEmitsChunkAndComplete
  with a fake transport that drives a known sequence.
- internal/api/acp_handler_test.go: TestAcpStreamHandlerRoutes
  smoke test using a fake opencode acp script.

Existing POST /api/workspaces/:id/acp/prompt is unchanged.

E2E: prompt 'say hi in exactly 3 words' -> 3 chunk events
('Hi',' there','!') + 1 complete {stopReason: 'end_turn'}.

Conversation: 019f3680-200f-79b0-860b-43302e60d0ea
2026-07-06 16:43:17 +08:00

118 lines
3.4 KiB
Go

package service
import (
"log/slog"
"codespace/internal/acp"
"codespace/internal/model"
"codespace/internal/process"
"codespace/internal/workspace"
)
// AcpService wraps the low-level acp.Service with logging and model mapping.
type AcpService struct {
svc acp.Service
lg *slog.Logger
}
// NewAcpService creates an AcpService.
// If lg is nil, slog.Default() is used.
func NewAcpService(processes process.Manager, workspaces workspace.Manager, lg *slog.Logger) *AcpService {
if lg == nil {
lg = slog.Default()
}
return &AcpService{svc: acp.NewService(processes, workspaces, lg), lg: lg}
}
// Status returns the ACP status for the workspace.
func (s *AcpService) Status(workspaceID string) (model.AcpStatusResponse, error) {
status, err := s.svc.Status(workspaceID)
if err != nil {
s.lg.Error("acp status failed", "workspace_id", workspaceID, "error", err)
return model.AcpStatusResponse{}, err
}
return model.AcpStatusResponse{
WorkspaceID: status.WorkspaceID,
Ready: status.Ready,
SessionID: status.SessionID,
Running: status.Running,
PID: status.PID,
Error: status.Error,
}, nil
}
// History returns the ACP conversation history for the workspace.
func (s *AcpService) History(workspaceID string) (model.AcpHistoryResponse, error) {
hist, err := s.svc.History(workspaceID)
if err != nil {
s.lg.Error("acp history failed", "workspace_id", workspaceID, "error", err)
return model.AcpHistoryResponse{}, err
}
messages := make([]model.AcpMessage, len(hist.Messages))
for i, m := range hist.Messages {
messages[i] = model.AcpMessage{Role: m.Role, Text: m.Text, Time: m.Time, MessageID: m.MessageID}
}
return model.AcpHistoryResponse{SessionID: hist.SessionID, Messages: messages}, nil
}
// Prompt sends a prompt to the agent and returns its response.
func (s *AcpService) Prompt(workspaceID string, content string) (model.AcpPromptResponse, error) {
res, err := s.svc.Prompt(workspaceID, content)
if err != nil {
s.lg.Error("acp prompt failed", "workspace_id", workspaceID, "error", err)
return model.AcpPromptResponse{}, err
}
return model.AcpPromptResponse{
SessionID: res.SessionID,
StopReason: res.StopReason,
Text: res.Text,
}, nil
}
// Cancel interrupts an in-flight prompt.
func (s *AcpService) Cancel(workspaceID string) error {
if err := s.svc.Cancel(workspaceID); err != nil {
s.lg.Error("acp cancel failed", "workspace_id", workspaceID, "error", err)
return err
}
return nil
}
// Stream sends a prompt to the agent and forwards streamed events to out.
func (s *AcpService) Stream(workspaceID string, content string, out chan<- model.AcpStreamEvent) error {
acpOut := make(chan acp.StreamEvent, 32)
var streamErr error
done := make(chan struct{})
go func() {
defer close(done)
streamErr = s.svc.Stream(workspaceID, content, acpOut)
close(acpOut)
}()
for ev := range acpOut {
out <- model.AcpStreamEvent{
Type: ev.Type,
MessageID: ev.MessageID,
Text: ev.Text,
StopReason: ev.StopReason,
Error: ev.Error,
}
}
<-done
if streamErr != nil {
s.lg.Error("acp stream failed", "workspace_id", workspaceID, "error", streamErr)
}
return streamErr
}
// EnsureReady ensures the ACP process is started and initialized.
func (s *AcpService) EnsureReady(workspaceID string) error {
if err := s.svc.EnsureReady(workspaceID); err != nil {
s.lg.Error("acp ensure ready failed", "workspace_id", workspaceID, "error", err)
return err
}
return nil
}