package history import ( "bufio" "encoding/json" "errors" "fmt" "os" "path/filepath" "strings" "time" ) // InterfaceSample is one per-interface network sample at timestamp. type InterfaceSample struct { Interface string `json:"interface"` RxMbps float64 `json:"rx_mbps"` TxMbps float64 `json:"tx_mbps"` RxDrops uint64 `json:"rx_drops"` TxDrops uint64 `json:"tx_drops"` } // NetworkSnapshot stores one full sample containing many interfaces. type NetworkSnapshot struct { Timestamp time.Time `json:"timestamp"` Interfaces []InterfaceSample `json:"interfaces"` } // NetworkStore appends and reads network snapshots (JSONL) per cluster. type NetworkStore struct { dir string } func NewNetworkStore(baseDir string) *NetworkStore { if strings.TrimSpace(baseDir) == "" { baseDir = "." } return &NetworkStore{ dir: filepath.Join(baseDir, "history", "network"), } } func (s *NetworkStore) Path(clusterID string) string { return filepath.Join(s.dir, sanitizeClusterID(clusterID)+".jsonl") } func (s *NetworkStore) Append(clusterID string, snapshot NetworkSnapshot) error { if strings.TrimSpace(clusterID) == "" { return errors.New("empty cluster id") } if snapshot.Timestamp.IsZero() { snapshot.Timestamp = time.Now().UTC() } if len(snapshot.Interfaces) == 0 { return nil } if err := os.MkdirAll(s.dir, 0o700); err != nil { return fmt.Errorf("create history dir: %w", err) } path := s.Path(clusterID) f, err := os.OpenFile(path, os.O_CREATE|os.O_WRONLY|os.O_APPEND, 0o600) if err != nil { return fmt.Errorf("open history file: %w", err) } defer f.Close() line, err := json.Marshal(snapshot) if err != nil { return fmt.Errorf("marshal history snapshot: %w", err) } if _, err := f.Write(append(line, '\n')); err != nil { return fmt.Errorf("append history snapshot: %w", err) } return nil } func (s *NetworkStore) Load(clusterID string, since time.Time) ([]NetworkSnapshot, error) { if strings.TrimSpace(clusterID) == "" { return nil, errors.New("empty cluster id") } path := s.Path(clusterID) f, err := os.Open(path) if err != nil { if errors.Is(err, os.ErrNotExist) { return []NetworkSnapshot{}, nil } return nil, fmt.Errorf("open history file: %w", err) } defer f.Close() sc := bufio.NewScanner(f) sc.Buffer(make([]byte, 0, 1024*64), 1024*1024*8) out := make([]NetworkSnapshot, 0, 256) for sc.Scan() { line := strings.TrimSpace(sc.Text()) if line == "" { continue } var snap NetworkSnapshot if err := json.Unmarshal([]byte(line), &snap); err != nil { continue } if !since.IsZero() && snap.Timestamp.Before(since) { continue } if len(snap.Interfaces) == 0 { continue } out = append(out, snap) } if err := sc.Err(); err != nil { return nil, fmt.Errorf("scan history file: %w", err) } return out, nil } func sanitizeClusterID(v string) string { v = strings.TrimSpace(v) if v == "" { return "unknown" } var b strings.Builder for _, r := range v { switch { case r >= 'a' && r <= 'z': b.WriteRune(r) case r >= 'A' && r <= 'Z': b.WriteRune(r) case r >= '0' && r <= '9': b.WriteRune(r) case r == '-' || r == '_': b.WriteRune(r) default: b.WriteByte('_') } } out := strings.Trim(b.String(), "_") if out == "" { return "unknown" } return out }