code refactor
This commit is contained in:
@@ -7,7 +7,7 @@ import (
|
|||||||
)
|
)
|
||||||
|
|
||||||
type MessageQueue struct {
|
type MessageQueue struct {
|
||||||
ch chan *Process
|
producerCh chan *Process
|
||||||
consumerCh chan struct{}
|
consumerCh chan struct{}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -23,7 +23,7 @@ func NewMessageQueue() *MessageQueue {
|
|||||||
}
|
}
|
||||||
|
|
||||||
return &MessageQueue{
|
return &MessageQueue{
|
||||||
ch: make(chan *Process, size),
|
producerCh: make(chan *Process, size),
|
||||||
consumerCh: make(chan struct{}, size),
|
consumerCh: make(chan struct{}, size),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -31,12 +31,13 @@ func NewMessageQueue() *MessageQueue {
|
|||||||
// Publish a message to the queue and set the task to a peding state.
|
// Publish a message to the queue and set the task to a peding state.
|
||||||
func (m *MessageQueue) Publish(p *Process) {
|
func (m *MessageQueue) Publish(p *Process) {
|
||||||
go p.SetPending()
|
go p.SetPending()
|
||||||
m.ch <- p
|
m.producerCh <- p
|
||||||
}
|
}
|
||||||
|
|
||||||
// Setup the consumer listened which "subscribes" to the queue events.
|
// Setup the consumer listener which subscribes to the changes to the producer
|
||||||
func (m *MessageQueue) SetupConsumer() {
|
// channel and triggers the "download" action.
|
||||||
for msg := range m.ch {
|
func (m *MessageQueue) Subscriber() {
|
||||||
|
for msg := range m.producerCh {
|
||||||
m.consumerCh <- struct{}{}
|
m.consumerCh <- struct{}{}
|
||||||
go func(p *Process) {
|
go func(p *Process) {
|
||||||
p.Start()
|
p.Start()
|
||||||
|
|||||||
@@ -28,10 +28,9 @@ func RunBlocking(port int, frontend fs.FS) {
|
|||||||
db.Restore()
|
db.Restore()
|
||||||
|
|
||||||
mq := internal.NewMessageQueue()
|
mq := internal.NewMessageQueue()
|
||||||
go mq.SetupConsumer()
|
go mq.Subscriber()
|
||||||
|
|
||||||
service := ytdlpRPC.Container(&db, mq)
|
service := ytdlpRPC.Container(&db, mq)
|
||||||
|
|
||||||
rpc.Register(service)
|
rpc.Register(service)
|
||||||
|
|
||||||
app := fiber.New()
|
app := fiber.New()
|
||||||
|
|||||||
Reference in New Issue
Block a user