From f361618f066670e4d6d6563362434b56e33d5116 Mon Sep 17 00:00:00 2001 From: Qiu Jian Date: Mon, 30 Mar 2020 19:25:00 +0800 Subject: [PATCH] fix: worker may block by worker_chan --- pkg/appsrv/appsrv.go | 6 +++--- pkg/appsrv/response.go | 7 ++++++- pkg/appsrv/workers.go | 3 +++ 3 files changed, 12 insertions(+), 4 deletions(-) diff --git a/pkg/appsrv/appsrv.go b/pkg/appsrv/appsrv.go index b2d7ad2e58..72f91951f0 100644 --- a/pkg/appsrv/appsrv.go +++ b/pkg/appsrv/appsrv.go @@ -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: diff --git a/pkg/appsrv/response.go b/pkg/appsrv/response.go index 0d102cf388..ec1e9a277d 100644 --- a/pkg/appsrv/response.go +++ b/pkg/appsrv/response.go @@ -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 { diff --git a/pkg/appsrv/workers.go b/pkg/appsrv/workers.go index a6defdcb7d..4747294b60 100644 --- a/pkg/appsrv/workers.go +++ b/pkg/appsrv/workers.go @@ -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 {