forked from vieter-v/vieter
feat(build): allowed invalidating entries in build queue
parent
3611123f45
commit
882a9a60a9
|
@ -8,6 +8,8 @@ import util
|
|||
|
||||
struct BuildJob {
|
||||
pub:
|
||||
// Time at which this build job was created/queued
|
||||
created time.Time
|
||||
// Next timestamp from which point this job is allowed to be executed
|
||||
timestamp time.Time
|
||||
// Required for calculating next timestamp after having pop'ed a job
|
||||
|
@ -22,6 +24,8 @@ fn (r1 BuildJob) < (r2 BuildJob) bool {
|
|||
return r1.timestamp < r2.timestamp
|
||||
}
|
||||
|
||||
// The build job queue is responsible for managing the list of scheduled builds
|
||||
// for each architecture. Agents receive jobs from this queue.
|
||||
pub struct BuildJobQueue {
|
||||
// Schedule to use for targets without explicitely defined cron expression
|
||||
default_schedule CronExpression
|
||||
|
@ -31,15 +35,17 @@ mut:
|
|||
mutex shared util.Dummy
|
||||
// For each architecture, a priority queue is tracked
|
||||
queues map[string]MinHeap<BuildJob>
|
||||
// Each queued build job is also stored in a map, with the keys being the
|
||||
// target IDs. This is used when removing or editing targets.
|
||||
// jobs map[int]BuildJob
|
||||
// When a target is removed from the server or edited, its previous build
|
||||
// configs will be invalid. This map allows for those to be simply skipped
|
||||
// by ignoring any build configs created before this timestamp.
|
||||
invalidated map[int]time.Time
|
||||
}
|
||||
|
||||
pub fn new_job_queue(default_schedule CronExpression, default_base_image string) BuildJobQueue {
|
||||
return BuildJobQueue{
|
||||
default_schedule: default_schedule
|
||||
default_base_image: default_base_image
|
||||
invalidated: map[int]time.Time{}
|
||||
}
|
||||
}
|
||||
|
||||
|
@ -63,6 +69,7 @@ pub fn (mut q BuildJobQueue) insert(target Target, arch string) ! {
|
|||
timestamp := ce.next_from_now()!
|
||||
|
||||
job := BuildJob{
|
||||
created: time.now()
|
||||
timestamp: timestamp
|
||||
ce: ce
|
||||
config: BuildConfig{
|
||||
|
@ -88,6 +95,7 @@ fn (mut q BuildJobQueue) reschedule(job BuildJob, arch string) ! {
|
|||
|
||||
new_job := BuildJob{
|
||||
...job
|
||||
created: time.now()
|
||||
timestamp: new_timestamp
|
||||
}
|
||||
|
||||
|
@ -96,18 +104,28 @@ fn (mut q BuildJobQueue) reschedule(job BuildJob, arch string) ! {
|
|||
|
||||
// peek shows the first job for the given architecture that's ready to be
|
||||
// executed, if present.
|
||||
pub fn (q &BuildJobQueue) peek(arch string) ?BuildJob {
|
||||
pub fn (mut q BuildJobQueue) peek(arch string) ?BuildJob {
|
||||
rlock q.mutex {
|
||||
if arch !in q.queues {
|
||||
return none
|
||||
}
|
||||
|
||||
for {
|
||||
job := q.queues[arch].peek() or { return none }
|
||||
|
||||
// Skip any invalidated jobs
|
||||
if job.config.target_id in q.invalidated
|
||||
&& job.created < q.invalidated[job.config.target_id] {
|
||||
// This pop *should* never fail according to the source code
|
||||
q.queues[arch].pop() or { return none }
|
||||
continue
|
||||
}
|
||||
|
||||
if job.timestamp < time.now() {
|
||||
return job
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
return none
|
||||
}
|
||||
|
@ -120,8 +138,17 @@ pub fn (mut q BuildJobQueue) pop(arch string) ?BuildJob {
|
|||
return none
|
||||
}
|
||||
|
||||
for {
|
||||
mut job := q.queues[arch].peek() or { return none }
|
||||
|
||||
// Skip any invalidated jobs
|
||||
if job.config.target_id in q.invalidated
|
||||
&& job.created < q.invalidated[job.config.target_id] {
|
||||
// This pop *should* never fail according to the source code
|
||||
q.queues[arch].pop() or { return none }
|
||||
continue
|
||||
}
|
||||
|
||||
if job.timestamp < time.now() {
|
||||
job = q.queues[arch].pop()?
|
||||
|
||||
|
@ -133,6 +160,7 @@ pub fn (mut q BuildJobQueue) pop(arch string) ?BuildJob {
|
|||
return job
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
return none
|
||||
}
|
||||
|
@ -146,18 +174,28 @@ pub fn (mut q BuildJobQueue) pop_n(arch string, n int) []BuildJob {
|
|||
|
||||
mut out := []BuildJob{}
|
||||
|
||||
for out.len < n {
|
||||
mut job := q.queues[arch].peek() or { break }
|
||||
outer: for out.len < n {
|
||||
for {
|
||||
mut job := q.queues[arch].peek() or { break outer }
|
||||
|
||||
// Skip any invalidated jobs
|
||||
if job.config.target_id in q.invalidated
|
||||
&& job.created < q.invalidated[job.config.target_id] {
|
||||
// This pop *should* never fail according to the source code
|
||||
q.queues[arch].pop() or { break outer }
|
||||
continue
|
||||
}
|
||||
|
||||
if job.timestamp < time.now() {
|
||||
job = q.queues[arch].pop() or { break }
|
||||
job = q.queues[arch].pop() or { break outer }
|
||||
|
||||
// TODO idem
|
||||
q.reschedule(job, arch) or {}
|
||||
|
||||
out << job
|
||||
} else {
|
||||
break
|
||||
break outer
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
@ -166,3 +204,8 @@ pub fn (mut q BuildJobQueue) pop_n(arch string, n int) []BuildJob {
|
|||
|
||||
return []
|
||||
}
|
||||
|
||||
// invalidate a target's old build jobs.
|
||||
pub fn (mut q BuildJobQueue) invalidate(target_id int) {
|
||||
q.invalidated[target_id] = time.now()
|
||||
}
|
||||
|
|
Loading…
Reference in New Issue