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/acp/client.go
T
tao.chen eb2bf3e684 fix(acp): parse session/update content as ContentBlock, not string
The ACP spec defines session/update content as a ContentBlock object
({type, text}) but we decoded it as a string. This dropped every
agent_message_chunk and left prompt responses empty. E2E against the
real opencode acp binary surfaced the bug.

- internal/acp/messages.go: add ContentBlock struct; Update.Content is
  now ContentBlock.
- internal/acp/client.go: handleNotification extracts Text when
  Content.Type == "text"; other types are logged at debug and dropped.
- internal/acp/client_test.go: NDJSON framing test sends content as a
  ContentBlock object and reads .Text.
- README.md: document /acp/* routes + CODESPACE_OPENCODE_ARGS env.

Verified end-to-end: POST /acp/prompt returns text="Hi there!", 4
messages in history (1 user + 3 agent chunks), 0 invalid session/update
log entries.

Conversation: 019f357d-4236-7010-949c-e61067618d42
2026-07-06 11:52:22 +08:00

278 lines
7.4 KiB
Go

package acp
import (
"context"
"encoding/json"
"log/slog"
"strings"
"sync"
"codespace/internal/fs"
"codespace/internal/process"
"codespace/internal/util"
)
// PromptResult is the outcome of a single prompt turn.
type PromptResult struct {
StopReason string
Text string
}
// Client manages one ACP session for a single workspace process.
type Client struct {
workspaceID string
root string
processes process.Manager
fs fs.FileSystem
lg *slog.Logger
mu sync.RWMutex
ready bool
sessionID string
transport *transport
history *History
}
// NewClient creates an ACP client for the given workspace.
func NewClient(workspaceID, root string, processes process.Manager, filesystem fs.FileSystem, lg *slog.Logger) *Client {
if lg == nil {
lg = slog.Default()
}
return &Client{
workspaceID: workspaceID,
root: root,
processes: processes,
fs: filesystem,
lg: lg,
history: NewHistory(),
}
}
// Initialize starts the process (if necessary), performs the ACP handshake,
// and creates a session. It is safe to call repeatedly; the client will reset
// and re-initialize if the process is not running.
func (c *Client) Initialize(ctx context.Context) error {
c.mu.RLock()
already := c.ready && c.sessionID != "" && c.processes.Status(c.workspaceID).Running
c.mu.RUnlock()
if already {
return nil
}
c.mu.Lock()
c.ready = false
c.sessionID = ""
c.history.Clear()
old := c.transport
c.transport = nil
c.mu.Unlock()
if old != nil {
old.onRequest = nil
old.onNotification = nil
old.onClose = nil
_ = old.Close()
}
if !c.processes.Status(c.workspaceID).Running {
if err := c.processes.Start(c.workspaceID, c.root); err != nil {
return util.Wrap(util.CodeInternal, "failed to start opencode", err)
}
}
stdin, err := c.processes.Stdin(c.workspaceID)
if err != nil {
return util.Wrap(util.CodeInternal, "failed to get opencode stdin", err)
}
sub, err := c.processes.Subscribe(c.workspaceID)
if err != nil {
return util.Wrap(util.CodeInternal, "failed to subscribe to opencode output", err)
}
tr := newTransport(stdin, sub, c.lg)
tr.onRequest = c.handleRequest
tr.onNotification = c.handleNotification
tr.onClose = c.onTransportClose
tr.start()
initParams := InitializeParams{
ProtocolVersion: 1,
ClientCapabilities: ClientCapabilities{
FS: ClientFSCapabilities{ReadTextFile: true, WriteTextFile: true},
},
ClientInfo: ClientInfo{Name: "codespace", Version: "0.1.0"},
}
var initRes InitializeResult
if err := tr.call(ctx, "initialize", initParams, &initRes); err != nil {
_ = tr.Close()
return util.Wrap(util.CodeInternal, "initialize failed", err)
}
sessParams := SessionNewParams{Cwd: c.root, MCPServers: []any{}}
var sessRes SessionNewResult
if err := tr.call(ctx, "session/new", sessParams, &sessRes); err != nil {
_ = tr.Close()
return util.Wrap(util.CodeInternal, "session/new failed", err)
}
c.mu.Lock()
c.transport = tr
c.sessionID = sessRes.SessionID
c.ready = true
c.mu.Unlock()
c.lg.Info("acp initialized", "workspace_id", c.workspaceID, "session_id", sessRes.SessionID)
return nil
}
// Prompt sends a user message and returns the agent's final text for this turn.
func (c *Client) Prompt(ctx context.Context, content string) (PromptResult, error) {
var result PromptResult
c.mu.RLock()
if !c.ready {
c.mu.RUnlock()
return result, util.New(util.CodeConflict, "client not ready")
}
sessionID := c.sessionID
tr := c.transport
startIdx := c.history.Len()
c.mu.RUnlock()
c.history.Add("user", content)
params := SessionPromptParams{
SessionID: sessionID,
Prompt: []PromptMessage{{Type: "text", Text: content}},
}
if err := tr.call(ctx, "session/prompt", params, &result); err != nil {
return result, util.Wrap(util.CodeInternal, "session/prompt failed", err)
}
msgs := c.history.Since(startIdx)
var sb strings.Builder
for _, m := range msgs {
if m.Role == "agent" {
sb.WriteString(m.Text)
}
}
result.Text = sb.String()
return result, nil
}
// Cancel sends a best-effort session/cancel notification.
func (c *Client) Cancel() error {
c.mu.RLock()
defer c.mu.RUnlock()
if !c.ready || c.transport == nil {
return nil
}
return c.transport.sendNotification("session/cancel", SessionCancelParams{SessionID: c.sessionID})
}
// Status returns whether the client has a ready session and its id.
func (c *Client) Status() (ready bool, sessionID string) {
c.mu.RLock()
defer c.mu.RUnlock()
return c.ready, c.sessionID
}
// History returns the in-memory conversation history.
func (c *Client) History() []Message {
return c.history.List()
}
// handleRequest dispatches agent-initiated JSON-RPC requests.
func (c *Client) handleRequest(method string, params json.RawMessage, id int) {
var result any
var rpcErr *rpcError
switch method {
case "fs/read_text_file":
var p FSReadParams
if err := json.Unmarshal(params, &p); err != nil {
rpcErr = &rpcError{Code: -32700, Message: "Parse error"}
break
}
data, err := c.fs.Read(p.Path)
if err != nil {
rpcErr = c.rpcErrorFromErr(err)
} else {
result = FSReadResult{Content: string(data)}
}
case "fs/write_text_file":
var p FSWriteParams
if err := json.Unmarshal(params, &p); err != nil {
rpcErr = &rpcError{Code: -32700, Message: "Parse error"}
break
}
if err := c.fs.Write(p.Path, []byte(p.Content)); err != nil {
rpcErr = c.rpcErrorFromErr(err)
} else {
result = FSWriteResult{}
}
default:
rpcErr = &rpcError{Code: -32601, Message: "Method not found: " + method}
}
c.mu.RLock()
tr := c.transport
c.mu.RUnlock()
if tr != nil {
if err := tr.sendResponse(id, result, rpcErr); err != nil {
c.lg.Error("acp failed to send response", "workspace_id", c.workspaceID, "method", method, "error", err)
}
}
}
// handleNotification processes agent-initiated notifications.
func (c *Client) handleNotification(method string, params json.RawMessage) {
if method != "session/update" {
c.lg.Debug("acp drop notification", "workspace_id", c.workspaceID, "method", method)
return
}
var up SessionUpdateParams
if err := json.Unmarshal(params, &up); err != nil {
c.lg.Error("acp invalid session/update", "workspace_id", c.workspaceID, "error", err)
return
}
switch up.Update.SessionUpdate {
case "agent_message_chunk":
if up.Update.Content.Type == "text" {
c.history.Add("agent", up.Update.Content.Text)
} else {
c.lg.Debug("acp drop non-text agent chunk", "workspace_id", c.workspaceID, "type", up.Update.Content.Type)
}
case "user_message_chunk":
if up.Update.Content.Type == "text" {
c.history.Add("user", up.Update.Content.Text)
} else {
c.lg.Debug("acp drop non-text user chunk", "workspace_id", c.workspaceID, "type", up.Update.Content.Type)
}
default:
c.lg.Debug("acp drop update", "workspace_id", c.workspaceID, "type", up.Update.SessionUpdate)
}
}
// onTransportClose is invoked when the process output stream closes.
func (c *Client) onTransportClose() {
c.mu.Lock()
defer c.mu.Unlock()
if c.ready {
c.ready = false
c.sessionID = ""
c.transport = nil
c.history.Clear()
c.lg.Info("acp transport closed", "workspace_id", c.workspaceID)
}
}
// rpcErrorFromErr maps a filesystem error to a JSON-RPC error.
func (c *Client) rpcErrorFromErr(err error) *rpcError {
switch util.CodeOf(err) {
case util.CodeNotFound, util.CodeBadRequest:
return &rpcError{Code: -32602, Message: err.Error()}
default:
return &rpcError{Code: -32603, Message: err.Error()}
}
}