From 62e4eb0fd3a268f6ea3b7b3efa37ef994aecd7f9 Mon Sep 17 00:00:00 2001 From: Robert Kopaczewski Date: Thu, 9 Jul 2020 18:22:49 +0200 Subject: [PATCH] feat: better codebox container caching --- app/lb/server.go | 58 +++++++++++++++++++++++--------------------- app/lb/worker.go | 8 ++++++ app/script/errors.go | 8 ------ app/script/runner.go | 2 +- 4 files changed, 39 insertions(+), 37 deletions(-) diff --git a/app/lb/server.go b/app/lb/server.go index b3aa835..fa4def1 100644 --- a/app/lb/server.go +++ b/app/lb/server.go @@ -282,14 +282,16 @@ func (s *Server) processRun(ctx context.Context, logger logrus.FieldLogger, stre break } - cont.Worker.ResetErrorCount() - response, err := s.relayReponse(logger, stream, cont, resCh) - if response != nil && response.Cached { - // Add container to worker cache if we got any response. - cont.ID = response.ContainerId - s.addCache(cont, def) + if response != nil { + cont.Worker.ResetErrorCount() + + if response.Cached { + // Add container to worker cache if we got cached response. + cont.ID = response.ContainerId + s.addCache(cont, def) + } } logger.WithFields(logrus.Fields{"took": time.Since(start), "mcpu": cont.mCPU}).Info("grpc:lb:Run") @@ -329,36 +331,36 @@ func (s *Server) removeCache(w *Worker, def *script.Definition, containerID stri s.mu.Lock() - if w.RemoveCache(def, containerID) { - // Remove worker container with specified definition from cache map. - if m, ok := s.workerContainersCached[defHash]; ok { - if _, ok = m[containerID]; ok { - removedContainer = true + w.RemoveCache(def, containerID) - delete(m, containerID) - } + // Remove worker container with specified definition from cache map. + if m, ok := s.workerContainersCached[defHash]; ok { + if _, ok = m[containerID]; ok { + removedContainer = true - if len(m) == 0 { - delete(s.workerContainersCached, defHash) - } + delete(m, containerID) } - // Decrease refcount of script index for given worker. - if m, ok := s.workersByIndex[idx]; ok { - if v, ok := m[w.ID]; ok { - removedIdx = true + if len(m) == 0 { + delete(s.workerContainersCached, defHash) + } + } - if v <= 1 { - delete(m, w.ID) - } else { - m[w.ID]-- - } - } + // Decrease refcount of script index for given worker. + if m, ok := s.workersByIndex[idx]; ok { + if v, ok := m[w.ID]; ok { + removedIdx = true - if len(m) == 0 { - delete(s.workersByIndex, idx) + if v <= 1 { + delete(m, w.ID) + } else { + m[w.ID]-- } } + + if len(m) == 0 { + delete(s.workersByIndex, idx) + } } s.mu.Unlock() diff --git a/app/lb/worker.go b/app/lb/worker.go index 44f70f3..dc2e853 100644 --- a/app/lb/worker.go +++ b/app/lb/worker.go @@ -320,6 +320,14 @@ type WorkerContainer struct { async uint32 } +func (w *WorkerContainer) String() string { + if w.ID == "" { + return fmt.Sprintf("{ID:Fresh, Worker:%s}", w.Worker.ID) + } + + return fmt.Sprintf("{ID:%s, Worker:%s}", w.ID, w.Worker.ID) +} + // Conns returns number of current connections. func (w *WorkerContainer) Conns() uint32 { return atomic.LoadUint32(&w.conns) diff --git a/app/script/errors.go b/app/script/errors.go index 1599c15..411a35c 100644 --- a/app/script/errors.go +++ b/app/script/errors.go @@ -2,17 +2,9 @@ package script import ( "errors" - - "google.golang.org/grpc/codes" - "google.golang.org/grpc/status" ) var ( - // gRPC errors. - - // ErrSourceNotAvailable signals that specified source hash was not found. - ErrSourceNotAvailable = status.Error(codes.FailedPrecondition, "source not available") - // Result errors. // ErrIncorrectData signals when output structure is malformed. diff --git a/app/script/runner.go b/app/script/runner.go index 5a9d024..c528329 100644 --- a/app/script/runner.go +++ b/app/script/runner.go @@ -428,7 +428,7 @@ func (r *DockerRunner) Run(ctx context.Context, logger logrus.FieldLogger, reque // Check and refresh source and environment. if r.fileRepo.Get(def.SourceHash) == "" || (def.Environment != "" && r.fileRepo.Get(def.Environment) == "") { logger.Error("grpc:script:Run source not available") - return nil, ErrSourceNotAvailable + return nil, common.ErrSourceNotAvailable } r.taskWaitGroup.Add(1)