diff options
| author | Chia <Chia@93.nz> | 2026-08-06 15:58:57 +1200 |
|---|---|---|
| committer | Chia <Chia@93.nz> | 2026-08-06 15:58:57 +1200 |
| commit | 3f702084d20b3c3a3ea916f3110e99b22bda60b3 (patch) | |
| tree | 517f76c51025ce1ee085ea4898c60f799e5c37ea /internal/providerhealth/tracker.go | |
| parent | 41e322c53d7b4b796eb377d0df9c29ecd10ba431 (diff) | |
feat: complete commercial developer workflowspublish-commercial-control-plane
Add tenant-safe usage observability, prepaid billing controls, API key lifecycle management, Embeddings metering, configurable billing alerts, and resilient provider health propagation. Harden Stripe failure handling, migrations, readiness, and the authenticated control-plane UI with end-to-end verification evidence.
Diffstat (limited to 'internal/providerhealth/tracker.go')
| -rw-r--r-- | internal/providerhealth/tracker.go | 162 |
1 files changed, 156 insertions, 6 deletions
diff --git a/internal/providerhealth/tracker.go b/internal/providerhealth/tracker.go index a5af2b2..6394def 100644 --- a/internal/providerhealth/tracker.go +++ b/internal/providerhealth/tracker.go @@ -8,31 +8,59 @@ import ( const recentWindow = 100 +type EventKind string + +const ( + EventOutcome EventKind = "outcome" + EventTTFT EventKind = "ttft" +) + type RouteKey struct { - ModelID string - ProviderID string - WireAPI string + ModelID string `json:"model_id"` + ProviderID string `json:"provider_id"` + WireAPI string `json:"wire_api"` } type Observation struct { StatusCode int Latency time.Duration Failed bool + Active bool ObservedAt time.Time } +type Event struct { + Kind EventKind `json:"kind"` + Key RouteKey `json:"key"` + StatusCode int `json:"status_code,omitempty"` + LatencyMillis int64 `json:"latency_ms,omitempty"` + Failed bool `json:"failed,omitempty"` + Active bool `json:"active,omitempty"` + ObservedAt time.Time `json:"observed_at"` +} + +type EventSink interface { + Enqueue(Event) +} + type Status struct { ModelID string `json:"model_id"` ProviderID string `json:"provider_id"` WireAPI string `json:"wire_api"` State string `json:"state"` Attempts uint64 `json:"attempts"` + ActiveProbes uint64 `json:"active_probes"` RecentSamples int `json:"recent_samples"` AvailabilityPercent float64 `json:"availability_percent"` HeaderLatencyEWMA int64 `json:"header_latency_ewma_ms"` + TTFTSamples uint64 `json:"ttft_samples"` + TTFTEWMA int64 `json:"ttft_ewma_ms"` + SharedAttempts uint64 `json:"shared_attempts"` + SharedTTFTSamples uint64 `json:"shared_ttft_samples"` ConsecutiveFailures uint64 `json:"consecutive_failures"` LastStatusCode int `json:"last_status_code,omitempty"` LastObservedAt *time.Time `json:"last_observed_at,omitempty"` + LastProbeAt *time.Time `json:"last_probe_at,omitempty"` LastHealthyAt *time.Time `json:"last_healthy_at,omitempty"` CircuitOpenUntil *time.Time `json:"circuit_open_until,omitempty"` } @@ -41,6 +69,7 @@ type Options struct { FailureThreshold uint64 OpenDuration time.Duration Now func() time.Time + Sink EventSink } type Tracker struct { @@ -48,23 +77,63 @@ type Tracker struct { failureThreshold uint64 openDuration time.Duration now func() time.Time + sinkMu sync.RWMutex + sink EventSink } type routeState struct { mu sync.RWMutex attempts uint64 + activeProbes uint64 consecutiveFailures uint64 lastStatusCode int lastObservedAt time.Time + lastProbeAt time.Time lastHealthyAt time.Time openUntil time.Time headerLatencyEWMA float64 + ttftSamples uint64 + ttftEWMA float64 + sharedAttempts uint64 + sharedTTFTSamples uint64 recent [recentWindow]bool recentCount int recentPosition int recentHealthy int } +// ObserveTTFT records user-visible response latency without counting a second +// request outcome. Forwarder observations already update availability and the +// circuit when response headers arrive. +func (t *Tracker) ObserveTTFT(key RouteKey, latency time.Duration) { + if t == nil || key.ProviderID == "" || latency <= 0 { + return + } + observedAt := t.now() + t.observeTTFT(key, latency, false) + t.publish(Event{Kind: EventTTFT, Key: key, LatencyMillis: durationMillis(latency), ObservedAt: observedAt}) +} + +func (t *Tracker) observeTTFT(key RouteKey, latency time.Duration, shared bool) { + value, _ := t.states.LoadOrStore(key, &routeState{}) + state := value.(*routeState) + state.mu.Lock() + defer state.mu.Unlock() + valueMS := float64(durationMillis(latency)) + if valueMS < 1 { + valueMS = 1 + } + state.ttftSamples++ + if shared { + state.sharedTTFTSamples++ + } + if state.ttftEWMA == 0 { + state.ttftEWMA = valueMS + } else { + state.ttftEWMA = state.ttftEWMA*0.8 + valueMS*0.2 + } +} + func New(options Options) *Tracker { if options.FailureThreshold == 0 { options.FailureThreshold = 3 @@ -75,7 +144,7 @@ func New(options Options) *Tracker { if options.Now == nil { options.Now = time.Now } - return &Tracker{failureThreshold: options.FailureThreshold, openDuration: options.OpenDuration, now: options.Now} + return &Tracker{failureThreshold: options.FailureThreshold, openDuration: options.OpenDuration, now: options.Now, sink: options.Sink} } func (t *Tracker) Observe(key RouteKey, observation Observation) { @@ -85,12 +154,26 @@ func (t *Tracker) Observe(key RouteKey, observation Observation) { if observation.ObservedAt.IsZero() { observation.ObservedAt = t.now() } + t.observe(key, observation, false) + t.publish(Event{Kind: EventOutcome, Key: key, StatusCode: observation.StatusCode, + LatencyMillis: durationMillis(observation.Latency), Failed: observation.Failed, + Active: observation.Active, ObservedAt: observation.ObservedAt}) +} + +func (t *Tracker) observe(key RouteKey, observation Observation, shared bool) { value, _ := t.states.LoadOrStore(key, &routeState{}) state := value.(*routeState) state.mu.Lock() defer state.mu.Unlock() state.attempts++ + if shared { + state.sharedAttempts++ + } + if observation.Active { + state.activeProbes++ + state.lastProbeAt = observation.ObservedAt + } state.lastStatusCode = observation.StatusCode state.lastObservedAt = observation.ObservedAt if observation.Latency > 0 { @@ -117,6 +200,53 @@ func (t *Tracker) Observe(key RouteKey, observation Observation) { state.lastHealthyAt = observation.ObservedAt } +func (t *Tracker) SetSink(sink EventSink) { + if t == nil { + return + } + t.sinkMu.Lock() + t.sink = sink + t.sinkMu.Unlock() +} + +func (t *Tracker) ApplyShared(event Event) { + if t == nil || event.Key.ProviderID == "" { + return + } + if event.ObservedAt.IsZero() { + event.ObservedAt = t.now() + } + switch event.Kind { + case EventOutcome: + t.observe(event.Key, Observation{StatusCode: event.StatusCode, Latency: time.Duration(event.LatencyMillis) * time.Millisecond, + Failed: event.Failed, Active: event.Active, ObservedAt: event.ObservedAt}, true) + case EventTTFT: + if event.LatencyMillis > 0 { + t.observeTTFT(event.Key, time.Duration(event.LatencyMillis)*time.Millisecond, true) + } + } +} + +func (t *Tracker) publish(event Event) { + t.sinkMu.RLock() + sink := t.sink + t.sinkMu.RUnlock() + if sink != nil { + sink.Enqueue(event) + } +} + +func durationMillis(value time.Duration) int64 { + if value <= 0 { + return 0 + } + milliseconds := value.Milliseconds() + if milliseconds < 1 { + return 1 + } + return milliseconds +} + func (s *routeState) addRecent(healthy bool) { if s.recentCount == recentWindow { if s.recent[s.recentPosition] { @@ -146,6 +276,21 @@ func (t *Tracker) CircuitOpen(key RouteKey) bool { return state.openUntil.After(t.now()) } +// StatusFor returns one immutable route-health snapshot for routing decisions. +func (t *Tracker) StatusFor(key RouteKey) (Status, bool) { + if t == nil { + return Status{}, false + } + value, ok := t.states.Load(key) + if !ok { + return Status{}, false + } + state := value.(*routeState) + state.mu.RLock() + defer state.mu.RUnlock() + return statusFromState(key, state, t.now()), true +} + func (t *Tracker) Snapshot() []Status { if t == nil { return []Status{} @@ -172,8 +317,9 @@ func (t *Tracker) Snapshot() []Status { func statusFromState(key RouteKey, state *routeState, now time.Time) Status { item := Status{ModelID: key.ModelID, ProviderID: key.ProviderID, WireAPI: key.WireAPI, - Attempts: state.attempts, RecentSamples: state.recentCount, HeaderLatencyEWMA: int64(state.headerLatencyEWMA + 0.5), - ConsecutiveFailures: state.consecutiveFailures, LastStatusCode: state.lastStatusCode} + Attempts: state.attempts, ActiveProbes: state.activeProbes, RecentSamples: state.recentCount, HeaderLatencyEWMA: int64(state.headerLatencyEWMA + 0.5), + TTFTSamples: state.ttftSamples, TTFTEWMA: int64(state.ttftEWMA + 0.5), SharedAttempts: state.sharedAttempts, + SharedTTFTSamples: state.sharedTTFTSamples, ConsecutiveFailures: state.consecutiveFailures, LastStatusCode: state.lastStatusCode} if state.recentCount > 0 { item.AvailabilityPercent = float64(state.recentHealthy) / float64(state.recentCount) * 100 } @@ -181,6 +327,10 @@ func statusFromState(key RouteKey, state *routeState, now time.Time) Status { value := state.lastObservedAt item.LastObservedAt = &value } + if !state.lastProbeAt.IsZero() { + value := state.lastProbeAt + item.LastProbeAt = &value + } if !state.lastHealthyAt.IsZero() { value := state.lastHealthyAt item.LastHealthyAt = &value |
