All checks were successful
release-tag / release-image (push) Successful in 1m36s
354 lines
8.8 KiB
Go
354 lines
8.8 KiB
Go
package state
|
|
|
|
import (
|
|
"bufio"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"os"
|
|
"path/filepath"
|
|
"sort"
|
|
"strconv"
|
|
"strings"
|
|
"sync"
|
|
|
|
"github.com/example/glpi-ai-agent/internal/model"
|
|
)
|
|
|
|
type Store struct {
|
|
mu sync.RWMutex
|
|
path string
|
|
indexPath string
|
|
processedVersions map[int64]string
|
|
escalationKeys map[string]struct{}
|
|
runs []model.RunRecord
|
|
runIndex map[string]int
|
|
analysisIndex map[string]model.AnalysisRun
|
|
maxRuns int
|
|
}
|
|
|
|
type durableIndex struct {
|
|
ProcessedVersions map[string]string `json:"processed_versions"`
|
|
EscalationKeys []string `json:"escalation_keys"`
|
|
}
|
|
|
|
func Open(dir string, maxRuns int) (*Store, error) {
|
|
if err := os.MkdirAll(dir, 0o750); err != nil {
|
|
return nil, err
|
|
}
|
|
s := &Store{
|
|
path: filepath.Join(dir, "runs.jsonl"),
|
|
indexPath: filepath.Join(dir, "state-index.json"),
|
|
processedVersions: map[int64]string{},
|
|
escalationKeys: map[string]struct{}{},
|
|
runIndex: map[string]int{},
|
|
analysisIndex: map[string]model.AnalysisRun{},
|
|
maxRuns: maxRuns,
|
|
}
|
|
// Fail fast during startup if the persistent data path is not writable.
|
|
// A read-only/root-owned Docker volume must not be discovered only after the first ticket.
|
|
f, err := os.OpenFile(s.path, os.O_CREATE|os.O_APPEND|os.O_WRONLY, 0o640)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("state directory %q is not writable: %w", dir, err)
|
|
}
|
|
if err := f.Close(); err != nil {
|
|
return nil, fmt.Errorf("close state write probe: %w", err)
|
|
}
|
|
if err := s.loadDurableIndex(); err != nil {
|
|
return nil, err
|
|
}
|
|
if err := s.load(); err != nil {
|
|
return nil, err
|
|
}
|
|
return s, nil
|
|
}
|
|
func (s *Store) Seen(id int64, version string) bool {
|
|
s.mu.RLock()
|
|
defer s.mu.RUnlock()
|
|
return version != "" && s.processedVersions[id] == version
|
|
}
|
|
func (s *Store) ProcessedVersionCount() int {
|
|
s.mu.RLock()
|
|
defer s.mu.RUnlock()
|
|
return len(s.processedVersions)
|
|
}
|
|
|
|
func (s *Store) Append(r model.RunRecord) error {
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
f, err := os.OpenFile(s.path, os.O_CREATE|os.O_APPEND|os.O_WRONLY, 0o640)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
enc := json.NewEncoder(f)
|
|
if err = enc.Encode(r); err == nil {
|
|
err = f.Sync()
|
|
}
|
|
cerr := f.Close()
|
|
if err == nil {
|
|
err = cerr
|
|
}
|
|
if err != nil {
|
|
return err
|
|
}
|
|
s.absorbDurableStateLocked(r)
|
|
s.runs = append(s.runs, r)
|
|
if len(s.runs) > s.maxRuns {
|
|
s.runs = s.runs[len(s.runs)-s.maxRuns:]
|
|
}
|
|
s.rebuildIndexesLocked()
|
|
if err := s.persistDurableIndexLocked(); err != nil {
|
|
return fmt.Errorf("persist durable state index: %w", err)
|
|
}
|
|
if info, statErr := os.Stat(s.path); statErr == nil && info.Size() > 64<<20 {
|
|
if compactErr := s.compactLocked(); compactErr != nil {
|
|
return fmt.Errorf("compact state: %w", compactErr)
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
func (s *Store) Recent(limit int) []model.RunRecord {
|
|
s.mu.RLock()
|
|
defer s.mu.RUnlock()
|
|
if limit <= 0 || limit > len(s.runs) {
|
|
limit = len(s.runs)
|
|
}
|
|
out := make([]model.RunRecord, limit)
|
|
for i := 0; i < limit; i++ {
|
|
out[i] = s.runs[len(s.runs)-1-i]
|
|
}
|
|
return out
|
|
}
|
|
|
|
func (s *Store) FindRun(runID string) (model.RunRecord, bool) {
|
|
s.mu.RLock()
|
|
defer s.mu.RUnlock()
|
|
i, ok := s.runIndex[runID]
|
|
if !ok || i < 0 || i >= len(s.runs) {
|
|
return model.RunRecord{}, false
|
|
}
|
|
return s.runs[i], true
|
|
}
|
|
|
|
func marksTicketVersionProcessed(r model.RunRecord) bool {
|
|
if r.Outcome == "error" || r.Error != "" || strings.EqualFold(r.Trigger, "scheduled_escalation") {
|
|
return false
|
|
}
|
|
return r.SourceVersion != ""
|
|
}
|
|
|
|
func (s *Store) FindAnalysis(analysisID string) (model.AnalysisRun, bool) {
|
|
s.mu.RLock()
|
|
defer s.mu.RUnlock()
|
|
a, ok := s.analysisIndex[analysisID]
|
|
return a, ok
|
|
}
|
|
|
|
func (s *Store) HasEscalationKey(key string) bool {
|
|
s.mu.RLock()
|
|
defer s.mu.RUnlock()
|
|
_, ok := s.escalationKeys[strings.TrimSpace(key)]
|
|
return ok
|
|
}
|
|
|
|
func (s *Store) load() error {
|
|
f, err := os.Open(s.path)
|
|
if errors.Is(err, os.ErrNotExist) {
|
|
return nil
|
|
}
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer f.Close()
|
|
sc := bufio.NewScanner(f)
|
|
buf := make([]byte, 64*1024)
|
|
sc.Buffer(buf, 2*1024*1024)
|
|
for sc.Scan() {
|
|
var r model.RunRecord
|
|
if json.Unmarshal(sc.Bytes(), &r) == nil {
|
|
s.absorbDurableStateLocked(r)
|
|
s.runs = append(s.runs, r)
|
|
}
|
|
}
|
|
if len(s.runs) > s.maxRuns {
|
|
s.runs = s.runs[len(s.runs)-s.maxRuns:]
|
|
}
|
|
sort.SliceStable(s.runs, func(i, j int) bool { return s.runs[i].FinishedAt.Before(s.runs[j].FinishedAt) })
|
|
s.rebuildIndexesLocked()
|
|
return sc.Err()
|
|
}
|
|
|
|
func (s *Store) rebuildIndexesLocked() {
|
|
s.runIndex = make(map[string]int, len(s.runs))
|
|
s.analysisIndex = make(map[string]model.AnalysisRun)
|
|
for i, r := range s.runs {
|
|
s.runIndex[r.RunID] = i
|
|
for _, a := range r.Analyses {
|
|
s.analysisIndex[a.AnalysisID] = a
|
|
}
|
|
}
|
|
}
|
|
|
|
func (s *Store) LatestTicketRun(ticketID int64, excludedTriggers ...string) (model.RunRecord, bool) {
|
|
s.mu.RLock()
|
|
defer s.mu.RUnlock()
|
|
excluded := make(map[string]struct{}, len(excludedTriggers))
|
|
for _, trigger := range excludedTriggers {
|
|
excluded[strings.ToLower(strings.TrimSpace(trigger))] = struct{}{}
|
|
}
|
|
for i := len(s.runs) - 1; i >= 0; i-- {
|
|
r := s.runs[i]
|
|
if r.TicketID != ticketID {
|
|
continue
|
|
}
|
|
if _, skip := excluded[strings.ToLower(strings.TrimSpace(r.Trigger))]; skip {
|
|
continue
|
|
}
|
|
return r, true
|
|
}
|
|
return model.RunRecord{}, false
|
|
}
|
|
|
|
func (s *Store) absorbDurableStateLocked(r model.RunRecord) {
|
|
if marksTicketVersionProcessed(r) {
|
|
s.processedVersions[r.TicketID] = r.SourceVersion
|
|
}
|
|
for _, a := range r.Analyses {
|
|
if a.AnalysisType != "escalation" {
|
|
continue
|
|
}
|
|
// New escalation plans persist every successfully executed step independently.
|
|
// This keeps idempotency intact even when a later step in the same plan fails.
|
|
for _, step := range a.Action.Steps {
|
|
if !step.Executed {
|
|
continue
|
|
}
|
|
if key := escalationKeyFromResult(step.Result); key != "" {
|
|
s.escalationKeys[key] = struct{}{}
|
|
}
|
|
}
|
|
// Preserve compatibility with historical single-action audit records.
|
|
if a.Action.Executed {
|
|
if key := escalationKeyFromResult(a.Action.Result); key != "" {
|
|
s.escalationKeys[key] = struct{}{}
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
func escalationKeyFromResult(result string) string {
|
|
parts := map[string]string{}
|
|
for _, part := range strings.Split(result, ";") {
|
|
part = strings.TrimSpace(part)
|
|
for _, name := range []string{"ticket", "level", "action", "target"} {
|
|
prefix := name + "="
|
|
if strings.HasPrefix(part, prefix) && len(part) > len(prefix) {
|
|
parts[name] = part
|
|
}
|
|
}
|
|
}
|
|
if parts["ticket"] == "" || parts["level"] == "" {
|
|
return ""
|
|
}
|
|
key := parts["ticket"] + ";" + parts["level"]
|
|
if parts["action"] != "" {
|
|
key += ";" + parts["action"]
|
|
}
|
|
if parts["target"] != "" {
|
|
key += ";" + parts["target"]
|
|
}
|
|
return key
|
|
}
|
|
|
|
func (s *Store) loadDurableIndex() error {
|
|
b, err := os.ReadFile(s.indexPath)
|
|
if errors.Is(err, os.ErrNotExist) {
|
|
return nil
|
|
}
|
|
if err != nil {
|
|
return fmt.Errorf("read durable state index: %w", err)
|
|
}
|
|
var idx durableIndex
|
|
if err := json.Unmarshal(b, &idx); err != nil {
|
|
return fmt.Errorf("decode durable state index: %w", err)
|
|
}
|
|
for rawID, version := range idx.ProcessedVersions {
|
|
id, err := strconv.ParseInt(rawID, 10, 64)
|
|
if err == nil && id > 0 && version != "" {
|
|
s.processedVersions[id] = version
|
|
}
|
|
}
|
|
for _, key := range idx.EscalationKeys {
|
|
key = strings.TrimSpace(key)
|
|
if key != "" {
|
|
s.escalationKeys[key] = struct{}{}
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (s *Store) persistDurableIndexLocked() error {
|
|
idx := durableIndex{ProcessedVersions: make(map[string]string, len(s.processedVersions)), EscalationKeys: make([]string, 0, len(s.escalationKeys))}
|
|
for id, version := range s.processedVersions {
|
|
if id > 0 && version != "" {
|
|
idx.ProcessedVersions[strconv.FormatInt(id, 10)] = version
|
|
}
|
|
}
|
|
for key := range s.escalationKeys {
|
|
idx.EscalationKeys = append(idx.EscalationKeys, key)
|
|
}
|
|
sort.Strings(idx.EscalationKeys)
|
|
b, err := json.MarshalIndent(idx, "", " ")
|
|
if err != nil {
|
|
return err
|
|
}
|
|
tmp := s.indexPath + ".tmp"
|
|
f, err := os.OpenFile(tmp, os.O_CREATE|os.O_TRUNC|os.O_WRONLY, 0o640)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
cleanup := func() {
|
|
_ = f.Close()
|
|
_ = os.Remove(tmp)
|
|
}
|
|
if _, err := f.Write(append(b, '\n')); err != nil {
|
|
cleanup()
|
|
return err
|
|
}
|
|
if err := f.Sync(); err != nil {
|
|
cleanup()
|
|
return err
|
|
}
|
|
if err := f.Close(); err != nil {
|
|
_ = os.Remove(tmp)
|
|
return err
|
|
}
|
|
return os.Rename(tmp, s.indexPath)
|
|
}
|
|
|
|
func (s *Store) compactLocked() error {
|
|
tmp := s.path + ".tmp"
|
|
f, err := os.OpenFile(tmp, os.O_CREATE|os.O_TRUNC|os.O_WRONLY, 0o640)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
enc := json.NewEncoder(f)
|
|
for _, r := range s.runs {
|
|
if err := enc.Encode(r); err != nil {
|
|
f.Close()
|
|
_ = os.Remove(tmp)
|
|
return err
|
|
}
|
|
}
|
|
if err := f.Sync(); err != nil {
|
|
f.Close()
|
|
_ = os.Remove(tmp)
|
|
return err
|
|
}
|
|
if err := f.Close(); err != nil {
|
|
_ = os.Remove(tmp)
|
|
return err
|
|
}
|
|
return os.Rename(tmp, s.path)
|
|
}
|