Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
26 changes: 26 additions & 0 deletions changelog.d/fixed/go-context-propagation-sweep.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,26 @@
- Go context-propagation sweep: callers that previously dropped a
`context.Context` at the subprocess / SQLite boundary now forward it via
`exec.CommandContext` / `(*sql.DB).ExecContext`. Affected sites:
- `pkg/libvmaf/scorer` — added `ScoreContext(ctx, …)` that wraps the
`vmaf` CLI invocation in `exec.CommandContext`. `Score(…)` is kept as a
backwards-compatible wrapper around `ScoreContext(context.Background(), …)`
and marked deprecated.
Comment on lines +4 to +7
- `cmd/vmafx-controller/{grpc_server,http_server}.go` and
`cmd/vmafx-server/{grpc_server,http_server}.go` — `/v1/score` and the
gRPC `Score` RPC now forward `r.Context()` / the gRPC context into
`ScoreContext`, so a client disconnect or graceful-shutdown signal aborts
the underlying `vmaf` subprocess instead of leaving a zombie.
- `cmd/vmafx-node/executor.go` — the node executor forwards the job's
`ctx` into `ScoreContext` so a controller-side cancellation reaches the
running scorer.
Comment on lines +11 to +15
- `cmd/vmafx-controller/queue/queue.go` — `Submit`, `PullWork`,
`ReportResult`, and `Cancel` now use `db.ExecContext` with the caller's
`ctx` (previously they took `_ context.Context` and dropped it).
- `cmd/vmafx-node/probe/probe.go` — `EncoderInventory(ctx, ffmpegBin)`
now accepts and forwards a context; `cmd/vmafx-node/main.go` binds the
startup probe to a 30 s timeout so a hung `ffmpeg` cannot stall node
boot.
- `cmd/vmafx-mcp/impl.go` — the `vmaf_score`, `probe_backend`,
`eval_model_on_split`, `compare_models`, and `describe_worst_frames`
handlers now plumb the MCP tool-call context through `runVmafScore` /
`delegateToPythonEval` and into `exec.CommandContext`.
25 changes: 16 additions & 9 deletions cmd/vmafx-controller/queue/queue.go
Original file line number Diff line number Diff line change
Expand Up @@ -214,7 +214,7 @@ func (q *SQLiteQueue) reload() error {

// Submit enqueues a new job. The job's ID is assigned here (UUID v4) and
// written to both SQLite and the in-memory FIFO.
func (q *SQLiteQueue) Submit(_ context.Context, job *Job) (string, error) {
func (q *SQLiteQueue) Submit(ctx context.Context, job *Job) (string, error) {
job.ID = uuid.New().String()
job.Status = StatusPending
now := time.Now()
Expand All @@ -226,7 +226,9 @@ func (q *SQLiteQueue) Submit(_ context.Context, job *Job) (string, error) {
return "", fmt.Errorf("queue: marshal scoring params: %w", err)
}

_, err = q.db.Exec(
// ExecContext propagates the caller's ctx so a cancelled gRPC SubmitJob
// can abort the INSERT instead of holding the SQLite write open.
_, err = q.db.ExecContext(ctx,
"INSERT INTO jobs (id, status, scoring, created_at, updated_at) VALUES (?,?,?,?,?)",
job.ID, StatusPending, string(scoringJSON), now.Unix(), now.Unix(),
)
Expand All @@ -245,7 +247,7 @@ func (q *SQLiteQueue) Submit(_ context.Context, job *Job) (string, error) {
// PullWork atomically dequeues the oldest PENDING job whose backend requirement
// (if any) is satisfied by the requesting node's capabilities. Returns
// (nil, nil) when no matching job is available.
func (q *SQLiteQueue) PullWork(_ context.Context, nodeID string, capacity NodeCapacity) (*Job, error) {
func (q *SQLiteQueue) PullWork(ctx context.Context, nodeID string, capacity NodeCapacity) (*Job, error) {
q.mu.Lock()
defer q.mu.Unlock()

Expand Down Expand Up @@ -287,9 +289,10 @@ func (q *SQLiteQueue) PullWork(_ context.Context, nodeID string, capacity NodeCa
// Remove from FIFO.
q.pendingFIFO = append(q.pendingFIFO[:matchIdx], q.pendingFIFO[matchIdx+1:]...)

// Transition to RUNNING in SQLite.
// Transition to RUNNING in SQLite. ExecContext propagates the caller's
// ctx so an aborted PullWork RPC does not leave the UPDATE in flight.
now := time.Now().Unix()
_, err := q.db.Exec(
_, err := q.db.ExecContext(ctx,
"UPDATE jobs SET status=?, assigned_node=?, updated_at=? WHERE id=?",
StatusRunning, nodeID, now, matchID,
)
Expand Down Expand Up @@ -345,7 +348,7 @@ func (q *SQLiteQueue) rollbackTopending(jobID string) error {

// ReportResult records the terminal outcome of a job. If result.Err is
// non-empty the job is marked FAILED; otherwise COMPLETED.
func (q *SQLiteQueue) ReportResult(_ context.Context, jobID string, result *JobResult) error {
func (q *SQLiteQueue) ReportResult(ctx context.Context, jobID string, result *JobResult) error {
status := StatusCompleted
if result.Err != "" {
status = StatusFailed
Expand All @@ -356,8 +359,10 @@ func (q *SQLiteQueue) ReportResult(_ context.Context, jobID string, result *JobR
featuresJSON = []byte("{}")
}

// ExecContext propagates the caller's ctx so the node's ReportResult RPC
// deadline / cancellation aborts the UPDATE cleanly.
now := time.Now().Unix()
_, err = q.db.Exec(
_, err = q.db.ExecContext(ctx,
"UPDATE jobs SET status=?, score=?, features=?, error=?, updated_at=? WHERE id=?",
status, result.Score, string(featuresJSON), result.Err, now, jobID,
)
Expand Down Expand Up @@ -427,9 +432,11 @@ func (q *SQLiteQueue) getUnlocked(jobID string) (*Job, error) {

// Cancel marks a PENDING or RUNNING job as CANCELLED. Returns nil if the job
// was already in a terminal state (idempotent).
func (q *SQLiteQueue) Cancel(_ context.Context, jobID string) error {
func (q *SQLiteQueue) Cancel(ctx context.Context, jobID string) error {
// ExecContext propagates the caller's ctx so an aborted CancelJob RPC
// does not leave the UPDATE in flight.
now := time.Now().Unix()
res, err := q.db.Exec(
res, err := q.db.ExecContext(ctx,
"UPDATE jobs SET status=?, updated_at=? WHERE id=? AND status IN (?,?)",
StatusCancelled, now, jobID, StatusPending, StatusRunning,
)
Expand Down
9 changes: 7 additions & 2 deletions cmd/vmafx-node/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@ import (
"os"
"os/signal"
"syscall"
"time"

"github.com/VMAFx/vmafx/cmd/vmafx-node/probe"
"github.com/VMAFx/vmafx/cmd/vmafx-node/server"
Expand Down Expand Up @@ -68,8 +69,12 @@ func run() error {
ffmpegBin := ffmpegPath()
slog.Info("ffmpeg discovery", "path", ffmpegBin)

// Startup probe: enumerate available encoders and cache.
encoders, err := probe.EncoderInventory(ffmpegBin)
// Startup probe: enumerate available encoders and cache. Bound the probe
// with a short timeout so a hung ffmpeg binary cannot stall node startup
// indefinitely.
probeCtx, probeCancel := context.WithTimeout(context.Background(), 30*time.Second)
encoders, err := probe.EncoderInventory(probeCtx, ffmpegBin)
probeCancel()
if err != nil {
// Non-fatal: node can still serve; encoders that are unavailable
// will fail at job-dispatch time with a clear error.
Expand Down
9 changes: 7 additions & 2 deletions cmd/vmafx-node/probe/probe.go
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@
package probe

import (
"context"
"fmt"
"log/slog"
"os/exec"
Expand Down Expand Up @@ -49,8 +50,12 @@ var expectedSoftwareCodecs = []string{

// EncoderInventory runs `ffmpegBin -encoders` and returns the parsed inventory.
// It also logs a WARN for each expected software codec that is absent.
func EncoderInventory(ffmpegBin string) (*Inventory, error) {
out, err := exec.Command(ffmpegBin, "-hide_banner", "-encoders").Output()
//
// ctx is forwarded to exec.CommandContext so that a cancelled caller (timeout,
// SIGINT during startup) terminates the ffmpeg subprocess instead of leaking
// it. Pass context.Background() if no caller-side context is available.
func EncoderInventory(ctx context.Context, ffmpegBin string) (*Inventory, error) {
out, err := exec.CommandContext(ctx, ffmpegBin, "-hide_banner", "-encoders").Output()
Comment on lines +57 to +58
if err != nil {
return nil, fmt.Errorf("ffmpeg -encoders: %w", err)
}
Expand Down
3 changes: 2 additions & 1 deletion cmd/vmafx-node/probe/probe_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@
package probe_test

import (
"context"
"os"
"os/exec"
"testing"
Expand Down Expand Up @@ -53,7 +54,7 @@ func TestEncoderInventory_RealBinary(t *testing.T) {
t.Skipf("ffmpeg not found on PATH (%s): %v", bin, err)
}

inv, err := probe.EncoderInventory(bin)
inv, err := probe.EncoderInventory(context.Background(), bin)
if err != nil {
t.Fatalf("EncoderInventory(%q): %v", bin, err)
}
Expand Down
Loading