Files

247 lines
5.3 KiB
Go
Raw Permalink Normal View History

package board
import (
"encoding/json"
"fmt"
"os"
"path/filepath"
"sync"
"time"
)
// TaskState represents the actual state of a task
type TaskState struct {
TaskID string `json:"task_id"`
Status string `json:"status"` // "pending", "in_progress", "completed", "failed"
CompletedAt time.Time `json:"completed_at,omitempty"`
FailedAt time.Time `json:"failed_at,omitempty"`
Error string `json:"error,omitempty"`
Branch string `json:"branch,omitempty"`
Metrics map[string]interface{} `json:"metrics,omitempty"`
}
// StateTracker tracks actual task states
type StateTracker struct {
mu sync.RWMutex
basePath string
states map[string]*TaskState
lastUpdate time.Time
}
// NewStateTracker creates a new state tracker
func NewStateTracker(basePath string) *StateTracker {
return &StateTracker{
basePath: basePath,
states: make(map[string]*TaskState),
}
}
// UpdateTaskState updates the state of a task
func (st *StateTracker) UpdateTaskState(taskID, status, branch string, err error) error {
st.mu.Lock()
defer st.mu.Unlock()
errorMsg := ""
if err != nil {
errorMsg = err.Error()
}
state := &TaskState{
TaskID: taskID,
Status: status,
Branch: branch,
Error: errorMsg,
Metrics: make(map[string]interface{}),
}
if status == "completed" {
state.CompletedAt = time.Now()
} else if status == "failed" {
state.FailedAt = time.Now()
}
st.states[taskID] = state
st.lastUpdate = time.Now()
return st.persistLocked()
}
// GetTaskState retrieves the state of a task
func (st *StateTracker) GetTaskState(taskID string) *TaskState {
st.mu.RLock()
defer st.mu.RUnlock()
return st.states[taskID]
}
// GetAllStates returns all task states
func (st *StateTracker) GetAllStates() map[string]*TaskState {
st.mu.RLock()
defer st.mu.RUnlock()
// Return a copy
copy := make(map[string]*TaskState)
for k, v := range st.states {
copy[k] = v
}
return copy
}
// GetCompletedTasks returns all completed tasks
func (st *StateTracker) GetCompletedTasks() []string {
st.mu.RLock()
defer st.mu.RUnlock()
completed := make([]string, 0)
for _, state := range st.states {
if state.Status == "completed" {
completed = append(completed, state.TaskID)
}
}
return completed
}
// GetFailedTasks returns all failed tasks
func (st *StateTracker) GetFailedTasks() []string {
st.mu.RLock()
defer st.mu.RUnlock()
failed := make([]string, 0)
for _, state := range st.states {
if state.Status == "failed" {
failed = append(failed, state.TaskID)
}
}
return failed
}
// GetPendingTasks returns all pending tasks
func (st *StateTracker) GetPendingTasks() []string {
st.mu.RLock()
defer st.mu.RUnlock()
pending := make([]string, 0)
for _, state := range st.states {
if state.Status == "pending" || state.Status == "in_progress" {
pending = append(pending, state.TaskID)
}
}
return pending
}
// AddMetric adds a metric to a task
func (st *StateTracker) AddMetric(taskID, metricName string, value interface{}) error {
st.mu.Lock()
defer st.mu.Unlock()
state, exists := st.states[taskID]
if !exists {
return fmt.Errorf("task state not found: %s", taskID)
}
state.Metrics[metricName] = value
st.lastUpdate = time.Now()
return st.persistLocked()
}
// Load loads state from disk
func (st *StateTracker) Load() error {
st.mu.Lock()
defer st.mu.Unlock()
statePath := filepath.Join(st.basePath, "board", "state.json")
data, err := os.ReadFile(statePath)
if err != nil {
if os.IsNotExist(err) {
return nil // File doesn't exist yet
}
return err
}
var states []TaskState
if err := json.Unmarshal(data, &states); err != nil {
return err
}
st.states = make(map[string]*TaskState)
for i := range states {
st.states[states[i].TaskID] = &states[i]
}
return nil
}
// persistLocked saves state to disk (must be called with lock held)
func (st *StateTracker) persistLocked() error {
states := make([]TaskState, 0)
for _, state := range st.states {
states = append(states, *state)
}
data, err := json.MarshalIndent(states, "", " ")
if err != nil {
return err
}
statePath := filepath.Join(st.basePath, "board", "state.json")
// Create directory if it doesn't exist
if err := os.MkdirAll(filepath.Dir(statePath), 0755); err != nil {
return err
}
return os.WriteFile(statePath, data, 0644)
}
// GetAsCompletionMap returns task completion status as a boolean map
func (st *StateTracker) GetAsCompletionMap() map[string]bool {
st.mu.RLock()
defer st.mu.RUnlock()
completion := make(map[string]bool)
for taskID, state := range st.states {
completion[taskID] = state.Status == "completed"
}
return completion
}
// GetLastUpdate returns the last time state was updated
func (st *StateTracker) GetLastUpdate() time.Time {
st.mu.RLock()
defer st.mu.RUnlock()
return st.lastUpdate
}
// GetStats returns statistics about task states
func (st *StateTracker) GetStats() map[string]interface{} {
st.mu.RLock()
defer st.mu.RUnlock()
stats := make(map[string]interface{})
counts := make(map[string]int)
for _, state := range st.states {
counts[state.Status]++
}
stats["total"] = len(st.states)
stats["counts"] = counts
stats["last_update"] = st.lastUpdate
return stats
}
// Reset clears all state
func (st *StateTracker) Reset() error {
st.mu.Lock()
defer st.mu.Unlock()
st.states = make(map[string]*TaskState)
st.lastUpdate = time.Time{}
return st.persistLocked()
}