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 }