Merge remote-tracking branch 'origin/polecat/mediocre-mk0vwdaf'
This commit is contained in:
76
internal/agent/state.go
Normal file
76
internal/agent/state.go
Normal file
@@ -0,0 +1,76 @@
|
|||||||
|
// Package agent provides shared types and utilities for Gas Town agents
|
||||||
|
// (witness, refinery, deacon, etc.).
|
||||||
|
package agent
|
||||||
|
|
||||||
|
import (
|
||||||
|
"encoding/json"
|
||||||
|
"os"
|
||||||
|
"path/filepath"
|
||||||
|
|
||||||
|
"github.com/steveyegge/gastown/internal/util"
|
||||||
|
)
|
||||||
|
|
||||||
|
// State represents an agent's running state.
|
||||||
|
type State string
|
||||||
|
|
||||||
|
const (
|
||||||
|
// StateStopped means the agent is not running.
|
||||||
|
StateStopped State = "stopped"
|
||||||
|
|
||||||
|
// StateRunning means the agent is actively operating.
|
||||||
|
StateRunning State = "running"
|
||||||
|
|
||||||
|
// StatePaused means the agent is paused (not operating but not stopped).
|
||||||
|
StatePaused State = "paused"
|
||||||
|
)
|
||||||
|
|
||||||
|
// StateManager handles loading and saving agent state to disk.
|
||||||
|
// It uses generics to work with any state type.
|
||||||
|
type StateManager[T any] struct {
|
||||||
|
stateFilePath string
|
||||||
|
defaultFactory func() *T
|
||||||
|
}
|
||||||
|
|
||||||
|
// NewStateManager creates a new StateManager for the given state file path.
|
||||||
|
// The defaultFactory function is called when the state file doesn't exist
|
||||||
|
// to create a new state with default values.
|
||||||
|
func NewStateManager[T any](rigPath, stateFileName string, defaultFactory func() *T) *StateManager[T] {
|
||||||
|
return &StateManager[T]{
|
||||||
|
stateFilePath: filepath.Join(rigPath, ".runtime", stateFileName),
|
||||||
|
defaultFactory: defaultFactory,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// StateFile returns the path to the state file.
|
||||||
|
func (m *StateManager[T]) StateFile() string {
|
||||||
|
return m.stateFilePath
|
||||||
|
}
|
||||||
|
|
||||||
|
// Load loads agent state from disk.
|
||||||
|
// If the file doesn't exist, returns a new state created by the default factory.
|
||||||
|
func (m *StateManager[T]) Load() (*T, error) {
|
||||||
|
data, err := os.ReadFile(m.stateFilePath)
|
||||||
|
if err != nil {
|
||||||
|
if os.IsNotExist(err) {
|
||||||
|
return m.defaultFactory(), nil
|
||||||
|
}
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
|
||||||
|
var state T
|
||||||
|
if err := json.Unmarshal(data, &state); err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
|
||||||
|
return &state, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// Save persists agent state to disk using atomic write.
|
||||||
|
func (m *StateManager[T]) Save(state *T) error {
|
||||||
|
dir := filepath.Dir(m.stateFilePath)
|
||||||
|
if err := os.MkdirAll(dir, 0755); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
|
return util.AtomicWriteJSON(m.stateFilePath, state)
|
||||||
|
}
|
||||||
@@ -2,7 +2,6 @@ package refinery
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"bytes"
|
"bytes"
|
||||||
"encoding/json"
|
|
||||||
"errors"
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
"io"
|
"io"
|
||||||
@@ -13,6 +12,7 @@ import (
|
|||||||
"strings"
|
"strings"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
|
"github.com/steveyegge/gastown/internal/agent"
|
||||||
"github.com/steveyegge/gastown/internal/beads"
|
"github.com/steveyegge/gastown/internal/beads"
|
||||||
"github.com/steveyegge/gastown/internal/claude"
|
"github.com/steveyegge/gastown/internal/claude"
|
||||||
"github.com/steveyegge/gastown/internal/config"
|
"github.com/steveyegge/gastown/internal/config"
|
||||||
@@ -33,9 +33,10 @@ var (
|
|||||||
|
|
||||||
// Manager handles refinery lifecycle and queue operations.
|
// Manager handles refinery lifecycle and queue operations.
|
||||||
type Manager struct {
|
type Manager struct {
|
||||||
rig *rig.Rig
|
rig *rig.Rig
|
||||||
workDir string
|
workDir string
|
||||||
output io.Writer // Output destination for user-facing messages
|
output io.Writer // Output destination for user-facing messages
|
||||||
|
stateManager *agent.StateManager[Refinery]
|
||||||
}
|
}
|
||||||
|
|
||||||
// NewManager creates a new refinery manager for a rig.
|
// NewManager creates a new refinery manager for a rig.
|
||||||
@@ -44,6 +45,12 @@ func NewManager(r *rig.Rig) *Manager {
|
|||||||
rig: r,
|
rig: r,
|
||||||
workDir: r.Path,
|
workDir: r.Path,
|
||||||
output: os.Stdout,
|
output: os.Stdout,
|
||||||
|
stateManager: agent.NewStateManager[Refinery](r.Path, "refinery.json", func() *Refinery {
|
||||||
|
return &Refinery{
|
||||||
|
RigName: r.Name,
|
||||||
|
State: StateStopped,
|
||||||
|
}
|
||||||
|
}),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -55,7 +62,7 @@ func (m *Manager) SetOutput(w io.Writer) {
|
|||||||
|
|
||||||
// stateFile returns the path to the refinery state file.
|
// stateFile returns the path to the refinery state file.
|
||||||
func (m *Manager) stateFile() string {
|
func (m *Manager) stateFile() string {
|
||||||
return filepath.Join(m.rig.Path, ".runtime", "refinery.json")
|
return m.stateManager.StateFile()
|
||||||
}
|
}
|
||||||
|
|
||||||
// sessionName returns the tmux session name for this refinery.
|
// sessionName returns the tmux session name for this refinery.
|
||||||
@@ -65,33 +72,12 @@ func (m *Manager) sessionName() string {
|
|||||||
|
|
||||||
// loadState loads refinery state from disk.
|
// loadState loads refinery state from disk.
|
||||||
func (m *Manager) loadState() (*Refinery, error) {
|
func (m *Manager) loadState() (*Refinery, error) {
|
||||||
data, err := os.ReadFile(m.stateFile())
|
return m.stateManager.Load()
|
||||||
if err != nil {
|
|
||||||
if os.IsNotExist(err) {
|
|
||||||
return &Refinery{
|
|
||||||
RigName: m.rig.Name,
|
|
||||||
State: StateStopped,
|
|
||||||
}, nil
|
|
||||||
}
|
|
||||||
return nil, err
|
|
||||||
}
|
|
||||||
|
|
||||||
var ref Refinery
|
|
||||||
if err := json.Unmarshal(data, &ref); err != nil {
|
|
||||||
return nil, err
|
|
||||||
}
|
|
||||||
|
|
||||||
return &ref, nil
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// saveState persists refinery state to disk using atomic write.
|
// saveState persists refinery state to disk using atomic write.
|
||||||
func (m *Manager) saveState(ref *Refinery) error {
|
func (m *Manager) saveState(ref *Refinery) error {
|
||||||
dir := filepath.Dir(m.stateFile())
|
return m.stateManager.Save(ref)
|
||||||
if err := os.MkdirAll(dir, 0755); err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
|
|
||||||
return util.AtomicWriteJSON(m.stateFile(), ref)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// Status returns the current refinery status.
|
// Status returns the current refinery status.
|
||||||
|
|||||||
@@ -5,20 +5,18 @@ import (
|
|||||||
"errors"
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
|
"github.com/steveyegge/gastown/internal/agent"
|
||||||
)
|
)
|
||||||
|
|
||||||
// State represents the refinery's running state.
|
// State is an alias for agent.State for backwards compatibility.
|
||||||
type State string
|
type State = agent.State
|
||||||
|
|
||||||
|
// State constants - re-exported from agent package for backwards compatibility.
|
||||||
const (
|
const (
|
||||||
// StateStopped means the refinery is not running.
|
StateStopped = agent.StateStopped
|
||||||
StateStopped State = "stopped"
|
StateRunning = agent.StateRunning
|
||||||
|
StatePaused = agent.StatePaused
|
||||||
// StateRunning means the refinery is actively processing.
|
|
||||||
StateRunning State = "running"
|
|
||||||
|
|
||||||
// StatePaused means the refinery is paused (not processing new items).
|
|
||||||
StatePaused State = "paused"
|
|
||||||
)
|
)
|
||||||
|
|
||||||
// Refinery represents a rig's merge queue processor.
|
// Refinery represents a rig's merge queue processor.
|
||||||
|
|||||||
@@ -1,12 +1,11 @@
|
|||||||
package witness
|
package witness
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"encoding/json"
|
|
||||||
"errors"
|
"errors"
|
||||||
"os"
|
"os"
|
||||||
"path/filepath"
|
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
|
"github.com/steveyegge/gastown/internal/agent"
|
||||||
"github.com/steveyegge/gastown/internal/rig"
|
"github.com/steveyegge/gastown/internal/rig"
|
||||||
"github.com/steveyegge/gastown/internal/util"
|
"github.com/steveyegge/gastown/internal/util"
|
||||||
)
|
)
|
||||||
@@ -19,8 +18,9 @@ var (
|
|||||||
|
|
||||||
// Manager handles witness lifecycle and monitoring operations.
|
// Manager handles witness lifecycle and monitoring operations.
|
||||||
type Manager struct {
|
type Manager struct {
|
||||||
rig *rig.Rig
|
rig *rig.Rig
|
||||||
workDir string
|
workDir string
|
||||||
|
stateManager *agent.StateManager[Witness]
|
||||||
}
|
}
|
||||||
|
|
||||||
// NewManager creates a new witness manager for a rig.
|
// NewManager creates a new witness manager for a rig.
|
||||||
@@ -28,43 +28,28 @@ func NewManager(r *rig.Rig) *Manager {
|
|||||||
return &Manager{
|
return &Manager{
|
||||||
rig: r,
|
rig: r,
|
||||||
workDir: r.Path,
|
workDir: r.Path,
|
||||||
|
stateManager: agent.NewStateManager[Witness](r.Path, "witness.json", func() *Witness {
|
||||||
|
return &Witness{
|
||||||
|
RigName: r.Name,
|
||||||
|
State: StateStopped,
|
||||||
|
}
|
||||||
|
}),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// stateFile returns the path to the witness state file.
|
// stateFile returns the path to the witness state file.
|
||||||
func (m *Manager) stateFile() string {
|
func (m *Manager) stateFile() string {
|
||||||
return filepath.Join(m.rig.Path, ".runtime", "witness.json")
|
return m.stateManager.StateFile()
|
||||||
}
|
}
|
||||||
|
|
||||||
// loadState loads witness state from disk.
|
// loadState loads witness state from disk.
|
||||||
func (m *Manager) loadState() (*Witness, error) {
|
func (m *Manager) loadState() (*Witness, error) {
|
||||||
data, err := os.ReadFile(m.stateFile())
|
return m.stateManager.Load()
|
||||||
if err != nil {
|
|
||||||
if os.IsNotExist(err) {
|
|
||||||
return &Witness{
|
|
||||||
RigName: m.rig.Name,
|
|
||||||
State: StateStopped,
|
|
||||||
}, nil
|
|
||||||
}
|
|
||||||
return nil, err
|
|
||||||
}
|
|
||||||
|
|
||||||
var w Witness
|
|
||||||
if err := json.Unmarshal(data, &w); err != nil {
|
|
||||||
return nil, err
|
|
||||||
}
|
|
||||||
|
|
||||||
return &w, nil
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// saveState persists witness state to disk using atomic write.
|
// saveState persists witness state to disk using atomic write.
|
||||||
func (m *Manager) saveState(w *Witness) error {
|
func (m *Manager) saveState(w *Witness) error {
|
||||||
dir := filepath.Dir(m.stateFile())
|
return m.stateManager.Save(w)
|
||||||
if err := os.MkdirAll(dir, 0755); err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
|
|
||||||
return util.AtomicWriteJSON(m.stateFile(), w)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// Status returns the current witness status.
|
// Status returns the current witness status.
|
||||||
|
|||||||
@@ -3,20 +3,18 @@ package witness
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
|
"github.com/steveyegge/gastown/internal/agent"
|
||||||
)
|
)
|
||||||
|
|
||||||
// State represents the witness's running state.
|
// State is an alias for agent.State for backwards compatibility.
|
||||||
type State string
|
type State = agent.State
|
||||||
|
|
||||||
|
// State constants - re-exported from agent package for backwards compatibility.
|
||||||
const (
|
const (
|
||||||
// StateStopped means the witness is not running.
|
StateStopped = agent.StateStopped
|
||||||
StateStopped State = "stopped"
|
StateRunning = agent.StateRunning
|
||||||
|
StatePaused = agent.StatePaused
|
||||||
// StateRunning means the witness is actively monitoring.
|
|
||||||
StateRunning State = "running"
|
|
||||||
|
|
||||||
// StatePaused means the witness is paused (not monitoring).
|
|
||||||
StatePaused State = "paused"
|
|
||||||
)
|
)
|
||||||
|
|
||||||
// Witness represents a rig's polecat monitoring agent.
|
// Witness represents a rig's polecat monitoring agent.
|
||||||
|
|||||||
Reference in New Issue
Block a user