Documentation
¶
Overview ¶
Package upstreams discovers models across the configured OpenAI-compatible endpoints and routes a chosen model id back to the client that serves it. It calls each upstream's /models endpoint and merges the results (e.g. an Ollama base + an OpenAI base expose a single combined list); discovery is best-effort so one unreachable upstream doesn't hide the others.
Index ¶
- Variables
- func HealthToAPI(health Health) api.UpstreamHealth
- func HealthsToAPI(health []Health) []api.UpstreamHealth
- func ModelToAPI(m Model) api.Model
- func WrapDriver(driver elelem.Driver, upstream string, limits RuntimeLimits, ...) (elelem.Driver, error)
- type Candidate
- type ClientFactory
- type Health
- type HealthState
- type HealthTracker
- func (t *HealthTracker) RecordFailure(upstream, operation string, latency time.Duration, err error)
- func (t *HealthTracker) RecordSuccess(upstream, operation string, latency time.Duration)
- func (t *HealthTracker) Snapshot(upstream string) (Health, bool)
- func (t *HealthTracker) Snapshots(upstreamNames []string) []Health
- type Model
- type Registry
- func (r *Registry) CandidatesFor(modelID string) []Candidate
- func (r *Registry) ClientFor(modelID string) (*elelem.Client, bool)
- func (r *Registry) Health() []Health
- func (r *Registry) Models() []Model
- func (r *Registry) ReadinessHealth() []operations.UpstreamHealth
- func (r *Registry) SetDefault(modelID string) error
- func (r *Registry) SetFallbacks(configured []config.Upstream) error
- type RuntimeLimits
- type TokenPrice
Constants ¶
This section is empty.
Variables ¶
var ( ErrInvalidLimits = errors.New("invalid upstream runtime limits") ErrFirstTokenTimeout = errors.New("upstream time-to-first-token exceeded") ErrTurnTimeout = errors.New("upstream total-turn timeout exceeded") ErrDiscoveryTimeout = errors.New( "upstream model discovery timeout exceeded", ) ErrInvalidDiscovery = errors.New( "invalid upstream discovery configuration", ) ErrNilDriver = errors.New("upstream driver is nil") ErrNilHealthTracker = errors.New("upstream health tracker is nil") ErrInvalidUpstreamName = errors.New("upstream name is empty") )
Functions ¶
func HealthToAPI ¶
func HealthToAPI(health Health) api.UpstreamHealth
HealthToAPI projects a redacted upstream-health snapshot to the wire shape.
func HealthsToAPI ¶
func HealthsToAPI(health []Health) []api.UpstreamHealth
HealthsToAPI projects the ordered registry snapshots to the wire shape.
func ModelToAPI ¶
ModelToAPI projects an available model to the wire shape.
func WrapDriver ¶
func WrapDriver( driver elelem.Driver, upstream string, limits RuntimeLimits, health *HealthTracker, ) (elelem.Driver, error)
WrapDriver adds queue limits, time-to-first-token and total-turn deadlines around a concurrency-safe Elelem driver. It never retries a request: Elelem owns provider-aware retries, and MCP tool execution is outside this wrapper.
Types ¶
type Candidate ¶
Candidate is one configured model and the provider client that serves it. CandidatesFor always returns the selected model first, followed by its explicitly configured fallback order.
type ClientFactory ¶
ClientFactory builds the elelem client for an upstream. Injected so tests can supply fakes instead of real HTTP clients.
type Health ¶
type Health struct {
Upstream string
State HealthState
LastOperation string
LastLatency time.Duration
LastSuccessAt time.Time
LastFailureAt time.Time
LastFailureClass string
ConsecutiveFailure int
}
Health is a snapshot of the latest operation outcome for one upstream.
type HealthState ¶
type HealthState string
HealthState is the current client-visible condition of one configured upstream. It records classes rather than provider text so diagnostics never retain prompts, headers, credentials, or response bodies.
const ( HealthStateUnknown HealthState = "unknown" HealthStateHealthy HealthState = "healthy" HealthStateDegraded HealthState = "degraded" )
type HealthTracker ¶
type HealthTracker struct {
// contains filtered or unexported fields
}
HealthTracker holds per-upstream snapshots. It is safe for concurrent stream callbacks and lets registry/discovery and runtime calls contribute to one source of truth.
func NewHealthTracker ¶
func NewHealthTracker(upstreamNames []string) *HealthTracker
NewHealthTracker creates unknown snapshots for every configured upstream.
func (*HealthTracker) RecordFailure ¶
func (t *HealthTracker) RecordFailure( upstream, operation string, latency time.Duration, err error, )
RecordFailure stores only a stable failure class, never raw provider text.
func (*HealthTracker) RecordSuccess ¶
func (t *HealthTracker) RecordSuccess( upstream, operation string, latency time.Duration, )
RecordSuccess records a completed idempotent or turn operation.
func (*HealthTracker) Snapshot ¶
func (t *HealthTracker) Snapshot(upstream string) (Health, bool)
Snapshot returns one immutable health value.
func (*HealthTracker) Snapshots ¶
func (t *HealthTracker) Snapshots(upstreamNames []string) []Health
Snapshots returns known upstream health in caller-selected order.
type Model ¶
type Model struct {
ID string `json:"id"`
Upstream string `json:"upstream"`
Alias string `json:"alias,omitempty"`
ContextWindow int `json:"contextWindow,omitempty"`
MaxOutputTokens int `json:"maxOutputTokens,omitempty"`
SupportsTools *bool `json:"supportsTools,omitempty"`
SupportsReasoning *bool `json:"supportsReasoning,omitempty"`
SupportsVision *bool `json:"supportsVision,omitempty"`
SupportsFiles *bool `json:"supportsFiles,omitempty"`
FirstTokenLatencyMs int64
InputTokenPrice *TokenPrice
OutputTokenPrice *TokenPrice
Default bool `json:"default"`
}
Model is a discovered model id and the upstream that serves it.
type Registry ¶
type Registry struct {
// contains filtered or unexported fields
}
Registry is the merged view of models across upstreams plus the per-upstream clients, so a chosen model routes to its owner.
func Discover ¶
func Discover( ctx context.Context, ups []config.Upstream, factory ClientFactory, connectTimeout time.Duration, health *HealthTracker, ) *Registry
Discover builds each upstream's client, lists its models, and merges them. The first upstream to claim a model id wins (later duplicates are skipped with a warning). An upstream whose /models call fails is logged and skipped — the rest still populate the registry.
func (*Registry) CandidatesFor ¶
CandidatesFor resolves a selected model into its primary client followed by its validated fallbacks. It returns no candidates when the selected model is not currently routable.
func (*Registry) Health ¶
Health returns current upstream status in the original configuration order.
func (*Registry) ReadinessHealth ¶
func (r *Registry) ReadinessHealth() []operations.UpstreamHealth
ReadinessHealth returns the redacted health subset for admin readiness.
func (*Registry) SetDefault ¶
SetDefault marks exactly one advertised model as the instance default. An empty model ID leaves the registry without an explicit default.
func (*Registry) SetFallbacks ¶
SetFallbacks applies model fallback lists after discovery, when the registry can prove every target model is currently routable. A configured source model that is not advertised is ignored: discovery is best-effort, so a temporarily unavailable upstream must not prevent the rest from starting. An advertised source with an unavailable target is rejected.
type RuntimeLimits ¶
type RuntimeLimits struct {
MaxConcurrent int
FirstTokenTimeout time.Duration
TurnTimeout time.Duration
}
RuntimeLimits bound one upstream's live provider traffic. The queue is caller-context-bound; first-token and turn timers start only after a slot is acquired, so queue time is not mislabeled as provider lag.
func (RuntimeLimits) Validate ¶
func (l RuntimeLimits) Validate() error
Validate ensures limits fail closed before any provider request is made.
type TokenPrice ¶
TokenPrice is a model's configured price per million input or output tokens. AmountSmallestUnit uses the configured currency's smallest unit.