Skip to content
Merged
Show file tree
Hide file tree
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
79 changes: 53 additions & 26 deletions internal/proxy/codex_capacity.go
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@ const (
CapacitySourceUsageCache CapacitySource = iota + 1
CapacitySourceHardLimit
CapacitySourceHTTPHeaders
CapacitySourceLiveUsage
CapacitySourceLiveRateLimits
)

Expand Down Expand Up @@ -80,6 +81,13 @@ type capacityBucketKey struct {
bucket CapacityBucket
}

type capacitySnapshotAggregate struct {
windows map[quota.WindowName]quota.Window
remaining int
reset time.Time
set bool
}

// CodexCapacityObservationStream orders facts from one upstream response or connection.
type CodexCapacityObservationStream struct {
generation uint64
Expand Down Expand Up @@ -191,7 +199,7 @@ func (l *CodexCapacityLedger) updateHardFenceState(key capacityFactKey, fact Cap
case CapacitySourceHardLimit:
live, ok := l.livePositiveHighWater[bucketKey]
l.suppressedHardFences[key] = ok && liveFactLiftsHardFence(live, fact)
case CapacitySourceLiveRateLimits:
case CapacitySourceLiveUsage, CapacitySourceLiveRateLimits:
if fact.Confidence != CapacityConfidenceAuthoritative || fact.RemainingPct <= 0 {
return
}
Expand Down Expand Up @@ -226,7 +234,7 @@ func validCapacityFact(fact CapacityFact) bool {
return fact.Confidence == CapacityConfidenceAdvisory && fact.ConnectionGeneration == 0
case CapacitySourceHardLimit:
return fact.Confidence == CapacityConfidenceAuthoritative && fact.ConnectionGeneration > 0 && fact.RemainingPct == 0
case CapacitySourceHTTPHeaders, CapacitySourceLiveRateLimits:
case CapacitySourceHTTPHeaders, CapacitySourceLiveUsage, CapacitySourceLiveRateLimits:
return fact.Confidence == CapacityConfidenceAuthoritative && fact.ConnectionGeneration > 0
default:
return false
Expand All @@ -238,13 +246,47 @@ func (l *CodexCapacityLedger) ObserveQuotaSnapshot(account codex.AccountKey, sna
if l == nil || account == "" || len(snap.Result.Windows) == 0 {
return
}
type aggregate struct {
windows map[quota.WindowName]quota.Window
remaining int
reset time.Time
set bool
aggregates := capacitySnapshotAggregates(snap)
for bucket, aggregate := range aggregates {
l.mu.Lock()
l.seq++
fact := CapacityFact{
AccountKey: account,
Windows: aggregate.windows,
Bucket: bucket,
RemainingPct: aggregate.remaining,
Source: CapacitySourceUsageCache,
Sequence: l.seq,
ObservedAt: snap.FetchedAt,
ResetAt: aggregate.reset,
Confidence: CapacityConfidenceAdvisory,
}
l.observeLocked(fact)
l.mu.Unlock()
}
}

// ObserveLivePositiveQuotaSnapshot records positive capacity from a fresh
// authenticated usage response. Positive live evidence can lift an older hard
// fence after a banked reset; zero usage remains advisory.
func (l *CodexCapacityLedger) ObserveLivePositiveQuotaSnapshot(stream *CodexCapacityObservationStream, account codex.AccountKey, snap QuotaSnapshot) {
if l == nil || stream == nil || account == "" || len(snap.Result.Windows) == 0 {
return
}
for bucket, aggregate := range capacitySnapshotAggregates(snap) {
if aggregate.remaining <= 0 {
continue
}
l.Observe(stream.Stamp(CapacityFact{
AccountKey: account, Windows: aggregate.windows, Bucket: bucket,
RemainingPct: aggregate.remaining, Source: CapacitySourceLiveUsage,
ObservedAt: snap.FetchedAt, ResetAt: aggregate.reset, Confidence: CapacityConfidenceAuthoritative,
}))
}
aggregates := make(map[CapacityBucket]aggregate)
}

func capacitySnapshotAggregates(snap QuotaSnapshot) map[CapacityBucket]capacitySnapshotAggregate {
aggregates := make(map[CapacityBucket]capacitySnapshotAggregate)
for name, window := range snap.Result.Windows {
bucket := CapacityBucketBase
if scoped := quota.WindowBucket(name); scoped != "" {
Expand All @@ -268,23 +310,7 @@ func (l *CodexCapacityLedger) ObserveQuotaSnapshot(account codex.AccountKey, sna
current.set = true
aggregates[bucket] = current
}
for bucket, aggregate := range aggregates {
l.mu.Lock()
l.seq++
fact := CapacityFact{
AccountKey: account,
Windows: aggregate.windows,
Bucket: bucket,
RemainingPct: aggregate.remaining,
Source: CapacitySourceUsageCache,
Sequence: l.seq,
ObservedAt: snap.FetchedAt,
ResetAt: aggregate.reset,
Confidence: CapacityConfidenceAdvisory,
}
l.observeLocked(fact)
l.mu.Unlock()
}
return aggregates
}

// Capacity returns exact bucket state, falling scoped requests back to shared
Expand Down Expand Up @@ -329,6 +355,7 @@ func (l *CodexCapacityLedger) capacityLocked(account codex.AccountKey, bucket Ca
CapacitySourceUsageCache,
CapacitySourceHardLimit,
CapacitySourceHTTPHeaders,
CapacitySourceLiveUsage,
CapacitySourceLiveRateLimits,
} {
fact, ok := l.facts[capacityFactKey{account: account, bucket: bucket, source: source}]
Expand Down Expand Up @@ -374,7 +401,7 @@ func (l *CodexCapacityLedger) factStale(fact CapacityFact, now time.Time) bool {
}

func liveFactLiftsHardFence(live, hard CapacityFact) bool {
if live.Source != CapacitySourceLiveRateLimits || live.Confidence != CapacityConfidenceAuthoritative || live.RemainingPct <= 0 {
if (live.Source != CapacitySourceLiveUsage && live.Source != CapacitySourceLiveRateLimits) || live.Confidence != CapacityConfidenceAuthoritative || live.RemainingPct <= 0 {
return false
}
if !capacityCursorAfter(live, hard) {
Expand Down
9 changes: 6 additions & 3 deletions internal/proxy/codex_capacity_refresh.go
Original file line number Diff line number Diff line change
Expand Up @@ -89,13 +89,14 @@ func (r *CodexRoutingCapacityRefresher) Refresh(ctx context.Context, accounts []
type result struct {
account codex.AccountKey
observation codex.UsageObservation
stream *CodexCapacityObservationStream
err error
panicked bool
}
results := make(chan result, len(eligible))
for _, account := range eligible {
go func() {
outcome := result{account: account}
outcome := result{account: account, stream: r.Capacity.NewObservationStream()}
defer func() {
if recover() != nil {
outcome.observation = codex.UsageObservation{}
Expand Down Expand Up @@ -125,10 +126,12 @@ func (r *CodexRoutingCapacityRefresher) Refresh(ctx context.Context, accounts []
retryAt = completedAt.Add(interval)
}
if valid {
r.Capacity.ObserveQuotaSnapshot(outcome.account, QuotaSnapshot{
snapshot := QuotaSnapshot{
Result: outcome.observation.Result,
FetchedAt: now,
})
}
r.Capacity.ObserveQuotaSnapshot(outcome.account, snapshot)
r.Capacity.ObserveLivePositiveQuotaSnapshot(outcome.stream, outcome.account, snapshot)
published = true
}
r.mu.Lock()
Expand Down
55 changes: 55 additions & 0 deletions internal/proxy/codex_capacity_refresh_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -95,6 +95,61 @@ func TestCodexRoutingCapacityRefresherPublishesUsageAndHonoursInterval(t *testin
}
}

func TestCodexRoutingCapacityRefresherLiftsHardFenceAfterReset(t *testing.T) {
now := time.Unix(1_800_000_000, 0)
resetAt := now.Add(7 * 24 * time.Hour)
ledger := NewCodexCapacityLedger(func() time.Time { return now }, 5*time.Minute)
stream := ledger.NewObservationStream()
if !ledger.Observe(stream.Stamp(CapacityFact{
AccountKey: "reset", Bucket: CapacityBucketBase, RemainingPct: 0,
Source: CapacitySourceHardLimit, ResetAt: resetAt, Confidence: CapacityConfidenceAuthoritative,
})) {
t.Fatal("hard limit was not observed")
}

now = now.Add(time.Second)
exact := 100.0
reader := &codexRoutingUsageReaderStub{
results: map[codex.AccountKey]codex.UsageObservation{
"reset": {Result: quota.Result{Status: quota.StatusOK, Windows: map[quota.WindowName]quota.Window{
quota.Window7Day: {RemainingPct: 100, RemainingPctExact: &exact, ResetAtUnix: resetAt.Unix()},
}}},
},
errors: make(map[codex.AccountKey]error), panics: make(map[codex.AccountKey]bool), calls: make(map[codex.AccountKey]int),
}
refresher := &CodexRoutingCapacityRefresher{Usage: reader, Capacity: ledger, Now: func() time.Time { return now }}
if !refresher.Refresh(context.Background(), []codex.AccountKey{"reset"}) {
t.Fatal("reset usage was not published")
}
if view := ledger.Capacity("reset", CapacityBucketBase); view.State != CapacityPositive || view.RemainingPct != 100 || view.Source != CapacitySourceLiveUsage {
t.Fatalf("capacity after reset = %+v, want authoritative positive 100%%", view)
}
}

func TestCodexLiveUsageStartedBeforeHardLimitDoesNotLiftFence(t *testing.T) {
now := time.Unix(1_800_000_000, 0)
resetAt := now.Add(7 * 24 * time.Hour)
ledger := NewCodexCapacityLedger(func() time.Time { return now }, 5*time.Minute)
usageStream := ledger.NewObservationStream()
hardStream := ledger.NewObservationStream()
if !ledger.Observe(hardStream.Stamp(CapacityFact{
AccountKey: "account", Bucket: CapacityBucketBase, RemainingPct: 0,
Source: CapacitySourceHardLimit, ResetAt: resetAt, Confidence: CapacityConfidenceAuthoritative,
})) {
t.Fatal("hard limit was not observed")
}
exact := 100.0
ledger.ObserveLivePositiveQuotaSnapshot(usageStream, "account", QuotaSnapshot{
FetchedAt: now,
Result: quota.Result{Status: quota.StatusOK, Windows: map[quota.WindowName]quota.Window{
quota.Window7Day: {RemainingPct: 100, RemainingPctExact: &exact, ResetAtUnix: resetAt.Unix()},
}},
})
if view := ledger.Capacity("account", CapacityBucketBase); view.State != CapacityZero || view.Source != CapacitySourceHardLimit {
t.Fatalf("capacity = %+v, want newer hard fence", view)
}
}

func TestCodexRoutingCapacityRefresherContainsFailureAndPanic(t *testing.T) {
now := time.Unix(1_800_000_000, 0)
reader := &codexRoutingUsageReaderStub{
Expand Down
Loading