dataplane

package
v0.3.1 Latest Latest
Warning

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

Go to latest
Published: Aug 2, 2026 License: MIT Imports: 44 Imported by: 0

Documentation

Index

Constants

View Source
const (
	ModalityChat            = "chat"
	ModalityMessages        = "messages"
	ModalityGenerateContent = "generate_content"
	ModalityTTS             = "tts"
	ModalitySTT             = "stt"
	ModalityEmbedding       = "embedding"
	ModalityImage           = "image"
	ModalityMCP             = "mcp"
	ModalityA2A             = "a2a"
)
View Source
const AFITagsHeader = "X-AFI-Tags"

AFITagsHeader is the request header for external user tags.

Variables

View Source
var ErrStreamUnsupported = errors.New("streaming is not supported for this provider")

ErrStreamUnsupported is returned when the provider capabilities disallow streaming.

Functions

func AuthenticateKey

func AuthenticateKey(snap *snapshot.Snapshot, rawKey string) (snapshot.APIKey, error)

AuthenticateKey is exported for unit tests.

func CopyResponse

func CopyResponse(w http.ResponseWriter, resp *http.Response) error

CopyResponse copies an upstream response to the client writer.

func HeadersForPolicy

func HeadersForPolicy(h http.Header) map[string]string

HeadersForPolicy copies inbound headers for CEL as lowercased key → first value. Sensitive headers (authorization, cookie, set-cookie) are omitted.

func ParseAFITags

func ParseAFITags(header string) map[string]string

ParseAFITags parses "key:value,key:value" tag headers. Pairs are comma-separated; each pair splits on the first ':'. Keys and values are trimmed; empty keys are skipped; last duplicate key wins.

func TagsFromRequest

func TagsFromRequest(r *http.Request) map[string]string

TagsFromRequest reads X-AFI-Tags from the request.

Types

type AfterCallHook

type AfterCallHook = sdkhook.AfterCallHook

Re-export SDK hook types so existing extensions can keep importing dataplane.

type AfterCallInfo

type AfterCallInfo = sdkhook.AfterCallInfo

Re-export SDK hook types so existing extensions can keep importing dataplane.

type AfterChatHook

type AfterChatHook = sdkhook.AfterChatHook

Re-export SDK hook types so existing extensions can keep importing dataplane.

type AfterChatInfo

type AfterChatInfo = sdkhook.AfterChatInfo

Re-export SDK hook types so existing extensions can keep importing dataplane.

type AnthropicTransport

type AnthropicTransport interface {
	PassThrough(ctx context.Context, provider snapshot.Provider, targetModel string, body []byte, stream bool) (*http.Response, error)
	Messages(ctx context.Context, provider snapshot.Provider, targetModel string, body []byte, stream bool) (*http.Response, error)
}

AnthropicTransport is the outbound Anthropic HTTP surface used by chat + /v1/messages.

type AnthropicTransportProvider

type AnthropicTransportProvider interface {
	AnthropicTransport() AnthropicTransport
}

AnthropicTransportProvider is implemented by ChatProvider adapters that expose Anthropic HTTP.

type AudioBackend

type AudioBackend interface {
	AudioSpeech(ctx context.Context, provider snapshot.Provider, targetModel string, body []byte) (*http.Response, error)
	AudioTranscriptions(ctx context.Context, provider snapshot.Provider, targetModel, contentType string, body io.Reader) (*http.Response, error)
}

AudioBackend is the modality port for TTS/STT.

type AudioTransportProvider

type AudioTransportProvider interface {
	AudioBackend() AudioBackend
}

AudioTransportProvider is implemented by ChatProvider adapters that expose TTS/STT without a full OpenAI transport (e.g. elevenlabs).

type BeforeCallHook

type BeforeCallHook = sdkhook.BeforeCallHook

Re-export SDK hook types so existing extensions can keep importing dataplane.

type BuiltinFactory

type BuiltinFactory func(sec secrets.Resolver) ChatProvider

BuiltinFactory builds a ChatProvider from a secret resolver.

type CallContext

type CallContext = sdkhook.CallContext

Re-export SDK hook types so existing extensions can keep importing dataplane.

type CallDecision

type CallDecision = sdkhook.CallDecision

Re-export SDK hook types so existing extensions can keep importing dataplane.

type ChatHook

type ChatHook = sdkhook.ChatHook

Re-export SDK hook types so existing extensions can keep importing dataplane.

type ChatProvider

type ChatProvider interface {
	Type() string
	Capabilities() ProviderCaps
	Chat(ctx context.Context, p snapshot.Provider, targetModel string, body []byte, stream bool) (*http.Response, error)
}

ChatProvider is the in-process adapter contract for gateway chat HTTP transport.

type CompositeCounters

type CompositeCounters struct {
	Total CounterStore // Postgres (window=total)
	Timed CounterStore // Redis (minute/hour/day)
}

CompositeCounters routes lifetime quotas to Postgres and timed windows to Redis.

func (CompositeCounters) Get

func (c CompositeCounters) Get(ctx context.Context, scopeType, scopeID, metric, window string) (int64, error)

func (CompositeCounters) Incr

func (c CompositeCounters) Incr(ctx context.Context, scopeType, scopeID, metric, window string, delta int64) (int64, error)

type CounterStore

type CounterStore interface {
	Get(ctx context.Context, scopeType, scopeID, metric, window string) (int64, error)
	Incr(ctx context.Context, scopeType, scopeID, metric, window string, delta int64) (int64, error)
}

CounterStore reads/writes durable quota counters (not config).

type EmbeddingsBackend

type EmbeddingsBackend interface {
	Embeddings(ctx context.Context, provider snapshot.Provider, targetModel string, body []byte) (*http.Response, error)
}

EmbeddingsBackend is the modality port for /v1/embeddings (OpenAI-compatible).

type Holder

type Holder struct {
	// contains filtered or unexported fields
}

Holder keeps the current immutable snapshot for the request path.

func NewHolder

func NewHolder() *Holder

func (*Holder) Get

func (h *Holder) Get() *snapshot.Snapshot

func (*Holder) Set

func (h *Holder) Set(s *snapshot.Snapshot)

type HookChain

type HookChain struct {
	// contains filtered or unexported fields
}

HookChain runs BeforeCall / AfterCall / BeforeChat / AfterChat hooks in registration order.

func NewHookChain

func NewHookChain() *HookChain

func (*HookChain) Infos

func (c *HookChain) Infos() []HookInfo

func (*HookChain) Names

func (c *HookChain) Names() []string

func (*HookChain) PrependBeforeCall

func (c *HookChain) PrependBeforeCall(h BeforeCallHook) *HookChain

PrependBeforeCall inserts a BeforeCall hook at the front of the chain.

func (*HookChain) Register

func (c *HookChain) Register(h ChatHook) *HookChain

Register adds a BeforeChat hook.

func (*HookChain) RegisterAfter

func (c *HookChain) RegisterAfter(h AfterChatHook) *HookChain

RegisterAfter adds an AfterChat hook.

func (*HookChain) RegisterAfterCall

func (c *HookChain) RegisterAfterCall(h AfterCallHook) *HookChain

RegisterAfterCall adds an AfterCall hook.

func (*HookChain) RegisterBeforeCall

func (c *HookChain) RegisterBeforeCall(h BeforeCallHook) *HookChain

RegisterBeforeCall adds a BeforeCall hook (appended; runs after earlier entries).

func (*HookChain) RegisterHook

func (c *HookChain) RegisterHook(h any) *HookChain

RegisterHook registers a value that may implement any of the hook interfaces.

func (*HookChain) RunAfterCall

func (c *HookChain) RunAfterCall(ctx context.Context, call *CallContext, info AfterCallInfo)

RunAfterCall runs AfterCall hooks (errors ignored).

func (*HookChain) RunAfterChat

func (c *HookChain) RunAfterChat(ctx context.Context, info AfterChatInfo)

func (*HookChain) RunBeforeCall

func (c *HookChain) RunBeforeCall(ctx context.Context, call *CallContext) (CallDecision, error)

RunBeforeCall runs BeforeCall hooks. First deny wins. Mutates call in place.

func (*HookChain) RunBeforeChat

func (c *HookChain) RunBeforeChat(ctx context.Context, req ir.ChatRequest) (ir.ChatRequest, error)

RunBeforeChat runs typed chat hooks in registration order.

type HookInfo

type HookInfo struct {
	Name       string `json:"name"`
	BeforeCall bool   `json:"before_call"`
	AfterCall  bool   `json:"after_call"`
	BeforeChat bool   `json:"before_chat"`
	AfterChat  bool   `json:"after_chat"`
}

HookInfo describes a registered hook for healthz / UI.

type IRChatProvider

type IRChatProvider interface {
	ChatIR(ctx context.Context, p snapshot.Provider, targetModel string, req ir.ChatRequest) (ir.ChatResult, error)
}

IRChatProvider is implemented by built-in adapters that speak chat IR.

type ImagesBackend

type ImagesBackend interface {
	Images(ctx context.Context, provider snapshot.Provider, targetModel string, body []byte) (*http.Response, error)
}

ImagesBackend is the modality port for /v1/images/generations (OpenAI-compatible).

type MessagesBackend

type MessagesBackend interface {
	PassThrough(ctx context.Context, provider snapshot.Provider, targetModel string, body []byte, stream bool) (*http.Response, error)
}

MessagesBackend is retained for Anthropic transport PassThrough (used by ChatIR). Client Anthropic dialect traffic goes through chat IR, not this port directly.

type OpenAITransport

type OpenAITransport interface {
	ChatCompletions(ctx context.Context, provider snapshot.Provider, targetModel string, body []byte, stream bool) (*http.Response, error)
	AudioSpeech(ctx context.Context, provider snapshot.Provider, targetModel string, body []byte) (*http.Response, error)
	AudioTranscriptions(ctx context.Context, provider snapshot.Provider, targetModel, contentType string, body io.Reader) (*http.Response, error)
	Embeddings(ctx context.Context, provider snapshot.Provider, targetModel string, body []byte) (*http.Response, error)
	Images(ctx context.Context, provider snapshot.Provider, targetModel string, body []byte) (*http.Response, error)
}

OpenAITransport is the outbound OpenAI-compatible HTTP surface used by chat + audio + embeddings + images.

type OpenAITransportProvider

type OpenAITransportProvider interface {
	OpenAITransport() OpenAITransport
}

OpenAITransportProvider is implemented by ChatProvider adapters that expose OpenAI HTTP.

type Pipeline

type Pipeline struct {
	Holder      *Holder
	Providers   *Registry
	Hooks       *HookChain
	Wasm        *WasmRunner
	Log         *slog.Logger
	Usage       func(UsageEvent)
	Counters    CounterStore
	Policies    *policy.Evaluator
	Credentials secrets.CredentialOpener
	Secrets     secrets.Resolver
	HTTP        *http.Client
	Replay      ReplayStore
	Metrics     *telemetry.GatewayMetrics
	// RouteRand optional RNG for weighted routing (tests); nil uses math/rand global.
	RouteRand *rand.Rand
	// RouteSignals optional gateway-local EWMA store for latency/cost adaptive routing.
	RouteSignals routing.SignalStore
}

func NewPipeline

func NewPipeline(holder *Holder, reg *Registry, log *slog.Logger) *Pipeline

NewPipeline builds a pipeline with an explicit provider registry. Built-in LLM adapters are registered from cmd/gateway via adapters/llm.

func NewPipelineWithRegistry

func NewPipelineWithRegistry(holder *Holder, reg *Registry, log *slog.Logger) *Pipeline

NewPipelineWithRegistry uses an explicit provider registry.

func (*Pipeline) Handler

func (p *Pipeline) Handler() http.Handler

type Principal

type Principal = sdkhook.Principal

Re-export SDK hook types so existing extensions can keep importing dataplane.

type ProviderCaps

type ProviderCaps struct {
	Chat      bool
	Stream    bool
	TTS       bool
	STT       bool
	Embedding bool
	Image     bool
}

ProviderCaps mirrors snapshot capabilities for adapters.

type Registry

type Registry struct {
	// contains filtered or unexported fields
}

Registry maps provider type strings to ChatProvider implementations.

func DefaultRegistry

func DefaultRegistry() *Registry

DefaultRegistry registers all in-tree adapters via registerBuiltin factories.

func NewRegistry

func NewRegistry() *Registry

func RegistryWithOpenAI

func RegistryWithOpenAI(openai *llm.OpenAIClient) *Registry

RegistryWithOpenAI builds DefaultRegistry but uses the given OpenAI client for type "openai" (tests inject mock HTTP transports).

func RegistryWithSecrets

func RegistryWithSecrets(sec secrets.Resolver) *Registry

RegistryWithSecrets builds the builtin registry with a custom secret resolver.

func (*Registry) AnthropicTransport

func (r *Registry) AnthropicTransport(typ string) (AnthropicTransport, bool)

AnthropicTransport looks up an Anthropic transport by provider type.

func (*Registry) AudioBackend

func (r *Registry) AudioBackend(typ string) (AudioBackend, bool)

AudioBackend returns the TTS/STT port for a provider type.

func (*Registry) EmbeddingsBackend

func (r *Registry) EmbeddingsBackend(typ string) (EmbeddingsBackend, bool)

EmbeddingsBackend returns the /v1/embeddings port for a provider type (OpenAI-compatible only).

func (*Registry) Get

func (r *Registry) Get(typ string) (ChatProvider, bool)

func (*Registry) ImagesBackend

func (r *Registry) ImagesBackend(typ string) (ImagesBackend, bool)

ImagesBackend returns the /v1/images/generations port for a provider type (OpenAI-compatible only).

func (*Registry) MessagesBackend

func (r *Registry) MessagesBackend(typ string) (MessagesBackend, bool)

MessagesBackend returns the native /v1/messages port for a provider type.

func (*Registry) OpenAITransport

func (r *Registry) OpenAITransport(typ string) (OpenAITransport, bool)

OpenAITransport looks up an OpenAI-compatible transport by provider type.

func (*Registry) Register

func (r *Registry) Register(p ChatProvider) *Registry

func (*Registry) RegisterSDK

func (r *Registry) RegisterSDK(p sdkprovider.ChatProvider) *Registry

RegisterSDK wraps an SDK ChatProvider into the gateway registry.

func (*Registry) Types

func (r *Registry) Types() []string

type ReplayStore added in v0.3.0

type ReplayStore interface {
	Use(ctx context.Context, key string, ttl time.Duration) (bool, error)
}

type RouteContext

type RouteContext = sdkhook.RouteContext

Re-export SDK hook types so existing extensions can keep importing dataplane.

type UsageEvent

type UsageEvent = usage.Event

UsageEvent is an alias for the canonical usage.Event emitted on the request path.

type WasmRunner

type WasmRunner struct {
	Cache *afiWasm.ModuleCache
	Log   *slog.Logger
}

WasmRunner executes org-scoped snapshot WASM bindings via a module cache.

func (*WasmRunner) RunAfterCall

func (r *WasmRunner) RunAfterCall(ctx context.Context, snap *snapshot.Snapshot, call *CallContext, info AfterCallInfo)

func (*WasmRunner) RunBeforeCall

func (r *WasmRunner) RunBeforeCall(ctx context.Context, snap *snapshot.Snapshot, call *CallContext) (CallDecision, error)

func (*WasmRunner) RunBeforeChat

func (r *WasmRunner) RunBeforeChat(ctx context.Context, snap *snapshot.Snapshot, orgID string, req chatir.Request) (chatir.Request, error)

RunBeforeChat executes org-scoped typed chat IR WASM hooks.

Directories

Path Synopsis
Package dialect encodes/decodes client wire formats against chat IR.
Package dialect encodes/decodes client wire formats against chat IR.
Package ir defines the gateway-owned chat internal representation.
Package ir defines the gateway-owned chat internal representation.

Jump to

Keyboard shortcuts

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