Make common cache typed with generics (#2896)
Making typed common cache so that we don't use wrong types across set/get methods and more higher-level methods can be defined for cache. Currently, we are not able to operate over all keys of the cache due to generic types. I also removed code comments around the cache. Signed-off-by: Sanket Sudake <sanketsudake@gmail.com>
This commit is contained in:
@@ -30,6 +30,7 @@ import (
|
||||
|
||||
fv1 "github.com/fission/fission/pkg/apis/core/v1"
|
||||
"github.com/fission/fission/pkg/cache"
|
||||
"github.com/fission/fission/pkg/crd"
|
||||
"github.com/fission/fission/pkg/generated/clientset/versioned"
|
||||
"github.com/fission/fission/pkg/utils"
|
||||
"github.com/fission/fission/pkg/utils/manager"
|
||||
@@ -45,7 +46,7 @@ type (
|
||||
podInformer map[string]k8sCache.SharedIndexInformer
|
||||
pkgInformer map[string]k8sCache.SharedIndexInformer
|
||||
storageSvcUrl string
|
||||
buildCache *cache.Cache
|
||||
buildCache *cache.Cache[crd.CacheKeyUR, *fv1.Package]
|
||||
}
|
||||
)
|
||||
|
||||
@@ -60,13 +61,13 @@ func makePackageWatcher(logger *zap.Logger, fissionClient versioned.Interface, k
|
||||
podInformer: podInformer,
|
||||
pkgInformer: pkgInformer,
|
||||
storageSvcUrl: storageSvcUrl,
|
||||
buildCache: cache.MakeCache(0, 0),
|
||||
buildCache: cache.MakeCache[crd.CacheKeyUR, *fv1.Package](0, 0),
|
||||
}
|
||||
return pkgw
|
||||
}
|
||||
|
||||
func (pkgw *packageWatcher) buildCacheKey(obj metav1.ObjectMeta) string {
|
||||
return fmt.Sprintf("%s-%s-%s", obj.Namespace, obj.Name, obj.ResourceVersion)
|
||||
func (pkgw *packageWatcher) buildCacheKey(obj metav1.ObjectMeta) crd.CacheKeyUR {
|
||||
return crd.CacheKeyURFromMeta(&obj)
|
||||
}
|
||||
|
||||
func (pkgw *packageWatcher) buildWithCache(ctx context.Context, srcpkg *fv1.Package) {
|
||||
@@ -90,33 +91,33 @@ func (pkgw *packageWatcher) buildWithCache(ctx context.Context, srcpkg *fv1.Pack
|
||||
// 6. Update package status to succeed state
|
||||
// *. Update package status to failed state,if any one of steps above failed/time out
|
||||
func (pkgw *packageWatcher) build(ctx context.Context, srcpkg *fv1.Package) {
|
||||
defer func() {
|
||||
key := pkgw.buildCacheKey(srcpkg.ObjectMeta)
|
||||
logger := pkgw.logger.With(zap.String("package", srcpkg.Name), zap.String("namespace", srcpkg.Namespace), zap.String("resource_version", srcpkg.ResourceVersion), zap.String("key", key.String()))
|
||||
|
||||
defer func() {
|
||||
err := pkgw.buildCache.Delete(key)
|
||||
if err != nil {
|
||||
pkgw.logger.Error("error deleting key from cache", zap.String("key", key), zap.Error(err))
|
||||
logger.Error("error deleting key from cache", zap.Any("key", key), zap.Error(err))
|
||||
}
|
||||
}()
|
||||
|
||||
pkgw.logger.Info("starting build for package", zap.String("package_name", srcpkg.ObjectMeta.Name), zap.String("resource_version", srcpkg.ObjectMeta.ResourceVersion))
|
||||
logger.Info("starting build for package")
|
||||
|
||||
pkg, err := updatePackage(ctx, pkgw.logger, pkgw.fissionClient, srcpkg, fv1.BuildStatusRunning, "", nil)
|
||||
pkg, err := updatePackage(ctx, logger, pkgw.fissionClient, srcpkg, fv1.BuildStatusRunning, "", nil)
|
||||
if err != nil {
|
||||
pkgw.logger.Error("error setting package pending state", zap.Error(err))
|
||||
logger.Error("error setting package pending state", zap.Error(err))
|
||||
return
|
||||
}
|
||||
|
||||
env, err := pkgw.fissionClient.CoreV1().Environments(pkg.Spec.Environment.Namespace).Get(ctx, pkg.Spec.Environment.Name, metav1.GetOptions{})
|
||||
if k8serrors.IsNotFound(err) {
|
||||
e := "environment does not exist"
|
||||
pkgw.logger.Error(e, zap.String("environment", pkg.Spec.Environment.Name))
|
||||
_, er := updatePackage(ctx, pkgw.logger, pkgw.fissionClient, pkg,
|
||||
logger.Error(e, zap.String("environment", pkg.Spec.Environment.Name))
|
||||
_, er := updatePackage(ctx, logger, pkgw.fissionClient, pkg,
|
||||
fv1.BuildStatusFailed, fmt.Sprintf("%s: %q", e, pkg.Spec.Environment.Name), nil)
|
||||
if er != nil {
|
||||
pkgw.logger.Error(
|
||||
logger.Error(
|
||||
"error updating package",
|
||||
zap.String("package_name", pkg.ObjectMeta.Name),
|
||||
zap.String("resource_version", pkg.ObjectMeta.ResourceVersion),
|
||||
zap.Error(er),
|
||||
)
|
||||
}
|
||||
@@ -127,6 +128,8 @@ func (pkgw *packageWatcher) build(ctx context.Context, srcpkg *fv1.Package) {
|
||||
healthCheckBackOff := utils.NewDefaultBackOff()
|
||||
builderNs := pkgw.nsResolver.GetBuilderNS(env.ObjectMeta.Namespace)
|
||||
|
||||
logger = logger.With(zap.String("environment", env.Name), zap.String("builder_namespace", builderNs), zap.String("environment_namespace", env.Namespace))
|
||||
|
||||
// if err != nil {
|
||||
// pkgw.logger.Error("Unable to create BackOff for Health Check", zap.Error(err))
|
||||
//}
|
||||
@@ -136,12 +139,12 @@ func (pkgw *packageWatcher) build(ctx context.Context, srcpkg *fv1.Package) {
|
||||
// iterate all available environment builders.
|
||||
items := pkgw.podInformer[builderNs].GetStore().List()
|
||||
if err != nil {
|
||||
pkgw.logger.Error("error retrieving pod information for environment", zap.Error(err), zap.String("environment", env.ObjectMeta.Name))
|
||||
logger.Error("error retrieving pod information for environment", zap.Error(err))
|
||||
return
|
||||
}
|
||||
|
||||
if len(items) == 0 {
|
||||
pkgw.logger.Info("builder pod does not exist for environment, will retry again later", zap.String("environment", pkg.Spec.Environment.Name))
|
||||
logger.Info("builder pod does not exist for environment, will retry again later")
|
||||
time.Sleep(healthCheckBackOff.GetCurrentBackoffDuration())
|
||||
continue
|
||||
}
|
||||
@@ -165,27 +168,22 @@ func (pkgw *packageWatcher) build(ctx context.Context, srcpkg *fv1.Package) {
|
||||
}
|
||||
|
||||
if !podIsReady {
|
||||
pkgw.logger.Info("builder pod is not ready for environment, will retry again later", zap.String("environment", pkg.Spec.Environment.Name))
|
||||
logger.Info("builder pod is not ready for environment, will retry again later")
|
||||
time.Sleep(healthCheckBackOff.GetCurrentBackoffDuration())
|
||||
break
|
||||
}
|
||||
|
||||
uploadResp, buildLogs, err := buildPackage(ctx, pkgw.logger, pkgw.fissionClient, builderNs, pkgw.storageSvcUrl, pkg)
|
||||
if err != nil {
|
||||
pkgw.logger.Error("error building package", zap.Error(err), zap.String("package_name", pkg.ObjectMeta.Name))
|
||||
_, er := updatePackage(ctx, pkgw.logger, pkgw.fissionClient, pkg, fv1.BuildStatusFailed, buildLogs, nil)
|
||||
logger.Error("error building package", zap.Error(err))
|
||||
_, er := updatePackage(ctx, logger, pkgw.fissionClient, pkg, fv1.BuildStatusFailed, buildLogs, nil)
|
||||
if er != nil {
|
||||
pkgw.logger.Error(
|
||||
"error updating package",
|
||||
zap.String("package_name", pkg.ObjectMeta.Name),
|
||||
zap.String("resource_version", pkg.ObjectMeta.ResourceVersion),
|
||||
zap.Error(er),
|
||||
)
|
||||
logger.Error("error updating package", zap.Error(er))
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
pkgw.logger.Info("starting package info update", zap.String("package_name", pkg.ObjectMeta.Name))
|
||||
logger.Info("starting package info update")
|
||||
|
||||
fnList, err := pkgw.fissionClient.CoreV1().
|
||||
Functions(pkg.Namespace).List(ctx, metav1.ListOptions{})
|
||||
@@ -197,8 +195,6 @@ func (pkgw *packageWatcher) build(ctx context.Context, srcpkg *fv1.Package) {
|
||||
if er != nil {
|
||||
pkgw.logger.Error(
|
||||
"error updating package",
|
||||
zap.String("package_name", pkg.ObjectMeta.Name),
|
||||
zap.String("resource_version", pkg.ObjectMeta.ResourceVersion),
|
||||
zap.Error(er),
|
||||
)
|
||||
}
|
||||
@@ -215,57 +211,41 @@ func (pkgw *packageWatcher) build(ctx context.Context, srcpkg *fv1.Package) {
|
||||
_, err = pkgw.fissionClient.CoreV1().Functions(fn.ObjectMeta.Namespace).Update(ctx, &fn, metav1.UpdateOptions{})
|
||||
if err != nil {
|
||||
e := "error updating function package resource version"
|
||||
pkgw.logger.Error(e, zap.Error(err))
|
||||
logger.Error(e, zap.Error(err))
|
||||
buildLogs += fmt.Sprintf("%s: %v\n", e, err)
|
||||
_, er := updatePackage(ctx, pkgw.logger, pkgw.fissionClient, pkg, fv1.BuildStatusFailed, buildLogs, nil)
|
||||
_, er := updatePackage(ctx, logger, pkgw.fissionClient, pkg, fv1.BuildStatusFailed, buildLogs, nil)
|
||||
if er != nil {
|
||||
pkgw.logger.Error(
|
||||
"error updating package",
|
||||
zap.String("package_name", pkg.ObjectMeta.Name),
|
||||
zap.String("resource_version", pkg.ObjectMeta.ResourceVersion),
|
||||
zap.Error(er),
|
||||
)
|
||||
logger.Error("error updating package", zap.Error(er))
|
||||
}
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
_, err = updatePackage(ctx, pkgw.logger, pkgw.fissionClient, pkg,
|
||||
_, err = updatePackage(ctx, logger, pkgw.fissionClient, pkg,
|
||||
fv1.BuildStatusSucceeded, buildLogs, uploadResp)
|
||||
if err != nil {
|
||||
pkgw.logger.Error("error updating package info", zap.Error(err), zap.String("package_name", pkg.ObjectMeta.Name))
|
||||
_, er := updatePackage(ctx, pkgw.logger, pkgw.fissionClient, pkg, fv1.BuildStatusFailed, buildLogs, nil)
|
||||
logger.Error("error updating package info", zap.Error(err))
|
||||
_, er := updatePackage(ctx, logger, pkgw.fissionClient, pkg, fv1.BuildStatusFailed, buildLogs, nil)
|
||||
if er != nil {
|
||||
pkgw.logger.Error(
|
||||
"error updating package",
|
||||
zap.String("package_name", pkg.ObjectMeta.Name),
|
||||
zap.String("resource_version", pkg.ObjectMeta.ResourceVersion),
|
||||
zap.Error(er),
|
||||
)
|
||||
logger.Error("error updating package", zap.Error(er))
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
pkgw.logger.Info("completed package build request", zap.String("package_name", pkg.ObjectMeta.Name))
|
||||
logger.Info("completed package build request")
|
||||
return
|
||||
}
|
||||
time.Sleep(healthCheckBackOff.GetNext())
|
||||
}
|
||||
// build timeout
|
||||
_, err = updatePackage(ctx, pkgw.logger, pkgw.fissionClient, pkg,
|
||||
_, err = updatePackage(ctx, logger, pkgw.fissionClient, pkg,
|
||||
fv1.BuildStatusFailed, "Build timeout due to environment builder not ready", nil)
|
||||
if err != nil {
|
||||
pkgw.logger.Error(
|
||||
"error updating package",
|
||||
zap.String("package_name", pkg.ObjectMeta.Name),
|
||||
zap.String("resource_version", pkg.ObjectMeta.ResourceVersion),
|
||||
zap.Error(err),
|
||||
)
|
||||
logger.Error("error updating package", zap.Error(err))
|
||||
}
|
||||
|
||||
pkgw.logger.Error("max retries exceeded in building source package, timeout due to environment builder not ready",
|
||||
zap.String("package", fmt.Sprintf("%s.%s", pkg.ObjectMeta.Name, pkg.ObjectMeta.Namespace)))
|
||||
logger.Error("max retries exceeded in building source package, timeout due to environment builder not ready")
|
||||
}
|
||||
|
||||
func (pkgw *packageWatcher) packageInformerHandler(ctx context.Context) k8sCache.ResourceEventHandlerFuncs {
|
||||
|
||||
Vendored
+36
-36
@@ -34,33 +34,33 @@ const (
|
||||
)
|
||||
|
||||
type (
|
||||
Value struct {
|
||||
Value[V any] struct {
|
||||
ctime time.Time
|
||||
atime time.Time
|
||||
value interface{}
|
||||
value V
|
||||
}
|
||||
Cache struct {
|
||||
cache map[interface{}]*Value
|
||||
Cache[K comparable, V any] struct {
|
||||
cache map[K]*Value[V]
|
||||
ctimeExpiry time.Duration
|
||||
atimeExpiry time.Duration
|
||||
requestChannel chan *request
|
||||
requestChannel chan *request[K, V]
|
||||
}
|
||||
|
||||
request struct {
|
||||
request[K comparable, V any] struct {
|
||||
requestType
|
||||
key interface{}
|
||||
value interface{}
|
||||
responseChannel chan *response
|
||||
key K
|
||||
value V
|
||||
responseChannel chan *response[K, V]
|
||||
}
|
||||
response struct {
|
||||
response[K comparable, V any] struct {
|
||||
error
|
||||
existingValue interface{}
|
||||
mapCopy map[interface{}]interface{}
|
||||
value interface{}
|
||||
existingValue V
|
||||
mapCopy map[K]V
|
||||
value V
|
||||
}
|
||||
)
|
||||
|
||||
func (c *Cache) IsOld(v *Value) bool {
|
||||
func (c *Cache[K, V]) IsOld(v *Value[V]) bool {
|
||||
if (c.ctimeExpiry != time.Duration(0)) && (time.Since(v.ctime) > c.ctimeExpiry) {
|
||||
return true
|
||||
}
|
||||
@@ -72,12 +72,12 @@ func (c *Cache) IsOld(v *Value) bool {
|
||||
return false
|
||||
}
|
||||
|
||||
func MakeCache(ctimeExpiry, atimeExpiry time.Duration) *Cache {
|
||||
c := &Cache{
|
||||
cache: make(map[interface{}]*Value),
|
||||
func MakeCache[K comparable, V any](ctimeExpiry, atimeExpiry time.Duration) *Cache[K, V] {
|
||||
c := &Cache[K, V]{
|
||||
cache: make(map[K]*Value[V]),
|
||||
ctimeExpiry: ctimeExpiry,
|
||||
atimeExpiry: atimeExpiry,
|
||||
requestChannel: make(chan *request),
|
||||
requestChannel: make(chan *request[K, V]),
|
||||
}
|
||||
go c.service()
|
||||
if ctimeExpiry != time.Duration(0) || atimeExpiry != time.Duration(0) {
|
||||
@@ -86,10 +86,10 @@ func MakeCache(ctimeExpiry, atimeExpiry time.Duration) *Cache {
|
||||
return c
|
||||
}
|
||||
|
||||
func (c *Cache) service() {
|
||||
func (c *Cache[K, V]) service() {
|
||||
for {
|
||||
req := <-c.requestChannel
|
||||
resp := &response{}
|
||||
resp := &response[K, V]{}
|
||||
switch req.requestType {
|
||||
case GET:
|
||||
val, ok := c.cache[req.key]
|
||||
@@ -115,7 +115,7 @@ func (c *Cache) service() {
|
||||
resp.existingValue = val.value
|
||||
resp.error = ferror.MakeError(ferror.ErrorNameExists, "key already exists")
|
||||
} else {
|
||||
c.cache[req.key] = &Value{
|
||||
c.cache[req.key] = &Value[V]{
|
||||
value: req.value,
|
||||
ctime: now,
|
||||
atime: now,
|
||||
@@ -133,7 +133,7 @@ func (c *Cache) service() {
|
||||
}
|
||||
// no response
|
||||
case COPY:
|
||||
resp.mapCopy = make(map[interface{}]interface{})
|
||||
resp.mapCopy = make(map[K]V)
|
||||
for k, v := range c.cache {
|
||||
resp.mapCopy[k] = v.value
|
||||
}
|
||||
@@ -146,9 +146,9 @@ func (c *Cache) service() {
|
||||
}
|
||||
}
|
||||
|
||||
func (c *Cache) Get(key interface{}) (interface{}, error) {
|
||||
respChannel := make(chan *response)
|
||||
c.requestChannel <- &request{
|
||||
func (c *Cache[K, V]) Get(key K) (V, error) {
|
||||
respChannel := make(chan *response[K, V])
|
||||
c.requestChannel <- &request[K, V]{
|
||||
requestType: GET,
|
||||
key: key,
|
||||
responseChannel: respChannel,
|
||||
@@ -159,9 +159,9 @@ func (c *Cache) Get(key interface{}) (interface{}, error) {
|
||||
|
||||
// if key exists in the cache, the new value is NOT set; instead an
|
||||
// error and the old value are returned
|
||||
func (c *Cache) Set(key interface{}, value interface{}) (interface{}, error) {
|
||||
respChannel := make(chan *response)
|
||||
c.requestChannel <- &request{
|
||||
func (c *Cache[K, V]) Set(key K, value V) (V, error) {
|
||||
respChannel := make(chan *response[K, V])
|
||||
c.requestChannel <- &request[K, V]{
|
||||
requestType: SET,
|
||||
key: key,
|
||||
value: value,
|
||||
@@ -171,9 +171,9 @@ func (c *Cache) Set(key interface{}, value interface{}) (interface{}, error) {
|
||||
return resp.existingValue, resp.error
|
||||
}
|
||||
|
||||
func (c *Cache) Delete(key interface{}) error {
|
||||
respChannel := make(chan *response)
|
||||
c.requestChannel <- &request{
|
||||
func (c *Cache[K, V]) Delete(key K) error {
|
||||
respChannel := make(chan *response[K, V])
|
||||
c.requestChannel <- &request[K, V]{
|
||||
requestType: DELETE,
|
||||
key: key,
|
||||
responseChannel: respChannel,
|
||||
@@ -182,9 +182,9 @@ func (c *Cache) Delete(key interface{}) error {
|
||||
return resp.error
|
||||
}
|
||||
|
||||
func (c *Cache) Copy() map[interface{}]interface{} {
|
||||
respChannel := make(chan *response)
|
||||
c.requestChannel <- &request{
|
||||
func (c *Cache[K, V]) Copy() map[K]V {
|
||||
respChannel := make(chan *response[K, V])
|
||||
c.requestChannel <- &request[K, V]{
|
||||
requestType: COPY,
|
||||
responseChannel: respChannel,
|
||||
}
|
||||
@@ -192,10 +192,10 @@ func (c *Cache) Copy() map[interface{}]interface{} {
|
||||
return resp.mapCopy
|
||||
}
|
||||
|
||||
func (c *Cache) expiryService() {
|
||||
func (c *Cache[K, V]) expiryService() {
|
||||
for {
|
||||
time.Sleep(time.Minute)
|
||||
c.requestChannel <- &request{
|
||||
c.requestChannel <- &request[K, V]{
|
||||
requestType: EXPIRE,
|
||||
}
|
||||
}
|
||||
|
||||
Vendored
+1
-1
@@ -29,7 +29,7 @@ func checkErr(err error) {
|
||||
}
|
||||
|
||||
func TestCache(t *testing.T) {
|
||||
c := MakeCache(100*time.Millisecond, 100*time.Millisecond)
|
||||
c := MakeCache[string, string](100*time.Millisecond, 100*time.Millisecond)
|
||||
|
||||
_, err := c.Set("a", "b")
|
||||
checkErr(err)
|
||||
|
||||
@@ -27,7 +27,7 @@ import (
|
||||
|
||||
type (
|
||||
canaryConfigCancelFuncMap struct {
|
||||
cache *cache.Cache // map[metadataKey]*context.Context
|
||||
cache *cache.Cache[metadataKey, *CanaryProcessingInfo]
|
||||
}
|
||||
|
||||
// metav1.ObjectMeta is not hashable, so we make a hashable copy
|
||||
@@ -45,7 +45,7 @@ type (
|
||||
|
||||
func makecanaryConfigCancelFuncMap() *canaryConfigCancelFuncMap {
|
||||
return &canaryConfigCancelFuncMap{
|
||||
cache: cache.MakeCache(0, 0),
|
||||
cache: cache.MakeCache[metadataKey, *CanaryProcessingInfo](0, 0),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -62,8 +62,7 @@ func (cancelFuncMap *canaryConfigCancelFuncMap) lookup(f *metav1.ObjectMeta) (*C
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
value := item.(*CanaryProcessingInfo)
|
||||
return value, nil
|
||||
return item, nil
|
||||
}
|
||||
|
||||
func (cancelFuncMap *canaryConfigCancelFuncMap) assign(f *metav1.ObjectMeta, value *CanaryProcessingInfo) error {
|
||||
|
||||
@@ -78,7 +78,7 @@ type (
|
||||
nsResolver *utils.NamespaceResolver
|
||||
|
||||
fissionClient versioned.Interface
|
||||
functionEnv *cache.Cache
|
||||
functionEnv *cache.Cache[crd.CacheKeyUR, *fv1.Environment]
|
||||
fsCache *fscache.FunctionServiceCache
|
||||
instanceID string
|
||||
requestChannel chan *request
|
||||
@@ -147,7 +147,7 @@ func MakeGenericPoolManager(ctx context.Context,
|
||||
nsResolver: utils.DefaultNSResolver(),
|
||||
metricsClient: metricsClient,
|
||||
fissionClient: fissionClient,
|
||||
functionEnv: cache.MakeCache(10*time.Second, 0),
|
||||
functionEnv: cache.MakeCache[crd.CacheKeyUR, *fv1.Environment](10*time.Second, 0),
|
||||
fsCache: fscache.MakeFunctionServiceCache(gpmLogger),
|
||||
instanceID: instanceID,
|
||||
requestChannel: make(chan *request),
|
||||
@@ -595,8 +595,7 @@ func (gpm *GenericPoolManager) getFunctionEnv(ctx context.Context, fn *fv1.Funct
|
||||
// TODO: the cache should be able to search by <env name, fn namespace> instead of function metadata.
|
||||
result, err := gpm.functionEnv.Get(crd.CacheKeyURFromMeta(&fn.ObjectMeta))
|
||||
if err == nil {
|
||||
env = result.(*fv1.Environment)
|
||||
return env, nil
|
||||
return result, nil
|
||||
}
|
||||
|
||||
// Get env from controller
|
||||
|
||||
@@ -68,9 +68,9 @@ type (
|
||||
// FunctionServiceCache represents the function service cache
|
||||
FunctionServiceCache struct {
|
||||
logger *zap.Logger
|
||||
byFunction *cache.Cache // function-key -> funcSvc : map[string]*funcSvc
|
||||
byAddress *cache.Cache // address -> function : map[string]metav1.ObjectMeta
|
||||
byFunctionUID *cache.Cache // function uid -> function : map[string]metav1.ObjectMeta
|
||||
byFunction *cache.Cache[crd.CacheKeyUR, *FuncSvc]
|
||||
byAddress *cache.Cache[string, metav1.ObjectMeta]
|
||||
byFunctionUID *cache.Cache[types.UID, metav1.ObjectMeta]
|
||||
connFunctionCache *PoolCache // function-key -> funcSvc : map[string]*funcSvc
|
||||
PodToFsvc sync.Map // pod-name -> funcSvc: map[string]*FuncSvc
|
||||
WebsocketFsvc sync.Map // funcSvc-name -> bool: map[string]bool
|
||||
@@ -110,9 +110,9 @@ func IsNameExistError(err error) bool {
|
||||
func MakeFunctionServiceCache(logger *zap.Logger) *FunctionServiceCache {
|
||||
fsc := &FunctionServiceCache{
|
||||
logger: logger.Named("function_service_cache"),
|
||||
byFunction: cache.MakeCache(0, 0),
|
||||
byAddress: cache.MakeCache(0, 0),
|
||||
byFunctionUID: cache.MakeCache(0, 0),
|
||||
byFunction: cache.MakeCache[crd.CacheKeyUR, *FuncSvc](0, 0),
|
||||
byAddress: cache.MakeCache[string, metav1.ObjectMeta](0, 0),
|
||||
byFunctionUID: cache.MakeCache[types.UID, metav1.ObjectMeta](0, 0),
|
||||
connFunctionCache: NewPoolCache(logger.Named("conn_function_cache")),
|
||||
requestChannel: make(chan *fscRequest),
|
||||
}
|
||||
@@ -132,14 +132,12 @@ func (fsc *FunctionServiceCache) service() {
|
||||
// get svcs idle for > req.age
|
||||
fscs := fsc.byFunctionUID.Copy()
|
||||
funcObjects := make([]*FuncSvc, 0)
|
||||
for _, funcSvc := range fscs {
|
||||
mI := funcSvc.(metav1.ObjectMeta)
|
||||
fsvcI, err := fsc.byFunction.Get(crd.CacheKeyURFromMeta(&mI))
|
||||
for _, m := range fscs {
|
||||
fsvc, err := fsc.byFunction.Get(crd.CacheKeyURFromMeta(&m))
|
||||
if err != nil {
|
||||
fsc.logger.Error("error while getting service", zap.String("error", err.Error()))
|
||||
return
|
||||
}
|
||||
fsvc := fsvcI.(*FuncSvc)
|
||||
if time.Since(fsvc.Atime) > req.age {
|
||||
funcObjects = append(funcObjects, fsvc)
|
||||
}
|
||||
@@ -149,8 +147,7 @@ func (fsc *FunctionServiceCache) service() {
|
||||
fsc.logger.Info("dumping function service cache")
|
||||
funcCopy := fsc.byFunction.Copy()
|
||||
info := []string{}
|
||||
for key, fsvcI := range funcCopy {
|
||||
fsvc := fsvcI.(*FuncSvc)
|
||||
for key, fsvc := range funcCopy {
|
||||
for _, kubeObj := range fsvc.KubernetesObjects {
|
||||
info = append(info, fmt.Sprintf("%v\t%v\t%v", key, kubeObj.Kind, kubeObj.Name))
|
||||
}
|
||||
@@ -196,13 +193,12 @@ func (fsc *FunctionServiceCache) DumpDebugInfo(ctx context.Context) error {
|
||||
func (fsc *FunctionServiceCache) GetByFunction(m *metav1.ObjectMeta) (*FuncSvc, error) {
|
||||
key := crd.CacheKeyURFromMeta(m)
|
||||
|
||||
fsvcI, err := fsc.byFunction.Get(key)
|
||||
fsvc, err := fsc.byFunction.Get(key)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// update atime
|
||||
fsvc := fsvcI.(*FuncSvc)
|
||||
fsvc.Atime = time.Now()
|
||||
|
||||
fsvcCopy := *fsvc
|
||||
@@ -228,20 +224,17 @@ func (fsc *FunctionServiceCache) GetFuncSvc(ctx context.Context, m *metav1.Objec
|
||||
|
||||
// GetByFunctionUID gets a function service from cache using function UUID.
|
||||
func (fsc *FunctionServiceCache) GetByFunctionUID(uid types.UID) (*FuncSvc, error) {
|
||||
mI, err := fsc.byFunctionUID.Get(uid)
|
||||
m, err := fsc.byFunctionUID.Get(uid)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
m := mI.(metav1.ObjectMeta)
|
||||
|
||||
fsvcI, err := fsc.byFunction.Get(crd.CacheKeyURFromMeta(&m))
|
||||
fsvc, err := fsc.byFunction.Get(crd.CacheKeyURFromMeta(&m))
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// update atime
|
||||
fsvc := fsvcI.(*FuncSvc)
|
||||
fsvc.Atime = time.Now()
|
||||
|
||||
fsvcCopy := *fsvc
|
||||
@@ -279,12 +272,11 @@ func (fsc *FunctionServiceCache) Add(fsvc FuncSvc) (*FuncSvc, error) {
|
||||
existing, err := fsc.byFunction.Set(crd.CacheKeyURFromMeta(fsvc.Function), &fsvc)
|
||||
if err != nil {
|
||||
if IsNameExistError(err) {
|
||||
f := existing.(*FuncSvc)
|
||||
err2 := fsc.TouchByAddress(f.Address)
|
||||
err2 := fsc.TouchByAddress(existing.Address)
|
||||
if err2 != nil {
|
||||
return nil, err2
|
||||
}
|
||||
fCopy := *f
|
||||
fCopy := *existing
|
||||
return &fCopy, nil
|
||||
}
|
||||
return nil, err
|
||||
@@ -333,16 +325,14 @@ func (fsc *FunctionServiceCache) TouchByAddress(address string) error {
|
||||
}
|
||||
|
||||
func (fsc *FunctionServiceCache) _touchByAddress(address string) error {
|
||||
mI, err := fsc.byAddress.Get(address)
|
||||
m, err := fsc.byAddress.Get(address)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
m := mI.(metav1.ObjectMeta)
|
||||
fsvcI, err := fsc.byFunction.Get(crd.CacheKeyURFromMeta(&m))
|
||||
fsvc, err := fsc.byFunction.Get(crd.CacheKeyURFromMeta(&m))
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
fsvc := fsvcI.(*FuncSvc)
|
||||
fsvc.Atime = time.Now()
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -34,7 +34,7 @@ type (
|
||||
// reference into a resolveResult
|
||||
functionReferenceResolver struct {
|
||||
// FunctionReference -> function metadata
|
||||
refCache *cache.Cache
|
||||
refCache *cache.Cache[namespacedTriggerReference, resolveResult]
|
||||
funcInformer map[string]k8sCache.SharedIndexInformer
|
||||
logger *zap.Logger
|
||||
// store k8sCache.Store
|
||||
@@ -73,7 +73,7 @@ const (
|
||||
|
||||
func makeFunctionReferenceResolver(logger *zap.Logger, funcInformer map[string]k8sCache.SharedIndexInformer) *functionReferenceResolver {
|
||||
frr := &functionReferenceResolver{
|
||||
refCache: cache.MakeCache(time.Minute, 0),
|
||||
refCache: cache.MakeCache[namespacedTriggerReference, resolveResult](time.Minute, 0),
|
||||
funcInformer: funcInformer,
|
||||
logger: logger.Named("function_ref_resolver"),
|
||||
}
|
||||
@@ -89,9 +89,8 @@ func (frr *functionReferenceResolver) resolve(trigger fv1.HTTPTrigger) (*resolve
|
||||
}
|
||||
|
||||
// check cache
|
||||
rrInt, err := frr.refCache.Get(nfr)
|
||||
result, err := frr.refCache.Get(nfr)
|
||||
if err == nil {
|
||||
result := rrInt.(resolveResult)
|
||||
return &result, nil
|
||||
}
|
||||
|
||||
@@ -216,11 +215,5 @@ func (frr *functionReferenceResolver) delete(namespace string, triggerName, trig
|
||||
}
|
||||
|
||||
func (frr *functionReferenceResolver) copy() map[namespacedTriggerReference]resolveResult {
|
||||
cache := make(map[namespacedTriggerReference]resolveResult)
|
||||
for k, v := range frr.refCache.Copy() {
|
||||
key := k.(namespacedTriggerReference)
|
||||
val := v.(resolveResult)
|
||||
cache[key] = val
|
||||
}
|
||||
return cache
|
||||
return frr.refCache.Copy()
|
||||
}
|
||||
|
||||
@@ -29,7 +29,7 @@ import (
|
||||
type (
|
||||
functionServiceMap struct {
|
||||
logger *zap.Logger
|
||||
cache *cache.Cache // map[metadataKey]*url.URL
|
||||
cache *cache.Cache[metadataKey, *url.URL]
|
||||
}
|
||||
|
||||
// metav1.ObjectMeta is not hashable, so we make a hashable copy
|
||||
@@ -44,7 +44,7 @@ type (
|
||||
func makeFunctionServiceMap(logger *zap.Logger, expiry time.Duration) *functionServiceMap {
|
||||
return &functionServiceMap{
|
||||
logger: logger.Named("function_service_map"),
|
||||
cache: cache.MakeCache(expiry, 0),
|
||||
cache: cache.MakeCache[metadataKey, *url.URL](expiry, 0),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -62,15 +62,14 @@ func (fmap *functionServiceMap) lookup(f *metav1.ObjectMeta) (*url.URL, error) {
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
u := item.(*url.URL)
|
||||
return u, nil
|
||||
return item, nil
|
||||
}
|
||||
|
||||
func (fmap *functionServiceMap) assign(f *metav1.ObjectMeta, serviceURL *url.URL) {
|
||||
mk := keyFromMetadata(f)
|
||||
old, err := fmap.cache.Set(*mk, serviceURL)
|
||||
if err != nil {
|
||||
if *serviceURL == *(old.(*url.URL)) {
|
||||
if *serviceURL == *old {
|
||||
return
|
||||
}
|
||||
fmap.logger.Error("error caching service url for function with a different value", zap.Error(err))
|
||||
|
||||
Reference in New Issue
Block a user