package web import ( "context" "sync" "gitgud.io/mike/mpv-manager/pkg/installer" ) // jobOutputRelay owns the channels used by one CommandRunner. Installers write // to the channels synchronously; the relay drains them into the job record and // is closed only after the installer call returns. Waiting for both drains // before a terminal job transition prevents tail output from being lost when // CompleteJob/SetError archives and removes the active job. type jobOutputRelay struct { output chan string errors chan error wg sync.WaitGroup once sync.Once } func newJobCommandRunner(ctx context.Context, manager *JobManager, jobID string, bufferSize int) (*installer.CommandRunner, *jobOutputRelay) { if bufferSize < 1 { bufferSize = 1 } relay := &jobOutputRelay{ output: make(chan string, bufferSize), errors: make(chan error, bufferSize), } relay.wg.Add(2) go func() { defer relay.wg.Done() for output := range relay.output { manager.AddOutput(jobID, output) } }() go func() { defer relay.wg.Done() for err := range relay.errors { if err != nil { manager.AddOutput(jobID, "ERROR: "+err.Error()) } } }() runner := installer.NewCommandRunnerWithContext(ctx, relay.output, relay.errors) runner.SetInteractiveSudoAllowed(false) runner.SetCommitGuard(func() error { return manager.BeginCommit(jobID) }) return runner, relay } // Close declares that the synchronous installer call has returned, drains all // buffered events, and waits for the relay goroutines to exit. It is idempotent // so callers can both defer it for panic safety and invoke it before setting a // terminal job state. func (r *jobOutputRelay) Close() { if r == nil { return } r.once.Do(func() { close(r.output) close(r.errors) r.wg.Wait() }) }