package providerhealth import ( "sort" "sync" "time" ) const recentWindow = 100 type EventKind string const ( EventOutcome EventKind = "outcome" EventTTFT EventKind = "ttft" ) type RouteKey struct { 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"` } type Options struct { FailureThreshold uint64 OpenDuration time.Duration Now func() time.Time Sink EventSink } type Tracker struct { states sync.Map 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 } if options.OpenDuration <= 0 { options.OpenDuration = 30 * time.Second } if options.Now == nil { options.Now = time.Now } return &Tracker{failureThreshold: options.FailureThreshold, openDuration: options.OpenDuration, now: options.Now, sink: options.Sink} } func (t *Tracker) Observe(key RouteKey, observation Observation) { if t == nil || key.ProviderID == "" { return } 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 { latency := float64(observation.Latency.Milliseconds()) if latency < 1 { latency = 1 } if state.headerLatencyEWMA == 0 { state.headerLatencyEWMA = latency } else { state.headerLatencyEWMA = state.headerLatencyEWMA*0.8 + latency*0.2 } } state.addRecent(!observation.Failed) if observation.Failed { state.consecutiveFailures++ if state.consecutiveFailures >= t.failureThreshold { state.openUntil = observation.ObservedAt.Add(t.openDuration) } return } state.consecutiveFailures = 0 state.openUntil = time.Time{} 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] { s.recentHealthy-- } } else { s.recentCount++ } s.recent[s.recentPosition] = healthy if healthy { s.recentHealthy++ } s.recentPosition = (s.recentPosition + 1) % recentWindow } func (t *Tracker) CircuitOpen(key RouteKey) bool { if t == nil { return false } value, ok := t.states.Load(key) if !ok { return false } state := value.(*routeState) state.mu.RLock() defer state.mu.RUnlock() 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{} } now := t.now() result := make([]Status, 0) t.states.Range(func(rawKey, rawState any) bool { key := rawKey.(RouteKey) state := rawState.(*routeState) state.mu.RLock() item := statusFromState(key, state, now) state.mu.RUnlock() result = append(result, item) return true }) sort.Slice(result, func(i, j int) bool { if result[i].ModelID != result[j].ModelID { return result[i].ModelID < result[j].ModelID } return result[i].ProviderID < result[j].ProviderID }) return result } func statusFromState(key RouteKey, state *routeState, now time.Time) Status { item := Status{ModelID: key.ModelID, ProviderID: key.ProviderID, WireAPI: key.WireAPI, 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 } if !state.lastObservedAt.IsZero() { 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 } if state.openUntil.After(now) { value := state.openUntil item.CircuitOpenUntil = &value item.State = "open" } else if state.attempts == 0 { item.State = "unknown" } else if state.consecutiveFailures > 0 || (state.recentCount >= 5 && item.AvailabilityPercent < 95) { item.State = "degraded" } else { item.State = "healthy" } return item }