Files

145 lines
3.7 KiB
Go

package updater
import (
"context"
"time"
"github.com/shcizo/package-updater/internal/discovery"
"github.com/shcizo/package-updater/internal/metrics"
)
// Status describes the outcome of executing a single Job.
type Status string
const (
StatusUpdated Status = "updated"
StatusFailed Status = "failed"
StatusRefused Status = "refused"
StatusTimeout Status = "timeout"
)
// Result captures a per-Job outcome to be surfaced in the HTTP response.
type Result struct {
Job discovery.Job
Status Status
Error string
DurationMs int64
}
// Queue serialises Job execution through a single worker so we never
// run two `docker compose` commands against the same stack concurrently.
// It also satisfies the "single global FIFO worker" design choice in
// spec section 5.7.
type Queue struct {
exec Executor
metrics *metrics.Metrics
ch chan submission
stop chan struct{}
done chan struct{}
}
type submission struct {
ctx context.Context
jobs []discovery.Job
results chan []Result
}
// NewQueue constructs a queue bound to the given executor. m may be nil,
// in which case metrics recording is silently skipped.
func NewQueue(exec Executor, m *metrics.Metrics) *Queue {
return &Queue{
exec: exec,
metrics: m,
ch: make(chan submission, 16),
stop: make(chan struct{}),
done: make(chan struct{}),
}
}
// Start launches the single background worker. Call Stop to terminate.
func (q *Queue) Start(ctx context.Context) {
go q.run(ctx)
}
// Stop signals the worker to exit and waits for it to finish.
func (q *Queue) Stop() {
close(q.stop)
<-q.done
}
// Submit blocks until all jobs in the batch have been processed by the
// worker, then returns per-job results in the same order. The ctx
// timeout (if any) applies to each individual Job's execution and is
// surfaced as StatusTimeout.
func (q *Queue) Submit(ctx context.Context, jobs []discovery.Job) []Result {
resCh := make(chan []Result, 1)
q.ch <- submission{ctx: ctx, jobs: jobs, results: resCh}
if q.metrics != nil {
q.metrics.QueueDepth.Set(float64(len(q.ch)))
}
return <-resCh
}
func (q *Queue) run(_ context.Context) {
defer close(q.done)
for {
select {
case <-q.stop:
return
case s := <-q.ch:
results := make([]Result, len(s.jobs))
for i, j := range s.jobs {
results[i] = q.runOne(s.ctx, j)
}
s.results <- results
}
}
}
func (q *Queue) runOne(ctx context.Context, job discovery.Job) Result {
start := time.Now()
r := Result{Job: job}
if job.Refused {
r.Status = StatusRefused
r.Error = job.RefusedReason
r.DurationMs = time.Since(start).Milliseconds()
if q.metrics != nil {
q.metrics.UpdateJobs.WithLabelValues(job.Project, job.Service, string(StatusRefused)).Inc()
q.metrics.UpdateDuration.WithLabelValues(job.Project, job.Service).Observe(time.Since(start).Seconds())
}
return r
}
err := q.exec.Execute(ctx, job)
r.DurationMs = time.Since(start).Milliseconds()
switch {
case err == nil:
r.Status = StatusUpdated
case errorsIsContextDeadline(err) || ctxDeadlineExceeded(ctx):
r.Status = StatusTimeout
r.Error = err.Error()
default:
r.Status = StatusFailed
r.Error = err.Error()
}
if q.metrics != nil {
elapsed := time.Since(start).Seconds()
q.metrics.UpdateJobs.WithLabelValues(job.Project, job.Service, string(r.Status)).Inc()
q.metrics.UpdateDuration.WithLabelValues(job.Project, job.Service).Observe(elapsed)
if r.Status == StatusUpdated {
q.metrics.LastUpdateTime.WithLabelValues(job.Project, job.Service).SetToCurrentTime()
}
}
return r
}
func errorsIsContextDeadline(err error) bool {
return err == context.DeadlineExceeded
}
func ctxDeadlineExceeded(ctx context.Context) bool {
return ctx.Err() == context.DeadlineExceeded
}