diff --git a/pkg/appsrv/appsrv.go b/pkg/appsrv/appsrv.go index c57af5a03b..38c17a7604 100644 --- a/pkg/appsrv/appsrv.go +++ b/pkg/appsrv/appsrv.go @@ -49,7 +49,7 @@ func NewApplication(name string, connMax int) *Application { app := Application{name: name, context: context.Background(), connMax: connMax, - session: NewWorkerManager("sessionMan", connMax, DEFAULT_BACKLOG), + session: NewWorkerManager("HttpRequestWorkerManager", connMax, DEFAULT_BACKLOG), roots: make(map[string]*RadixNode), rootLock: &sync.Mutex{}, idleTimeout: DEFAULT_IDLE_TIMEOUT, @@ -239,6 +239,7 @@ func (app *Application) addDefaultHandler() { app.AddHandler("POST", "/ping", PingHandler) app.AddHandler("GET", "/ping", PingHandler) // app.AddHandler("OPTIONS", "/", CORSHandler) + app.AddHandler("GET", "/worker_stats", WorkerStatsHandler) } func timeoutHandle(h http.Handler) http.HandlerFunc { diff --git a/pkg/appsrv/workers.go b/pkg/appsrv/workers.go index bffed2f678..ed3b940eee 100644 --- a/pkg/appsrv/workers.go +++ b/pkg/appsrv/workers.go @@ -1,11 +1,14 @@ package appsrv import ( + "container/list" + "context" + "fmt" + "net/http" "runtime/debug" "sync" - "container/list" - "fmt" + "yunion.io/x/jsonutils" "yunion.io/x/log" ) @@ -20,6 +23,12 @@ func enableDebug() { isDebug = true } +var workerManagers []*SWorkerManager + +func init() { + workerManagers = make([]*SWorkerManager, 0) +} + type SWorker struct { id uint64 state int @@ -133,6 +142,8 @@ func NewWorkerManager(name string, workerCount int, backlog int) *SWorkerManager detachedWorker: newWorkerList(), workerLock: &sync.Mutex{}, workerId: 0} + + workerManagers = append(workerManagers, &manager) return &manager } @@ -201,6 +212,34 @@ func (wm *SWorkerManager) DetachedWorkerCount() int { return wm.detachedWorker.size() } +type SWorkerManagerStates struct { + Name string + Backlog int + ActiveWorkerCnt int + DetachWorkerCnt int +} + +func (wm *SWorkerManager) getState() SWorkerManagerStates { + state := SWorkerManagerStates{} + + state.Name = wm.name + state.Backlog = wm.queue.Size() + state.ActiveWorkerCnt = wm.activeWorker.size() + state.DetachWorkerCnt = wm.detachedWorker.size() + + return state +} + +func WorkerStatsHandler(ctx context.Context, w http.ResponseWriter, r *http.Request) { + stats := make([]SWorkerManagerStates, 0) + for i := 0; i < len(workerManagers); i += 1 { + stats = append(stats, workerManagers[i].getState()) + } + result := jsonutils.NewDict() + result.Add(jsonutils.Marshal(&stats), "workers") + fmt.Fprintf(w, result.String()) +} + func WaitChannel(ch chan interface{}) interface{} { var ret interface{} stop := false