- tmux som sanningskälla (agentoberoende, överlever omstart via adoption) - API: sessions, prompt, question/answer, keys, mode, config, shares, notify, events - frågedetektor testad mot riktiga agy 1.1.9-dumpar + mock - integrationstest: 26 tester i isolerad debian-container (scripts/run-integration.sh) - e2e-verifierad mot riktig agy (trust-fråga -> svar -> prompt) - .gitea/workflows/helmd-release.yaml -> rullande helmd-latest (x64+arm64) Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
324 lines
8.5 KiB
Go
324 lines
8.5 KiB
Go
package main
|
|
|
|
import (
|
|
"crypto/sha256"
|
|
"encoding/hex"
|
|
"fmt"
|
|
"log"
|
|
"strings"
|
|
"sync"
|
|
"time"
|
|
)
|
|
|
|
const sessPrefix = "helm-"
|
|
|
|
// SessionState is what clients see for one managed tmux session.
|
|
type SessionState struct {
|
|
Name string `json:"name"` // API name (without tmux prefix)
|
|
Tmux string `json:"tmux"` // tmux session name
|
|
Cmd string `json:"cmd"` // command line launched in the pane
|
|
Alive bool `json:"alive"` // pane process still running
|
|
Seq uint64 `json:"seq"` // bumped on every screen change
|
|
Question *Question `json:"question"` // pending question, if any
|
|
Changed time.Time `json:"changed"` // last screen change
|
|
Created time.Time `json:"created"`
|
|
}
|
|
|
|
type Event struct {
|
|
Type string `json:"type"` // screen | question | session | share
|
|
Session string `json:"session,omitempty"`
|
|
Data interface{} `json:"data,omitempty"`
|
|
}
|
|
|
|
type Manager struct {
|
|
cfg *Config
|
|
ntfy *Ntfy
|
|
|
|
mu sync.Mutex
|
|
sessions map[string]*SessionState
|
|
notified map[string]string // session -> last question hash sent to ntfy
|
|
|
|
subMu sync.Mutex
|
|
subs map[chan Event]struct{}
|
|
}
|
|
|
|
func NewManager(cfg *Config, ntfy *Ntfy) *Manager {
|
|
m := &Manager{
|
|
cfg: cfg,
|
|
ntfy: ntfy,
|
|
sessions: map[string]*SessionState{},
|
|
notified: map[string]string{},
|
|
subs: map[chan Event]struct{}{},
|
|
}
|
|
m.adoptExisting()
|
|
go m.pollLoop()
|
|
return m
|
|
}
|
|
|
|
// adoptExisting re-registers helm-* tmux sessions after a helmd restart:
|
|
// tmux is the source of truth, helmd only mirrors it.
|
|
func (m *Manager) adoptExisting() {
|
|
for _, t := range tmuxListSessions() {
|
|
if strings.HasPrefix(t, sessPrefix) {
|
|
name := strings.TrimPrefix(t, sessPrefix)
|
|
m.sessions[name] = &SessionState{
|
|
Name: name, Tmux: t, Cmd: "(adopterad)",
|
|
Created: time.Now(), Changed: time.Now(),
|
|
}
|
|
log.Printf("adopterade befintlig tmux-session %s", t)
|
|
}
|
|
}
|
|
}
|
|
|
|
func (m *Manager) Subscribe() chan Event {
|
|
ch := make(chan Event, 32)
|
|
m.subMu.Lock()
|
|
m.subs[ch] = struct{}{}
|
|
m.subMu.Unlock()
|
|
return ch
|
|
}
|
|
|
|
func (m *Manager) Unsubscribe(ch chan Event) {
|
|
m.subMu.Lock()
|
|
delete(m.subs, ch)
|
|
m.subMu.Unlock()
|
|
}
|
|
|
|
func (m *Manager) emit(ev Event) {
|
|
m.subMu.Lock()
|
|
for ch := range m.subs {
|
|
select {
|
|
case ch <- ev:
|
|
default: // slow client: drop rather than block the poller
|
|
}
|
|
}
|
|
m.subMu.Unlock()
|
|
}
|
|
|
|
func (m *Manager) List() []*SessionState {
|
|
m.mu.Lock()
|
|
defer m.mu.Unlock()
|
|
out := make([]*SessionState, 0, len(m.sessions))
|
|
for _, s := range m.sessions {
|
|
out = append(out, s)
|
|
}
|
|
return out
|
|
}
|
|
|
|
func (m *Manager) Get(name string) (*SessionState, bool) {
|
|
m.mu.Lock()
|
|
defer m.mu.Unlock()
|
|
s, ok := m.sessions[name]
|
|
return s, ok
|
|
}
|
|
|
|
// Create starts a new tmux session running the agent (or a custom
|
|
// command). An empty cmd uses the configured agent + args.
|
|
func (m *Manager) Create(name, cmdline, cwd string) (*SessionState, error) {
|
|
if name == "" || strings.ContainsAny(name, " \t/.:") {
|
|
return nil, fmt.Errorf("ogiltigt sessionsnamn %q", name)
|
|
}
|
|
m.mu.Lock()
|
|
defer m.mu.Unlock()
|
|
if _, exists := m.sessions[name]; exists {
|
|
return nil, fmt.Errorf("sessionen %s finns redan", name)
|
|
}
|
|
tmuxName := sessPrefix + name
|
|
if tmuxHasSession(tmuxName) {
|
|
return nil, fmt.Errorf("tmux-sessionen %s finns redan (adoptera med POST {\"adopt\":true})", tmuxName)
|
|
}
|
|
if cmdline == "" {
|
|
cmdline = m.cfg.AgentCmd
|
|
if len(m.cfg.AgentArgs) > 0 {
|
|
cmdline += " " + strings.Join(m.cfg.AgentArgs, " ")
|
|
}
|
|
}
|
|
if cwd == "" {
|
|
cwd = m.cfg.Workdir
|
|
}
|
|
if err := tmuxNewSession(tmuxName, cwd, m.cfg.Cols, m.cfg.Rows, cmdline); err != nil {
|
|
return nil, err
|
|
}
|
|
s := &SessionState{
|
|
Name: name, Tmux: tmuxName, Cmd: cmdline, Alive: true,
|
|
Created: time.Now(), Changed: time.Now(),
|
|
}
|
|
m.sessions[name] = s
|
|
m.emit(Event{Type: "session", Session: name, Data: "created"})
|
|
return s, nil
|
|
}
|
|
|
|
// Adopt registers an already-running tmux session under helmd control.
|
|
func (m *Manager) Adopt(name, tmuxName string) (*SessionState, error) {
|
|
if tmuxName == "" {
|
|
tmuxName = sessPrefix + name
|
|
}
|
|
if !tmuxHasSession(tmuxName) {
|
|
return nil, fmt.Errorf("ingen tmux-session %q", tmuxName)
|
|
}
|
|
m.mu.Lock()
|
|
defer m.mu.Unlock()
|
|
if _, exists := m.sessions[name]; exists {
|
|
return nil, fmt.Errorf("sessionen %s finns redan", name)
|
|
}
|
|
s := &SessionState{
|
|
Name: name, Tmux: tmuxName, Cmd: "(adopterad)", Alive: true,
|
|
Created: time.Now(), Changed: time.Now(),
|
|
}
|
|
m.sessions[name] = s
|
|
m.emit(Event{Type: "session", Session: name, Data: "adopted"})
|
|
return s, nil
|
|
}
|
|
|
|
// Kill terminates the tmux session and forgets it.
|
|
func (m *Manager) Kill(name string) error {
|
|
m.mu.Lock()
|
|
s, ok := m.sessions[name]
|
|
if ok {
|
|
delete(m.sessions, name)
|
|
delete(m.notified, name)
|
|
}
|
|
m.mu.Unlock()
|
|
if !ok {
|
|
return fmt.Errorf("ingen session %q", name)
|
|
}
|
|
m.emit(Event{Type: "session", Session: name, Data: "killed"})
|
|
if tmuxHasSession(s.Tmux) {
|
|
return tmuxKillSession(s.Tmux)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// Prompt pastes text into the session and (by default) submits it.
|
|
func (m *Manager) Prompt(name, text string, submit bool) error {
|
|
s, ok := m.Get(name)
|
|
if !ok {
|
|
return fmt.Errorf("ingen session %q", name)
|
|
}
|
|
if err := tmuxPasteText(s.Tmux, text); err != nil {
|
|
return err
|
|
}
|
|
if submit {
|
|
time.Sleep(150 * time.Millisecond) // let the TUI ingest the paste
|
|
return tmuxSendKeys(s.Tmux, "Enter")
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// Answer picks an option in the pending question. Numbered lists get
|
|
// the digit key; unnumbered lists are navigated relative to the
|
|
// current selection.
|
|
func (m *Manager) Answer(name string, option int) error {
|
|
s, ok := m.Get(name)
|
|
if !ok {
|
|
return fmt.Errorf("ingen session %q", name)
|
|
}
|
|
q := s.Question
|
|
if q == nil {
|
|
return fmt.Errorf("ingen väntande fråga i %s", name)
|
|
}
|
|
if option < 1 || option > len(q.Options) {
|
|
return fmt.Errorf("option %d utanför 1..%d", option, len(q.Options))
|
|
}
|
|
idx := option - 1
|
|
if q.Numbered {
|
|
num := q.Options[idx].Num
|
|
if err := tmuxSendKeys(s.Tmux, fmt.Sprint(num)); err != nil {
|
|
return err
|
|
}
|
|
time.Sleep(200 * time.Millisecond)
|
|
// Digit selection auto-confirms in agy; if the prompt is still
|
|
// up (pure highlight moved), Enter confirms. Extra Enter on an
|
|
// already-gone prompt lands in the input box and is harmless
|
|
// only if empty — so only confirm when the question persists.
|
|
if cur, err := tmuxCapture(s.Tmux, false, 0); err == nil {
|
|
if nq := DetectQuestion(cur); nq != nil && nq.Hash == q.Hash {
|
|
return tmuxSendKeys(s.Tmux, "Enter")
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
if q.Selected < 0 {
|
|
return fmt.Errorf("okänd markering — svara med raw keys i stället")
|
|
}
|
|
diff := idx - q.Selected
|
|
key, n := "Down", diff
|
|
if diff < 0 {
|
|
key, n = "Up", -diff
|
|
}
|
|
for i := 0; i < n; i++ {
|
|
if err := tmuxSendKeys(s.Tmux, key); err != nil {
|
|
return err
|
|
}
|
|
time.Sleep(60 * time.Millisecond)
|
|
}
|
|
return tmuxSendKeys(s.Tmux, "Enter")
|
|
}
|
|
|
|
// Keys sends raw tmux key names (Escape, BTab, C-c, y, ...).
|
|
func (m *Manager) Keys(name string, keys []string) error {
|
|
s, ok := m.Get(name)
|
|
if !ok {
|
|
return fmt.Errorf("ingen session %q", name)
|
|
}
|
|
return tmuxSendKeys(s.Tmux, keys...)
|
|
}
|
|
|
|
func (m *Manager) Screen(name string, ansi bool, history int) (string, error) {
|
|
s, ok := m.Get(name)
|
|
if !ok {
|
|
return "", fmt.Errorf("ingen session %q", name)
|
|
}
|
|
return tmuxCapture(s.Tmux, ansi, history)
|
|
}
|
|
|
|
// pollLoop watches every session for screen changes and questions.
|
|
func (m *Manager) pollLoop() {
|
|
interval := time.Duration(m.cfg.PollMs) * time.Millisecond
|
|
if interval <= 0 {
|
|
interval = 500 * time.Millisecond
|
|
}
|
|
hashes := map[string]string{}
|
|
for {
|
|
time.Sleep(interval)
|
|
for _, s := range m.List() {
|
|
alive := tmuxHasSession(s.Tmux) && tmuxSessionAlive(s.Tmux)
|
|
if !alive {
|
|
if s.Alive {
|
|
s.Alive = false
|
|
m.emit(Event{Type: "session", Session: s.Name, Data: "exited"})
|
|
m.ntfy.Publish("", "helm: "+s.Name+" avslutades", "Agentprocessen i "+s.Tmux+" kör inte längre.", "default")
|
|
}
|
|
continue
|
|
}
|
|
s.Alive = true
|
|
screen, err := tmuxCapture(s.Tmux, false, 0)
|
|
if err != nil {
|
|
continue
|
|
}
|
|
sum := sha256.Sum256([]byte(screen))
|
|
h := hex.EncodeToString(sum[:8])
|
|
if h == hashes[s.Name] {
|
|
continue
|
|
}
|
|
hashes[s.Name] = h
|
|
s.Seq++
|
|
s.Changed = time.Now()
|
|
s.Question = DetectQuestion(screen)
|
|
m.emit(Event{Type: "screen", Session: s.Name, Data: s.Seq})
|
|
if s.Question != nil {
|
|
m.emit(Event{Type: "question", Session: s.Name, Data: s.Question})
|
|
m.mu.Lock()
|
|
already := m.notified[s.Name] == s.Question.Hash
|
|
if !already {
|
|
m.notified[s.Name] = s.Question.Hash
|
|
}
|
|
m.mu.Unlock()
|
|
if !already {
|
|
m.ntfy.Question(s.Name, s.Question)
|
|
}
|
|
}
|
|
}
|
|
}
|
|
}
|