add: processing controls — start/hold, pause/resume, retry
Three lifecycle controls on the dashboard, sharing one /api surface and the status feed: - Start/hold gate (#11): the watcher tracks the queue but holds processing until POST /api/start. Held by default; AV1DAE_AUTOSTART=1 restores start- on-boot. NOTE: flips the previous auto-start default. - Pause/resume the active encode (#10): SIGSTOP/SIGCONT on the ffmpeg process (libsvtav1 is in-process, so one signal suspends all its threads). Resumes exactly where it left off; the tracker excludes paused time from elapsed. - Retry a failed file (#3): POST /api/retry?file=NAME moves it from failed/ back to input/, with a base-name guard against path traversal. /status now reports running + the failed list; snapshots carry a paused flag. Verified live: held queue, retry move, traversal -> 400, and a real encode suspending to process state T on pause and S on resume. Closes #3 Closes #10 Closes #11
This commit is contained in:
+31
-2
@@ -91,8 +91,29 @@ func main() {
|
|||||||
ctx, cancel := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
|
ctx, cancel := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
|
||||||
defer cancel()
|
defer cancel()
|
||||||
|
|
||||||
|
// Start/hold gate: held by default so the user clicks Start; AV1DAE_AUTOSTART
|
||||||
|
// (truthy) restores start-on-boot.
|
||||||
|
autostart := envTruthy(os.Getenv("AV1DAE_AUTOSTART"))
|
||||||
|
w := watcher.New(cfg.Paths.Input, 15, autostart)
|
||||||
|
if !autostart {
|
||||||
|
log.Info("Queue held on startup — click Start (set AV1DAE_AUTOSTART=1 to auto-start)")
|
||||||
|
}
|
||||||
|
|
||||||
|
controls := server.Controls{
|
||||||
|
Running: w.Running,
|
||||||
|
SetRunning: w.SetRunning,
|
||||||
|
Pause: enc.Pause,
|
||||||
|
Resume: enc.Resume,
|
||||||
|
RetryFailed: func(name string) error {
|
||||||
|
if name == "" || name != filepath.Base(name) {
|
||||||
|
return fmt.Errorf("invalid file name")
|
||||||
|
}
|
||||||
|
return mover.Rename(filepath.Join(cfg.Paths.Failed, name), filepath.Join(cfg.Paths.Input, name))
|
||||||
|
},
|
||||||
|
}
|
||||||
|
|
||||||
if addr := *cfg.HTTPAddr; addr != "" {
|
if addr := *cfg.HTTPAddr; addr != "" {
|
||||||
srv := &http.Server{Addr: addr, Handler: server.New(tracker, log, store, cfg.Paths.Input).Handler()}
|
srv := &http.Server{Addr: addr, Handler: server.New(tracker, log, store, controls, cfg.Paths.Input, cfg.Paths.Failed).Handler()}
|
||||||
go func() {
|
go func() {
|
||||||
log.Info(fmt.Sprintf("Status server listening on %s", addr))
|
log.Info(fmt.Sprintf("Status server listening on %s", addr))
|
||||||
if err := srv.ListenAndServe(); err != nil && err != http.ErrServerClosed {
|
if err := srv.ListenAndServe(); err != nil && err != http.ErrServerClosed {
|
||||||
@@ -107,7 +128,6 @@ func main() {
|
|||||||
}()
|
}()
|
||||||
}
|
}
|
||||||
|
|
||||||
w := watcher.New(cfg.Paths.Input, 15)
|
|
||||||
w.Start(ctx, func(ctx context.Context, inputPath string) error {
|
w.Start(ctx, func(ctx context.Context, inputPath string) error {
|
||||||
return processFile(ctx, inputPath, cfg, enc, metaClient, log, tracker, store)
|
return processFile(ctx, inputPath, cfg, enc, metaClient, log, tracker, store)
|
||||||
})
|
})
|
||||||
@@ -269,6 +289,15 @@ func failToFailed(path, failedDir string, log *logger.Logger) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// envTruthy reports whether an env var is set to a truthy value.
|
||||||
|
func envTruthy(s string) bool {
|
||||||
|
switch strings.ToLower(strings.TrimSpace(s)) {
|
||||||
|
case "1", "true", "yes", "on":
|
||||||
|
return true
|
||||||
|
}
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
|
||||||
// toStatusStreams maps probed source streams to the display shape, keeping only
|
// toStatusStreams maps probed source streams to the display shape, keeping only
|
||||||
// audio and subtitle streams (the video stream isn't shown).
|
// audio and subtitle streams (the video stream isn't shown).
|
||||||
func toStatusStreams(streams []encoder.StreamMetadata) []status.Stream {
|
func toStatusStreams(streams []encoder.StreamMetadata) []status.Stream {
|
||||||
|
|||||||
+4
-1
@@ -12,7 +12,10 @@ services:
|
|||||||
cpus: "8.0"
|
cpus: "8.0"
|
||||||
memory: 4g
|
memory: 4g
|
||||||
ports:
|
ports:
|
||||||
- "8080:8080" # status server: GET http://host:8080/status (needs http_addr: ":8080" in config)
|
- "8080:8080" # status server + dashboard at http://host:8080 (needs http_addr: ":8080" in config)
|
||||||
|
# By default the queue is HELD on boot — click Start in the UI. Uncomment to auto-start:
|
||||||
|
# environment:
|
||||||
|
# - AV1DAE_AUTOSTART=1
|
||||||
volumes:
|
volumes:
|
||||||
- ./data/config.yaml:/config/config.yaml:ro # your config — paths inside must point at /data/*
|
- ./data/config.yaml:/config/config.yaml:ro # your config — paths inside must point at /data/*
|
||||||
- ./data/media:/data # holds input/ output/ originals/ failed/ work/ + logs
|
- ./data/media:/data # holds input/ output/ originals/ failed/ work/ + logs
|
||||||
|
|||||||
@@ -5,12 +5,15 @@ import (
|
|||||||
"context"
|
"context"
|
||||||
"encoding/json"
|
"encoding/json"
|
||||||
"fmt"
|
"fmt"
|
||||||
|
"os"
|
||||||
"os/exec"
|
"os/exec"
|
||||||
"path/filepath"
|
"path/filepath"
|
||||||
"regexp"
|
"regexp"
|
||||||
"sort"
|
"sort"
|
||||||
"strconv"
|
"strconv"
|
||||||
"strings"
|
"strings"
|
||||||
|
"sync"
|
||||||
|
"syscall"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"av1dae/internal/logger"
|
"av1dae/internal/logger"
|
||||||
@@ -26,6 +29,59 @@ type Encoder struct {
|
|||||||
opusencPath string
|
opusencPath string
|
||||||
log *logger.Logger
|
log *logger.Logger
|
||||||
tracker *status.Tracker
|
tracker *status.Tracker
|
||||||
|
|
||||||
|
procMu sync.Mutex
|
||||||
|
cur *os.Process // the in-flight video encode, for pause/resume
|
||||||
|
}
|
||||||
|
|
||||||
|
// Pause freezes the active video encode in place via SIGSTOP. libsvtav1 runs in
|
||||||
|
// the ffmpeg process (no forked children), so one signal suspends all its
|
||||||
|
// threads. No-op if nothing is encoding. The process keeps its memory and
|
||||||
|
// partial output and resumes exactly where it left off.
|
||||||
|
func (e *Encoder) Pause() error {
|
||||||
|
e.procMu.Lock()
|
||||||
|
defer e.procMu.Unlock()
|
||||||
|
if e.cur == nil {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
if err := e.cur.Signal(syscall.SIGSTOP); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
if e.tracker != nil {
|
||||||
|
e.tracker.SetPaused(true)
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// Resume thaws a paused encode via SIGCONT. No-op if nothing is encoding.
|
||||||
|
func (e *Encoder) Resume() error {
|
||||||
|
e.procMu.Lock()
|
||||||
|
defer e.procMu.Unlock()
|
||||||
|
if e.cur == nil {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
if err := e.cur.Signal(syscall.SIGCONT); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
if e.tracker != nil {
|
||||||
|
e.tracker.SetPaused(false)
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (e *Encoder) setProc(p *os.Process) {
|
||||||
|
e.procMu.Lock()
|
||||||
|
e.cur = p
|
||||||
|
e.procMu.Unlock()
|
||||||
|
}
|
||||||
|
|
||||||
|
func (e *Encoder) clearProc() {
|
||||||
|
e.procMu.Lock()
|
||||||
|
e.cur = nil
|
||||||
|
e.procMu.Unlock()
|
||||||
|
if e.tracker != nil {
|
||||||
|
e.tracker.SetPaused(false)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
type StreamInfo struct {
|
type StreamInfo struct {
|
||||||
@@ -94,6 +150,9 @@ func (e *Encoder) runCmdProgress(ctx context.Context, label, file, name string,
|
|||||||
if err := cmd.Start(); err != nil {
|
if err := cmd.Start(); err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
|
// Publish the process so Pause/Resume can signal it; clear on exit.
|
||||||
|
e.setProc(cmd.Process)
|
||||||
|
defer e.clearProc()
|
||||||
|
|
||||||
// Reads stdout to EOF (when ffmpeg exits), so Wait below is safe afterwards.
|
// Reads stdout to EOF (when ffmpeg exits), so Wait below is safe afterwards.
|
||||||
var lastLog time.Time
|
var lastLog time.Time
|
||||||
|
|||||||
@@ -114,6 +114,22 @@
|
|||||||
.day:first-child { margin-top:0; }
|
.day:first-child { margin-top:0; }
|
||||||
|
|
||||||
@media (prefers-reduced-motion:reduce){ .dot.ok{animation:none;} .bar.indet>i{animation:none; width:100%; opacity:.4;} }
|
@media (prefers-reduced-motion:reduce){ .dot.ok{animation:none;} .bar.indet>i{animation:none; width:100%; opacity:.4;} }
|
||||||
|
|
||||||
|
/* control bar */
|
||||||
|
.controls { display:flex; align-items:center; gap:.6rem; flex-wrap:wrap; margin-bottom:1.25rem; min-height:1px; }
|
||||||
|
.ctl { font-family:var(--mono); font-size:.82rem; cursor:pointer; border-radius:7px; padding:.45rem .9rem;
|
||||||
|
background:var(--panel-2); color:var(--text); border:1px solid var(--line); }
|
||||||
|
.ctl:hover { border-color:var(--amber); color:var(--amber); }
|
||||||
|
.ctl.primary { background:var(--amber); color:#1a1200; border-color:var(--amber); font-weight:600; }
|
||||||
|
.ctl.primary:hover { background:var(--amber-soft); color:#1a1200; }
|
||||||
|
.ctl.small { font-size:.72rem; padding:.25rem .65rem; }
|
||||||
|
.ctl-state { font-family:var(--mono); font-size:.72rem; color:var(--muted); margin-left:.3rem; text-transform:uppercase; letter-spacing:.08em; }
|
||||||
|
.ctl-state.paused { color:var(--amber); }
|
||||||
|
.count { font-family:var(--mono); color:var(--red); }
|
||||||
|
/* failed list */
|
||||||
|
#failed .frow { display:flex; align-items:center; gap:.6rem; padding:.35rem 0; border-bottom:1px solid var(--line); font-family:var(--mono); font-size:.82rem; }
|
||||||
|
#failed .frow:last-child { border-bottom:0; }
|
||||||
|
#failed .frow .fn { flex:1; color:var(--text); word-break:break-word; }
|
||||||
</style>
|
</style>
|
||||||
</head>
|
</head>
|
||||||
<body>
|
<body>
|
||||||
@@ -124,10 +140,17 @@
|
|||||||
<span class="live"><span class="dot" id="dot"></span><span id="livetext">connecting</span></span>
|
<span class="live"><span class="dot" id="dot"></span><span id="livetext">connecting</span></span>
|
||||||
</header>
|
</header>
|
||||||
|
|
||||||
|
<div class="controls" id="controls"></div>
|
||||||
|
|
||||||
<section class="card" id="job">
|
<section class="card" id="job">
|
||||||
<div id="job-body"></div>
|
<div id="job-body"></div>
|
||||||
</section>
|
</section>
|
||||||
|
|
||||||
|
<section class="card" id="failedCard" hidden>
|
||||||
|
<p class="eyebrow">Failed <span class="count" id="fcount"></span></p>
|
||||||
|
<div id="failed"></div>
|
||||||
|
</section>
|
||||||
|
|
||||||
<div class="grid">
|
<div class="grid">
|
||||||
<section class="card">
|
<section class="card">
|
||||||
<p class="eyebrow">Queue</p>
|
<p class="eyebrow">Queue</p>
|
||||||
@@ -356,17 +379,60 @@
|
|||||||
txt.textContent = state === "ok" ? "live" : state === "idle" ? "idle" : state === "err" ? "offline" : "connecting";
|
txt.textContent = state === "ok" ? "live" : state === "idle" ? "idle" : state === "err" ? "offline" : "connecting";
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// ---- runtime controls: start/hold, pause/resume, retry ----
|
||||||
|
async function post(path) {
|
||||||
|
try { await fetch(path, { method: "POST" }); } catch (e) {}
|
||||||
|
poll(); // reflect the new state immediately
|
||||||
|
}
|
||||||
|
|
||||||
|
function renderControls(d) {
|
||||||
|
const c = d.current || {};
|
||||||
|
const encoding = c.file && c.phase !== "idle";
|
||||||
|
let html = d.running
|
||||||
|
? `<button class="ctl" data-act="/api/hold">⏸ Hold queue</button>`
|
||||||
|
: `<button class="ctl primary" data-act="/api/start">▶ Start queue</button>`;
|
||||||
|
if (encoding) {
|
||||||
|
html += c.paused
|
||||||
|
? `<button class="ctl primary" data-act="/api/resume">▶ Resume encode</button>`
|
||||||
|
: `<button class="ctl" data-act="/api/pause">⏸ Pause encode</button>`;
|
||||||
|
}
|
||||||
|
const state = !d.running ? "held" : c.paused ? "running · encode paused" : "running";
|
||||||
|
html += `<span class="ctl-state${c.paused ? " paused" : ""}">${state}</span>`;
|
||||||
|
$("controls").innerHTML = html;
|
||||||
|
}
|
||||||
|
|
||||||
|
function renderFailed(failed) {
|
||||||
|
const card = $("failedCard");
|
||||||
|
if (!failed || !failed.length) { card.hidden = true; $("failed").innerHTML = ""; return; }
|
||||||
|
card.hidden = false;
|
||||||
|
$("fcount").textContent = "(" + failed.length + ")";
|
||||||
|
$("failed").innerHTML = failed.map(f =>
|
||||||
|
`<div class="frow"><span class="fn">${esc(f)}</span><button class="ctl small" data-retry="${esc(f)}">retry</button></div>`
|
||||||
|
).join("");
|
||||||
|
}
|
||||||
|
|
||||||
|
$("controls").addEventListener("click", e => {
|
||||||
|
const b = e.target.closest("[data-act]");
|
||||||
|
if (b) post(b.getAttribute("data-act"));
|
||||||
|
});
|
||||||
|
$("failed").addEventListener("click", e => {
|
||||||
|
const b = e.target.closest("[data-retry]");
|
||||||
|
if (b) post("/api/retry?file=" + encodeURIComponent(b.getAttribute("data-retry")));
|
||||||
|
});
|
||||||
|
|
||||||
async function poll() {
|
async function poll() {
|
||||||
try {
|
try {
|
||||||
const r = await fetch("status", { cache: "no-store" });
|
const r = await fetch("status", { cache: "no-store" });
|
||||||
if (!r.ok) throw new Error(r.status);
|
if (!r.ok) throw new Error(r.status);
|
||||||
const d = await r.json();
|
const d = await r.json();
|
||||||
pushSpeed(d.current);
|
pushSpeed(d.current);
|
||||||
|
renderControls(d);
|
||||||
renderJob(d.current);
|
renderJob(d.current);
|
||||||
|
renderFailed(d.failed);
|
||||||
renderQueue(d.queue, d.current);
|
renderQueue(d.queue, d.current);
|
||||||
renderEvents(d.recent);
|
renderEvents(d.recent);
|
||||||
const active = d.current && d.current.file && d.current.phase !== "idle";
|
const active = d.current && d.current.file && d.current.phase !== "idle";
|
||||||
setLive(active ? "ok" : "idle");
|
setLive(d.current && d.current.paused ? "idle" : active ? "ok" : "idle");
|
||||||
} catch (e) {
|
} catch (e) {
|
||||||
setLive("err");
|
setLive("err");
|
||||||
} finally {
|
} finally {
|
||||||
|
|||||||
+94
-13
@@ -20,15 +20,27 @@ var indexHTML []byte
|
|||||||
//go:embed settings.html
|
//go:embed settings.html
|
||||||
var settingsHTML []byte
|
var settingsHTML []byte
|
||||||
|
|
||||||
type Server struct {
|
// Controls is the runtime-control surface the UI drives, wired in main from the
|
||||||
tracker *status.Tracker
|
// watcher (start/hold gate), encoder (pause/resume), and mover (retry).
|
||||||
log *logger.Logger
|
type Controls struct {
|
||||||
store *settings.Store
|
Running func() bool
|
||||||
inputDir string
|
SetRunning func(bool)
|
||||||
|
Pause func() error
|
||||||
|
Resume func() error
|
||||||
|
RetryFailed func(name string) error
|
||||||
}
|
}
|
||||||
|
|
||||||
func New(tracker *status.Tracker, log *logger.Logger, store *settings.Store, inputDir string) *Server {
|
type Server struct {
|
||||||
return &Server{tracker: tracker, log: log, store: store, inputDir: inputDir}
|
tracker *status.Tracker
|
||||||
|
log *logger.Logger
|
||||||
|
store *settings.Store
|
||||||
|
controls Controls
|
||||||
|
inputDir string
|
||||||
|
failedDir string
|
||||||
|
}
|
||||||
|
|
||||||
|
func New(tracker *status.Tracker, log *logger.Logger, store *settings.Store, controls Controls, inputDir, failedDir string) *Server {
|
||||||
|
return &Server{tracker: tracker, log: log, store: store, controls: controls, inputDir: inputDir, failedDir: failedDir}
|
||||||
}
|
}
|
||||||
|
|
||||||
// Handler returns the mux for all status routes. Phase 3 adds "/" (the HTML
|
// Handler returns the mux for all status routes. Phase 3 adds "/" (the HTML
|
||||||
@@ -38,10 +50,72 @@ func (s *Server) Handler() http.Handler {
|
|||||||
mux.HandleFunc("/status", s.handleStatus)
|
mux.HandleFunc("/status", s.handleStatus)
|
||||||
mux.HandleFunc("/settings", s.handleSettingsPage)
|
mux.HandleFunc("/settings", s.handleSettingsPage)
|
||||||
mux.HandleFunc("/api/settings", s.handleAPISettings)
|
mux.HandleFunc("/api/settings", s.handleAPISettings)
|
||||||
|
mux.HandleFunc("/api/start", s.gateHandler(true))
|
||||||
|
mux.HandleFunc("/api/hold", s.gateHandler(false))
|
||||||
|
mux.HandleFunc("/api/pause", s.actionHandler(func() error { return s.controls.Pause() }, "Encode paused"))
|
||||||
|
mux.HandleFunc("/api/resume", s.actionHandler(func() error { return s.controls.Resume() }, "Encode resumed"))
|
||||||
|
mux.HandleFunc("/api/retry", s.handleRetry)
|
||||||
mux.HandleFunc("/", s.handleIndex)
|
mux.HandleFunc("/", s.handleIndex)
|
||||||
return mux
|
return mux
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// gateHandler flips the start/hold gate. POST only.
|
||||||
|
func (s *Server) gateHandler(run bool) http.HandlerFunc {
|
||||||
|
return func(w http.ResponseWriter, r *http.Request) {
|
||||||
|
if r.Method != http.MethodPost {
|
||||||
|
methodNotAllowed(w)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
s.controls.SetRunning(run)
|
||||||
|
if run {
|
||||||
|
s.log.Info("Queue started via web UI")
|
||||||
|
} else {
|
||||||
|
s.log.Info("Queue held via web UI")
|
||||||
|
}
|
||||||
|
w.WriteHeader(http.StatusNoContent)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// actionHandler wraps a no-arg control action (pause/resume). POST only.
|
||||||
|
func (s *Server) actionHandler(fn func() error, logMsg string) http.HandlerFunc {
|
||||||
|
return func(w http.ResponseWriter, r *http.Request) {
|
||||||
|
if r.Method != http.MethodPost {
|
||||||
|
methodNotAllowed(w)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
if err := fn(); err != nil {
|
||||||
|
http.Error(w, err.Error(), http.StatusInternalServerError)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
s.log.Info(logMsg + " via web UI")
|
||||||
|
w.WriteHeader(http.StatusNoContent)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// handleRetry moves a named file from failed/ back to input/. POST ?file=NAME.
|
||||||
|
func (s *Server) handleRetry(w http.ResponseWriter, r *http.Request) {
|
||||||
|
if r.Method != http.MethodPost {
|
||||||
|
methodNotAllowed(w)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
name := r.URL.Query().Get("file")
|
||||||
|
if name == "" {
|
||||||
|
http.Error(w, "missing file", http.StatusBadRequest)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
if err := s.controls.RetryFailed(name); err != nil {
|
||||||
|
http.Error(w, err.Error(), http.StatusBadRequest)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
s.log.Info("Retry requested via web UI: " + name)
|
||||||
|
w.WriteHeader(http.StatusNoContent)
|
||||||
|
}
|
||||||
|
|
||||||
|
func methodNotAllowed(w http.ResponseWriter) {
|
||||||
|
w.Header().Set("Allow", "POST")
|
||||||
|
http.Error(w, "method not allowed", http.StatusMethodNotAllowed)
|
||||||
|
}
|
||||||
|
|
||||||
func (s *Server) handleSettingsPage(w http.ResponseWriter, r *http.Request) {
|
func (s *Server) handleSettingsPage(w http.ResponseWriter, r *http.Request) {
|
||||||
w.Header().Set("Content-Type", "text/html; charset=utf-8")
|
w.Header().Set("Content-Type", "text/html; charset=utf-8")
|
||||||
_, _ = w.Write(settingsHTML)
|
_, _ = w.Write(settingsHTML)
|
||||||
@@ -83,17 +157,14 @@ func (s *Server) handleIndex(w http.ResponseWriter, r *http.Request) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
type statusResponse struct {
|
type statusResponse struct {
|
||||||
|
Running bool `json:"running"`
|
||||||
Current status.Snapshot `json:"current"`
|
Current status.Snapshot `json:"current"`
|
||||||
Queue []string `json:"queue"`
|
Queue []string `json:"queue"`
|
||||||
|
Failed []string `json:"failed"`
|
||||||
Recent []logger.LogEntry `json:"recent"`
|
Recent []logger.LogEntry `json:"recent"`
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *Server) handleStatus(w http.ResponseWriter, r *http.Request) {
|
func (s *Server) handleStatus(w http.ResponseWriter, r *http.Request) {
|
||||||
queue := []string{}
|
|
||||||
for _, f := range watcher.InputFiles(s.inputDir) {
|
|
||||||
queue = append(queue, filepath.Base(f))
|
|
||||||
}
|
|
||||||
|
|
||||||
recent, err := s.log.RecentLogs(50)
|
recent, err := s.log.RecentLogs(50)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
http.Error(w, "reading logs", http.StatusInternalServerError)
|
http.Error(w, "reading logs", http.StatusInternalServerError)
|
||||||
@@ -102,8 +173,18 @@ func (s *Server) handleStatus(w http.ResponseWriter, r *http.Request) {
|
|||||||
|
|
||||||
w.Header().Set("Content-Type", "application/json")
|
w.Header().Set("Content-Type", "application/json")
|
||||||
_ = json.NewEncoder(w).Encode(statusResponse{
|
_ = json.NewEncoder(w).Encode(statusResponse{
|
||||||
|
Running: s.controls.Running(),
|
||||||
Current: s.tracker.Snapshot(),
|
Current: s.tracker.Snapshot(),
|
||||||
Queue: queue,
|
Queue: baseNames(watcher.InputFiles(s.inputDir)),
|
||||||
|
Failed: baseNames(watcher.InputFiles(s.failedDir)),
|
||||||
Recent: recent,
|
Recent: recent,
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func baseNames(paths []string) []string {
|
||||||
|
out := []string{}
|
||||||
|
for _, p := range paths {
|
||||||
|
out = append(out, filepath.Base(p))
|
||||||
|
}
|
||||||
|
return out
|
||||||
|
}
|
||||||
|
|||||||
+46
-11
@@ -43,16 +43,19 @@ type Stream struct {
|
|||||||
// Tracker holds live progress for the active job. Safe for concurrent use:
|
// Tracker holds live progress for the active job. Safe for concurrent use:
|
||||||
// the encode goroutine writes, HTTP/log readers call Snapshot.
|
// the encode goroutine writes, HTTP/log readers call Snapshot.
|
||||||
type Tracker struct {
|
type Tracker struct {
|
||||||
mu sync.RWMutex
|
mu sync.RWMutex
|
||||||
file string
|
file string
|
||||||
phase string
|
phase string
|
||||||
meta *JobMeta
|
meta *JobMeta
|
||||||
streams []Stream
|
streams []Stream
|
||||||
totalSec float64 // source duration; 0 until known
|
totalSec float64 // source duration; 0 until known
|
||||||
outTime float64 // encoded position in seconds
|
outTime float64 // encoded position in seconds
|
||||||
fps float64
|
fps float64
|
||||||
speed float64
|
speed float64
|
||||||
startedAt time.Time
|
startedAt time.Time
|
||||||
|
paused bool
|
||||||
|
pausedAt time.Time // when the current pause began
|
||||||
|
pausedTotal time.Duration // accumulated paused time this job
|
||||||
}
|
}
|
||||||
|
|
||||||
// Snapshot is an immutable view of the tracker for readers.
|
// Snapshot is an immutable view of the tracker for readers.
|
||||||
@@ -64,6 +67,7 @@ type Snapshot struct {
|
|||||||
Percent float64 `json:"percent"`
|
Percent float64 `json:"percent"`
|
||||||
FPS float64 `json:"fps"`
|
FPS float64 `json:"fps"`
|
||||||
Speed float64 `json:"speed"`
|
Speed float64 `json:"speed"`
|
||||||
|
Paused bool `json:"paused"`
|
||||||
ElapsedSec int `json:"elapsed_sec"`
|
ElapsedSec int `json:"elapsed_sec"`
|
||||||
ETASec int `json:"eta_sec"`
|
ETASec int `json:"eta_sec"`
|
||||||
StartedAt time.Time `json:"started_at"`
|
StartedAt time.Time `json:"started_at"`
|
||||||
@@ -83,6 +87,26 @@ func (t *Tracker) Begin(file string) {
|
|||||||
t.streams = nil
|
t.streams = nil
|
||||||
t.totalSec, t.outTime, t.fps, t.speed = 0, 0, 0, 0
|
t.totalSec, t.outTime, t.fps, t.speed = 0, 0, 0, 0
|
||||||
t.startedAt = time.Now()
|
t.startedAt = time.Now()
|
||||||
|
t.paused = false
|
||||||
|
t.pausedAt = time.Time{}
|
||||||
|
t.pausedTotal = 0
|
||||||
|
}
|
||||||
|
|
||||||
|
// SetPaused records pause/resume transitions so elapsed time excludes the
|
||||||
|
// paused span. Idempotent on repeated same-state calls.
|
||||||
|
func (t *Tracker) SetPaused(p bool) {
|
||||||
|
t.mu.Lock()
|
||||||
|
defer t.mu.Unlock()
|
||||||
|
if p == t.paused {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
if p {
|
||||||
|
t.pausedAt = time.Now()
|
||||||
|
} else if !t.pausedAt.IsZero() {
|
||||||
|
t.pausedTotal += time.Since(t.pausedAt)
|
||||||
|
t.pausedAt = time.Time{}
|
||||||
|
}
|
||||||
|
t.paused = p
|
||||||
}
|
}
|
||||||
|
|
||||||
func (t *Tracker) SetPhase(p string) {
|
func (t *Tracker) SetPhase(p string) {
|
||||||
@@ -129,6 +153,9 @@ func (t *Tracker) Idle() {
|
|||||||
t.streams = nil
|
t.streams = nil
|
||||||
t.totalSec, t.outTime, t.fps, t.speed = 0, 0, 0, 0
|
t.totalSec, t.outTime, t.fps, t.speed = 0, 0, 0, 0
|
||||||
t.startedAt = time.Time{}
|
t.startedAt = time.Time{}
|
||||||
|
t.paused = false
|
||||||
|
t.pausedAt = time.Time{}
|
||||||
|
t.pausedTotal = 0
|
||||||
}
|
}
|
||||||
|
|
||||||
func (t *Tracker) Snapshot() Snapshot {
|
func (t *Tracker) Snapshot() Snapshot {
|
||||||
@@ -141,10 +168,18 @@ func (t *Tracker) Snapshot() Snapshot {
|
|||||||
Streams: t.streams,
|
Streams: t.streams,
|
||||||
FPS: t.fps,
|
FPS: t.fps,
|
||||||
Speed: t.speed,
|
Speed: t.speed,
|
||||||
|
Paused: t.paused,
|
||||||
StartedAt: t.startedAt,
|
StartedAt: t.startedAt,
|
||||||
}
|
}
|
||||||
if !t.startedAt.IsZero() {
|
if !t.startedAt.IsZero() {
|
||||||
s.ElapsedSec = int(time.Since(t.startedAt).Seconds())
|
elapsed := time.Since(t.startedAt) - t.pausedTotal
|
||||||
|
if t.paused && !t.pausedAt.IsZero() {
|
||||||
|
elapsed -= time.Since(t.pausedAt)
|
||||||
|
}
|
||||||
|
if elapsed < 0 {
|
||||||
|
elapsed = 0
|
||||||
|
}
|
||||||
|
s.ElapsedSec = int(elapsed.Seconds())
|
||||||
}
|
}
|
||||||
if t.totalSec > 0 {
|
if t.totalSec > 0 {
|
||||||
s.Percent = t.outTime / t.totalSec * 100
|
s.Percent = t.outTime / t.totalSec * 100
|
||||||
|
|||||||
@@ -81,6 +81,27 @@ func TestSnapshotPercentAndETA(t *testing.T) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func TestSetPausedFlag(t *testing.T) {
|
||||||
|
tr := New()
|
||||||
|
tr.Begin("x.mkv")
|
||||||
|
if tr.Snapshot().Paused {
|
||||||
|
t.Fatal("should not start paused")
|
||||||
|
}
|
||||||
|
tr.SetPaused(true)
|
||||||
|
if !tr.Snapshot().Paused {
|
||||||
|
t.Error("expected paused after SetPaused(true)")
|
||||||
|
}
|
||||||
|
tr.SetPaused(true) // idempotent — must not double-count
|
||||||
|
tr.SetPaused(false)
|
||||||
|
if tr.Snapshot().Paused {
|
||||||
|
t.Error("expected not paused after SetPaused(false)")
|
||||||
|
}
|
||||||
|
tr.Idle()
|
||||||
|
if tr.Snapshot().Paused {
|
||||||
|
t.Error("Idle should clear paused")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
func TestSnapshotIdleNoDivByZero(t *testing.T) {
|
func TestSnapshotIdleNoDivByZero(t *testing.T) {
|
||||||
s := New().Snapshot() // no Begin, totalSec 0
|
s := New().Snapshot() // no Begin, totalSec 0
|
||||||
if s.Percent != 0 || s.ETASec != 0 || s.Phase != PhaseIdle {
|
if s.Percent != 0 || s.ETASec != 0 || s.Phase != PhaseIdle {
|
||||||
|
|||||||
@@ -6,6 +6,7 @@ import (
|
|||||||
"os"
|
"os"
|
||||||
"path/filepath"
|
"path/filepath"
|
||||||
"sort"
|
"sort"
|
||||||
|
"sync"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"av1dae/pkg/types"
|
"av1dae/pkg/types"
|
||||||
@@ -59,6 +60,10 @@ type failureRecord struct {
|
|||||||
type Watcher struct {
|
type Watcher struct {
|
||||||
inputDir string
|
inputDir string
|
||||||
interval time.Duration
|
interval time.Duration
|
||||||
|
// runMu guards running: the start/hold gate. While held, settled files are
|
||||||
|
// tracked (so the queue is reported) but not handed to processFn.
|
||||||
|
runMu sync.RWMutex
|
||||||
|
running bool
|
||||||
// seen tracks the last-observed (mtime, size) for every file currently in
|
// seen tracks the last-observed (mtime, size) for every file currently in
|
||||||
// the input dir. A file is only handed to processFn once two consecutive
|
// the input dir. A file is only handed to processFn once two consecutive
|
||||||
// ticks agree on both fields (partial-write protection).
|
// ticks agree on both fields (partial-write protection).
|
||||||
@@ -68,7 +73,7 @@ type Watcher struct {
|
|||||||
failed map[string]failureRecord
|
failed map[string]failureRecord
|
||||||
}
|
}
|
||||||
|
|
||||||
func New(inputDir string, intervalSeconds int) *Watcher {
|
func New(inputDir string, intervalSeconds int, autostart bool) *Watcher {
|
||||||
interval := time.Duration(intervalSeconds) * time.Second
|
interval := time.Duration(intervalSeconds) * time.Second
|
||||||
if interval < 10*time.Second {
|
if interval < 10*time.Second {
|
||||||
interval = 10 * time.Second
|
interval = 10 * time.Second
|
||||||
@@ -79,11 +84,26 @@ func New(inputDir string, intervalSeconds int) *Watcher {
|
|||||||
return &Watcher{
|
return &Watcher{
|
||||||
inputDir: inputDir,
|
inputDir: inputDir,
|
||||||
interval: interval,
|
interval: interval,
|
||||||
|
running: autostart,
|
||||||
seen: make(map[string]fileStat),
|
seen: make(map[string]fileStat),
|
||||||
failed: make(map[string]failureRecord),
|
failed: make(map[string]failureRecord),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// SetRunning toggles the start/hold gate.
|
||||||
|
func (w *Watcher) SetRunning(v bool) {
|
||||||
|
w.runMu.Lock()
|
||||||
|
w.running = v
|
||||||
|
w.runMu.Unlock()
|
||||||
|
}
|
||||||
|
|
||||||
|
// Running reports whether the queue is being processed.
|
||||||
|
func (w *Watcher) Running() bool {
|
||||||
|
w.runMu.RLock()
|
||||||
|
defer w.runMu.RUnlock()
|
||||||
|
return w.running
|
||||||
|
}
|
||||||
|
|
||||||
func (w *Watcher) Start(ctx context.Context, processFn func(context.Context, string) error) {
|
func (w *Watcher) Start(ctx context.Context, processFn func(context.Context, string) error) {
|
||||||
ticker := time.NewTicker(w.interval)
|
ticker := time.NewTicker(w.interval)
|
||||||
defer ticker.Stop()
|
defer ticker.Stop()
|
||||||
@@ -146,6 +166,12 @@ func (w *Watcher) scanAndProcess(ctx context.Context, processFn func(context.Con
|
|||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
|
|
||||||
|
if !w.Running() {
|
||||||
|
// Held: the file stays in the queue (already in nextSeen) and is
|
||||||
|
// picked up once the user starts processing.
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
|
||||||
if err := processFn(ctx, file); err != nil {
|
if err := processFn(ctx, file); err != nil {
|
||||||
fmt.Printf("Error processing %s: %v\n", file, err)
|
fmt.Printf("Error processing %s: %v\n", file, err)
|
||||||
nextFailed[file] = failureRecord{
|
nextFailed[file] = failureRecord{
|
||||||
|
|||||||
Reference in New Issue
Block a user