mirror of
https://github.com/fnproject/fn.git
synced 2022-10-28 21:29:17 +03:00
@@ -30,6 +30,8 @@ type Runner struct {
|
||||
availableMem int64
|
||||
usedMem int64
|
||||
usedMemMutex sync.RWMutex
|
||||
|
||||
stats
|
||||
}
|
||||
|
||||
var (
|
||||
|
||||
46
api/runner/stats.go
Normal file
46
api/runner/stats.go
Normal file
@@ -0,0 +1,46 @@
|
||||
package runner
|
||||
|
||||
import "sync"
|
||||
|
||||
type stats struct {
|
||||
mu sync.Mutex
|
||||
queue uint64
|
||||
running uint64
|
||||
complete uint64
|
||||
}
|
||||
|
||||
type Stats struct {
|
||||
Queue uint64
|
||||
Running uint64
|
||||
Complete uint64
|
||||
}
|
||||
|
||||
func (s *stats) Enqueue() {
|
||||
s.mu.Lock()
|
||||
s.queue++
|
||||
s.mu.Unlock()
|
||||
}
|
||||
|
||||
func (s *stats) Start() {
|
||||
s.mu.Lock()
|
||||
s.queue--
|
||||
s.running++
|
||||
s.mu.Unlock()
|
||||
}
|
||||
|
||||
func (s *stats) Complete() {
|
||||
s.mu.Lock()
|
||||
s.running--
|
||||
s.complete++
|
||||
s.mu.Unlock()
|
||||
}
|
||||
|
||||
func (s *stats) Stats() Stats {
|
||||
var stats Stats
|
||||
s.mu.Lock()
|
||||
stats.Running = s.running
|
||||
stats.Complete = s.complete
|
||||
stats.Queue = s.queue
|
||||
s.mu.Unlock()
|
||||
return stats
|
||||
}
|
||||
@@ -96,10 +96,12 @@ func StartWorkers(ctx context.Context, rnr *Runner, tasks <-chan task.Request) {
|
||||
continue
|
||||
}
|
||||
|
||||
rnr.Start()
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return
|
||||
case p <- task:
|
||||
rnr.Complete()
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -344,6 +346,8 @@ func (hc *htfn) serve(ctx context.Context) {
|
||||
|
||||
func runTaskReq(rnr *Runner, wg *sync.WaitGroup, t task.Request) {
|
||||
defer wg.Done()
|
||||
rnr.Start()
|
||||
defer rnr.Complete()
|
||||
result, err := rnr.Run(t.Ctx, t.Config)
|
||||
select {
|
||||
case t.Response <- task.Response{result, err}:
|
||||
|
||||
Reference in New Issue
Block a user