27a52d2986
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
297 lines
6.9 KiB
Go
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()
|
|
}
|