sink_messages.go

 1package acp
 2
 3import (
 4	"log/slog"
 5
 6	"github.com/charmbracelet/crush/internal/message"
 7	"github.com/charmbracelet/crush/internal/pubsub"
 8	"github.com/coder/acp-go-sdk"
 9)
10
11// HandleMessage translates a Crush message event to ACP session updates.
12func (s *Sink) HandleMessage(event pubsub.Event[message.Message]) {
13	msg := event.Payload
14
15	// Only handle messages for our session.
16	if msg.SessionID != s.sessionID {
17		return
18	}
19
20	for _, part := range msg.Parts {
21		update := s.translatePart(msg.ID, msg.Role, part)
22		if update == nil {
23			continue
24		}
25
26		if err := s.conn.SessionUpdate(s.ctx, acp.SessionNotification{
27			SessionId: acp.SessionId(s.sessionID),
28			Update:    *update,
29		}); err != nil {
30			slog.Error("Failed to send session update", "error", err)
31		}
32	}
33}
34
35// translatePart converts a message part to an ACP session update.
36func (s *Sink) translatePart(msgID string, role message.MessageRole, part message.ContentPart) *acp.SessionUpdate {
37	switch p := part.(type) {
38	case message.TextContent:
39		return s.translateText(msgID, role, p)
40
41	case message.ReasoningContent:
42		return s.translateReasoning(msgID, p)
43
44	case message.ToolCall:
45		return s.translateToolCall(p)
46
47	case message.ToolResult:
48		return s.translateToolResult(p)
49
50	case message.Finish:
51		// Reset offsets on message finish.
52		delete(s.textOffsets, msgID)
53		delete(s.reasoningOffsets, msgID)
54		return nil
55
56	default:
57		return nil
58	}
59}
60
61func (s *Sink) translateText(msgID string, role message.MessageRole, text message.TextContent) *acp.SessionUpdate {
62	// Skip user messages - the client already knows what it sent via the
63	// prompt request.
64	if role != message.Assistant {
65		return nil
66	}
67
68	offset := s.textOffsets[msgID]
69	if len(text.Text) <= offset {
70		return nil
71	}
72
73	delta := text.Text[offset:]
74	s.textOffsets[msgID] = len(text.Text)
75
76	if delta == "" {
77		return nil
78	}
79
80	update := acp.UpdateAgentMessageText(delta)
81	return &update
82}
83
84func (s *Sink) translateReasoning(msgID string, reasoning message.ReasoningContent) *acp.SessionUpdate {
85	offset := s.reasoningOffsets[msgID]
86	if len(reasoning.Thinking) <= offset {
87		return nil
88	}
89
90	delta := reasoning.Thinking[offset:]
91	s.reasoningOffsets[msgID] = len(reasoning.Thinking)
92
93	if delta == "" {
94		return nil
95	}
96
97	update := acp.UpdateAgentThoughtText(delta)
98	return &update
99}