upstreams

package
v0.7.9 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Aug 23, 2026 License: MIT Imports: 14 Imported by: 0

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

Constants

This section is empty.

Variables

View Source
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

func ModelToAPI(m Model) api.Model

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

type Candidate struct {
	ModelID string
	Client  *elelem.Client
}

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

type ClientFactory func(u config.Upstream) *elelem.Client

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

func (r *Registry) CandidatesFor(modelID string) []Candidate

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) ClientFor

func (r *Registry) ClientFor(
	modelID string,
) (*elelem.Client, bool)

func (*Registry) Health

func (r *Registry) Health() []Health

Health returns current upstream status in the original configuration order.

func (*Registry) Models

func (r *Registry) Models() []Model

Models returns the merged model list, sorted by id.

func (*Registry) ReadinessHealth

func (r *Registry) ReadinessHealth() []operations.UpstreamHealth

ReadinessHealth returns the redacted health subset for admin readiness.

func (*Registry) SetDefault

func (r *Registry) SetDefault(modelID string) error

SetDefault marks exactly one advertised model as the instance default. An empty model ID leaves the registry without an explicit default.

func (*Registry) SetFallbacks

func (r *Registry) SetFallbacks(configured []config.Upstream) error

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

type TokenPrice struct {
	AmountSmallestUnit int64
	Currency           string
}

TokenPrice is a model's configured price per million input or output tokens. AmountSmallestUnit uses the configured currency's smallest unit.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL