106 lines
2.6 KiB
Go
106 lines
2.6 KiB
Go
package agent
|
|
|
|
import (
|
|
"fmt"
|
|
"sync"
|
|
)
|
|
|
|
// --- Chat log methods ---
|
|
|
|
const maxChatLogBytes = 16 * 1024 * 1024
|
|
|
|
// 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
|
|
}
|
|
|
|
// --- 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()
|
|
}
|