247 lines
5.3 KiB
Go
247 lines
5.3 KiB
Go
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()
|
||
|
|
}
|