mirror of
https://github.com/yunionio/cloudpods.git
synced 2026-09-24 16:03:43 +08:00
fix: worker may block by worker_chan
This commit is contained in:
@@ -268,7 +268,7 @@ func (app *Application) defaultHandle(w http.ResponseWriter, r *http.Request, ri
|
||||
hand, ok := handler.(*SHandlerInfo)
|
||||
if ok {
|
||||
fw := newResponseWriterChannel(w)
|
||||
worker := make(chan *SWorker)
|
||||
currentWorker := make(chan *SWorker, 1) // make it a buffered channel
|
||||
to := hand.FetchProcessTimeout(r)
|
||||
if to == 0 {
|
||||
to = app.processTimeout
|
||||
@@ -320,13 +320,13 @@ func (app *Application) defaultHandle(w http.ResponseWriter, r *http.Request, ri
|
||||
} // otherwise, the task has been timeout
|
||||
fw.closeChannels()
|
||||
},
|
||||
worker,
|
||||
currentWorker,
|
||||
func(err error) {
|
||||
httperrors.InternalServerError(&fw, "Internal server error: %s", err)
|
||||
fw.closeChannels()
|
||||
},
|
||||
)
|
||||
runErr := fw.wait(ctx, worker)
|
||||
runErr := fw.wait(ctx, currentWorker)
|
||||
if runErr != nil {
|
||||
switch runErr.(type) {
|
||||
case *httputils.JSONClientError:
|
||||
|
||||
@@ -100,7 +100,12 @@ func (w *responseWriterChannel) wait(ctx context.Context, workerChan chan *SWork
|
||||
stop := false
|
||||
for !stop {
|
||||
select {
|
||||
case worker = <-workerChan:
|
||||
case curWorker, more := <-workerChan:
|
||||
if more {
|
||||
worker = curWorker
|
||||
} else {
|
||||
// ignore, worker is responsible for close the channel
|
||||
}
|
||||
case <-ctx.Done():
|
||||
// ctx deadline reached, timeout
|
||||
if worker != nil {
|
||||
|
||||
@@ -83,6 +83,9 @@ func (worker *SWorker) run() {
|
||||
task := req.(*sWorkerTask)
|
||||
if task.worker != nil {
|
||||
task.worker <- worker
|
||||
// worker channel is buffered
|
||||
// close the worker channel
|
||||
close(task.worker)
|
||||
}
|
||||
execCallback(task)
|
||||
} else {
|
||||
|
||||
Reference in New Issue
Block a user