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

151 lines
3.6 KiB
Go

// state.go — Agent state machine and state-change notifications.
//
// The agent state tracks execution progress: idle, running, completed, error.
// WaitChange blocks until a field changes, enabling 9P clients to observe
// state transitions via the event stream.
package agent
import (
"context"
"fmt"
"strconv"
"strings"
"time"
)
// WatchField names supported by Agent.WaitChange.
const (
WatchState = "state"
WatchFeed = "feed"
)
// State returns the agent's current execution state.
func (ag *Agent) State() string {
ag.stateMu.RLock()
s := ag.state
ag.stateMu.RUnlock()
return s
}
// SetState sets the agent's execution state and notifies waiters.
func (ag *Agent) SetState(state string) {
ag.stateMu.Lock()
ag.state = state
ag.stateSince = time.Now()
ag.stateMu.Unlock()
ag.notifyChange()
if ag.onStateChange != nil {
ag.onStateChange(ag.id, state)
}
}
// Status renders a human-readable one-line status: the current state and how
// long it has been in that state, e.g. "thinking · 12s" or "calling shell · 3s".
// The raw state string is machine-facing; this is the human-facing view.
func (ag *Agent) Status() string {
ag.stateMu.RLock()
state := ag.state
since := ag.stateSince
ag.stateMu.RUnlock()
if state == "" {
state = "idle"
}
// "calling: shell" reads better as "calling shell".
state = strings.Replace(state, "calling: ", "calling ", 1)
if state == "idle" || since.IsZero() {
return state
}
return state + " · " + humanDuration(time.Since(since))
}
// humanDuration renders a short elapsed duration: "3s", "2m10s", "1h04m".
func humanDuration(d time.Duration) string {
secs := int(d.Seconds())
if secs < 60 {
return strconv.Itoa(secs) + "s"
}
if secs < 3600 {
return strconv.Itoa(secs/60) + "m" + strconv.Itoa(secs%60) + "s"
}
return strconv.Itoa(secs/3600) + "h" + fmt.Sprintf("%02dm", (secs%3600)/60)
}
// Reply returns the agent's last assistant response.
func (ag *Agent) Reply() string {
ag.stateMu.RLock()
r := ag.reply
ag.stateMu.RUnlock()
return r
}
// SetOnStateChange sets a callback invoked whenever agent state changes.
func (ag *Agent) SetOnStateChange(fn func(agentID, state string)) {
ag.onStateChange = fn
}
// setReply sets the agent's last response.
func (ag *Agent) setReply(reply string) {
ag.stateMu.Lock()
ag.reply = reply
ag.stateMu.Unlock()
}
// notifyChange wakes all goroutines waiting on state changes.
func (ag *Agent) notifyChange() {
ag.signalMu.Lock()
close(ag.signalCh)
ag.signalCh = make(chan struct{})
ag.signalMu.Unlock()
}
// SignalCh returns the current signal channel (closed on any change).
func (ag *Agent) SignalCh() <-chan struct{} {
ag.signalMu.Lock()
ch := ag.signalCh
ag.signalMu.Unlock()
return ch
}
// WaitChange blocks until the agent's state differs from current.
// Returns the new value and true, or ("", false) if ctx is cancelled.
func (ag *Agent) WaitChange(ctx context.Context, field, current string) (string, bool) {
for {
// Snapshot the signal channel before reading state.
ag.signalMu.Lock()
ch := ag.signalCh
ag.signalMu.Unlock()
var val string
switch field {
case WatchState:
val = ag.State()
case WatchFeed:
val = ag.Feed.Hash()
if val == "" {
val = current // no data yet — block
}
default:
return "", false
}
if val != current {
return val, true
}
// Wait for either a state change or context cancellation.
select {
case <-ch:
// Changed — loop to re-check.
case <-ctx.Done():
return "", false
}
}
}
// emit sends an event to the agent's output handler.
func (ag *Agent) emit(ev Event) {
if ag.output != nil {
ag.output(ev)
}
}