From caf0e3c36fa7b7773d8cd0da786d050d611b2f3e Mon Sep 17 00:00:00 2001 From: Qiu Jian Date: Fri, 9 Nov 2018 11:17:14 +0800 Subject: [PATCH 1/2] =?UTF-8?q?=E4=BF=AE=E6=AD=A3=EF=BC=9A1.=20panic?= =?UTF-8?q?=E5=90=8E=E4=B8=8D=E8=BF=94=E5=9B=9E500=202.=20=E8=B6=85?= =?UTF-8?q?=E6=97=B6=E8=BF=98=E4=BC=9A=E9=98=BB=E5=A1=9Eworker?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- pkg/appsrv/appsrv.go | 80 ++++++++++++++++++++++++-------------- pkg/appsrv/handlerinfo.go | 6 +++ pkg/appsrv/response.go | 36 +++++++++++------ pkg/appsrv/test/testapp.go | 29 ++++++++++++++ pkg/appsrv/workers.go | 55 +++++++++++++++++--------- 5 files changed, 146 insertions(+), 60 deletions(-) create mode 100644 pkg/appsrv/test/testapp.go diff --git a/pkg/appsrv/appsrv.go b/pkg/appsrv/appsrv.go index 9e56f105d9..3720420051 100644 --- a/pkg/appsrv/appsrv.go +++ b/pkg/appsrv/appsrv.go @@ -24,6 +24,7 @@ type Application struct { name string context context.Context session *SWorkerManager + systemSession *SWorkerManager roots map[string]*RadixNode rootLock *sync.Mutex connMax int @@ -51,6 +52,7 @@ func NewApplication(name string, connMax int) *Application { context: context.Background(), connMax: connMax, session: NewWorkerManager("HttpRequestWorkerManager", connMax, DEFAULT_BACKLOG), + systemSession: NewWorkerManager("InternalHttpRequestWorkerManager", 1, DEFAULT_BACKLOG), roots: make(map[string]*RadixNode), rootLock: &sync.Mutex{}, idleTimeout: DEFAULT_IDLE_TIMEOUT, @@ -170,10 +172,15 @@ func (app *Application) ServeHTTP(w http.ResponseWriter, r *http.Request) { duration := float64(time.Since(start).Nanoseconds()) / 1000000 counter.hit += 1 counter.duration += duration - log.Infof("%d %s %s %s (%s) %.2fms", lrw.status, rid, r.Method, r.URL, r.RemoteAddr, duration) + if !hi.skipLog { + log.Infof("%d %s %s %s (%s) %.2fms", lrw.status, rid, r.Method, r.URL, r.RemoteAddr, duration) + } } func (app *Application) handleCORS(w http.ResponseWriter, r *http.Request) bool { + if app.cors == nil { + return false + } if r.Method == "OPTIONS" && r.Header.Get("Access-Control-Request-Method") != "" { app.cors.handlePreflight(w, r) return true @@ -196,7 +203,6 @@ func (app *Application) defaultHandle(w http.ResponseWriter, r *http.Request, ri if ok { fw := newResponseWriterChannel(w) worker := make(chan *SWorker) - errChan := make(chan interface{}) to := hand.processTimeout if to == 0 { to = app.processTimeout @@ -207,25 +213,32 @@ func (app *Application) defaultHandle(w http.ResponseWriter, r *http.Request, ri if session == nil { session = app.session } - session.Run(func() { - defer fw.closeChannels() - if ctx.Err() == nil { - ctx = context.WithValue(ctx, appctx.APP_CONTEXT_KEY_REQUEST_ID, rid) - ctx = context.WithValue(ctx, appctx.APP_CONTEXT_KEY_CUR_ROOT, hand.path) - ctx = context.WithValue(ctx, appctx.APP_CONTEXT_KEY_CUR_PATH, segs[len(hand.path):]) - ctx = context.WithValue(ctx, appctx.APP_CONTEXT_KEY_PARAMS, params) - if hand.metadata != nil { - ctx = context.WithValue(ctx, appctx.APP_CONTEXT_KEY_METADATA, hand.metadata) - } - func() { - span := trace.StartServerTrace(w, r, hand.GetName(params), app.GetName(), hand.GetTags()) - defer span.EndTrace() - ctx = context.WithValue(ctx, appctx.APP_CONTEXT_KEY_TRACE, span) - hand.handler(ctx, &fw, r) - }() - } // otherwise, the task has been timeout - }, worker, errChan) - runErr := fw.wait(ctx, worker, errChan) + session.Run( + func() { + if ctx.Err() == nil { + ctx = context.WithValue(ctx, appctx.APP_CONTEXT_KEY_REQUEST_ID, rid) + ctx = context.WithValue(ctx, appctx.APP_CONTEXT_KEY_CUR_ROOT, hand.path) + ctx = context.WithValue(ctx, appctx.APP_CONTEXT_KEY_CUR_PATH, segs[len(hand.path):]) + ctx = context.WithValue(ctx, appctx.APP_CONTEXT_KEY_PARAMS, params) + if hand.metadata != nil { + ctx = context.WithValue(ctx, appctx.APP_CONTEXT_KEY_METADATA, hand.metadata) + } + func() { + span := trace.StartServerTrace(&fw, r, hand.GetName(params), app.GetName(), hand.GetTags()) + defer span.EndTrace() + ctx = context.WithValue(ctx, appctx.APP_CONTEXT_KEY_TRACE, span) + hand.handler(ctx, &fw, r) + }() + } // otherwise, the task has been timeout + fw.closeChannels() + }, + worker, + func(err error) { + httperrors.InternalServerError(&fw, "Internal server error: %s", err) + fw.closeChannels() + }, + ) + runErr := fw.wait(ctx, worker) if runErr != nil { switch runErr.(type) { case *httputils.JSONClientError: @@ -235,6 +248,7 @@ func (app *Application) defaultHandle(w http.ResponseWriter, r *http.Request, ri httperrors.InternalServerError(w, "Internal server error") } } + fw.closeChannels() return hand } else { log.Errorf("Invalid handler for %s", r.URL) @@ -247,13 +261,19 @@ func (app *Application) defaultHandle(w http.ResponseWriter, r *http.Request, ri return nil } -func (app *Application) addDefaultHandler() { - app.AddHandler("GET", "/version", VersionHandler) - app.AddHandler("GET", "/stats", StatisticHandler) - app.AddHandler("POST", "/ping", PingHandler) - app.AddHandler("GET", "/ping", PingHandler) - // app.AddHandler("OPTIONS", "/", CORSHandler) - app.AddHandler("GET", "/worker_stats", WorkerStatsHandler) +func (app *Application) addDefaultHandler(method string, prefix string, handler func(context.Context, http.ResponseWriter, *http.Request), name string) { + segs := SplitPath(prefix) + hi := newHandlerInfo(method, segs, handler, nil, name, nil) + hi.SetSkipLog(true).SetWorkerManager(app.systemSession) + app.AddHandler3(hi) +} + +func (app *Application) addDefaultHandlers() { + app.addDefaultHandler("GET", "/version", VersionHandler, "version") + app.addDefaultHandler("GET", "/stats", StatisticHandler, "stats") + app.addDefaultHandler("POST", "/ping", PingHandler, "ping") + app.addDefaultHandler("GET", "/ping", PingHandler, "ping") + app.addDefaultHandler("GET", "/worker_stats", WorkerStatsHandler, "worker_stats") } func timeoutHandle(h http.Handler) http.HandlerFunc { @@ -274,10 +294,10 @@ func (app *Application) initServer(addr string) *http.Server { db.SetMaxIdleConns(app.connMax + 1) db.SetMaxOpenConns(app.connMax + 1) } - app.addDefaultHandler() + app.addDefaultHandlers() s := &http.Server{ Addr: addr, - Handler: timeoutHandle(app), + Handler: app, IdleTimeout: app.idleTimeout, ReadTimeout: app.readTimeout, ReadHeaderTimeout: app.readHeaderTimeout, diff --git a/pkg/appsrv/handlerinfo.go b/pkg/appsrv/handlerinfo.go index e6f07c7cbb..0eb87f6ff2 100644 --- a/pkg/appsrv/handlerinfo.go +++ b/pkg/appsrv/handlerinfo.go @@ -24,6 +24,7 @@ type SHandlerInfo struct { counter5XX handlerRequestCounter processTimeout time.Duration workerMan *SWorkerManager + skipLog bool } func (this *SHandlerInfo) GetName(params map[string]string) string { @@ -96,3 +97,8 @@ func (hi *SHandlerInfo) SetWorkerManager(workerMan *SWorkerManager) *SHandlerInf hi.workerMan = workerMan return hi } + +func (hi *SHandlerInfo) SetSkipLog(skip bool) *SHandlerInfo { + hi.skipLog = skip + return hi +} diff --git a/pkg/appsrv/response.go b/pkg/appsrv/response.go index 6bf18eea4c..718149909d 100644 --- a/pkg/appsrv/response.go +++ b/pkg/appsrv/response.go @@ -6,6 +6,7 @@ import ( "yunion.io/x/log" + "fmt" "yunion.io/x/onecloud/pkg/httperrors" ) @@ -15,19 +16,25 @@ type responseWriterResponse struct { } type responseWriterChannel struct { - backend http.ResponseWriter + backend http.ResponseWriter + bodyChan chan []byte bodyResp chan responseWriterResponse statusChan chan int statusResp chan bool + + isClosed bool } func newResponseWriterChannel(backend http.ResponseWriter) responseWriterChannel { - return responseWriterChannel{backend: backend, + return responseWriterChannel{ + backend: backend, bodyChan: make(chan []byte), bodyResp: make(chan responseWriterResponse), statusChan: make(chan int), - statusResp: make(chan bool)} + statusResp: make(chan bool), + isClosed: false, + } } func (w *responseWriterChannel) Header() http.Header { @@ -35,24 +42,30 @@ func (w *responseWriterChannel) Header() http.Header { } func (w *responseWriterChannel) Write(bytes []byte) (int, error) { + if w.isClosed { + return 0, fmt.Errorf("response stream has been closed") + } w.bodyChan <- bytes v := <-w.bodyResp return v.count, v.err } func (w *responseWriterChannel) WriteHeader(status int) { + if w.isClosed { + return + } w.statusChan <- status <-w.statusResp } -func (w *responseWriterChannel) wait(ctx context.Context, workerChan chan *SWorker, errChan chan interface{}) interface{} { - var err interface{} +func (w *responseWriterChannel) wait(ctx context.Context, workerChan chan *SWorker) interface{} { + var err error var worker *SWorker stop := false for !stop { select { case worker = <-workerChan: - log.Infof("request is being handled by worker %s", worker) + log.Debugf("request is being handled by worker %s", worker) case <-ctx.Done(): // ctx deadline reached, timeout if worker != nil { @@ -60,12 +73,6 @@ func (w *responseWriterChannel) wait(ctx context.Context, workerChan chan *SWork } err = httperrors.NewTimeoutError("request process timeout") stop = true - case e, more := <-errChan: - if more { - err = e - } else { - stop = true - } case bytes, more := <-w.bodyChan: // log.Print("Recive body ", len(bytes), " more ", more) if more { @@ -88,6 +95,11 @@ func (w *responseWriterChannel) wait(ctx context.Context, workerChan chan *SWork } func (w *responseWriterChannel) closeChannels() { + if w.isClosed { + return + } + w.isClosed = true + close(w.bodyChan) close(w.bodyResp) close(w.statusChan) diff --git a/pkg/appsrv/test/testapp.go b/pkg/appsrv/test/testapp.go new file mode 100644 index 0000000000..b424cba689 --- /dev/null +++ b/pkg/appsrv/test/testapp.go @@ -0,0 +1,29 @@ +package main + +import ( + "context" + "net/http" + "time" + + "yunion.io/x/log" + "yunion.io/x/onecloud/pkg/appsrv" +) + +func main() { + app := appsrv.NewApplication("test", 4) + app.AddHandler("GET", "/delay", func(ctx context.Context, w http.ResponseWriter, r *http.Request) { + time.Sleep(time.Second * 20) + log.Debugf("end of delay sleep....") + appsrv.Send(w, "pong") + }) + app.AddHandler("GET", "/panic", func(ctx context.Context, w http.ResponseWriter, r *http.Request) { + panic("the handler is panic") + appsrv.Send(w, "pong") + }) + app.AddHandler("GET", "/delaypanic", func(ctx context.Context, w http.ResponseWriter, r *http.Request) { + time.Sleep(time.Second * 1) + panic("the handler is panic") + appsrv.Send(w, "pong") + }) + app.ListenAndServe("0.0.0.0:44444") +} diff --git a/pkg/appsrv/workers.go b/pkg/appsrv/workers.go index 45c48b9a6e..6f328bde28 100644 --- a/pkg/appsrv/workers.go +++ b/pkg/appsrv/workers.go @@ -53,6 +53,8 @@ func (worker *SWorker) isDetached() bool { } func (worker *SWorker) run() { + defer log.Debugf("no more job, exit worker %s", worker) + defer worker.manager.removeWorker(worker) for { if worker.isDetached() { if isDebug { @@ -80,7 +82,6 @@ func (worker *SWorker) run() { break } } - worker.manager.removeWorker(worker) } func (worker *SWorker) Detach(reason string) { @@ -92,18 +93,28 @@ func (worker *SWorker) Detach(reason string) { worker.manager.detachedWorker.addWithLock(worker) log.Warningf("detach worker %s due to reason %s", worker, reason) + + worker.manager.scheduleWithLock() +} + +func (worker *SWorker) StateStr() string { + if worker.state == WORKER_STATE_ACTIVE { + return "active" + } else { + return "detach" + } } func (worker *SWorker) String() string { - return fmt.Sprintf("#%d(%d)", worker.id, worker.state) + return fmt.Sprintf("#%d(%p, %s)", worker.id, worker, worker.StateStr()) } type SWorkerList struct { list *list.List } -func newWorkerList() SWorkerList { - return SWorkerList{ +func newWorkerList() *SWorkerList { + return &SWorkerList{ list: list.New(), } } @@ -127,8 +138,8 @@ type SWorkerManager struct { queue *Ring workerCount int backlog int - activeWorker SWorkerList - detachedWorker SWorkerList + activeWorker *SWorkerList + detachedWorker *SWorkerList workerLock *sync.Mutex workerId uint64 } @@ -148,15 +159,21 @@ func NewWorkerManager(name string, workerCount int, backlog int) *SWorkerManager } type sWorkerTask struct { - task func() - worker chan *SWorker - err chan interface{} + task func() + worker chan *SWorker + onError func(error) } -func (wm *SWorkerManager) Run(task func(), worker chan *SWorker, err chan interface{}) bool { - ret := wm.queue.Push(&sWorkerTask{task: task, worker: worker, err: err}) +func (wm *SWorkerManager) String() string { + return wm.name +} + +func (wm *SWorkerManager) Run(task func(), worker chan *SWorker, onErr func(error)) bool { + ret := wm.queue.Push(&sWorkerTask{task: task, worker: worker, onError: onErr}) if ret { wm.schedule() + } else { + log.Warningf("[%] BUSY fail to push task, queue is FULL") } return ret } @@ -176,23 +193,23 @@ func execCallback(task *sWorkerTask) { defer func() { if r := recover(); r != nil { log.Errorf("WorkerManager exec callback error: %s", r) - debug.PrintStack() - if task.err != nil { - task.err <- r - close(task.err) + if task.onError != nil { + task.onError(fmt.Errorf("%s", r)) } + debug.PrintStack() } }() task.task() - if task.err != nil { - close(task.err) - } } func (wm *SWorkerManager) schedule() { wm.workerLock.Lock() defer wm.workerLock.Unlock() + wm.scheduleWithLock() +} + +func (wm *SWorkerManager) scheduleWithLock() { if wm.activeWorker.size() < wm.workerCount && wm.queue.Size() > 0 { wm.workerId += 1 worker := newWorker(wm.workerId, wm) @@ -201,6 +218,8 @@ func (wm *SWorkerManager) schedule() { log.Debugf("no enough worker, add new worker %s", worker) } go worker.run() + } else { + log.Warningf("[%s] BUSY activeWork %d max %d queue: %d", wm, wm.activeWorker.size(), wm.workerCount, wm.queue.Size()) } } From 77ed2af387c5863f9d8447150ea916f704330983 Mon Sep 17 00:00:00 2001 From: Qiu Jian Date: Fri, 9 Nov 2018 12:28:18 +0800 Subject: [PATCH 2/2] minor fixes --- pkg/appsrv/test/testapp.go | 2 -- 1 file changed, 2 deletions(-) diff --git a/pkg/appsrv/test/testapp.go b/pkg/appsrv/test/testapp.go index b424cba689..540e43db87 100644 --- a/pkg/appsrv/test/testapp.go +++ b/pkg/appsrv/test/testapp.go @@ -18,12 +18,10 @@ func main() { }) app.AddHandler("GET", "/panic", func(ctx context.Context, w http.ResponseWriter, r *http.Request) { panic("the handler is panic") - appsrv.Send(w, "pong") }) app.AddHandler("GET", "/delaypanic", func(ctx context.Context, w http.ResponseWriter, r *http.Request) { time.Sleep(time.Second * 1) panic("the handler is panic") - appsrv.Send(w, "pong") }) app.ListenAndServe("0.0.0.0:44444") }