diff --git a/pkg/executor/fscache/poolcache.go b/pkg/executor/fscache/poolcache.go index 9cdd4469..984591bd 100644 --- a/pkg/executor/fscache/poolcache.go +++ b/pkg/executor/fscache/poolcache.go @@ -96,7 +96,6 @@ type ( ) // NewPoolCache create a Cache object - func NewPoolCache(logger *zap.Logger) *PoolCache { c := &PoolCache{ cache: make(map[crd.CacheKeyURG]*funcSvcGroup), @@ -122,6 +121,7 @@ func (c *PoolCache) service() { case getValue: funcSvcGroup, ok := c.cache[req.function] if !ok { + // first request for this function, create a new group c.cache[req.function] = NewFuncSvcGroup() c.cache[req.function].svcWaiting++ resp.error = ferror.MakeError(ferror.ErrorNotFound, @@ -131,6 +131,7 @@ func (c *PoolCache) service() { } found := false totalActiveRequests := 0 + // check if any specialized pod is available for addr := range funcSvcGroup.svcs { totalActiveRequests += funcSvcGroup.svcs[addr].activeRequests if funcSvcGroup.svcs[addr].activeRequests < req.requestsPerPod && @@ -145,12 +146,22 @@ func (c *PoolCache) service() { break } } + // if specialized pod is available then return svc if found { req.responseChannel <- resp continue } - specializationInProgress := funcSvcGroup.svcWaiting - funcSvcGroup.queue.Len() - capacity := ((specializationInProgress + len(funcSvcGroup.svcs)) * req.requestsPerPod) - (totalActiveRequests + funcSvcGroup.svcWaiting) + concurrencyUsed := len(funcSvcGroup.svcs) + (funcSvcGroup.svcWaiting - funcSvcGroup.queue.Len()) + // if concurrency is available then be aggressive and use it as we are not sure if specialization will complete for other requests + if req.concurrency > 0 && concurrencyUsed < req.concurrency { + funcSvcGroup.svcWaiting++ + resp.error = ferror.MakeError(ferror.ErrorNotFound, fmt.Sprintf("function '%s' not found", req.function)) + req.responseChannel <- resp + continue + } + // if no concurrency is available then check if there is any virtual capacity in the existing pods to serve the request in future + // if specialization doesnt complete within request then request will be timeout + capacity := (concurrencyUsed * req.requestsPerPod) - (totalActiveRequests + funcSvcGroup.svcWaiting) if capacity > 0 { funcSvcGroup.svcWaiting++ svcWait := &svcWait{ @@ -164,8 +175,8 @@ func (c *PoolCache) service() { } // concurrency should not be set to zero and - //sum of specialization in progress and specialized pods should be less then req.concurrency - if req.concurrency > 0 && (specializationInProgress+len(funcSvcGroup.svcs)) >= req.concurrency { + // sum of specialization in progress and specialized pods should be less then req.concurrency + if req.concurrency > 0 && concurrencyUsed >= req.concurrency { resp.error = ferror.MakeError(ferror.ErrorTooManyRequests, fmt.Sprintf("function '%s' concurrency '%d' limit reached.", req.function, req.concurrency)) } else { funcSvcGroup.svcWaiting++ diff --git a/pkg/executor/fscache/poolcache_test.go b/pkg/executor/fscache/poolcache_test.go index db11bbbc..892bbc33 100644 --- a/pkg/executor/fscache/poolcache_test.go +++ b/pkg/executor/fscache/poolcache_test.go @@ -224,6 +224,7 @@ func TestPoolCacheRequests(t *testing.T) { simultaneous = 1 } for i := 1; i <= tt.requests; i++ { + reqno := i wg.Add(1) go func(reqno int) { defer wg.Done() @@ -231,10 +232,11 @@ func TestPoolCacheRequests(t *testing.T) { if err != nil { code, _ := ferror.GetHTTPError(err) if code == http.StatusNotFound { - p.SetSvcValue(context.Background(), key, fmt.Sprintf("svc-%d", svcCounter), &FuncSvc{ - Name: "value", - }, resource.MustParse("45m"), tt.rpp, tt.retainPods) atomic.AddUint64(&svcCounter, 1) + address := fmt.Sprintf("svc-%d", atomic.LoadUint64(&svcCounter)) + p.SetSvcValue(context.Background(), key, address, &FuncSvc{ + Name: address, + }, resource.MustParse("45m"), tt.rpp, tt.retainPods) } else { t.Log(reqno, "=>", err) atomic.AddUint64(&failedRequests, 1) @@ -244,9 +246,12 @@ func TestPoolCacheRequests(t *testing.T) { t.Log(reqno, "=>", "svc is nil") atomic.AddUint64(&failedRequests, 1) } + // } else { + // t.Log(reqno, "=>", svc.Name) + // } } - }(i) - if i%simultaneous == 0 { + }(reqno) + if reqno%simultaneous == 0 { wg.Wait() } } @@ -258,10 +263,11 @@ func TestPoolCacheRequests(t *testing.T) { for i := 0; i < tt.concurrency; i++ { for j := 0; j < tt.rpp; j++ { wg.Add(1) - go func(i int) { + svcno := i + go func(svcno int) { defer wg.Done() - p.MarkAvailable(key, fmt.Sprintf("svc-%d", i)) - }(i) + p.MarkAvailable(key, fmt.Sprintf("svc-%d", svcno+1)) + }(svcno) } } wg.Wait() @@ -270,8 +276,9 @@ func TestPoolCacheRequests(t *testing.T) { UID: "func", Generation: 2, } - p.SetSvcValue(context.Background(), newKey, fmt.Sprintf("svc-%d", svcCounter), &FuncSvc{ - Name: "value", + address := fmt.Sprintf("svc-%d", svcCounter) + p.SetSvcValue(context.Background(), newKey, address, &FuncSvc{ + Name: address, }, resource.MustParse("45m"), tt.rpp, tt.retainPods) funcSvc := p.ListAvailableValue() require.Equal(t, tt.concurrency, len(funcSvc))