2
0
mirror of https://github.com/hibiken/asynq.git synced 2024-11-14 11:31:18 +08:00

Fix bug around releasing semaphore token

This commit is contained in:
Ken Hibino 2019-11-18 07:42:26 -08:00
parent c6f482d4f8
commit 3daef02632

View File

@ -154,9 +154,10 @@ func (w *Workers) Run(handler TaskHandler) {
log.Printf("[Servere Error] could not parse json encoded message %s: %v", data, err) log.Printf("[Servere Error] could not parse json encoded message %s: %v", data, err)
continue continue
} }
w.poolTokens <- struct{}{} // acquire a token
t := &Task{Type: msg.Type, Payload: msg.Payload} t := &Task{Type: msg.Type, Payload: msg.Payload}
w.poolTokens <- struct{}{} // acquire token
go func(task *Task) { go func(task *Task) {
defer func() { <-w.poolTokens }() // release token
err := handler(task) err := handler(task)
if err != nil { if err != nil {
if msg.Retried >= msg.Retry { if msg.Retried >= msg.Retry {
@ -178,7 +179,6 @@ func (w *Workers) Run(handler TaskHandler) {
} }
} }
<-w.poolTokens // release the token
}(t) }(t)
} }
} }