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
This commit is contained in:
@@ -0,0 +1,296 @@
|
||||
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()
|
||||
}
|
||||
Reference in New Issue
Block a user