From 5b459c0b2eadeeb45ee1e8d4ec3a1ba1b69cf69f Mon Sep 17 00:00:00 2001 From: Daniel Arroyo Date: Tue, 7 Jul 2026 09:07:08 -0400 Subject: [PATCH] fix(storage): stream job output in real-time with StdoutPipe/StderrPipe - 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 --- Makefile | 2 +- internal/storage/jobs.go | 89 ++++++++++++++++++++++++---------------- 2 files changed, 54 insertions(+), 37 deletions(-) diff --git a/Makefile b/Makefile index 30e2f8f..590f52a 100644 --- a/Makefile +++ b/Makefile @@ -1,5 +1,5 @@ BINARY=nasctl -VERSION?=0.8.1 +VERSION?=0.8.2 GO?=go LDFLAGS=-s -w -X github.com/darroyo/nasctl/internal/web.Version=$(VERSION) -X github.com/darroyo/nasctl/internal/web.Commit=$(shell git rev-parse --short HEAD 2>/dev/null || echo unknown) BUILD_FLAGS=CGO_ENABLED=0 diff --git a/internal/storage/jobs.go b/internal/storage/jobs.go index 1444a4f..2cdfd3b 100644 --- a/internal/storage/jobs.go +++ b/internal/storage/jobs.go @@ -2,14 +2,14 @@ package storage import ( "bufio" - "bytes" "context" "encoding/json" "fmt" + "io" "log" "os/exec" + "strings" "sync" - "time" "github.com/darroyo/nasctl/internal/db" ) @@ -197,57 +197,58 @@ func (jm *JobManager) runJob(job db.StorageJob) { } 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 + stdout, err := cmd.StdoutPipe() + if err != nil { + jm.failJob(job.ID, -1, fmt.Sprintf("stdout pipe: %v", err)) + return } - _ = jm.db.MarkJobRunning(jid, pid) + 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"}) - err = cmd.Run() + var wg sync.WaitGroup + wg.Add(2) - // 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() + go jm.scanAndEmit(stdout, jid, &wg) + go jm.scanAndEmit(stderr, jid, &wg) - for scanner.Scan() { - line := scanner.Text() - lineCount++ - jm.logLine(jid, line) - if lineCount%25 == 0 { - _ = jm.db.TrimJobOutput(jid, jm.maxLines) - } - } + err = cmd.Start() - exitCode := 0 + var exitCode int 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 - } + 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 exitCode != 0 { + 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("exit code %d", exitCode) + errMsg = fmt.Sprintf("execution error: %v", err) + exitCode = -1 } - if err := jm.db.MarkJobFinished(jid, status, exitCode, combined, errMsg); err != nil { + if err := jm.db.MarkJobFinished(jid, status, exitCode, "", errMsg); err != nil { log.Printf("mark job finished: %v", err) } @@ -264,6 +265,22 @@ func (jm *JobManager) failJob(jobID int64, exitCode int, msg string) { 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) }