Initial Backupper Unraid plugin
This commit is contained in:
@@ -0,0 +1,157 @@
|
||||
package queue
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"sort"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/backupper-unraid/backupper/internal/model"
|
||||
)
|
||||
|
||||
var ErrDuplicate = errors.New("task already queued or running")
|
||||
|
||||
type Handler func(context.Context, model.Run) model.Run
|
||||
type OnChange func([]model.Run)
|
||||
|
||||
type Queue struct {
|
||||
mu sync.Mutex
|
||||
cond *sync.Cond
|
||||
pending []model.Run
|
||||
history []model.Run
|
||||
active *model.Run
|
||||
cancel context.CancelFunc
|
||||
stopping bool
|
||||
handler Handler
|
||||
onChange OnChange
|
||||
}
|
||||
|
||||
func New(handler Handler, onChange OnChange) *Queue {
|
||||
q := &Queue{handler: handler, onChange: onChange}
|
||||
q.cond = sync.NewCond(&q.mu)
|
||||
return q
|
||||
}
|
||||
|
||||
func (q *Queue) RestoreHistory(runs []model.Run) {
|
||||
q.mu.Lock()
|
||||
defer q.mu.Unlock()
|
||||
for _, run := range runs {
|
||||
if run.Status == "queued" || run.Status == "running" {
|
||||
run.Status = "failed"
|
||||
run.ErrorCode = "environment"
|
||||
run.Message = "daemon restarted while task was active"
|
||||
now := time.Now().UTC()
|
||||
run.FinishedAt = &now
|
||||
}
|
||||
q.history = append(q.history, run)
|
||||
}
|
||||
}
|
||||
|
||||
func (q *Queue) Enqueue(run model.Run) error {
|
||||
q.mu.Lock()
|
||||
defer q.mu.Unlock()
|
||||
if q.stopping {
|
||||
return errors.New("queue is stopping")
|
||||
}
|
||||
if q.active != nil && run.JobID != "" && q.active.JobID == run.JobID {
|
||||
return ErrDuplicate
|
||||
}
|
||||
for _, item := range q.pending {
|
||||
if run.JobID != "" && item.JobID == run.JobID {
|
||||
return ErrDuplicate
|
||||
}
|
||||
}
|
||||
run.Status = "queued"
|
||||
q.pending = append(q.pending, run)
|
||||
sort.SliceStable(q.pending, func(i, j int) bool { return q.pending[i].Priority > q.pending[j].Priority })
|
||||
q.emitLocked()
|
||||
q.cond.Signal()
|
||||
return nil
|
||||
}
|
||||
|
||||
func (q *Queue) Run(ctx context.Context) {
|
||||
for {
|
||||
q.mu.Lock()
|
||||
for len(q.pending) == 0 && !q.stopping {
|
||||
q.cond.Wait()
|
||||
}
|
||||
if q.stopping {
|
||||
q.mu.Unlock()
|
||||
return
|
||||
}
|
||||
run := q.pending[0]
|
||||
q.pending = q.pending[1:]
|
||||
now := time.Now().UTC()
|
||||
run.Status, run.StartedAt = "running", &now
|
||||
workCtx, cancel := context.WithCancel(ctx)
|
||||
q.cancel, q.active = cancel, &run
|
||||
q.emitLocked()
|
||||
q.mu.Unlock()
|
||||
|
||||
result := q.handler(workCtx, run)
|
||||
finished := time.Now().UTC()
|
||||
result.FinishedAt = &finished
|
||||
|
||||
q.mu.Lock()
|
||||
q.active, q.cancel = nil, nil
|
||||
q.history = append(q.history, result)
|
||||
if len(q.history) > 500 {
|
||||
q.history = q.history[len(q.history)-500:]
|
||||
}
|
||||
q.emitLocked()
|
||||
q.mu.Unlock()
|
||||
}
|
||||
}
|
||||
|
||||
func (q *Queue) Cancel(id string) bool {
|
||||
q.mu.Lock()
|
||||
defer q.mu.Unlock()
|
||||
if q.active != nil && q.active.ID == id && q.cancel != nil {
|
||||
q.cancel()
|
||||
return true
|
||||
}
|
||||
for i := range q.pending {
|
||||
if q.pending[i].ID == id {
|
||||
run := q.pending[i]
|
||||
q.pending = append(q.pending[:i], q.pending[i+1:]...)
|
||||
now := time.Now().UTC()
|
||||
run.Status, run.ErrorCode, run.FinishedAt = "cancelled", "cancelled", &now
|
||||
q.history = append(q.history, run)
|
||||
q.emitLocked()
|
||||
return true
|
||||
}
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
func (q *Queue) Stop() {
|
||||
q.mu.Lock()
|
||||
q.stopping = true
|
||||
if q.cancel != nil {
|
||||
q.cancel()
|
||||
}
|
||||
q.cond.Broadcast()
|
||||
q.mu.Unlock()
|
||||
}
|
||||
|
||||
func (q *Queue) Snapshot() []model.Run {
|
||||
q.mu.Lock()
|
||||
defer q.mu.Unlock()
|
||||
return q.snapshotLocked()
|
||||
}
|
||||
|
||||
func (q *Queue) snapshotLocked() []model.Run {
|
||||
result := append([]model.Run{}, q.history...)
|
||||
if q.active != nil {
|
||||
result = append(result, *q.active)
|
||||
}
|
||||
result = append(result, q.pending...)
|
||||
return result
|
||||
}
|
||||
|
||||
func (q *Queue) emitLocked() {
|
||||
if q.onChange != nil {
|
||||
q.onChange(q.snapshotLocked())
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,30 @@
|
||||
package queue
|
||||
|
||||
import (
|
||||
"context"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/backupper-unraid/backupper/internal/model"
|
||||
)
|
||||
|
||||
func TestQueueDeduplicatesJobAndRunsTask(t *testing.T) {
|
||||
done := make(chan struct{})
|
||||
q := New(func(_ context.Context, run model.Run) model.Run { run.Status = "success"; close(done); return run }, nil)
|
||||
run := model.Run{ID: "one", JobID: "job", TaskType: "backup", CreatedAt: time.Now()}
|
||||
if err := q.Enqueue(run); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := q.Enqueue(model.Run{ID: "two", JobID: "job"}); err != ErrDuplicate {
|
||||
t.Fatalf("duplicate error = %v", err)
|
||||
}
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
defer cancel()
|
||||
go q.Run(ctx)
|
||||
select {
|
||||
case <-done:
|
||||
case <-time.After(time.Second):
|
||||
t.Fatal("queue did not execute task")
|
||||
}
|
||||
q.Stop()
|
||||
}
|
||||
Reference in New Issue
Block a user