Files
move-data-nas/internal/syncengine/queue.go
T

73 lines
1.4 KiB
Go

package syncengine
import (
"errors"
"sync"
)
var ErrAlreadyRunning = errors.New("job already running for this sync pair")
type Queue struct {
mu sync.Mutex
runs map[int64]*RunInfo
}
type RunInfo struct {
JobID int64
Cancel func()
ByUser bool
}
func NewQueue() *Queue {
return &Queue{runs: make(map[int64]*RunInfo)}
}
func (q *Queue) Enqueue(syncPairID, jobID int64, cancel func()) error {
q.mu.Lock()
defer q.mu.Unlock()
if _, exists := q.runs[syncPairID]; exists {
return ErrAlreadyRunning
}
q.runs[syncPairID] = &RunInfo{JobID: jobID, Cancel: cancel}
return nil
}
func (q *Queue) Dequeue(syncPairID int64) {
q.mu.Lock()
defer q.mu.Unlock()
delete(q.runs, syncPairID)
}
func (q *Queue) IsRunning(syncPairID int64) bool {
q.mu.Lock()
defer q.mu.Unlock()
_, exists := q.runs[syncPairID]
return exists
}
func (q *Queue) GetJobID(syncPairID int64) (int64, bool) {
q.mu.Lock()
defer q.mu.Unlock()
info, exists := q.runs[syncPairID]
if !exists {
return 0, false
}
return info.JobID, true
}
func (q *Queue) Cancel(syncPairID int64, byUser bool) {
q.mu.Lock()
defer q.mu.Unlock()
if info, exists := q.runs[syncPairID]; exists && info.Cancel != nil {
info.ByUser = byUser
info.Cancel()
}
}
func (q *Queue) IsCancelledByUser(syncPairID int64) bool {
q.mu.Lock()
defer q.mu.Unlock()
info, exists := q.runs[syncPairID]
return exists && info.ByUser
}