From 85e1baf59ade79c621f5676995dea03e1af2c531 Mon Sep 17 00:00:00 2001 From: mrhid6 Date: Mon, 20 Jul 2026 15:18:35 +0100 Subject: [PATCH] feat(server): workflow log-writer registry, retention setting, sweeper --- server/cmd/main.go | 2 + server/internal/api/handlers.go | 7 +- server/internal/models/settings.go | 2 + server/internal/services/settings.go | 21 ++- server/internal/services/steplogs.go | 210 +++++++++++++++++++++++++++ 5 files changed, 237 insertions(+), 5 deletions(-) create mode 100644 server/internal/services/steplogs.go diff --git a/server/cmd/main.go b/server/cmd/main.go index 6f71e21..f5b12c2 100644 --- a/server/cmd/main.go +++ b/server/cmd/main.go @@ -31,6 +31,8 @@ func main() { log.Printf("warning: failed to ensure workflow indexes: %v", err) } + services.StartLogSweeper() + redisAddr := getEnv("REDIS_ADDR", "localhost:6379") if err := auth.InitRedis(redisAddr); err != nil { log.Fatalf("failed to connect to Redis: %v", err) diff --git a/server/internal/api/handlers.go b/server/internal/api/handlers.go index 03b182c..52a0174 100644 --- a/server/internal/api/handlers.go +++ b/server/internal/api/handlers.go @@ -450,14 +450,15 @@ func getSettings(c *gin.Context) { func saveSettings(c *gin.Context) { var body struct { - Alerts models.AlertSettings `json:"alerts"` - Email models.EmailSettings `json:"email"` + Alerts models.AlertSettings `json:"alerts"` + Email models.EmailSettings `json:"email"` + WorkflowLogRetentionDays *int `json:"workflow_log_retention_days"` } if err := c.ShouldBindJSON(&body); err != nil { c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()}) return } - if err := services.SaveSettings(body.Alerts, body.Email); err != nil { + if err := services.SaveSettings(body.Alerts, body.Email, body.WorkflowLogRetentionDays); err != nil { c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) return } diff --git a/server/internal/models/settings.go b/server/internal/models/settings.go index c7e25a9..60bdaf2 100644 --- a/server/internal/models/settings.go +++ b/server/internal/models/settings.go @@ -36,4 +36,6 @@ type Settings struct { Alerts AlertSettings `bson:"alerts" json:"alerts"` Email EmailSettings `bson:"email" json:"email"` Secrets SecretsSettings `bson:"secrets" json:"secrets"` + // WorkflowLogRetentionDays: nil = default 30, 0 = keep forever, N = N days. + WorkflowLogRetentionDays *int `bson:"workflow_log_retention_days,omitempty" json:"workflow_log_retention_days,omitempty"` } diff --git a/server/internal/services/settings.go b/server/internal/services/settings.go index 149eafd..9895192 100644 --- a/server/internal/services/settings.go +++ b/server/internal/services/settings.go @@ -100,7 +100,7 @@ func VerifySecretsReadToken(token string) bool { return subtle.ConstantTimeCompare(expected, got[:]) == 1 } -func SaveSettings(alerts models.AlertSettings, email models.EmailSettings) error { +func SaveSettings(alerts models.AlertSettings, email models.EmailSettings, retentionDays *int) error { ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) defer cancel() @@ -111,14 +111,31 @@ func SaveSettings(alerts models.AlertSettings, email models.EmailSettings) error email.SMTPPort = 587 } + set := bson.M{"alerts": alerts, "email": email} + if retentionDays != nil { + set["workflow_log_retention_days"] = *retentionDays + } _, err := db.Col("settings").UpdateOne(ctx, bson.M{}, - bson.M{"$set": bson.M{"alerts": alerts, "email": email}}, + bson.M{"$set": set}, options.UpdateOne().SetUpsert(true), ) return err } +// GetWorkflowLogRetentionDays returns the log retention in days: 30 when unset, +// 0 for keep-forever, or the configured value. +func GetWorkflowLogRetentionDays() (int, error) { + s, err := GetSettings() + if err != nil { + return 30, err + } + if s.WorkflowLogRetentionDays == nil { + return 30, nil + } + return *s.WorkflowLogRetentionDays, nil +} + func SendOfflineWebhook(webhookURL, hostname, serverID, ipAddress string) { payload := map[string]any{ "event": "server.offline", diff --git a/server/internal/services/steplogs.go b/server/internal/services/steplogs.go new file mode 100644 index 0000000..f70d8df --- /dev/null +++ b/server/internal/services/steplogs.go @@ -0,0 +1,210 @@ +package services + +import ( + "os" + "path/filepath" + "strings" + "sync" + "time" + + "github.com/mrhid6/vantage/server/internal/db" + "go.mongodb.org/mongo-driver/v2/bson" +) + +// WorkflowLogDir returns the base directory for workflow step logs, creating it. +func WorkflowLogDir() string { + dir := os.Getenv("VANTAGE_WORKFLOW_LOG_DIR") + if dir == "" { + dir = filepath.Join("data", "workflow-logs") + } + _ = os.MkdirAll(dir, 0700) + return dir +} + +// ServerRunLogPath is the per-server-run log file path. +func ServerRunLogPath(runID, serverID string) string { + return filepath.Join(WorkflowLogDir(), runID, serverID+".log") +} + +// AppendMarker appends a line to the server-run log and returns the byte offset +// at which the write began (used as a step's log_offset). +func AppendMarker(runID, serverID, line string) (int64, error) { + path := ServerRunLogPath(runID, serverID) + if err := os.MkdirAll(filepath.Dir(path), 0700); err != nil { + return 0, err + } + f, err := os.OpenFile(path, os.O_APPEND|os.O_CREATE|os.O_WRONLY, 0600) + if err != nil { + return 0, err + } + defer f.Close() + off, _ := f.Seek(0, 2) // current end = offset before write + if _, err := f.WriteString(line); err != nil { + return off, err + } + return off, nil +} + +// ---- streamed chunk writer, boundary-safe secret masking ---- + +type stepLogWriter struct { + mu sync.Mutex + f *os.File + carry []byte + secrets []string + maxSecret int +} + +type stepLogRegistry struct { + mu sync.Mutex + writers map[string]*stepLogWriter +} + +var StepLogs = &stepLogRegistry{writers: make(map[string]*stepLogWriter)} + +// Open opens (append) the server-run file for a step's streamed chunks. +func (r *stepLogRegistry) Open(commandID, path string, secrets []string) error { + if err := os.MkdirAll(filepath.Dir(path), 0700); err != nil { + return err + } + f, err := os.OpenFile(path, os.O_APPEND|os.O_CREATE|os.O_WRONLY, 0600) + if err != nil { + return err + } + max := 0 + for _, s := range secrets { + if len(s) > max { + max = len(s) + } + } + w := &stepLogWriter{f: f, secrets: secrets, maxSecret: max} + r.mu.Lock() + r.writers[commandID] = w + r.mu.Unlock() + return nil +} + +func (r *stepLogRegistry) get(commandID string) *stepLogWriter { + r.mu.Lock() + defer r.mu.Unlock() + return r.writers[commandID] +} + +// Append masks and writes a chunk, holding back the last maxSecret-1 bytes so a +// secret split across a chunk boundary is still masked on the next append/close. +func (r *stepLogRegistry) Append(commandID string, data []byte) { + w := r.get(commandID) + if w == nil { + return + } + w.mu.Lock() + defer w.mu.Unlock() + if len(w.secrets) == 0 || w.maxSecret <= 1 { + _, _ = w.f.Write(data) + return + } + buf := append(w.carry, data...) + hold := w.maxSecret - 1 + if len(buf) <= hold { + w.carry = buf + return + } + flush := buf[:len(buf)-hold] + w.carry = append([]byte{}, buf[len(buf)-hold:]...) + _, _ = w.f.Write(maskBytes(flush, w.secrets)) +} + +// Close flushes the carry (masked) and closes the file. +func (r *stepLogRegistry) Close(commandID string) { + r.mu.Lock() + w := r.writers[commandID] + delete(r.writers, commandID) + r.mu.Unlock() + if w == nil { + return + } + w.mu.Lock() + defer w.mu.Unlock() + if len(w.carry) > 0 { + _, _ = w.f.Write(maskBytes(w.carry, w.secrets)) + w.carry = nil + } + _ = w.f.Close() +} + +func maskBytes(b []byte, secrets []string) []byte { + s := string(b) + for _, v := range secrets { + if v == "" { + continue + } + s = strings.ReplaceAll(s, v, "***") + } + return []byte(s) +} + +// ---- retention sweeper ---- + +// StartLogSweeper sweeps expired run-log dirs hourly (and once now). +func StartLogSweeper() { + go func() { + sweepLogs() + t := time.NewTicker(time.Hour) + defer t.Stop() + for range t.C { + sweepLogs() + } + }() +} + +func sweepLogs() { + days := retentionDays() + if days <= 0 { + return + } + cutoff := time.Now().AddDate(0, 0, -days) + base := WorkflowLogDir() + entries, err := os.ReadDir(base) + if err != nil { + return + } + for _, e := range entries { + if !e.IsDir() { + continue + } + runID := e.Name() + dir := filepath.Join(base, runID) + if runExpired(runID, dir, cutoff) { + _ = os.RemoveAll(dir) + } + } +} + +// runExpired is true when the run finished before cutoff (falling back to dir +// mtime when the run doc is gone). +func runExpired(runID, dir string, cutoff time.Time) bool { + ctx, cancel := wfCtx() + defer cancel() + var run struct { + FinishedAt *time.Time `bson:"finished_at"` + } + err := db.Col("workflow_runs").FindOne(ctx, bson.M{"run_id": runID}).Decode(&run) + if err == nil { + if run.FinishedAt == nil { + return false // still running / never finished — keep + } + return run.FinishedAt.Before(cutoff) + } + // run doc gone: use dir mtime + if fi, e := os.Stat(dir); e == nil { + return fi.ModTime().Before(cutoff) + } + return false +} + +func retentionDays() int { + if v, err := GetWorkflowLogRetentionDays(); err == nil { + return v + } + return 30 +}