feat: [wip] - implement storage

This commit is contained in:
sujit
2024-10-17 10:03:07 +05:45
parent 53b68572dd
commit ea266be846
13 changed files with 730 additions and 482 deletions

View File

@@ -1,6 +1,8 @@
package mq
import (
"context"
"github.com/oarkflow/mq/storage"
"github.com/oarkflow/mq/storage/memory"
)
@@ -29,3 +31,32 @@ func (b *Broker) NewQueue(qName string) *Queue {
go b.dispatchWorker(q)
return q
}
type QueueTask struct {
ctx context.Context
payload *Task
priority int
}
type PriorityQueue []*QueueTask
func (pq PriorityQueue) Len() int { return len(pq) }
func (pq PriorityQueue) Less(i, j int) bool {
return pq[i].priority > pq[j].priority
}
func (pq PriorityQueue) Swap(i, j int) { pq[i], pq[j] = pq[j], pq[i] }
func (pq *PriorityQueue) Push(x interface{}) {
item := x.(*QueueTask)
*pq = append(*pq, item)
}
func (pq *PriorityQueue) Pop() interface{} {
old := *pq
n := len(old)
item := old[n-1]
*pq = old[0 : n-1]
return item
}