5b459c0b2e
- Replace blocking outBuf/errBuf capture with StdoutPipe+StderrPipe + goroutines - scanAndEmit runs in parallel for stdout and stderr, emitting each line via SSE broadcast and DB append immediately (not after cmd.Run completes) - Custom scanner split: bytes.TrimRight strips trailing \r so snapraid progress lines (e.g. '2%, 57396 MB\r') are stored as clean lines - scanner.Buffer increased to 128KB max to handle large outputs - Use sync.WaitGroup to ensure both scanners finish before MarkJobFinished - Fix PID capture: cmd.Process is nil before cmd.Start(), pass 0 instead
314 lines
7.2 KiB
Go
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 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()
|
|
}
|