Files
baby-nas/internal/storage/jobs.go
T
darroyo 27a52d2986 feat: add mergerfs mover and snapraid integration
New 'Almacenamiento' page with:
- Auto-detection of rsync, mergerfs, snapraid binaries and mergerfs mount
- Configurable pool settings (source/dest, macOS cleanup, rsync flags)
- Mergerfs mover with dry-run preview and live SSE output streaming
- SnapRAID diff/sync/scrub/check with live SSE output
- Async job system (1 concurrent job) with SSE streaming
- Job history table

Backend:
- internal/storage/ package with capabilities, mergerfs, snapraid, jobs
- storage_config and storage_jobs DB tables (migration 0008)
- GET/PUT /api/storage/config, GET /api/storage/capabilities
- POST/GET /api/storage/jobs, GET /api/storage/jobs/{id}/stream
- Storage operations disabled when NASCTL_EXEC_SYSTEM=false

Closes #new-feature
2026-07-06 23:23:22 -04:00

297 lines
6.9 KiB
Go

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()
}