ollie/cmd/olliesrv/internal/agent/chat.go

158 lines
4.1 KiB
Go

// chat.go — Chat log storage and streaming.
//
// The chat log is an append-only buffer of formatted output. Readers can
// stream from any offset; StreamChat returns a read function that blocks
// until data is available.
package agent
import (
"bytes"
"fmt"
"sync"
"sync/atomic"
)
// --- Chat log methods ---
const maxChatLogBytes = 16 * 1024 * 1024
// NextBlockID returns the next block ID and increments the counter.
// Block IDs are sequential integers formatted as hex strings.
func (ag *Agent) NextBlockID() string {
id := atomic.AddUint64(&ag.blockCounter, 1) - 1
return fmt.Sprintf("%x", id)
}
// AppendChat appends data to the chat log and notifies stream readers.
func (ag *Agent) AppendChat(data []byte) {
if len(data) == 0 {
return
}
ag.chatMu.Lock()
ag.chatLog = append(ag.chatLog, data...)
if len(ag.chatLog) > maxChatLogBytes {
drop := len(ag.chatLog) - maxChatLogBytes
ag.chatLog = append([]byte(nil), ag.chatLog[drop:]...)
ag.chatStart += drop
}
ag.chatVers++
ag.chatMu.Unlock()
ag.chatCond.Broadcast()
ag.chatSignalMu.Lock()
close(ag.chatSignalCh)
ag.chatSignalCh = make(chan struct{})
ag.chatSignalMu.Unlock()
}
// EnsureTrailingNewline appends a newline if the log doesn't end with one.
func (ag *Agent) EnsureTrailingNewline() {
ag.chatMu.Lock()
if len(ag.chatLog) > 0 && ag.chatLog[len(ag.chatLog)-1] != '\n' {
ag.chatLog = append(ag.chatLog, '\n')
}
ag.chatVers++
ag.chatMu.Unlock()
}
// ChatMu returns the chat log mutex for external locking (streaming).
func (ag *Agent) ChatMu() *sync.RWMutex { return &ag.chatMu }
// ChatCond returns the condvar for blocking chat readers.
func (ag *Agent) ChatCond() *sync.Cond { return ag.chatCond }
// ChatSignal returns the current chat signal channel (closed on new chat data).
func (ag *Agent) ChatSignal() <-chan struct{} {
ag.chatSignalMu.Lock()
ch := ag.chatSignalCh
ag.chatSignalMu.Unlock()
return ch
}
// ChatLog returns the raw chat log bytes (caller must hold ChatMu.RLock).
func (ag *Agent) ChatLog() []byte { return ag.chatLog }
// ChatRead returns new chat data since the given offset (base).
// Returns (data, nextBase, error). If no new data, returns empty data.
func (ag *Agent) ChatRead(base string) ([]byte, string, error) {
var offset int
if base != "" {
fmt.Sscanf(base, "%d", &offset)
} else {
ag.chatMu.RLock()
offset = len(ag.chatLog) + ag.chatStart
ag.chatMu.RUnlock()
}
ag.chatMu.RLock()
if offset < ag.chatStart {
offset = ag.chatStart
}
localOffset := offset - ag.chatStart
log := ag.chatLog
if len(log) <= localOffset {
ag.chatMu.RUnlock()
return nil, fmt.Sprintf("%d", offset), nil
}
data := make([]byte, len(log)-localOffset)
copy(data, log[localOffset:])
newOffset := ag.chatStart + len(log)
ag.chatMu.RUnlock()
return data, fmt.Sprintf("%d", newOffset), nil
}
// ChatSearchByID returns the block content for the given block ID.
// Returns the block text (without markers) and true if found, empty and false otherwise.
func (ag *Agent) ChatSearchByID(blockID string) ([]byte, bool) {
if blockID == "" {
return nil, false
}
ag.chatMu.RLock()
log := ag.chatLog
ag.chatMu.RUnlock()
// Parse blocks looking for the target ID
// Format: [[[role:name#id]]] ... [[[end]]]
target := "#" + blockID + "]]]"
idx := bytes.Index(log, []byte(target))
if idx < 0 {
return nil, false
}
// Find the start of this block header
headerStart := bytes.LastIndex(log[:idx], []byte("[[["))
if headerStart < 0 {
return nil, false
}
// Find [[[end]]] after the header
endMarker := []byte("[[[end]]]")
endIdx := bytes.Index(log[idx:], endMarker)
if endIdx < 0 {
return nil, false
}
blockEnd := idx + endIdx + len(endMarker)
// Extract block including header and end marker
block := log[headerStart:blockEnd]
return block, true
}
// --- Plan methods ---
// Plan returns a copy of the plan.
func (ag *Agent) Plan() []byte {
ag.chatMu.RLock()
p := make([]byte, len(ag.plan))
copy(p, ag.plan)
ag.chatMu.RUnlock()
return p
}
// SetPlan replaces the plan.
func (ag *Agent) SetPlan(data []byte) {
ag.chatMu.Lock()
ag.plan = make([]byte, len(data))
copy(ag.plan, data)
ag.chatMu.Unlock()
}