Merge pull request #444 in YUNIONIO/onecloud from ~QIUJIAN/onecloud:hotfix/qj-fix-workermanager-bugs to release/2.3.0

* commit '77ed2af387c5863f9d8447150ea916f704330983':
  minor fixes
  修正:1. panic后不返回500 2. 超时还会阻塞worker
This commit is contained in:
邱剑
2018-11-09 17:27:50 +08:00
5 changed files with 144 additions and 60 deletions
+50 -30
View File
@@ -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,
+6
View File
@@ -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
}
+24 -12
View File
@@ -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)
+27
View File
@@ -0,0 +1,27 @@
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")
})
app.AddHandler("GET", "/delaypanic", func(ctx context.Context, w http.ResponseWriter, r *http.Request) {
time.Sleep(time.Second * 1)
panic("the handler is panic")
})
app.ListenAndServe("0.0.0.0:44444")
}
+37 -18
View File
@@ -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())
}
}