summaryrefslogtreecommitdiff
path: root/internal/providerhealth/tracker.go
diff options
context:
space:
mode:
authorChia <Chia@93.nz>2026-08-06 15:58:57 +1200
committerChia <Chia@93.nz>2026-08-06 15:58:57 +1200
commit3f702084d20b3c3a3ea916f3110e99b22bda60b3 (patch)
tree517f76c51025ce1ee085ea4898c60f799e5c37ea /internal/providerhealth/tracker.go
parent41e322c53d7b4b796eb377d0df9c29ecd10ba431 (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.go162
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