Skip to content

Commit 097bbc5

Browse files
authored
fix: optimize customer usage attribution lookup (#4684)
1 parent 79501a8 commit 097bbc5

4 files changed

Lines changed: 449 additions & 17 deletions

File tree

openmeter/customer/adapter/customer.go

Lines changed: 76 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -512,22 +512,11 @@ func (a *adapter) GetCustomerByUsageAttribution(ctx context.Context, input custo
512512
now := clock.Now().UTC()
513513

514514
query := repo.db.Customer.Query().
515-
Where(customerdb.Namespace(input.Namespace)).
516515
Where(
517-
customerdb.Or(
518-
// We lookup the customer by subject key in the subjects table
519-
customerdb.HasSubjectsWith(
520-
customersubjectsdb.SubjectKey(input.Key),
521-
customersubjectsdb.Or(
522-
customersubjectsdb.DeletedAtIsNil(),
523-
customersubjectsdb.DeletedAtGT(now),
524-
),
525-
),
526-
// Or else we lookup the customer by key in the customers table
527-
customerdb.Key(input.Key),
528-
),
529-
).
530-
Where(customerdb.DeletedAtIsNil())
516+
customerdb.Namespace(input.Namespace),
517+
customerdb.DeletedAtIsNil(),
518+
customerMatchesUsageAttributionKey(input.Namespace, input.Key, now),
519+
)
531520
query = WithSubjects(query, now)
532521
if slices.Contains(input.Expands, customer.ExpandSubscriptions) {
533522
query = WithActiveSubscriptions(query, now)
@@ -552,6 +541,78 @@ func (a *adapter) GetCustomerByUsageAttribution(ctx context.Context, input custo
552541
})
553542
}
554543

544+
// customerMatchesUsageAttributionKey resolves a customer key before a subject
545+
// key while keeping both lookup branches independently indexable. The returned
546+
// candidate is applied to the outer customer query as:
547+
//
548+
// WHERE customers.id IN (
549+
// SELECT matches.id
550+
// FROM (
551+
// SELECT c.id, 0 AS lookup_priority
552+
// FROM customers AS c
553+
// WHERE c.namespace = $1
554+
// AND c.key = $2
555+
// AND c.deleted_at IS NULL
556+
//
557+
// UNION ALL
558+
//
559+
// SELECT c.id, 1 AS lookup_priority
560+
// FROM customer_subjects AS cs
561+
// JOIN customers AS c ON c.id = cs.customer_id
562+
// WHERE cs.namespace = $1
563+
// AND cs.subject_key = $2
564+
// AND (cs.deleted_at IS NULL OR cs.deleted_at > $3)
565+
// AND c.namespace = $1
566+
// AND c.deleted_at IS NULL
567+
// ) AS matches
568+
// ORDER BY matches.lookup_priority
569+
// LIMIT 1
570+
// )
571+
func customerMatchesUsageAttributionKey(namespace, key string, at time.Time) predicate.Customer {
572+
return func(s *sql.Selector) {
573+
keyCustomerTable := sql.Table(customerdb.Table).As("customer_by_key")
574+
customerKeyMatch := sql.Select(keyCustomerTable.C(customerdb.FieldID)).
575+
AppendSelectExprAs(sql.Expr("0"), "lookup_priority").
576+
From(keyCustomerTable).
577+
Where(sql.And(
578+
sql.EQ(keyCustomerTable.C(customerdb.FieldNamespace), namespace),
579+
sql.EQ(keyCustomerTable.C(customerdb.FieldKey), key),
580+
sql.IsNull(keyCustomerTable.C(customerdb.FieldDeletedAt)),
581+
))
582+
583+
customerSubjectsTable := sql.Table(customersubjectsdb.Table).As("customer_subjects")
584+
subjectCustomerTable := sql.Table(customerdb.Table).As("customer_by_subject")
585+
subjectKeyMatch := sql.Select(subjectCustomerTable.C(customerdb.FieldID)).
586+
AppendSelectExprAs(sql.Expr("1"), "lookup_priority").
587+
From(customerSubjectsTable).
588+
Join(subjectCustomerTable).
589+
On(
590+
subjectCustomerTable.C(customerdb.FieldID),
591+
customerSubjectsTable.C(customersubjectsdb.FieldCustomerID),
592+
).
593+
Where(sql.And(
594+
sql.EQ(customerSubjectsTable.C(customersubjectsdb.FieldNamespace), namespace),
595+
sql.EQ(customerSubjectsTable.C(customersubjectsdb.FieldSubjectKey), key),
596+
sql.Or(
597+
sql.IsNull(customerSubjectsTable.C(customersubjectsdb.FieldDeletedAt)),
598+
sql.GT(customerSubjectsTable.C(customersubjectsdb.FieldDeletedAt), at),
599+
),
600+
sql.EQ(subjectCustomerTable.C(customerdb.FieldNamespace), namespace),
601+
sql.IsNull(subjectCustomerTable.C(customerdb.FieldDeletedAt)),
602+
))
603+
604+
candidates := customerKeyMatch.
605+
UnionAll(subjectKeyMatch).
606+
As("usage_attribution_matches")
607+
preferredCustomerID := sql.Select(candidates.C(customerdb.FieldID)).
608+
From(candidates).
609+
OrderBy(candidates.C("lookup_priority")).
610+
Limit(1)
611+
612+
s.Where(sql.In(s.C(customerdb.FieldID), preferredCustomerID))
613+
}
614+
}
615+
555616
// GetCustomersByUsageAttribution resolves multiple customers by usage attribution keys in a single query.
556617
// A key matches a customer either by the customer's own key or by one of its subject keys, mirroring
557618
// the single-key GetCustomerByUsageAttribution. Keys that match no customer are simply absent from the

openmeter/customer/adapter/customer_test.go

Lines changed: 188 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -391,6 +391,194 @@ func customerIDs(customers []customer.Customer) []string {
391391
})
392392
}
393393

394+
func TestGetCustomerByUsageAttribution(t *testing.T) {
395+
t.Run("MatchByCustomerKey", func(t *testing.T) {
396+
env := newTestEnv(t)
397+
ns := ulid.Make().String()
398+
id := env.seedCustomerWithKey(ns, "cust-key", "subj-1")
399+
400+
got, err := env.adapter.GetCustomerByUsageAttribution(t.Context(), customer.GetCustomerByUsageAttributionInput{
401+
Namespace: ns,
402+
Key: "cust-key",
403+
})
404+
require.NoError(t, err)
405+
require.NotNil(t, got)
406+
assert.Equal(t, id, got.ID)
407+
require.NotNil(t, got.UsageAttribution)
408+
assert.Equal(t, []string{"subj-1"}, got.UsageAttribution.SubjectKeys)
409+
})
410+
411+
t.Run("MatchBySubjectKey", func(t *testing.T) {
412+
env := newTestEnv(t)
413+
ns := ulid.Make().String()
414+
id := env.seedCustomerWithKey(ns, "cust-key", "subj-1", "subj-2")
415+
416+
got, err := env.adapter.GetCustomerByUsageAttribution(t.Context(), customer.GetCustomerByUsageAttributionInput{
417+
Namespace: ns,
418+
Key: "subj-2",
419+
})
420+
require.NoError(t, err)
421+
require.NotNil(t, got)
422+
assert.Equal(t, id, got.ID)
423+
require.NotNil(t, got.UsageAttribution)
424+
assert.Equal(t, []string{"subj-1", "subj-2"}, got.UsageAttribution.SubjectKeys)
425+
})
426+
427+
t.Run("CustomerKeyTakesPriorityOverAnotherCustomerSubjectKey", func(t *testing.T) {
428+
env := newTestEnv(t)
429+
ns := ulid.Make().String()
430+
customerKeyID := env.seedCustomerWithKey(ns, "shared-key", "customer-key-subject")
431+
_ = env.seedCustomerWithKey(ns, "subject-owner", "shared-key")
432+
433+
got, err := env.adapter.GetCustomerByUsageAttribution(t.Context(), customer.GetCustomerByUsageAttributionInput{
434+
Namespace: ns,
435+
Key: "shared-key",
436+
})
437+
require.NoError(t, err)
438+
require.NotNil(t, got)
439+
assert.Equal(t, customerKeyID, got.ID, "a direct customer-key match must take priority over a subject-key match")
440+
})
441+
442+
t.Run("CustomerMatchedByOwnKeyAndSubjectKeyReturnedOnce", func(t *testing.T) {
443+
env := newTestEnv(t)
444+
ns := ulid.Make().String()
445+
id := env.seedCustomerWithKey(ns, "shared-key", "shared-key")
446+
447+
got, err := env.adapter.GetCustomerByUsageAttribution(t.Context(), customer.GetCustomerByUsageAttributionInput{
448+
Namespace: ns,
449+
Key: "shared-key",
450+
})
451+
require.NoError(t, err)
452+
require.NotNil(t, got)
453+
assert.Equal(t, id, got.ID)
454+
})
455+
456+
t.Run("SoftDeletedCustomerExcluded", func(t *testing.T) {
457+
env := newTestEnv(t)
458+
ns := ulid.Make().String()
459+
id := env.seedCustomerWithKey(ns, "cust-key", "subj-1")
460+
_ = freezeTime(t, time.Now())
461+
462+
require.NoError(t, env.adapter.DeleteCustomer(t.Context(), customer.DeleteCustomerInput{
463+
Namespace: ns,
464+
ID: id,
465+
}))
466+
467+
for _, key := range []string{"cust-key", "subj-1"} {
468+
_, err := env.adapter.GetCustomerByUsageAttribution(t.Context(), customer.GetCustomerByUsageAttributionInput{
469+
Namespace: ns,
470+
Key: key,
471+
})
472+
require.Error(t, err)
473+
assert.True(t, models.IsGenericNotFoundError(err))
474+
}
475+
})
476+
477+
t.Run("SubjectDeletedInPastExcluded", func(t *testing.T) {
478+
env := newTestEnv(t)
479+
ns := ulid.Make().String()
480+
id := env.seedCustomerWithKey(ns, "cust-key", "subj-1")
481+
now := freezeTime(t, time.Now())
482+
483+
_, err := env.db.CustomerSubjects.Update().
484+
Where(
485+
customersubjectsdb.Namespace(ns),
486+
customersubjectsdb.CustomerID(id),
487+
customersubjectsdb.SubjectKey("subj-1"),
488+
).
489+
SetDeletedAt(now.Add(-time.Minute)).
490+
Save(t.Context())
491+
require.NoError(t, err)
492+
493+
_, err = env.adapter.GetCustomerByUsageAttribution(t.Context(), customer.GetCustomerByUsageAttributionInput{
494+
Namespace: ns,
495+
Key: "subj-1",
496+
})
497+
require.Error(t, err)
498+
assert.True(t, models.IsGenericNotFoundError(err))
499+
})
500+
501+
t.Run("SubjectDeletedAtLookupTimeExcluded", func(t *testing.T) {
502+
env := newTestEnv(t)
503+
ns := ulid.Make().String()
504+
id := env.seedCustomerWithKey(ns, "cust-key", "subj-1")
505+
now := freezeTime(t, time.Now())
506+
507+
_, err := env.db.CustomerSubjects.Update().
508+
Where(
509+
customersubjectsdb.Namespace(ns),
510+
customersubjectsdb.CustomerID(id),
511+
customersubjectsdb.SubjectKey("subj-1"),
512+
).
513+
SetDeletedAt(now).
514+
Save(t.Context())
515+
require.NoError(t, err)
516+
517+
_, err = env.adapter.GetCustomerByUsageAttribution(t.Context(), customer.GetCustomerByUsageAttributionInput{
518+
Namespace: ns,
519+
Key: "subj-1",
520+
})
521+
require.Error(t, err)
522+
assert.True(t, models.IsGenericNotFoundError(err))
523+
})
524+
525+
t.Run("SubjectDeletedInFutureIncluded", func(t *testing.T) {
526+
env := newTestEnv(t)
527+
ns := ulid.Make().String()
528+
id := env.seedCustomerWithKey(ns, "cust-key", "subj-1")
529+
now := freezeTime(t, time.Now())
530+
531+
_, err := env.db.CustomerSubjects.Update().
532+
Where(
533+
customersubjectsdb.Namespace(ns),
534+
customersubjectsdb.CustomerID(id),
535+
customersubjectsdb.SubjectKey("subj-1"),
536+
).
537+
SetDeletedAt(now.Add(time.Minute)).
538+
Save(t.Context())
539+
require.NoError(t, err)
540+
541+
got, err := env.adapter.GetCustomerByUsageAttribution(t.Context(), customer.GetCustomerByUsageAttributionInput{
542+
Namespace: ns,
543+
Key: "subj-1",
544+
})
545+
require.NoError(t, err)
546+
require.NotNil(t, got)
547+
assert.Equal(t, id, got.ID)
548+
})
549+
550+
t.Run("NamespaceIsolation", func(t *testing.T) {
551+
env := newTestEnv(t)
552+
nsA := ulid.Make().String()
553+
nsB := ulid.Make().String()
554+
idA := env.seedCustomerWithKey(nsA, "shared-customer-key", "shared-subject-key")
555+
_ = env.seedCustomerWithKey(nsB, "shared-customer-key", "shared-subject-key")
556+
557+
for _, key := range []string{"shared-customer-key", "shared-subject-key"} {
558+
got, err := env.adapter.GetCustomerByUsageAttribution(t.Context(), customer.GetCustomerByUsageAttributionInput{
559+
Namespace: nsA,
560+
Key: key,
561+
})
562+
require.NoError(t, err)
563+
require.NotNil(t, got)
564+
assert.Equal(t, idA, got.ID)
565+
}
566+
})
567+
568+
t.Run("UnknownKeyReturnsNotFound", func(t *testing.T) {
569+
env := newTestEnv(t)
570+
ns := ulid.Make().String()
571+
_ = env.seedCustomerWithKey(ns, "cust-key", "subj-1")
572+
573+
_, err := env.adapter.GetCustomerByUsageAttribution(t.Context(), customer.GetCustomerByUsageAttributionInput{
574+
Namespace: ns,
575+
Key: "missing",
576+
})
577+
require.Error(t, err)
578+
assert.True(t, models.IsGenericNotFoundError(err))
579+
})
580+
}
581+
394582
func TestGetCustomersByUsageAttribution(t *testing.T) {
395583
t.Run("MatchByCustomerKey", func(t *testing.T) {
396584
env := newTestEnv(t)

0 commit comments

Comments
 (0)