Files
mesh-catalog/modules/docker-compose/cmd/docker-compose-tools/jobs.go
T
jochen 45befbbd02 docker-compose: the package, and its projects as tools (hq to-be 42 phase 2.9)
Compose and nothing else, for the two workstations, where it is already
installed by hand; the runtime, buildx and the group stay the docker
module's. Nine Go tools: projects (with their directories from the
containers' labels), ps, logs, a rendered config with secret-looking values
redacted, and up, down, restart and pull by directory or name. Acts run as
jobs inside the bundle, waited on for 18 s and followed with
docker_compose_job, because an up that pulls outlasts a call. down never
removes volumes.
2026-10-04 13:02:23 +02:00

152 lines
4.1 KiB
Go

package main
// jobs.go is the same file in the bundles whose acts can outlast one call (flatpak, docker-compose):
// an install or an `up` that pulls images takes minutes, and the runtime gives a call 30 s. Such an
// act is started as a job inside this process, waited on for a while, and answered either finished
// or with the job's id for the module's `_job` tool to follow. A job ends with this process: if the
// runtime restarts the bundle, a running job is cut off, and its id is then unknown.
import (
"fmt"
"sort"
"strings"
"sync"
"time"
)
// JobLimit is the longest a job may run; JobWait how long an act waits before answering a job id.
const (
JobLimit = 15 * time.Minute
JobWait = 18 * time.Second
keptJobs = 50
)
// Job is one long act, as its tool answers it.
type Job struct {
ID string `json:"job"`
Command string `json:"command"`
Started time.Time `json:"started"`
Finished *time.Time `json:"finished,omitempty"`
Running bool `json:"running"`
Status *int `json:"status,omitempty"`
Error string `json:"error,omitempty"`
Output string `json:"output,omitempty"`
Truncated bool `json:"truncated,omitempty"`
done chan struct{}
}
type jobBook struct {
mu sync.Mutex
seq int
jobs map[string]*Job
}
var jobs = &jobBook{jobs: map[string]*Job{}}
// startJob runs c in the background, held to JobLimit.
func startJob(c Cmd) *Job {
c.Timeout = JobLimit
name, args := argv(c)
jobs.mu.Lock()
jobs.seq++
j := &Job{ID: fmt.Sprintf("%d-%d", time.Now().Unix(), jobs.seq), Command: strings.TrimSpace(name + " " + strings.Join(args, " ")),
Started: time.Now().UTC(), Running: true, done: make(chan struct{})}
jobs.jobs[j.ID] = j
jobs.forgetOldest()
jobs.mu.Unlock()
go func() {
r := run(c)
var err error
if r.Status != 0 || r.Error != "" {
err = failure(c, r)
}
jobs.mu.Lock()
now := time.Now().UTC()
j.Finished, j.Running = &now, false
status := r.Status
j.Status = &status
if err != nil {
j.Error = err.Error()
}
j.Output = tail(strings.TrimSpace(r.Stdout+"\n"+r.Stderr), 16<<10)
j.Truncated = r.Truncated || len(r.Stdout)+len(r.Stderr) > 16<<10
jobs.mu.Unlock()
close(j.done)
}()
return j
}
// forgetOldest keeps the book bounded; finished jobs go first. Called with the lock held.
func (b *jobBook) forgetOldest() {
if len(b.jobs) <= keptJobs {
return
}
all := make([]*Job, 0, len(b.jobs))
for _, j := range b.jobs {
all = append(all, j)
}
sort.Slice(all, func(i, k int) bool { return all[i].Started.Before(all[k].Started) })
for _, j := range all {
if len(b.jobs) <= keptJobs {
return
}
if !j.Running {
delete(b.jobs, j.ID)
}
}
}
// awaitJob waits up to d for a job to finish and answers a copy of it as it then stands.
func awaitJob(j *Job, d time.Duration) Job {
select {
case <-j.done:
case <-time.After(d):
}
return snapshot(j)
}
func snapshot(j *Job) Job {
jobs.mu.Lock()
defer jobs.mu.Unlock()
c := *j
c.done = nil
return c
}
// jobByID answers a job by its id, or says it is not known to this process.
func jobByID(id string) (Job, error) {
jobs.mu.Lock()
j, ok := jobs.jobs[id]
jobs.mu.Unlock()
if !ok {
return Job{}, fmt.Errorf("no job %s in this process: it was never started here, was forgotten after %d newer ones, or the bundle has restarted since", id, keptJobs)
}
return snapshot(j), nil
}
// listJobs answers every job this process knows, newest first.
func listJobs() []Job {
jobs.mu.Lock()
all := make([]*Job, 0, len(jobs.jobs))
for _, j := range jobs.jobs {
all = append(all, j)
}
jobs.mu.Unlock()
sort.Slice(all, func(i, k int) bool { return all[i].Started.After(all[k].Started) })
out := make([]Job, 0, len(all))
for _, j := range all {
out = append(out, snapshot(j))
}
return out
}
// actAsJob starts c and answers the job once it finishes or JobWait passes, whichever is first.
// A finished job that failed is answered as an error, so a failed act is never read as success.
func actAsJob(c Cmd) (Job, error) {
j := awaitJob(startJob(c), JobWait)
if !j.Running && j.Error != "" {
return j, fmt.Errorf("%s (job %s)", j.Error, j.ID)
}
return j, nil
}