package storage import ( "bufio" "bytes" "context" "encoding/json" "fmt" "log" "os/exec" "sync" "time" "github.com/darroyo/nasctl/internal/db" ) type Event struct { Type string `json:"type"` Line string `json:"line,omitempty"` Status string `json:"status,omitempty"` ExitCode int `json:"exit_code,omitempty"` Message string `json:"message,omitempty"` } type JobManager struct { db *db.DB mu sync.Mutex running map[int64]context.CancelFunc subs map[int64]map[int]chan Event nextSub int maxLines int execSystem bool } func NewJobManager(database *db.DB, execSystem bool) *JobManager { return &JobManager{ db: database, running: make(map[int64]context.CancelFunc), subs: make(map[int64]map[int]chan Event), nextSub: 1, maxLines: 500, execSystem: execSystem, } } func (jm *JobManager) Ok() bool { return jm.execSystem } func (jm *JobManager) Subscribe(jobID int64) (<-chan Event, func()) { jm.mu.Lock() defer jm.mu.Unlock() if jm.subs[jobID] == nil { jm.subs[jobID] = make(map[int]chan Event) } ch := make(chan Event, 64) id := jm.nextSub jm.nextSub++ jm.subs[jobID][id] = ch unsubscribe := func() { jm.mu.Lock() defer jm.mu.Unlock() delete(jm.subs[jobID], id) close(ch) } return ch, unsubscribe } func (jm *JobManager) broadcast(jobID int64, ev Event) { jm.mu.Lock() defer jm.mu.Unlock() if subs, ok := jm.subs[jobID]; ok { for _, ch := range subs { select { case ch <- ev: default: } } } } func (jm *JobManager) Start(kind string, args map[string]any) (*db.StorageJob, error) { jm.mu.Lock() running, err := jm.db.GetRunningJob() if err != nil { jm.mu.Unlock() return nil, err } if running != nil { jm.mu.Unlock() return nil, fmt.Errorf("a job is already running (id=%d, kind=%s)", running.ID, running.Kind) } argsJSON, err := json.Marshal(args) if err != nil { jm.mu.Unlock() return nil, fmt.Errorf("marshal args: %w", err) } job, err := jm.db.CreateStorageJob(kind, string(argsJSON)) if err != nil { jm.mu.Unlock() return nil, err } jm.mu.Unlock() go jm.runJob(job) return &job, nil } func (jm *JobManager) runJob(job db.StorageJob) { ctx, cancel := context.WithCancel(context.Background()) jm.mu.Lock() jm.running[job.ID] = cancel jm.mu.Unlock() defer func() { jm.mu.Lock() delete(jm.running, job.ID) jm.mu.Unlock() }() var args []string var envcfg db.StorageConfig var err error // Build command based on kind switch job.Kind { case "mergerfs_preview", "mergerfs_move": envcfg, err = jm.db.GetStorageConfig() if err != nil { jm.failJob(job.ID, -1, fmt.Sprintf("get config: %v", err)) return } cfg := MoverConfig{ Source: envcfg.MoverSource, Dest: envcfg.MoverDest, CleanMacOS: envcfg.MoverCleanMacOS, RemoveSrc: envcfg.MoverRemoveSource, Inplace: envcfg.MoverInplace, DryRun: job.Kind == "mergerfs_preview", } if envcfg.MoverRsyncOptions != "" { flags, fe := ParseExtraRsyncFlags(envcfg.MoverRsyncOptions) if fe != nil { jm.failJob(job.ID, -1, fmt.Sprintf("invalid rsync flags: %v", fe)) return } cfg.ExtraFlags = flags } if err := ValidateMoverPaths(cfg); err != nil { jm.failJob(job.ID, -1, fmt.Sprintf("path validation: %v", err)) return } if job.Kind == "mergerfs_move" && cfg.CleanMacOS { if err := CleanMacOSJunk(cfg.Source); err != nil { jm.logLine(job.ID, fmt.Sprintf("warning: clean macos junk: %v", err)) } } args, err = BuildMoverArgs(cfg) if err != nil { jm.failJob(job.ID, -1, fmt.Sprintf("build rsync args: %v", err)) return } case "snapraid_diff", "snapraid_sync", "snapraid_scrub", "snapraid_check": envcfg, err = jm.db.GetStorageConfig() if err != nil { jm.failJob(job.ID, -1, fmt.Sprintf("get config: %v", err)) return } sc := SnapraidConfig{ Content: envcfg.SnapraidContent, DataDirs: ParseSnapraidDataDirs(envcfg.SnapraidDataDirs), ParityDir: envcfg.SnapraidParityDir, ScrubPlan: envcfg.SnapraidScrubPlan, } if err := ValidateSnapraidContent(sc.Content); err != nil { jm.failJob(job.ID, -1, fmt.Sprintf("content validation: %v", err)) return } if job.Kind == "snapraid_scrub" { if err := ValidateScrubPlan(sc.ScrubPlan); err != nil { jm.failJob(job.ID, -1, fmt.Sprintf("scrub plan validation: %v", err)) return } } args, err = BuildSnapraidArgs(job.Kind, sc) if err != nil { jm.failJob(job.ID, -1, fmt.Sprintf("build snapraid args: %v", err)) return } default: jm.failJob(job.ID, -1, fmt.Sprintf("unknown job kind: %s", job.Kind)) return } cmd := exec.CommandContext(ctx, args[0], args[1:]...) var outBuf, errBuf bytes.Buffer cmd.Stdout = &outBuf cmd.Stderr = &errBuf // Mark running and get PID jid := job.ID pid := 0 if cmd.Process != nil { pid = cmd.Process.Pid } _ = jm.db.MarkJobRunning(jid, pid) jm.broadcast(jid, Event{Type: "status", Status: "running"}) err = cmd.Run() // Collect all output combined := outBuf.String() + errBuf.String() scanner := bufio.NewScanner(bufio.NewReader(bytes.NewReader([]byte(combined)))) lineCount := 0 flushTicker := time.NewTicker(1 * time.Second) defer flushTicker.Stop() for scanner.Scan() { line := scanner.Text() lineCount++ jm.logLine(jid, line) if lineCount%25 == 0 { _ = jm.db.TrimJobOutput(jid, jm.maxLines) } } exitCode := -1 if err != nil { if exitErr, ok := err.(*exec.ExitError); ok { exitCode = exitErr.ExitCode() } else { jm.failJob(jid, exitCode, fmt.Sprintf("execution error: %v", err)) return } } _ = jm.db.TrimJobOutput(jid, jm.maxLines) status := "success" errMsg := "" if exitCode != 0 { status = "failed" errMsg = fmt.Sprintf("exit code %d", exitCode) } if err := jm.db.MarkJobFinished(jid, status, exitCode, combined, errMsg); err != nil { log.Printf("mark job finished: %v", err) } jm.broadcast(jid, Event{Type: "end", Status: status, ExitCode: exitCode}) } func (jm *JobManager) logLine(jobID int64, line string) { _ = jm.db.AppendJobOutput(jobID, line) jm.broadcast(jobID, Event{Type: "line", Line: line}) } func (jm *JobManager) failJob(jobID int64, exitCode int, msg string) { _ = jm.db.MarkJobFinished(jobID, "failed", exitCode, "", msg) jm.broadcast(jobID, Event{Type: "end", Status: "failed", ExitCode: exitCode, Message: msg}) } func (jm *JobManager) Get(id int64) (db.StorageJob, error) { return jm.db.GetStorageJob(id) } func (jm *JobManager) List(limit int) ([]db.StorageJob, error) { return jm.db.ListStorageJobs(limit) } func (jm *JobManager) Cancel(id int64) error { jm.mu.Lock() cancel, ok := jm.running[id] jm.mu.Unlock() if !ok { job, err := jm.db.GetStorageJob(id) if err != nil { return err } if job.Status == "queued" { return jm.db.SetJobCancelled(id) } return fmt.Errorf("job %d is not running", id) } cancel() _ = jm.db.SetJobCancelled(id) return nil } func (jm *JobManager) ResetOrphans() error { return jm.db.ResetOrphanedJobs() }