Skip to content
Draft
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
89 changes: 87 additions & 2 deletions pkg/services/orgresolver/linking.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,10 +2,12 @@ package orgresolver

import (
"context"
"database/sql"
"errors"
"fmt"
"time"

"go.opentelemetry.io/otel/attribute"
"go.opentelemetry.io/otel/metric"
"google.golang.org/grpc"
"google.golang.org/grpc/credentials"
Expand All @@ -32,6 +34,15 @@ type OrgResolver interface {
Get(ctx context.Context, owner string) (string, error)
}

// CacheStore persists owner->orgID mappings to durable storage (e.g. Postgres).
// Implementations are provided by core and must be safe for concurrent use.
type CacheStore interface {
// GetOrg returns the cached orgID for owner, or sql.ErrNoRows if absent.
GetOrg(ctx context.Context, owner string) (string, error)
// UpsertOrg stores the owner->orgID mapping.
UpsertOrg(ctx context.Context, owner, orgID string) error
}

type Config struct {
URL string
TLSEnabled bool
Expand All @@ -43,6 +54,11 @@ type Config struct {

Client linkingclient.LinkingServiceClient // optional
Meter metric.Meter // optional

// CacheEnabled turns on durable caching of owner->orgID mappings via CacheStore.
CacheEnabled bool
// CacheStore is required when CacheEnabled is true.
CacheStore CacheStore
}

// orgResolver makes direct calls to the linking service to resolve organization IDs from workflow owners.
Expand All @@ -57,10 +73,26 @@ type orgResolver struct {
jwtGenerator JWTGenerator
requestTimeout time.Duration

passCount metric.Int64Counter
failCount metric.Int64Counter
cacheEnabled bool
cacheStore CacheStore

passCount metric.Int64Counter
failCount metric.Int64Counter
cacheLookups metric.Int64Counter // tagged with result=hit|miss|error
}

// cacheLookupResult is the attribute key for cache lookup outcomes.
var cacheResultAttr = "result"

const (
cacheResultHit = "hit"
cacheResultMiss = "miss"
)

// ErrCacheMiss is returned by CacheStore.GetOrg when no mapping exists for owner.
// Stores backed by SQL may return sql.ErrNoRows; both are treated as a miss.
var ErrCacheMiss = errors.New("org not found in cache")

// NewOrgResolver creates a new org resolver with the specified configuration
// Deprecated: Use Config.New
//
Expand All @@ -84,12 +116,18 @@ func (cfg *Config) New(logger log.Logger) (*orgResolver, error) {
requestTimeout = defaultRequestTimeout
}

if cfg.CacheEnabled && cfg.CacheStore == nil {
return nil, errors.New("CacheStore is required when CacheEnabled is true")
}

resolver := &orgResolver{
workflowRegistryAddress: cfg.WorkflowRegistryAddress,
workflowRegistryChainSelector: cfg.WorkflowRegistryChainSelector,
logger: log.Sugared(logger).Named("OrgResolver"),
jwtGenerator: cfg.JWTGenerator,
requestTimeout: requestTimeout,
cacheEnabled: cfg.CacheEnabled,
cacheStore: cfg.CacheStore,
}

if cfg.Client != nil {
Expand Down Expand Up @@ -125,6 +163,12 @@ func (cfg *Config) New(logger log.Logger) (*orgResolver, error) {
if err != nil {
return nil, fmt.Errorf("failed to create failure count metric: %w", err)
}
if resolver.cacheEnabled {
resolver.cacheLookups, err = cfg.Meter.Int64Counter("org_resolver_cache_lookups")
if err != nil {
return nil, fmt.Errorf("failed to create cache lookups metric: %w", err)
}
}
}

return resolver, nil
Expand All @@ -148,6 +192,12 @@ func (o *orgResolver) addJWTAuth(ctx context.Context, req any) (context.Context,
}

func (o *orgResolver) Get(ctx context.Context, owner string) (string, error) {
if o.cacheEnabled {
if orgID, ok := o.checkCache(ctx, owner); ok {
return orgID, nil
}
}

ctx, cancel := context.WithTimeout(ctx, o.requestTimeout)
defer cancel()

Expand All @@ -174,9 +224,44 @@ func (o *orgResolver) Get(ctx context.Context, owner string) (string, error) {
if o.passCount != nil {
o.passCount.Add(ctx, 1)
}

if o.cacheEnabled {
o.storeInCache(ctx, owner, resp.OrganizationId)
}
return resp.OrganizationId, nil
}

// checkCache looks up owner in the durable cache. Returns (orgID, true) on hit.
// A cache store error is logged and treated as a miss so lookups remain resilient.
func (o *orgResolver) checkCache(ctx context.Context, owner string) (string, bool) {
orgID, err := o.cacheStore.GetOrg(ctx, owner)
if err != nil {
if errors.Is(err, ErrCacheMiss) || errors.Is(err, sql.ErrNoRows) {
o.recordCacheLookup(ctx, cacheResultMiss)
} else {
o.logger.Warnw("Failed to read org from cache store, falling back to linking service", "owner", owner, "error", err)
o.recordCacheLookup(ctx, "error")
}
return "", false
}
o.recordCacheLookup(ctx, cacheResultHit)
return orgID, true
}

// storeInCache persists the owner->orgID mapping. Failures are logged but not
// propagated so a store hiccup does not break resolution.
func (o *orgResolver) storeInCache(ctx context.Context, owner, orgID string) {
if err := o.cacheStore.UpsertOrg(ctx, owner, orgID); err != nil {
o.logger.Warnw("Failed to persist org to cache store", "owner", owner, "error", err)
}
}

func (o *orgResolver) recordCacheLookup(ctx context.Context, result string) {
if o.cacheLookups != nil {
o.cacheLookups.Add(ctx, 1, metric.WithAttributes(attribute.String(cacheResultAttr, result)))
}
}

func (o *orgResolver) Start(_ context.Context) error {
return nil
}
Expand Down
Loading