Files
darroyo 56e5fabe8d feat(storage): paginate operations history with prev/next controls
Backend:
- ListStorageJobs(limit, offset) adds OFFSET for pagination
- CountStorageJobs() returns total row count for UI
- handler returns { jobs, total, limit, offset }
- JobManager.List(limit, offset) updated signature

Frontend:
- loadJobs(offset) with default 0
- Pagination UI: 'Mostrando X-Y de Z' + Anterior/Siguiente buttons
- After job start/end, reloads from offset 0
- listStorageJobs(limit, offset) API updated

Tests: fix List/ListStorageJobs calls to include offset=0
2026-07-07 09:27:06 -04:00

314 lines
7.2 KiB
Go

package storage
import (
"bufio"
"context"
"encoding/json"
"fmt"
"io"
"log"
"os/exec"
"strings"
"sync"
"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{
Conf: envcfg.SnapraidConf,
DataDirs: ParseSnapraidDataDirs(envcfg.SnapraidDataDirs),
ParityDir: envcfg.SnapraidParityDir,
ScrubPlan: envcfg.SnapraidScrubPlan,
}
if err := ValidateSnapraidConf(sc.Conf); err != nil {
jm.failJob(job.ID, -1, fmt.Sprintf("conf 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:]...)
stdout, err := cmd.StdoutPipe()
if err != nil {
jm.failJob(job.ID, -1, fmt.Sprintf("stdout pipe: %v", err))
return
}
stderr, err := cmd.StderrPipe()
if err != nil {
jm.failJob(job.ID, -1, fmt.Sprintf("stderr pipe: %v", err))
return
}
jid := job.ID
_ = jm.db.MarkJobRunning(jid, 0)
jm.broadcast(jid, Event{Type: "status", Status: "running"})
var wg sync.WaitGroup
wg.Add(2)
go jm.scanAndEmit(stdout, jid, &wg)
go jm.scanAndEmit(stderr, jid, &wg)
err = cmd.Start()
var exitCode int
if err != nil {
exitCode = -1
jm.failJob(jid, exitCode, fmt.Sprintf("start: %v", err))
wg.Wait()
return
}
_ = cmd.Wait()
wg.Wait()
_ = jm.db.TrimJobOutput(jid, jm.maxLines)
status := "success"
errMsg := ""
if exitErr, ok := err.(*exec.ExitError); ok {
exitCode = exitErr.ExitCode()
if exitCode != 0 {
status = "failed"
errMsg = fmt.Sprintf("exit code %d", exitCode)
}
} else if err != nil {
status = "failed"
errMsg = fmt.Sprintf("execution error: %v", err)
exitCode = -1
}
if err := jm.db.MarkJobFinished(jid, status, exitCode, "", 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) scanAndEmit(r io.Reader, jobID int64, wg *sync.WaitGroup) {
defer wg.Done()
scanner := bufio.NewScanner(r)
scanner.Buffer(make([]byte, 1024), 128*1024)
lineCount := 0
for scanner.Scan() {
line := scanner.Text()
line = strings.TrimRight(line, "\r")
lineCount++
jm.logLine(jobID, line)
if lineCount%25 == 0 {
_ = jm.db.TrimJobOutput(jobID, jm.maxLines)
}
}
}
func (jm *JobManager) Get(id int64) (db.StorageJob, error) {
return jm.db.GetStorageJob(id)
}
func (jm *JobManager) List(limit, offset int) ([]db.StorageJob, error) {
return jm.db.ListStorageJobs(limit, offset)
}
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()
}