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 }