Documentation
¶
Overview ¶
Package ai owns provider transports, native transcripts, stream events, OAuth/account pools and model catalogs. The native engine consumes its AssistantMessageEventStream and lossless Message types directly.
The native providers include Codex Responses, Antigravity, Gemini-compatible Completions and Claude Code/Anthropic Messages. Provider-neutral client contracts, provider session state and retry/error classification live in ai/protocol, below both this package and agentcore. Neither ai nor protocol imports the agent runtime.
NativeProvider dispatches native transports; FallbackProvider composes them with bounded retries and ordered fallback under the same StreamFn contract. Durable hosts bind selection and usage callbacks through FallbackProvider.Run; retry/fallback never replays a request after visible content has escaped.
NewClient retains the Chat/Stream client API used by SDK integrations and connectivity probes. New adds managed-provider identity and model discovery; Collection routes those clients by model. These client contracts are distinct from the native transcript and stream types used by engine.
Conversation-scoped provider state can cache endpoint capabilities or account-specific transport state. It is an optimization, not durable agent history; an empty cache must remain correct, and account rotation selectively clears account-scoped records.
Index ¶
- Constants
- func APIKeyOptional(vendor string) bool
- func AbortableSleep(ms float64, ctx context.Context, cancelMessage string) error
- func AppendGrammarToolInputJSONDelta(buffer *GrammarToolInputJSONBuffer, property, next string, close bool) (*string, error)
- func ApplyNativeControls(api string, raw json.RawMessage, opts NativeGenerationControls) (json.RawMessage, error)
- func AssertChatModel(model any) error
- func AssertClassifierModel(model any) error
- func AssertImageModel(model any) error
- func BuildAnthropicParams(rawModel json.RawMessage, context TranscriptContext, oauth bool, ...) (json.RawMessage, error)
- func BuildAnthropicSimpleOptions(rawModel json.RawMessage, transcript TranscriptContext, ...) (json.RawMessage, error)
- func BuildAzureResponsesParams(rawModel json.RawMessage, transcript TranscriptContext, ...) (json.RawMessage, error)
- func BuildAzureResponsesSimpleOptions(model json.RawMessage, transcript TranscriptContext, options json.RawMessage) (json.RawMessage, error)
- func BuildCodexResponsesParams(rawModel json.RawMessage, context TranscriptContext, ...) (json.RawMessage, error)
- func BuildCodexResponsesSimpleOptions(rawModel json.RawMessage, transcript TranscriptContext, ...) (json.RawMessage, error)
- func BuildOpenAICompletionsParams(rawModel json.RawMessage, context TranscriptContext, ...) (json.RawMessage, error)
- func BuildOpenAICompletionsSimpleOptions(rawModel json.RawMessage, transcript TranscriptContext, ...) (json.RawMessage, error)
- func BuildOpenAIResponsesParams(rawModel json.RawMessage, context TranscriptContext, ...) (json.RawMessage, error)
- func BuildOpenAIResponsesSimpleOptions(rawModel json.RawMessage, transcript TranscriptContext, ...) (json.RawMessage, error)
- func CalculateContextTokens(usage Usage) float64
- func CapabilitiesFor(vendor, model string) protocol.ModelCapabilities
- func ClampMaxTokensToContext(contextWindow float64, transcript TranscriptContext, maxTokens float64) float64
- func ClampOpenAIPromptCacheKey(raw json.RawMessage) (json.RawMessage, error)
- func ClampReasoning(effort string) string
- func ClampThinkingBudgetToAnswerRoom(budget, ceiling float64) float64
- func ClampThinkingLevel(rawModel json.RawMessage, level string) (string, error)
- func CloseCodexResponsesSessions(session string)
- func ContentText(content MessageContent, separator ...string) string
- func ContextWindowFor(vendor, model string) int
- func ConvertAnthropicTools(tools []Tool, options AnthropicToolsOptions) (json.RawMessage, error)
- func ConvertOpenAICompletionsMessages(rawModel json.RawMessage, context TranscriptContext, ...) (json.RawMessage, error)
- func ConvertOpenAICompletionsTools(tools []Tool, compat OpenAICompletionsCompat) (json.RawMessage, error)
- func ConvertResponsesMessages(rawModel json.RawMessage, context TranscriptContext, ...) (json.RawMessage, error)
- func ConvertResponsesTools(tools []Tool, options ResponsesToolsOptions) (json.RawMessage, error)
- func CreateGrammarToolInputProperties(tools []Tool, supported bool) (map[string]string, error)
- func DeclarationsEqual(left, right Tool) bool
- func EstimateMessageTokens(message Message) float64
- func EstimateTextAndImageContentTokens(content MessageContent) float64
- func EstimateTextTokens(text string) float64
- func GetCurrentSystemPrompt(messages []Message) string
- func GetGrammarToolInput(toolName string, arguments json.RawMessage, property string) (string, error)
- func GetJSONSchemaToolParameters(tool Tool, strict *bool) (json.RawMessage, error)
- func GetModelType(model any) (any, error)
- func GetSupportedThinkingLevels(rawModel json.RawMessage) ([]string, error)
- func GetSystemMessageText(message Message) string
- func HasAPI(model, api any) (bool, error)
- func HasAnthropicFederationConfig(provider string, rawOptions json.RawMessage) bool
- func HasNonAdditiveToolChanges(messages []Message) bool
- func HasToolRedefinitions(messages []Message) bool
- func IsContextOverflow(err error) bool
- func IsModelType(model, kind any) (bool, error)
- func IsOAuthVendor(vendor string) bool
- func MakeStrictJSONSchema(schema json.RawMessage, check UnsupportedStrictSchemaKeywordCheck) (json.RawMessage, error)
- func ModelsAreEqual(a, b any) (bool, error)
- func NewChatHTTPClient(timeout time.Duration) *http.Client
- func NewClient(spec ClientSpec) (protocol.LLMProvider, error)
- func NewStreamHTTPClient(timeout time.Duration) *http.Client
- func NormalizeOAuthVendor(v string) string
- func NormalizeRadiusGatewayURL(value string) string
- func NormalizeVendor(v string) string
- func OAuthErrorHTML(message string, details ...string) string
- func OAuthSuccessHTML(message string) string
- func ParseJSNumber(value string) float64
- func ParseJSONWithRepair(input string) (json.RawMessage, error)
- func ParseStreamingJSON(input string) json.RawMessage
- func PollOAuthDeviceCodeFlow(options *OAuthDeviceCodePollOptions) (any, error)
- func ProviderDefaultURL(id string) string
- func ProviderPurpose(id string) string
- func RenderSystemMessageUpdate(message Message) string
- func RepairJSON(input string) string
- func ResetCodexWebSocketDebugStats(session string)
- func ResolveJSONSchemaStrictSampling(tool Tool, supported bool, check UnsupportedStrictSchemaKeywordCheck) (*bool, error)
- func ResolveProviderAuth(ctx context.Context, provider *ModelProvider, ...) (any, error)
- func ResolveSamplingParams(rawModel json.RawMessage, thinkingLevel string, requestParams json.RawMessage) (json.RawMessage, error)
- func SanitizeSurrogates(text string) string
- func StringifyJSON(raw []byte) ([]byte, error)
- func Synthesize(ctx context.Context, s SpeechSpec) ([]byte, string, error)
- func ThinkingBudgetForLevel(level string, custom map[string]json.RawMessage) float64
- func ValidateNativeModel(raw json.RawMessage) error
- func ValidateNativeToolChoice(choice protocol.ToolChoice, names []string) error
- func ValidatePiMessagesControls(opts NativeGenerationControls) error
- func WithAssistantStreamSynchronization(ctx context.Context, relay *AssistantMessageEventStream) context.Context
- type APIKeyAuth
- type APIKeyAuthInput
- type AnthropicClientError
- type AnthropicCompat
- type AnthropicConvertedMessages
- type AnthropicMessageClient
- type AnthropicMessagesOptions
- type AnthropicOAuthOptions
- type AnthropicProvider
- func (p *AnthropicProvider) Chat(ctx context.Context, req protocol.ChatRequest) (protocol.ChatResponse, error)
- func (p *AnthropicProvider) ModelCapabilities(model string) protocol.ModelCapabilities
- func (p *AnthropicProvider) Name() string
- func (p *AnthropicProvider) Stream(ctx context.Context, req protocol.ChatRequest) (<-chan protocol.ChatDelta, error)
- func (p *AnthropicProvider) SupportsTools() bool
- func (p *AnthropicProvider) UpdateAPIKey(key string)
- type AnthropicRequestOptions
- type AnthropicStreamOptions
- type AnthropicToolsOptions
- type AntigravityOAuthOptions
- type AntigravityProvider
- func (p *AntigravityProvider) Chat(ctx context.Context, req protocol.ChatRequest) (protocol.ChatResponse, error)
- func (p *AntigravityProvider) ModelCapabilities(model string) protocol.ModelCapabilities
- func (p *AntigravityProvider) Name() string
- func (p *AntigravityProvider) Stream(ctx context.Context, req protocol.ChatRequest) (<-chan protocol.ChatDelta, error)
- func (p *AntigravityProvider) SupportsTools() bool
- type Array
- type AssistantMessageEvent
- type AssistantMessageEventStream
- func LazyStream(ctx context.Context, model any, ...) *AssistantMessageEventStream
- func NewAssistantMessageEventStream() *AssistantMessageEventStream
- func NewAssistantMessageEventStreamFor(ctx context.Context) *AssistantMessageEventStream
- func NewAssistantMessageEventStreamWithHooks(hooks AssistantStreamHooks) *AssistantMessageEventStream
- func StreamAnthropic(ctx context.Context, rawModel json.RawMessage, transcript TranscriptContext, ...) *AssistantMessageEventStream
- func StreamAnthropicSimple(ctx context.Context, rawModel json.RawMessage, transcript TranscriptContext, ...) (*AssistantMessageEventStream, error)
- func StreamAnthropicWithClient(ctx context.Context, rawModel json.RawMessage, transcript TranscriptContext, ...) *AssistantMessageEventStream
- func StreamAntigravityPooled(ctx context.Context, model json.RawMessage, transcript TranscriptContext, ...) (*AssistantMessageEventStream, error)
- func StreamAzureResponses(ctx context.Context, rawModel json.RawMessage, transcript TranscriptContext, ...) *AssistantMessageEventStream
- func StreamAzureResponsesSimple(ctx context.Context, rawModel json.RawMessage, transcript TranscriptContext, ...) (*AssistantMessageEventStream, error)
- func StreamClaudeCodePooled(ctx context.Context, model json.RawMessage, transcript TranscriptContext, ...) (*AssistantMessageEventStream, error)
- func StreamCodexResponses(ctx context.Context, rawModel json.RawMessage, transcript TranscriptContext, ...) *AssistantMessageEventStream
- func StreamCodexResponsesPooled(ctx context.Context, rawModel json.RawMessage, transcript TranscriptContext, ...) (*AssistantMessageEventStream, error)
- func StreamCodexResponsesSSE(ctx context.Context, rawModel json.RawMessage, transcript TranscriptContext, ...) *AssistantMessageEventStream
- func StreamCodexResponsesSimple(ctx context.Context, rawModel json.RawMessage, transcript TranscriptContext, ...) (*AssistantMessageEventStream, error)
- func StreamCodexResponsesSimpleSSE(ctx context.Context, rawModel json.RawMessage, transcript TranscriptContext, ...) (*AssistantMessageEventStream, error)
- func StreamDevinPooled(ctx context.Context, rawModel json.RawMessage, transcript TranscriptContext, ...) (*AssistantMessageEventStream, error)
- func StreamOpenAICompletions(ctx context.Context, rawModel json.RawMessage, transcript TranscriptContext, ...) *AssistantMessageEventStream
- func StreamOpenAICompletionsSimple(ctx context.Context, rawModel json.RawMessage, transcript TranscriptContext, ...) (*AssistantMessageEventStream, error)
- func StreamOpenAIResponses(ctx context.Context, rawModel json.RawMessage, transcript TranscriptContext, ...) *AssistantMessageEventStream
- func StreamOpenAIResponsesSimple(ctx context.Context, rawModel json.RawMessage, transcript TranscriptContext, ...) (*AssistantMessageEventStream, error)
- func StreamPiMessages(ctx context.Context, model *Object, transcript TranscriptContext, ...) *AssistantMessageEventStream
- func StreamPiMessagesJSON(ctx context.Context, rawModel json.RawMessage, transcript TranscriptContext, ...) (*AssistantMessageEventStream, error)
- func StreamSimplePiMessages(ctx context.Context, model *Object, transcript TranscriptContext, ...) *AssistantMessageEventStream
- func (s *AssistantMessageEventStream) End(result ...*Message)
- func (s *AssistantMessageEventStream) Next(ctx context.Context) (AssistantMessageEvent, bool, error)
- func (s *AssistantMessageEventStream) Result(ctx context.Context) (*Message, error)
- func (s *AssistantMessageEventStream) SnapshotEvent(event AssistantMessageEvent) (AssistantMessageEvent, error)
- func (s *AssistantMessageEventStream) SnapshotResult(ctx context.Context) (*Message, error)
- func (s *AssistantMessageEventStream) Synchronize(fn func())
- func (s *AssistantMessageEventStream) SynchronizeYielding(callback func(PayloadYield) error) error
- func (s *AssistantMessageEventStream) WaitForEnd(ctx context.Context) error
- type AssistantStreamHooks
- type AttemptOutcome
- type AttemptTrace
- type AuthContext
- type AuthResolutionOverrides
- type AzureResponsesConfig
- type AzureResponsesStreamOptions
- type BlockList
- func (l *BlockList) Append(blocks ...*ContentBlock) int
- func (l *BlockList) Delete(index int)
- func (l *BlockList) Get(index int) *ContentBlock
- func (l *BlockList) Len() int
- func (l *BlockList) MarshalJSON() ([]byte, error)
- func (l *BlockList) Pop() *ContentBlock
- func (l *BlockList) Set(index int, block *ContentBlock)
- func (l *BlockList) SetLength(length int)
- func (l *BlockList) UnmarshalJSON(data []byte) error
- func (l *BlockList) Values() []*ContentBlock
- type ChatGPTOAuthAuthorization
- type ChatGPTOAuthCallback
- type ChatGPTOAuthCallbackOptions
- type ClientSpec
- type CodexProvider
- func (p *CodexProvider) Chat(ctx context.Context, req protocol.ChatRequest) (protocol.ChatResponse, error)
- func (p *CodexProvider) ModelCapabilities(model string) protocol.ModelCapabilities
- func (p *CodexProvider) Name() string
- func (p *CodexProvider) Stream(ctx context.Context, req protocol.ChatRequest) (<-chan protocol.ChatDelta, error)
- func (p *CodexProvider) SupportsTools() bool
- type CodexResponsesStreamOptions
- type CodexWebSocketDebugStats
- type Collection
- func (c *Collection) Chat(ctx context.Context, req protocol.ChatRequest) (protocol.ChatResponse, error)
- func (c *Collection) ChatOn(ctx context.Context, providerID string, req protocol.ChatRequest) (protocol.ChatResponse, error)
- func (c *Collection) Get(id string) (Provider, bool)
- func (c *Collection) ListModels(ctx context.Context) (ListResult, error)
- func (c *Collection) ModelCapabilities(model string) protocol.ModelCapabilities
- func (c *Collection) Name() string
- func (c *Collection) Owner(ctx context.Context, modelID string) (Provider, error)
- func (c *Collection) Providers() []Provider
- func (c *Collection) Register(p Provider)
- func (c *Collection) Stream(ctx context.Context, req protocol.ChatRequest) (<-chan protocol.ChatDelta, error)
- func (c *Collection) StreamOn(ctx context.Context, providerID string, req protocol.ChatRequest) (<-chan protocol.ChatDelta, error)
- func (c *Collection) SupportsTools() bool
- type Compat
- type CompletionsResponse
- type ContentBlock
- type Context
- type ContextUsageEstimate
- type CredentialPersistence
- type DataCloneError
- type DeferredStreamFunc
- type DevinOAuthOptions
- type DevinProvider
- func (p *DevinProvider) Chat(ctx context.Context, req protocol.ChatRequest) (protocol.ChatResponse, error)
- func (p *DevinProvider) ModelCapabilities(model string) protocol.ModelCapabilities
- func (p *DevinProvider) Name() string
- func (p *DevinProvider) Stream(ctx context.Context, req protocol.ChatRequest) (<-chan protocol.ChatDelta, error)
- func (p *DevinProvider) SupportsTools() bool
- type EventStream
- type FallbackAttempt
- type FallbackCandidate
- type FallbackProvider
- type FallbackRequest
- type FallbackSelection
- type GitHubCopilotOAuthOptions
- type GrammarConstrainedSampling
- type GrammarToolInputJSONBuffer
- type HTTPDoer
- type InMemoryCredentialStore
- func (s *InMemoryCredentialStore) Delete(ctx context.Context, providerID string) error
- func (s *InMemoryCredentialStore) List(ctx context.Context) (*Array, error)
- func (s *InMemoryCredentialStore) Modify(ctx context.Context, providerID string, modify func(any) (any, error)) (any, error)
- func (s *InMemoryCredentialStore) Read(ctx context.Context, providerID string) (any, error)
- type InMemoryModelsStore
- type JSONMethod
- type KimiCodingOAuthOptions
- type LazyAPICapabilities
- type LazyOAuthOptions
- type ListError
- type ListResult
- type Message
- func CreateInitialSystemMessage(systemPrompt string, tools []Tool) *Message
- func GetCurrentSystemMessage(messages []Message) *Message
- func GetInitialSystemMessage(messages []Message) *Message
- func TransformMessageReferences(messages []*Message, model *Model, normalize ToolCallIDNormalizer) []*Message
- func TransformMessages(messages []Message, model Model, normalize ToolCallIDNormalizer) []Message
- func WithToolChanges(message Message, changes ToolStateChanges) Message
- func WithoutInitialSystemMessage(messages []Message) []Message
- type MessageContent
- type MetaOAuthOptions
- type Model
- type ModelPrice
- type ModelProvider
- type ModelRefreshContext
- type ModelStreamFunc
- type Models
- func (m *Models) CancelDeferred(ctx context.Context, model, handle any, options ...*ModelsRequestOptions) error
- func (m *Models) CheckAuth(ctx context.Context, providerID string) (any, error)
- func (m *Models) Classify(ctx context.Context, model, input any, options ...*ModelsRequestOptions) (any, error)
- func (m *Models) ClearProviders()
- func (m *Models) Complete(ctx context.Context, model any, input Context, ...) (*Message, error)
- func (m *Models) CompleteSimple(ctx context.Context, model any, input Context, ...) (*Message, error)
- func (m *Models) DeleteProvider(id string)
- func (m *Models) FetchDeferred(ctx context.Context, model, handle any, options ...*ModelsRequestOptions) (*Message, error)
- func (m *Models) GenerateImages(ctx context.Context, model, input any, options ...*ModelsRequestOptions) (any, error)
- func (m *Models) GetAllAvailable(ctx context.Context, providerID ...string) (*Array, error)
- func (m *Models) GetAllModels(providerID ...string) *Array
- func (m *Models) GetAuth(ctx context.Context, providerOrModel any, options ...AuthResolutionOverrides) (any, error)
- func (m *Models) GetAvailable(ctx context.Context, providerID ...string) (*Array, error)
- func (m *Models) GetAvailableOfType(ctx context.Context, kind any, providerID ...string) (*Array, error)
- func (m *Models) GetModel(providerID string, id any) (any, error)
- func (m *Models) GetModelOfType(kind any, providerID string, id any) (any, error)
- func (m *Models) GetModels(providerID ...string) *Array
- func (m *Models) GetModelsOfType(kind any, providerID ...string) (*Array, error)
- func (m *Models) GetProvider(id string) *ModelProvider
- func (m *Models) GetProviders() []*ModelProvider
- func (m *Models) Login(providerID, authType string, interaction ProviderAuthInteraction, ...) (any, error)
- func (m *Models) Logout(ctx context.Context, providerID string) error
- func (m *Models) Refresh(ctx context.Context, options ...ModelsRefreshOptions) ModelsRefreshResult
- func (m *Models) SetProvider(provider *ModelProvider)
- func (m *Models) Stream(ctx context.Context, model any, input Context, ...) *AssistantMessageEventStream
- func (m *Models) StreamDeferred(ctx context.Context, model, handle any, options ...*ModelsRequestOptions) *AssistantMessageEventStream
- func (m *Models) StreamSimple(ctx context.Context, model any, input Context, ...) *AssistantMessageEventStream
- type ModelsError
- type ModelsOptions
- type ModelsPersistence
- type ModelsPublication
- type ModelsPublicationScope
- type ModelsPublisher
- type ModelsRefreshFailure
- type ModelsRefreshOptions
- type ModelsRefreshResult
- type ModelsRequestOptions
- type NativeClient
- type NativeGenerationControls
- type NativeProvider
- type NativeProviderFailure
- type OAuthAuth
- func AnthropicOAuth(settings ...AnthropicOAuthOptions) *OAuthAuth
- func AntigravityOAuth(settings ...AntigravityOAuthOptions) *OAuthAuth
- func DevinOAuth(settings ...DevinOAuthOptions) *OAuthAuth
- func GitHubCopilotOAuth(options GitHubCopilotOAuthOptions) (*OAuthAuth, error)
- func KimiCodingOAuth(settings ...KimiCodingOAuthOptions) *OAuthAuth
- func LazyOAuth(input *LazyOAuthOptions) *OAuthAuth
- func MetaOAuth(settings ...MetaOAuthOptions) *OAuthAuth
- func OpenAIChatGPTOAuth(settings ...OpenAIChatGPTOAuthOptions) *OAuthAuth
- func OpenAICodexOAuth(settings ...OpenAICodexOAuthOptions) *OAuthAuth
- func OpenRouterOAuth(settings ...OpenRouterOAuthOptions) *OAuthAuth
- func RadiusOAuth(input *RadiusOAuthOptions, settings ...RadiusOAuthRuntimeOptions) *OAuthAuth
- func XaiOAuth(settings ...XaiOAuthOptions) *OAuthAuth
- type OAuthCallbackServer
- type OAuthCallbackServerOptions
- type OAuthDeviceCodePollOptions
- type OAuthDeviceCodePollResult
- type OAuthDiagnosticError
- type OAuthLoginOptions
- type OAuthManualPrompt
- type OAuthToken
- type Object
- func FlattenChatModelCatalog(_ string, groups any) (*Object, error)
- func FlattenClassifierModelCatalog(_ string, groups any) (*Object, error)
- func FlattenImageModelCatalog(_ string, groups any) (*Object, error)
- func GetRadiusCredentialConfig(credential any) *Object
- func LoadRadiusGatewayConfig(ctx context.Context, gateway string, apiKey any, clients ...*http.Client) (*Object, error)
- func NewObject(properties ...Property) *Object
- func WaitForCallbackOrManualInput(interaction ProviderAuthInteraction, callback *OAuthCallbackServer, ...) (*Object, error)
- type OpenAIChatGPTOAuthOptions
- type OpenAICodexOAuthOptions
- type OpenAICompletionsCompat
- type OpenAICompletionsStreamOptions
- type OpenAIEmbedder
- type OpenAIProvider
- func (p *OpenAIProvider) Chat(ctx context.Context, req protocol.ChatRequest) (protocol.ChatResponse, error)
- func (p *OpenAIProvider) ModelCapabilities(model string) protocol.ModelCapabilities
- func (p *OpenAIProvider) Name() string
- func (p *OpenAIProvider) Stream(ctx context.Context, req protocol.ChatRequest) (<-chan protocol.ChatDelta, error)
- func (p *OpenAIProvider) SupportsTools() bool
- func (p *OpenAIProvider) UpdateAPIKey(key string)
- type OpenAIResponsesCompat
- type OpenAIResponsesProvider
- func (p *OpenAIResponsesProvider) Chat(ctx context.Context, req protocol.ChatRequest) (protocol.ChatResponse, error)
- func (p *OpenAIResponsesProvider) ModelCapabilities(model string) protocol.ModelCapabilities
- func (p *OpenAIResponsesProvider) Name() string
- func (p *OpenAIResponsesProvider) Stream(ctx context.Context, req protocol.ChatRequest) (<-chan protocol.ChatDelta, error)
- func (p *OpenAIResponsesProvider) SupportsTools() bool
- func (p *OpenAIResponsesProvider) UpdateAPIKey(key string)
- type OpenAIResponsesStreamOptions
- type OpenAIWire
- type OpenRouterOAuthOptions
- type PKCE
- type PayloadYield
- type PiMessagesResponseError
- type PiMessagesStreamOptions
- type PreparationError
- type Pricing
- type Property
- type Provider
- type ProviderAuth
- type ProviderAuthInteraction
- type ProviderClassifier
- type ProviderEventSource
- type ProviderFactoryOptions
- type ProviderImages
- type ProviderModelCatalog
- type ProviderOption
- type ProviderStreams
- func AnthropicMessagesAPI(defaults OpenAICompletionsStreamOptions) *ProviderStreams
- func AzureOpenAIResponsesAPI(defaults OpenAICompletionsStreamOptions) *ProviderStreams
- func LazyAPI(load func(context.Context) (*ProviderStreams, error), ...) *ProviderStreams
- func OpenAICodexResponsesAPI(defaults OpenAICompletionsStreamOptions) *ProviderStreams
- func OpenAICompletionsAPI(defaults OpenAICompletionsStreamOptions) *ProviderStreams
- func OpenAIResponsesAPI(defaults OpenAICompletionsStreamOptions) *ProviderStreams
- func PiMessagesAPI(defaults PiMessagesStreamOptions) *ProviderStreams
- type RadiusOAuthOptions
- type RadiusOAuthResponseError
- type RadiusOAuthRuntimeOptions
- type RadiusProviderDependencies
- type RadiusProviderOptions
- type ResponsesMessagesOptions
- type ResponsesToolsOptions
- type Spec
- type SpeechSpec
- type StreamFn
- type SystemSection
- type SystemSections
- type ThinkingTokenBudget
- type TokenSource
- type Tool
- type ToolCallIDNormalizer
- type ToolReference
- type ToolStateChanges
- type TranscriptContext
- type TranscriptTools
- type UnsupportedStrictJSONSchemaError
- type UnsupportedStrictSchemaKeywordCheck
- type Usage
- type UsageCost
- type UsageObservation
- type XaiOAuthOptions
Constants ¶
const ( // DefaultChatTimeout bounds one chat exchange end to end, body included. DefaultChatTimeout = 10 * time.Minute // DefaultResponseHeaderTimeout bounds the wait for response headers on a // STREAMED request — the "is this gateway alive" check the absolute cap was // mistakenly doing. It is meaningless on a buffered completion. DefaultResponseHeaderTimeout = 60 * time.Second )
Timeouts for the chat providers. There are two knobs, they do different jobs, and which of them applies depends on whether the request is STREAMED.
Go's http.Client.Timeout is absolute: it covers dialing, the response headers, AND the full body read. For a chat provider the body IS the model generating — a turn that produces a complete HTML page can legitimately spend minutes there. Sizing that cap like a handshake (the old 120s) does not protect against a dead gateway; it kills the slowest *healthy* turns, and because the kill is deterministic the retry ladder then reproduces it on every attempt.
ResponseHeaderTimeout looks like the safe replacement — time to first response header is a real liveness signal — but only on an SSE request. There the gateway writes its 200 immediately and the tokens follow, so headers precede generation. On a NON-streamed completion the response is a single buffered JSON document: headers do not leave the gateway until generation has finished, so a "header" deadline is a generation deadline wearing a disguise. Bench #4 proved it — the first turn generated 9.5k tokens, the header deadline fired at exactly 60.3s, and the (successful) retry then spent 135s regenerating the same answer: 60s of dead wall clock and a doubled bill for a healthy call.
So: the absolute cap applies to both paths and is sized for the slowest legitimate turn; the header deadline applies to the streamed path ONLY.
const ( VendorClaudeCode = "claude-code" VendorOpenAICodex = "openai-codex" VendorGoogleAntigravity = "google-antigravity" VendorDevin = "devin" VendorXaiOAuth = "xai-oauth" )
OAuth vendor ids: subscription providers whose credential is a pool of OAuth accounts (many logins per provider row) rather than a single API key. The vendor id is the provider row's `vendor` column and the wire client's identity — it is NOT the API-key vendor it resembles on the wire (claude-code speaks the Anthropic Messages API but authenticates with a Bearer grant, not x-api-key).
const ( // AntigravityUserAgent mirrors the real antigravity/hub client; the backend // gates model availability on the version it carries. Exported because the // OAuth login path (internal/oauth) must present the same fingerprint — a // version bump in one place only would desynchronize login from chat. AntigravityUserAgent = "antigravity/hub/2.8.0 (aidev_client; os_type=darwin; arch=arm64; cl=963137146)" )
const DefaultRadiusGateway = "https://radius.pi.dev"
const (
// DevinDefaultBaseURL is the Cascade origin used when a provider row sets none.
DevinDefaultBaseURL = "https://server.codeium.com"
)
Devin (Cognition SWE / Codeium Cascade) speaks a Connect-RPC protobuf surface at server.codeium.com, not an OpenAI-shaped wire. This file holds the protocol.LLMProvider adapter (Chat/Stream over the native event stream) and the model listing/uid resolution that the OAuth vendor registry expects; the wire producer itself lives in devin_native.go.
const MinAnswerTokens = 1024
const Null = jsonjs.Null
const OAuthMinimumValidityMS = 5 * 60 * 1000
const OAuthPoolKey = "oauth-pool"
OAuthPoolKey is the sentinel ResolveWorkspaceRun puts in the per-tier key map for an OAuth vendor. It is never sent on the wire — pooled providers ignore their static key and pull a live access token from their TokenSource per request — but it must be non-empty so the run path's "flash key configured" gate treats a pooled provider as configured.
const OAuthRefreshTimeout = 15 * time.Second
const OpenAIPromptCacheKeyMaxLength = 64
const Undefined = jsonjs.Undefined
const VendorAzureResponses = "azure-openai-responses"
const (
VendorOpenAIResponses = "openai-responses"
)
const VendorPiMessages = "pi-messages"
Variables ¶
This section is empty.
Functions ¶
func APIKeyOptional ¶
APIKeyOptional reports OpenAI-compatible engines whose normal local/self- hosted deployment accepts unauthenticated requests. These are explicit vendor identities rather than a hostname heuristic so the same configuration works on a laptop (localhost) and a server (a private service/DNS name). A supplied key is still sent, which supports secured deployments of the same engines.
func AbortableSleep ¶
AbortableSleep replaces the abort cause with the requested cancellation message, as Pi does. Delay conversion follows the native timer boundary.
func AppendGrammarToolInputJSONDelta ¶
func AppendGrammarToolInputJSONDelta(buffer *GrammarToolInputJSONBuffer, property, next string, close bool) (*string, error)
AppendGrammarToolInputJSONDelta returns nil for an unchanged open input or a duplicate close. A non-nil delta is an exact JSON string fragment.
func ApplyNativeControls ¶
func ApplyNativeControls(api string, raw json.RawMessage, opts NativeGenerationControls) (json.RawMessage, error)
ApplyNativeControls adds host generation controls through Pi's onPayload hook. Pi still owns all message encoding, signatures, caching, and streaming. Raw fields preserve extension values and numbers without a float64 roundtrip.
func AssertChatModel ¶
func AssertClassifierModel ¶
func AssertImageModel ¶
func BuildAnthropicParams ¶
func BuildAnthropicParams(rawModel json.RawMessage, context TranscriptContext, oauth bool, rawOptions json.RawMessage) (json.RawMessage, error)
BuildAnthropicParams ports Pi's request builder. OAuth is supplied by the client/auth selection; this function does not resolve tokens or federation. The caller resolves mid-conversation system-message support beforehand.
func BuildAnthropicSimpleOptions ¶
func BuildAnthropicSimpleOptions(rawModel json.RawMessage, transcript TranscriptContext, rawOptions json.RawMessage) (json.RawMessage, error)
BuildAnthropicSimpleOptions projects Pi's simple controls, including adaptive effort and the two context clamps around budget-based thinking.
func BuildAzureResponsesParams ¶
func BuildAzureResponsesParams(rawModel json.RawMessage, transcript TranscriptContext, rawOptions json.RawMessage) (json.RawMessage, error)
BuildAzureResponsesParams ports the Azure request builder, including its distinct strict-tool defaults, deployment name and prompt-cache behavior. The input is a normalized transcript; no SDK or TypeScript code is involved.
func BuildAzureResponsesSimpleOptions ¶
func BuildAzureResponsesSimpleOptions(model json.RawMessage, transcript TranscriptContext, options json.RawMessage) (json.RawMessage, error)
BuildAzureResponsesSimpleOptions retains only Pi's shared simple controls; Azure-specific settings must come through the scoped/process environment.
func BuildCodexResponsesParams ¶
func BuildCodexResponsesParams(rawModel json.RawMessage, context TranscriptContext, rawOptions json.RawMessage, cacheSessionID *string) (json.RawMessage, error)
BuildCodexResponsesParams ports Pi's Codex request builder. The transcript must already be normalized/resolved, and cacheSessionID is the transport's chosen cache identity (nil preserves an absent prompt_cache_key).
func BuildCodexResponsesSimpleOptions ¶
func BuildCodexResponsesSimpleOptions(rawModel json.RawMessage, transcript TranscriptContext, rawOptions json.RawMessage) (json.RawMessage, error)
BuildCodexResponsesSimpleOptions uses the source's shared base-option and thinking-level projection. Provider-only effort/summary/tier options do not bypass the simple-mode reasoning contract.
func BuildOpenAICompletionsParams ¶
func BuildOpenAICompletionsParams(rawModel json.RawMessage, context TranscriptContext, rawOptions json.RawMessage) (json.RawMessage, error)
BuildOpenAICompletionsParams ports Pi's request-body builder. Options are the provider-specific (not streamSimple) options. Credentials, callbacks and HTTP configuration never enter the body unless explicitly set by samplingParams, whose final override semantics match the original provider.
func BuildOpenAICompletionsSimpleOptions ¶
func BuildOpenAICompletionsSimpleOptions(rawModel json.RawMessage, transcript TranscriptContext, rawOptions json.RawMessage) (json.RawMessage, error)
BuildOpenAICompletionsSimpleOptions ports streamSimple's serializable option projection. Concrete Go callbacks, HTTP client and signal remain outside JSON.
func BuildOpenAIResponsesParams ¶
func BuildOpenAIResponsesParams(rawModel json.RawMessage, context TranscriptContext, rawOptions json.RawMessage) (json.RawMessage, error)
BuildOpenAIResponsesParams builds the native Responses request body. Callers supply the resolved transcript, as Pi's stream does before invoking its builder.
func BuildOpenAIResponsesSimpleOptions ¶
func BuildOpenAIResponsesSimpleOptions(rawModel json.RawMessage, transcript TranscriptContext, rawOptions json.RawMessage) (json.RawMessage, error)
BuildOpenAIResponsesSimpleOptions follows Responses' streamSimple projection: thinkingBudgets, reasoningSummary and serviceTier are not forwarded.
func CalculateContextTokens ¶
func CapabilitiesFor ¶
func CapabilitiesFor(vendor, model string) protocol.ModelCapabilities
CapabilitiesFor returns adapter-level knowledge and conservative wire-policy limits for a provider/model pair. It intentionally stays sparse: this is not a model catalog, and unknown support preserves the optimistic request path. Live discovery may overlay these defaults with explicit model-level facts.
func ClampMaxTokensToContext ¶
func ClampMaxTokensToContext(contextWindow float64, transcript TranscriptContext, maxTokens float64) float64
func ClampOpenAIPromptCacheKey ¶
func ClampOpenAIPromptCacheKey(raw json.RawMessage) (json.RawMessage, error)
ClampOpenAIPromptCacheKey counts Unicode code points as Pi does with Array.from. A nil raw input represents an absent key. For serialized non-string inputs it retains the original value when the array-like length does not exceed the limit.
func ClampReasoning ¶
func ClampThinkingLevel ¶
func ClampThinkingLevel(rawModel json.RawMessage, level string) (string, error)
func CloseCodexResponsesSessions ¶
func CloseCodexResponsesSessions(session string)
CloseCodexResponsesSessions releases cached sockets for one session, or all sessions when session is empty. As in Pi, closing does not reset fallback memory.
func ContentText ¶
func ContentText(content MessageContent, separator ...string) string
ContentText follows Pi's filter/map/join reads. Malformed content panics with its property-read error; valid sparse arrays skip absent slots.
func ContextWindowFor ¶
ContextWindowFor returns the input context window for a model id, or 0 when this package has nothing better than a guess. Zero is a real answer — "no opinion" — and callers must treat it as such rather than as a window of size zero.
Vendor is accepted for future per-vendor disambiguation and is currently only normalized; ids are matched on their own because a workspace reaches the same model through several vendor ids (openai, openai-compat, a router).
func ConvertAnthropicTools ¶
func ConvertAnthropicTools(tools []Tool, options AnthropicToolsOptions) (json.RawMessage, error)
ConvertAnthropicTools retains Pi's legacy schema projection unless a supported strict constraint is selected. Only the last declaration gets a cache marker.
func ConvertOpenAICompletionsMessages ¶
func ConvertOpenAICompletionsMessages(rawModel json.RawMessage, context TranscriptContext, compat OpenAICompletionsCompat, grammarProperties map[string]string) (json.RawMessage, error)
ConvertOpenAICompletionsMessages ports Pi's exported convertMessages. The input remains a native transcript; no legacy agentcore.Message projection is used. Grammar properties come from CreateGrammarToolInputProperties.
func ConvertOpenAICompletionsTools ¶
func ConvertOpenAICompletionsTools(tools []Tool, compat OpenAICompletionsCompat) (json.RawMessage, error)
ConvertOpenAICompletionsTools shares Pi's strict and grammar decisions with both top-level declarations and mid-conversation tool additions.
func ConvertResponsesMessages ¶
func ConvertResponsesMessages(rawModel json.RawMessage, context TranscriptContext, allowedToolCallProviders []string, options ResponsesMessagesOptions) (json.RawMessage, error)
ConvertResponsesMessages ports the shared Responses transcript conversion. allowedToolCallProviders is supplied by the concrete provider; it controls cross-provider call-ID normalization, not signature replay permissions.
func ConvertResponsesTools ¶
func ConvertResponsesTools(tools []Tool, options ResponsesToolsOptions) (json.RawMessage, error)
ConvertResponsesTools converts declarations for the Responses API, including strict schemas, grammar tools and deferred tool-search results.
func DeclarationsEqual ¶
func EstimateMessageTokens ¶
func EstimateTextAndImageContentTokens ¶
func EstimateTextAndImageContentTokens(content MessageContent) float64
func EstimateTextTokens ¶
func GetCurrentSystemPrompt ¶
func GetGrammarToolInput ¶
func GetJSONSchemaToolParameters ¶
func GetJSONSchemaToolParameters(tool Tool, strict *bool) (json.RawMessage, error)
GetJSONSchemaToolParameters only applies strict conversion when explicitly selected. Otherwise the original parameter bytes remain authoritative.
func GetModelType ¶
GetModelType uses nullish fallback, retaining unknown and malformed explicit type values. Unlike catalog flattening, an absent/null type means chat.
func GetSupportedThinkingLevels ¶
func GetSupportedThinkingLevels(rawModel json.RawMessage) ([]string, error)
func GetSystemMessageText ¶
func HasAnthropicFederationConfig ¶
func HasAnthropicFederationConfig(provider string, rawOptions json.RawMessage) bool
HasAnthropicFederationConfig reports whether Pi would select federation for these controls. It only resolves configuration; it never reads or exchanges an identity token. Hosts can use it before allowing an absent API key.
func HasToolRedefinitions ¶
func IsContextOverflow ¶
IsContextOverflow recognizes explicit context rejections, not transient errors, arbitrary payload limits or an ambiguous transport interruption.
func IsModelType ¶
func IsOAuthVendor ¶
IsOAuthVendor reports whether vendor is one of the subscription/OAuth pool vendors. Such a provider row has no api_key; it owns a pool of accounts in workspace_provider_accounts.
func MakeStrictJSONSchema ¶
func MakeStrictJSONSchema(schema json.RawMessage, check UnsupportedStrictSchemaKeywordCheck) (json.RawMessage, error)
MakeStrictJSONSchema ports Pi's constrained-sampling conversion. It returns a detached schema, requires every object property, and makes optional fields nullable. Unsupported constructs fail instead of being silently discarded.
func ModelsAreEqual ¶
func NewChatHTTPClient ¶
NewChatHTTPClient builds the HTTP client for NON-streamed chat completions: a generous absolute deadline and no header deadline, because the headers of a buffered completion arrive only once the model has finished. A zero timeout uses DefaultChatTimeout.
Exported so a caller that knows its own latency envelope can size the cap itself and assign it to the provider's HTTP field, rather than inheriting a default chosen for the slowest case.
func NewClient ¶
func NewClient(spec ClientSpec) (protocol.LLMProvider, error)
NewClient resolves a ClientSpec into an protocol.LLMProvider. Adding a vendor is additive here — a new case (or, for OpenAI-compatible vendors, just a compat entry + base_url) — and never requires touching the agent loop.
func NewStreamHTTPClient ¶
NewStreamHTTPClient builds the HTTP client for SSE chat requests: the same absolute cap plus the header deadline, which is honest here — an SSE gateway flushes its headers before the first token, so a silent wait past DefaultResponseHeaderTimeout means the endpoint is not answering at all.
func NormalizeOAuthVendor ¶
NormalizeOAuthVendor folds aliases onto the canonical vendor ids. Exported because the oauth package's descriptor lookup must agree with this table — a second copy has already drifted once (missing case-fold, wrong zero value).
func NormalizeRadiusGatewayURL ¶
NormalizeRadiusGatewayURL retains the source text, adding a scheme only when no HTTP(S) prefix is present and removing trailing slashes.
func NormalizeVendor ¶
NormalizeVendor maps aliases onto the built-in vendor ids.
func OAuthErrorHTML ¶
func OAuthSuccessHTML ¶
func ParseJSNumber ¶
ParseJSNumber implements Number(string) for Go ports of Pi's coercion paths. Empty/whitespace input becomes zero; invalid syntax becomes NaN. It accepts unsigned hex, binary and octal integers of arbitrary width and rounds them to binary64. Overflow becomes infinity. Callers apply their own finite/integer and empty-input policies, as the upstream call sites do.
func ParseJSONWithRepair ¶
func ParseJSONWithRepair(input string) (json.RawMessage, error)
ParseJSONWithRepair returns the parsed JSON representation, retaining object member order. Callers can unmarshal it into their desired Go type. Invalid syntax still fails; this function does not complete unfinished documents.
func ParseStreamingJSON ¶
func ParseStreamingJSON(input string) json.RawMessage
ParseStreamingJSON ports Pi's parseStreamingJson, including the permissive partial-json 0.1.7 fallback. Its result is JSON, not necessarily an object: complete scalar/null/array inputs retain their values. Unparseable input and a partial null resolve to {}. No errors escape this streaming helper.
func PollOAuthDeviceCodeFlow ¶
func PollOAuthDeviceCodeFlow(options *OAuthDeviceCodePollOptions) (any, error)
func ProviderDefaultURL ¶
func ProviderPurpose ¶
func RepairJSON ¶
RepairJSON ports Pi's repairJson: it only repairs control characters and invalid escapes inside quoted strings, leaving syntax outside them intact.
func ResetCodexWebSocketDebugStats ¶
func ResetCodexWebSocketDebugStats(session string)
ResetCodexWebSocketDebugStats clears counters and sticky fallback for one session (or all sessions if empty). It does not close cached connections.
func ResolveJSONSchemaStrictSampling ¶
func ResolveJSONSchemaStrictSampling(tool Tool, supported bool, check UnsupportedStrictSchemaKeywordCheck) (*bool, error)
func ResolveProviderAuth ¶
func ResolveProviderAuth(ctx context.Context, provider *ModelProvider, credentials *CredentialPersistence, authContext *AuthContext, overrides AuthResolutionOverrides) (any, error)
func ResolveSamplingParams ¶
func ResolveSamplingParams(rawModel json.RawMessage, thinkingLevel string, requestParams json.RawMessage) (json.RawMessage, error)
ResolveSamplingParams applies model defaults, effective-thinking-level defaults, then request overrides. A nil result means no sampling object was supplied.
func SanitizeSurrogates ¶
SanitizeSurrogates ports Pi's provider-text filter. Transcripts retain lone UTF-16 units as WTF-8; only the provider fields selected by Pi drop them. Valid pairs, including pairs assembled by concatenating WTF-8 strings, survive.
func StringifyJSON ¶
StringifyJSON normalizes serialized JSON using JavaScript value semantics, retaining numeric index-key ordering, binary64 numbers and UTF-16 surrogates.
func Synthesize ¶
Synthesize buffers bounded audio, so an unsuccessful attempt can safely fall back before any audio is exposed. Redirects must not forward provider keys.
func ThinkingBudgetForLevel ¶
func ThinkingBudgetForLevel(level string, custom map[string]json.RawMessage) float64
func ValidateNativeModel ¶
func ValidateNativeModel(raw json.RawMessage) error
ValidateNativeModel rejects APIs without a native Go transport.
func ValidateNativeToolChoice ¶
func ValidateNativeToolChoice(choice protocol.ToolChoice, names []string) error
func ValidatePiMessagesControls ¶
func ValidatePiMessagesControls(opts NativeGenerationControls) error
Gateway controls are nested and tools are declared by transcript system messages. Do not add OpenAI fields which this protocol does not forward.
func WithAssistantStreamSynchronization ¶
func WithAssistantStreamSynchronization(ctx context.Context, relay *AssistantMessageEventStream) context.Context
WithAssistantStreamSynchronization lets a host relay native provider events without copying Pi's live message pointers. Provider streams constructed with NewAssistantMessageEventStreamFor share the relay's payload lock. The context carries only synchronization, never a queue, result, or cancellation policy.
Types ¶
type APIKeyAuth ¶
type APIKeyAuth struct {
Name string
Login func(ProviderAuthInteraction, ...*OAuthLoginOptions) (any, error)
Check func(APIKeyAuthInput) (any, error)
Resolve func(APIKeyAuthInput) (any, error)
}
func EnvAPIKeyAuth ¶
func EnvAPIKeyAuth(name string, envVars *Array) *APIKeyAuth
EnvAPIKeyAuth retains Pi's stored-key-first resolution and interactive login. envVars is a live array of string variable names. Environment values and credential fields remain unmodified, including whitespace and scoped config.
type APIKeyAuthInput ¶
type APIKeyAuthInput struct {
Context context.Context
AuthContext *AuthContext
Credential any
}
type AnthropicClientError ¶
AnthropicClientError retains the status/headers used by Pi's retry policy. Status zero denotes an unspecified status (for example a connection failure).
func (*AnthropicClientError) Error ¶
func (e *AnthropicClientError) Error() string
type AnthropicCompat ¶
type AnthropicCompat struct {
SupportsEagerToolInputStreaming bool `json:"supportsEagerToolInputStreaming"`
SupportsLongCacheRetention bool `json:"supportsLongCacheRetention"`
SendSessionAffinityHeaders bool `json:"sendSessionAffinityHeaders"`
SessionAffinityFormat *string `json:"sessionAffinityFormat,omitempty"`
SupportsCacheControlOnTools bool `json:"supportsCacheControlOnTools"`
SupportsTemperature bool `json:"supportsTemperature"`
AllowEmptySignature bool `json:"allowEmptySignature"`
SupportsStrictTools bool `json:"supportsStrictTools"`
SupportsMidConvoSystemMessages bool `json:"supportsMidConvoSystemMessages"`
SupportsMidConvoToolChanges bool `json:"supportsMidConvoToolChanges"`
}
func ResolveAnthropicCompat ¶
func ResolveAnthropicCompat(rawModel json.RawMessage) (AnthropicCompat, error)
type AnthropicConvertedMessages ¶
type AnthropicConvertedMessages struct {
Messages json.RawMessage `json:"messages"`
AssistantLevels map[int]string `json:"assistantLevels"`
}
func ConvertAnthropicMessages ¶
func ConvertAnthropicMessages(messages []Message, options AnthropicMessagesOptions) (AnthropicConvertedMessages, error)
ConvertAnthropicMessages consumes an already transformed conversation (the initial system message is sent separately). Later system updates are delayed until the next assistant so tool results remain adjacent to their tool use.
type AnthropicMessageClient ¶
type AnthropicMessageClient struct {
CreateResponse func(context.Context, json.RawMessage, AnthropicRequestOptions) (*http.Response, error)
}
AnthropicMessageClient is the Go equivalent of Pi's injected messaging client. It owns authentication, request encoding and endpoint adaptation. CreateResponse must honor ctx and return an unread response body, or an error. The provider closes returned bodies and owns callbacks, SSE decoding and accumulation.
type AnthropicMessagesOptions ¶
type AnthropicMessagesOptions struct {
OAuth bool `json:"oauth"`
CacheControl json.RawMessage `json:"cacheControl,omitempty"`
AllowEmptySignature bool `json:"allowEmptySignature"`
ManagedProvider *string `json:"managedProvider,omitempty"`
ConvertToolDefinitions func([]Tool) (json.RawMessage, error) `json:"-"`
}
type AnthropicOAuthOptions ¶
type AnthropicOAuthOptions struct {
Client *http.Client
PKCE func() (PKCE, error)
StartCallback func(*OAuthCallbackServerOptions) (*OAuthCallbackServer, error)
CallbackHost *string
Now func() float64
RequestTimeout time.Duration
ErrorStack func(name, message string) string
}
AnthropicOAuthOptions supplies concrete native dependencies. CallbackHost is captured when the flow is constructed, like Pi's loaded module constant. ErrorStack controls stacks for errors created by this flow; the default is a native Go stack. Supplied OAuthDiagnosticError values retain their own stacks.
type AnthropicProvider ¶
type AnthropicProvider struct {
APIKey string
BaseURL string
// OAuth marks the claude-code subscription mode: the credential is an OAuth
// access token sent as Authorization: Bearer (never x-api-key), the request
// carries the Claude Code client fingerprint (User-Agent, x-app, the fixed
// beta set), and the system prompt gains the Claude Code identity block.
OAuth bool
// HTTP serves the non-streamed Chat path (absolute cap, no header deadline).
HTTP *http.Client
// StreamHTTP serves the SSE path, where a header deadline is meaningful. Nil
// falls back to HTTP, so a caller that overrides only HTTP still works.
StreamHTTP *http.Client
// contains filtered or unexported fields
}
AnthropicProvider speaks the Anthropic Messages API. The wire format diverges from OpenAI completions (top-level system, tool_use/tool_result content blocks, x-api-key auth), so per the spec it gets its own implementation rather than a compat entry — and proves the protocol.LLMProvider seam (§12 AC: a second provider needs no edits to agent.go).
func NewAnthropicProvider ¶
func NewAnthropicProvider(apiKey, baseURL string) *AnthropicProvider
NewAnthropicProvider builds a provider; an empty baseURL uses the vendor default.
func (*AnthropicProvider) Chat ¶
func (p *AnthropicProvider) Chat(ctx context.Context, req protocol.ChatRequest) (protocol.ChatResponse, error)
Chat performs one non-streaming Messages call.
func (*AnthropicProvider) ModelCapabilities ¶
func (p *AnthropicProvider) ModelCapabilities(model string) protocol.ModelCapabilities
func (*AnthropicProvider) Name ¶
func (p *AnthropicProvider) Name() string
func (*AnthropicProvider) Stream ¶
func (p *AnthropicProvider) Stream(ctx context.Context, req protocol.ChatRequest) (<-chan protocol.ChatDelta, error)
Stream performs a real token-streaming Messages call: text_delta blocks are forwarded as content deltas; tool_use blocks accumulate their input_json_delta fragments (keyed by content-block index) into whole ToolCalls flushed before the terminal Done delta. Proves the streaming seam holds for a second provider with no edits to the loop (§12 AC).
func (*AnthropicProvider) SupportsTools ¶
func (p *AnthropicProvider) SupportsTools() bool
func (*AnthropicProvider) UpdateAPIKey ¶
func (p *AnthropicProvider) UpdateAPIKey(key string)
UpdateAPIKey swaps the key used for subsequent requests (protocol.KeyUpdater), letting the loop refresh an expiring BYO token between turns.
type AnthropicRequestOptions ¶
type AnthropicRequestOptions struct {
// Raw JSON preserves absent/null/fractional timeout values. Validation belongs
// to the injected client, just as it does for a prebuilt SDK client in Pi.
TimeoutMS json.RawMessage `json:"timeout,omitempty"`
MaxRetries int `json:"maxRetries"`
}
type AnthropicStreamOptions ¶
type AnthropicStreamOptions = OpenAICompletionsStreamOptions
AnthropicStreamOptions separates concrete Go transport/callback dependencies from serializable provider controls. Client replaces fetch, not SDK auth.
type AnthropicToolsOptions ¶
type AnthropicToolsOptions struct {
OAuth bool `json:"oauth"`
EagerInputStreaming bool `json:"eagerInputStreaming"`
StrictTools bool `json:"strictTools"`
CacheControl json.RawMessage `json:"cacheControl,omitempty"`
}
type AntigravityOAuthOptions ¶
type AntigravityOAuthOptions struct {
Client *http.Client
RandomValue func() (string, error)
CallbackHost func() string
StartCallback func(*OAuthCallbackServerOptions) (*OAuthCallbackServer, error)
Now func() float64
}
type AntigravityProvider ¶
type AntigravityProvider struct {
// BaseURL, when set, replaces the endpoint failover list with a single
// endpoint (tests, proxies). Empty tries daily then sandbox.
BaseURL string
// HTTP serves the list-models path.
HTTP *http.Client
// StreamHTTP serves the SSE path. Nil falls back to HTTP.
StreamHTTP *http.Client
// contains filtered or unexported fields
}
AntigravityProvider speaks Google's Cloud Code Assist internal API — the backend behind the Antigravity subscription — authenticated with a Google OAuth access token plus the account's Cloud Code project id. The wire is a Gemini-shaped generateContent envelope wrapped in an agent request (project/requestId/labels/session bookkeeping the real client sends).
func NewAntigravityProvider ¶
func NewAntigravityProvider() *AntigravityProvider
NewAntigravityProvider builds the Antigravity wire client. Credentials arrive per request via applyOAuthToken — there is no static key.
func (*AntigravityProvider) Chat ¶
func (p *AntigravityProvider) Chat(ctx context.Context, req protocol.ChatRequest) (protocol.ChatResponse, error)
Chat consumes the SSE stream to completion and returns the assembled response — the backend only streams, so there is no second wire path.
func (*AntigravityProvider) ModelCapabilities ¶
func (p *AntigravityProvider) ModelCapabilities(model string) protocol.ModelCapabilities
func (*AntigravityProvider) Name ¶
func (p *AntigravityProvider) Name() string
func (*AntigravityProvider) Stream ¶
func (p *AntigravityProvider) Stream(ctx context.Context, req protocol.ChatRequest) (<-chan protocol.ChatDelta, error)
Stream POSTs the envelope to {endpoint}/v1internal:streamGenerateContent?alt=sse, failing over to the sandbox endpoint when the daily host fails at transport level or with a 5xx before the first event. Text parts (thought != true) stream as content deltas; functionCall parts arrive whole; usageMetadata carries the token accounting.
func (*AntigravityProvider) SupportsTools ¶
func (p *AntigravityProvider) SupportsTools() bool
type AssistantMessageEvent ¶
type AssistantMessageEvent struct {
Type string `json:"type"`
ContentIndex int `json:"contentIndex,omitempty"`
Delta string `json:"delta,omitempty"`
Content string `json:"content,omitempty"`
ToolCall *ContentBlock `json:"toolCall,omitempty"`
Partial *Message `json:"partial,omitempty"`
Reason string `json:"reason,omitempty"`
Message *Message `json:"message,omitempty"`
Error *Message `json:"error,omitempty"`
Extra map[string]json.RawMessage `json:"-"`
// contains filtered or unexported fields
}
AssistantMessageEvent is Pi's event union. Type selects the fields emitted on the wire; Partial is the live accumulator, not an event-time snapshot.
func (AssistantMessageEvent) MarshalJSON ¶
func (e AssistantMessageEvent) MarshalJSON() ([]byte, error)
type AssistantMessageEventStream ¶
type AssistantMessageEventStream struct {
*EventStream[AssistantMessageEvent, *Message]
// contains filtered or unexported fields
}
AssistantMessageEventStream retains Pi's live message pointers. Producers updating published payloads must use Synchronize; consumers either inspect them in Synchronize or use SnapshotEvent/SnapshotResult. Queue synchronization alone cannot protect a message being changed after Push returns.
func LazyStream ¶
func LazyStream(ctx context.Context, model any, setup func(context.Context) (*ProviderEventSource, error), now func() int64) *AssistantMessageEventStream
LazyStream returns immediately. Only setup/provider callbacks own request cancellation; forwarding drains the source independently of reader waits. A terminal event settles Result before forwarding and source.Result finish.
func NewAssistantMessageEventStream ¶
func NewAssistantMessageEventStream() *AssistantMessageEventStream
func NewAssistantMessageEventStreamFor ¶
func NewAssistantMessageEventStreamFor(ctx context.Context) *AssistantMessageEventStream
NewAssistantMessageEventStreamFor constructs an independent queue/result using a host relay's payload lock when supplied, or its own lock otherwise. Set the context before starting a producer; locks are immutable once published.
func NewAssistantMessageEventStreamWithHooks ¶
func NewAssistantMessageEventStreamWithHooks(hooks AssistantStreamHooks) *AssistantMessageEventStream
func StreamAnthropic ¶
func StreamAnthropic(ctx context.Context, rawModel json.RawMessage, transcript TranscriptContext, options AnthropicStreamOptions) *AssistantMessageEventStream
StreamAnthropic executes the native Messages provider without a worker.
func StreamAnthropicSimple ¶
func StreamAnthropicSimple(ctx context.Context, rawModel json.RawMessage, transcript TranscriptContext, options AnthropicStreamOptions) (*AssistantMessageEventStream, error)
StreamAnthropicSimple rejects missing credentials before admitting a stream. Federation is admitted when the provider's scoped/process env is complete.
func StreamAnthropicWithClient ¶
func StreamAnthropicWithClient(ctx context.Context, rawModel json.RawMessage, transcript TranscriptContext, options AnthropicStreamOptions, client *AnthropicMessageClient) *AssistantMessageEventStream
StreamAnthropicWithClient skips built-in auth/client construction and OAuth request identity. A nil client selects the ordinary native HTTP provider. Simple mode intentionally does not accept this override, matching Pi.
func StreamAntigravityPooled ¶
func StreamAntigravityPooled(ctx context.Context, model json.RawMessage, transcript TranscriptContext, options OpenAICompletionsStreamOptions, source TokenSource) (*AssistantMessageEventStream, error)
StreamAntigravityPooled keeps account credentials out of the transcript and carries the selected account's project through each native request attempt.
func StreamAzureResponses ¶
func StreamAzureResponses(ctx context.Context, rawModel json.RawMessage, transcript TranscriptContext, options AzureResponsesStreamOptions) *AssistantMessageEventStream
StreamAzureResponses executes Pi's Azure Responses stream directly in Go. Inputs are frozen before the asynchronous producer starts; callbacks receive the original model, while deployment and endpoint settings affect only HTTP.
func StreamAzureResponsesSimple ¶
func StreamAzureResponsesSimple(ctx context.Context, rawModel json.RawMessage, transcript TranscriptContext, options AzureResponsesStreamOptions) (*AssistantMessageEventStream, error)
StreamAzureResponsesSimple rejects missing API keys synchronously. Unlike the OpenAI providers, an authorization header cannot satisfy this admission.
func StreamClaudeCodePooled ¶
func StreamClaudeCodePooled(ctx context.Context, model json.RawMessage, transcript TranscriptContext, options AnthropicStreamOptions, source TokenSource) (*AssistantMessageEventStream, error)
StreamClaudeCodePooled uses the native Anthropic transcript and SSE machinery. It never calls the legacy Chat/Stream adapter or serializes the selected token into the model, transcript or returned events.
func StreamCodexResponses ¶
func StreamCodexResponses(ctx context.Context, rawModel json.RawMessage, transcript TranscriptContext, options CodexResponsesStreamOptions) *AssistantMessageEventStream
StreamCodexResponses selects the native WebSocket/SSE path using Pi's transport option and session fallback memory. apiKey is an OAuth access token.
func StreamCodexResponsesPooled ¶
func StreamCodexResponsesPooled(ctx context.Context, rawModel json.RawMessage, transcript TranscriptContext, options CodexResponsesStreamOptions, source TokenSource) (*AssistantMessageEventStream, error)
StreamCodexResponsesPooled binds an existing host account source to the native simple stream. Credential rotation is limited to typed authentication failures before visible content; the selected token is never stored on a shared client.
func StreamCodexResponsesSSE ¶
func StreamCodexResponsesSSE(ctx context.Context, rawModel json.RawMessage, transcript TranscriptContext, options CodexResponsesStreamOptions) *AssistantMessageEventStream
StreamCodexResponsesSSE executes Pi's explicit SSE transport in Go. The combined transport selector is available through StreamCodexResponses. Credentials are a per-call OAuth access token in apiKey.
func StreamCodexResponsesSimple ¶
func StreamCodexResponsesSimple(ctx context.Context, rawModel json.RawMessage, transcript TranscriptContext, options CodexResponsesStreamOptions) (*AssistantMessageEventStream, error)
StreamCodexResponsesSimple applies the simple-mode contract with native transport selection.
func StreamCodexResponsesSimpleSSE ¶
func StreamCodexResponsesSimpleSSE(ctx context.Context, rawModel json.RawMessage, transcript TranscriptContext, options CodexResponsesStreamOptions) (*AssistantMessageEventStream, error)
StreamCodexResponsesSimpleSSE preserves streamSimple's synchronous API-key check while executing the explicit SSE transport. It never obtains a token from environment fallback or a different account's cached credential.
func StreamDevinPooled ¶
func StreamDevinPooled(ctx context.Context, rawModel json.RawMessage, transcript TranscriptContext, options OpenAICompletionsStreamOptions, source TokenSource) (*AssistantMessageEventStream, error)
StreamDevinPooled binds an account source to the native Devin stream. It mirrors StreamCodexResponsesPooled: a typed 401/403 before visible content rotates the account and replays; after the first delta the failure is forwarded and Report feeds the outcome back to the pool.
func StreamOpenAICompletions ¶
func StreamOpenAICompletions(ctx context.Context, rawModel json.RawMessage, transcript TranscriptContext, options OpenAICompletionsStreamOptions) *AssistantMessageEventStream
StreamOpenAICompletions executes the native provider against its HTTP endpoint. It returns the live stream immediately; preparation, callback and HTTP errors settle it with an error event. It does not invoke a JavaScript runtime or the legacy agentcore transcript adapter.
func StreamOpenAICompletionsSimple ¶
func StreamOpenAICompletionsSimple(ctx context.Context, rawModel json.RawMessage, transcript TranscriptContext, options OpenAICompletionsStreamOptions) (*AssistantMessageEventStream, error)
StreamOpenAICompletionsSimple retains the original synchronous credential check. Provider/request failures after admission settle the returned stream.
func StreamOpenAIResponses ¶
func StreamOpenAIResponses(ctx context.Context, rawModel json.RawMessage, transcript TranscriptContext, options OpenAIResponsesStreamOptions) *AssistantMessageEventStream
StreamOpenAIResponses runs Responses directly in Go. It freezes caller inputs and returns the live stream immediately; producer failures emit error events.
func StreamOpenAIResponsesSimple ¶
func StreamOpenAIResponsesSimple(ctx context.Context, rawModel json.RawMessage, transcript TranscriptContext, options OpenAIResponsesStreamOptions) (*AssistantMessageEventStream, error)
StreamOpenAIResponsesSimple validates credentials before admitting the stream, then runs the native provider with Pi's clamped simple options.
func StreamPiMessages ¶
func StreamPiMessages(ctx context.Context, model *Object, transcript TranscriptContext, options PiMessagesStreamOptions) *AssistantMessageEventStream
StreamPiMessages implements the gateway protocol entirely in Go. A terminal provider event wins even if the request context has since been canceled.
func StreamPiMessagesJSON ¶
func StreamPiMessagesJSON(ctx context.Context, rawModel json.RawMessage, transcript TranscriptContext, input OpenAICompletionsStreamOptions) (*AssistantMessageEventStream, error)
StreamPiMessagesJSON adapts the host's serialized model/options and callback boundary to pi-messages. It owns decoded input objects; callbacks can replace raw events (including primitive JSON values) before protocol conversion. Nil payload output preserves the payload, while JSON null replaces it.
func StreamSimplePiMessages ¶
func StreamSimplePiMessages(ctx context.Context, model *Object, transcript TranscriptContext, options PiMessagesStreamOptions) *AssistantMessageEventStream
Both upstream entry points use the same transport and controls.
func (*AssistantMessageEventStream) End ¶
func (s *AssistantMessageEventStream) End(result ...*Message)
End marks producer settlement separately from the first terminal event. A proxy can still mutate its live result after that event while draining the response. WaitForEnd is useful before inspecting retained raw pointers.
func (*AssistantMessageEventStream) Next ¶
func (s *AssistantMessageEventStream) Next(ctx context.Context) (AssistantMessageEvent, bool, error)
func (*AssistantMessageEventStream) Result ¶
func (s *AssistantMessageEventStream) Result(ctx context.Context) (*Message, error)
func (*AssistantMessageEventStream) SnapshotEvent ¶
func (s *AssistantMessageEventStream) SnapshotEvent(event AssistantMessageEvent) (AssistantMessageEvent, error)
SnapshotEvent copies the payload as observed now, not at the original Push. Shared message, block and usage objects remain shared within the snapshot. Retained raw events still observe subsequent mutations, as in Pi.
func (*AssistantMessageEventStream) SnapshotResult ¶
func (s *AssistantMessageEventStream) SnapshotResult(ctx context.Context) (*Message, error)
func (*AssistantMessageEventStream) Synchronize ¶
func (s *AssistantMessageEventStream) Synchronize(fn func())
Synchronize protects live payload access. The callback may Push, but must not wait on the stream or call Synchronize/Snapshot methods recursively.
func (*AssistantMessageEventStream) SynchronizeYielding ¶
func (s *AssistantMessageEventStream) SynchronizeYielding(callback func(PayloadYield) error) error
SynchronizeYielding protects a callback's live message access with the same lock as Synchronize. Unlike Synchronize, its callback can wait for a provider using yield without preventing that provider from updating the payload. Use yield synchronously in the callback. Nested/overlapping yields and calls retained after the callback returns are rejected without invoking their work.
func (*AssistantMessageEventStream) WaitForEnd ¶
func (s *AssistantMessageEventStream) WaitForEnd(ctx context.Context) error
type AssistantStreamHooks ¶
type AssistantStreamHooks struct {
Result func(context.Context) (*Message, error)
AfterEnd func(context.Context) error
}
AssistantStreamHooks supplies the overridable result/iteration operations used by Pi's host callback adapter. Its callback promise can reject after a terminal event was already pushed. Ordinary provider streams need no hooks. Hooks are immutable after construction and must honor reader cancellation.
type AttemptOutcome ¶
type AttemptOutcome struct {
AdmissionError error
Terminal AssistantMessageEvent
Committed bool
// contains filtered or unexported fields
}
AttemptOutcome keeps the original terminal event and any withheld metadata. The caller decides retry/escalation and publishes a final terminal; a failed attempt must not settle the logical request's stream.
func (AttemptOutcome) Message ¶
func (r AttemptOutcome) Message() *Message
Message returns the original settled message without copying native payloads.
type AttemptTrace ¶
type AttemptTrace struct {
// contains filtered or unexported fields
}
AttemptTrace records detached JSON after producer settlement. It never relays events or feeds a display projection back into the native engine.
func NewAttemptTrace ¶
func NewAttemptTrace(span *telemetry.Span) *AttemptTrace
func (*AttemptTrace) Finish ¶
func (t *AttemptTrace) Finish(attempt FallbackAttempt)
Finish records one settled attempt, including admission and terminal errors.
func (*AttemptTrace) Start ¶
func (t *AttemptTrace) Start(model json.RawMessage, view TranscriptContext)
Start snapshots the request immediately before opening an attempt.
type AuthContext ¶
type AuthContext struct {
Env func(context.Context, string) (any, error)
FileExists func(context.Context, string) (bool, error)
}
AuthContext is the concrete environment boundary for provider-owned auth. Env returns Undefined when absent. The callback is read on every resolution.
func DefaultProviderAuthContext ¶
func DefaultProviderAuthContext() *AuthContext
DefaultProviderAuthContext reads process environment lazily. Whitespace-only values are absent, while nonempty values retain their original whitespace. FileExists follows Pi's literal leading-tilde expansion, including ~suffix.
type AuthResolutionOverrides ¶
nil/Undefined mean omitted overrides; Null retains an explicit null. Now is a native clock dependency and is not forwarded to provider callbacks.
type AzureResponsesConfig ¶
type AzureResponsesConfig struct {
BaseURL string `json:"baseUrl"`
APIVersion string `json:"apiVersion"`
DeploymentName string `json:"deploymentName"`
}
AzureResponsesConfig resolves Pi's Azure endpoint and deployment settings. It performs no HTTP requests or credential lookup.
func ResolveAzureResponsesConfig ¶
func ResolveAzureResponsesConfig(rawModel, rawOptions json.RawMessage) (AzureResponsesConfig, error)
type AzureResponsesStreamOptions ¶
type AzureResponsesStreamOptions = OpenAICompletionsStreamOptions
type BlockList ¶
type BlockList struct {
// contains filtered or unexported fields
}
BlockList retains content-array identity across shallow message copies. A nil entry is a hole; NullContentBlock represents an explicit null. Its typed view uses the shared JavaScript array implementation for length and sparse slots. Callers synchronize live stream reads and writes through the owning stream.
func NewBlockList ¶
func NewBlockList(blocks ...*ContentBlock) *BlockList
func (*BlockList) Append ¶
func (l *BlockList) Append(blocks ...*ContentBlock) int
func (*BlockList) Get ¶
func (l *BlockList) Get(index int) *ContentBlock
func (*BlockList) MarshalJSON ¶
func (*BlockList) Pop ¶
func (l *BlockList) Pop() *ContentBlock
func (*BlockList) Set ¶
func (l *BlockList) Set(index int, block *ContentBlock)
func (*BlockList) UnmarshalJSON ¶
func (*BlockList) Values ¶
func (l *BlockList) Values() []*ContentBlock
Values returns a detached outer slice; the block objects remain shared. Mutate membership through Set/Append/Delete/SetLength, not this projection.
type ChatGPTOAuthAuthorization ¶
type ChatGPTOAuthAuthorization struct{ Code, ClientID string }
type ChatGPTOAuthCallback ¶
type ChatGPTOAuthCallback struct {
URL string
Ready <-chan struct{}
Result func() (ChatGPTOAuthAuthorization, error)
Close func()
}
Ready is shared by all observers. Result is read after Ready closes. Closing an unused listener does not settle the authorization result, matching Pi.
func StartChatGPTOAuthCallback ¶
func StartChatGPTOAuthCallback(options ChatGPTOAuthCallbackOptions) (*ChatGPTOAuthCallback, error)
type ClientSpec ¶
type ClientSpec struct {
Name string // "openai" | "anthropic" | an OAuth vendor | any OpenAI-compatible vendor
APIKey string
BaseURL string
Compat Compat // optional; zero value falls back to the vendor default
// OpenAIWire chooses the transport for an OpenAI provider identity. Empty
// and OpenAIWireChat use Chat Completions; OpenAIWireResponses uses the
// public Responses API while Name() remains "openai". Keeping provider
// identity separate from wire format lets one workspace row serve models
// with different API contracts, matching the model.api split in OMP.
OpenAIWire OpenAIWire
// SessionScope distinguishes independently configured provider rows that use
// the same vendor/base URL. OAuth account-rotation state must never confuse
// two sibling pools just because both speak (for example) openai-codex.
SessionScope string
// TokenSource is the OAuth account pool a subscription vendor
// (claude-code | openai-codex | google-antigravity) draws per-request
// credentials from. Required for those vendors, ignored by the rest.
TokenSource TokenSource
}
ClientSpec is the resolved configuration for constructing a wire client: vendor name, decrypted key, and an optional base-URL override.
It is deliberately smaller than Spec: a ClientSpec yields a bare protocol.LLMProvider (Chat/Stream and nothing else), which is what a run needs. Spec adds the identity and live model list that a workspace-managed provider needs, and builds on this.
type CodexProvider ¶
type CodexProvider struct {
// BaseURL overrides the backend origin; empty uses codexBaseURL. The
// /codex/responses suffix is appended by responsesURL, so a test server
// can point BaseURL at itself.
BaseURL string
// HTTP serves the list-models path.
HTTP *http.Client
// StreamHTTP serves the SSE path (the only chat transport this backend
// offers). Nil falls back to HTTP.
StreamHTTP *http.Client
// contains filtered or unexported fields
}
CodexProvider speaks the ChatGPT Codex backend: the Responses API at chatgpt.com/backend-api/codex/responses, authenticated with a ChatGPT OAuth access token rather than an OpenAI API key. The wire is Responses-shaped (input items, function_call/function_call_output pairs, SSE event types under response.*), but the backend is stricter than the public API: it rejects sampling parameters outright, requires stream:true, and wants the codex client headers (originator, version, OpenAI-Beta) on every call.
func NewCodexProvider ¶
func NewCodexProvider() *CodexProvider
NewCodexProvider builds the Codex wire client. Credentials arrive per request via applyOAuthToken — there is no static key.
func (*CodexProvider) Chat ¶
func (p *CodexProvider) Chat(ctx context.Context, req protocol.ChatRequest) (protocol.ChatResponse, error)
Chat consumes the SSE stream to completion and returns the assembled response — the backend has no non-streaming mode, so this is the simplest correct implementation rather than a second wire path.
func (*CodexProvider) ModelCapabilities ¶
func (p *CodexProvider) ModelCapabilities(model string) protocol.ModelCapabilities
func (*CodexProvider) Name ¶
func (p *CodexProvider) Name() string
func (*CodexProvider) Stream ¶
func (p *CodexProvider) Stream(ctx context.Context, req protocol.ChatRequest) (<-chan protocol.ChatDelta, error)
Stream performs the only chat call this backend offers: a streamed Responses request. Text arrives as response.output_text.delta events; tool calls land whole on response.output_item.done; usage and the stop reason ride the terminal response.completed/done/incomplete event.
func (*CodexProvider) SupportsTools ¶
func (p *CodexProvider) SupportsTools() bool
type CodexResponsesStreamOptions ¶
type CodexResponsesStreamOptions = OpenAICompletionsStreamOptions
type CodexWebSocketDebugStats ¶
type CodexWebSocketDebugStats struct {
Requests int `json:"requests"`
ConnectionsCreated int `json:"connectionsCreated"`
ConnectionsReused int `json:"connectionsReused"`
CachedContextRequests int `json:"cachedContextRequests"`
StoreTrueRequests int `json:"storeTrueRequests"`
FullContextRequests int `json:"fullContextRequests"`
DeltaRequests int `json:"deltaRequests"`
LastInputItems int `json:"lastInputItems"`
LastDeltaInputItems *int `json:"lastDeltaInputItems,omitempty"`
LastPreviousResponseID json.RawMessage `json:"lastPreviousResponseId,omitempty"`
WebSocketFailures int `json:"websocketFailures"`
SSEFallbacks int `json:"sseFallbacks"`
WebSocketFallbackActive *bool `json:"websocketFallbackActive,omitempty"`
LastWebSocketError *string `json:"lastWebSocketError,omitempty"`
}
CodexWebSocketDebugStats is a snapshot; optional fields retain Pi's distinction between a value that has never been recorded and an explicit false/zero.
func GetCodexWebSocketDebugStats ¶
func GetCodexWebSocketDebugStats(session string) *CodexWebSocketDebugStats
type Collection ¶
type Collection struct {
// contains filtered or unexported fields
}
Collection holds registered providers and routes Chat/Stream to the owner of the requested model. ListModels unions the live catalogs of every registered provider (failures are returned per-provider, not swallowed).
func CollectionFromSpecs ¶
func CollectionFromSpecs(specs []Spec) (*Collection, error)
CollectionFromSpecs registers one provider per Spec. Used by the workspace API and persist tests so listing and Chat share the same construction path.
func (*Collection) Chat ¶
func (c *Collection) Chat(ctx context.Context, req protocol.ChatRequest) (protocol.ChatResponse, error)
Chat implements protocol.LLMProvider: the owning provider serves the request.
func (*Collection) ChatOn ¶
func (c *Collection) ChatOn(ctx context.Context, providerID string, req protocol.ChatRequest) (protocol.ChatResponse, error)
ChatOn sends the request to a specific registered provider (used when a tier names both a provider id and a model, so an ID collision cannot mis-route).
func (*Collection) Get ¶
func (c *Collection) Get(id string) (Provider, bool)
Get returns a registered provider by id.
func (*Collection) ListModels ¶
func (c *Collection) ListModels(ctx context.Context) (ListResult, error)
ListModels refreshes every registered provider's catalog via its vendor list-models HTTP API and returns the union. The last-known map is updated so a subsequent Chat/Stream can find the owner.
func (*Collection) ModelCapabilities ¶
func (c *Collection) ModelCapabilities(model string) protocol.ModelCapabilities
func (*Collection) Name ¶
func (c *Collection) Name() string
func (*Collection) Owner ¶
Owner returns the provider that listed modelID. If the cache is empty it refreshes first. First registered owner wins on an ID collision.
func (*Collection) Providers ¶
func (c *Collection) Providers() []Provider
Providers returns registered providers in registration order.
func (*Collection) Register ¶
func (c *Collection) Register(p Provider)
Register adds or replaces a provider. The provider's ID is the map key.
func (*Collection) Stream ¶
func (c *Collection) Stream(ctx context.Context, req protocol.ChatRequest) (<-chan protocol.ChatDelta, error)
Stream implements protocol.LLMProvider.
func (*Collection) StreamOn ¶
func (c *Collection) StreamOn(ctx context.Context, providerID string, req protocol.ChatRequest) (<-chan protocol.ChatDelta, error)
StreamOn is Stream scoped to one registered provider.
func (*Collection) SupportsTools ¶
func (c *Collection) SupportsTools() bool
type Compat ¶
type Compat struct {
// MaxTokensField is "max_tokens" or "max_completion_tokens".
MaxTokensField string
// SupportsTools indicates the vendor accepts the tools/tool_calls fields.
SupportsTools bool
}
Compat is the per-vendor capability table. Most "providers" are the OpenAI completions API at a different base URL plus a few flags (pi-ai pattern), so an OpenAI-compatible vendor is config, not new code.
type CompletionsResponse ¶
type ContentBlock ¶
type ContentBlock struct {
Type string `json:"type"`
Text string `json:"text,omitempty"`
TextSignature *string `json:"textSignature,omitempty"`
Thinking string `json:"thinking,omitempty"`
ThinkingSignature *string `json:"thinkingSignature,omitempty"`
Redacted *bool `json:"redacted,omitempty"`
Data string `json:"data,omitempty"`
MIMEType string `json:"mimeType,omitempty"`
ID string `json:"id,omitempty"`
Name string `json:"name,omitempty"`
Arguments json.RawMessage `json:"arguments,omitempty"`
ThoughtSignature *string `json:"thoughtSignature,omitempty"`
Namespace *string `json:"namespace,omitempty"`
Extra map[string]json.RawMessage `json:"-"`
// contains filtered or unexported fields
}
ContentBlock is Pi's text/thinking/image/toolCall union. Optional signatures use pointers so an explicit empty signature survives JSON round trips. A nil slot represents an array hole; decoded null entries retain a distinct sentinel. Both serialize as null within an array, but filtering distinguishes them.
func NullContentBlock ¶
func NullContentBlock() *ContentBlock
NullContentBlock constructs an explicitly present null entry, distinct from a nil slot (hole) in MessageContent.Blocks.
func (*ContentBlock) IsNull ¶
func (b *ContentBlock) IsNull() bool
IsNull reports an explicit decoded null. A nil block pointer is a sparse content-array hole, which Array.filter skips without reading its properties.
func (ContentBlock) MarshalJSON ¶
func (b ContentBlock) MarshalJSON() ([]byte, error)
func (*ContentBlock) UnmarshalJSON ¶
func (b *ContentBlock) UnmarshalJSON(data []byte) error
type ContextUsageEstimate ¶
type ContextUsageEstimate struct {
Tokens float64 `json:"tokens"`
UsageTokens float64 `json:"usageTokens"`
TrailingTokens float64 `json:"trailingTokens"`
LastUsageIndex *int `json:"lastUsageIndex"`
}
ContextUsageEstimate follows Pi's usage-aware estimate. Usage applies only when its response is at least as new as every preceding transcript message.
func EstimateContextTokens ¶
func EstimateContextTokens(messages []Message) ContextUsageEstimate
type CredentialPersistence ¶
type CredentialPersistence struct {
Read func(context.Context, string) (any, error)
List func(context.Context) (*Array, error)
Modify func(context.Context, string, func(any) (any, error)) (any, error)
Delete func(context.Context, string) error
}
CredentialPersistence is the app-owned callback boundary used by Models. Undefined means absent/no replacement; Null is an explicitly stored null. Credentials and metadata use Object/Array values to preserve field presence and identity. No new provider interface is required.
type DataCloneError ¶
type DataCloneError struct{}
DataCloneError identifies values that cannot be stored by structured clone.
func (*DataCloneError) Error ¶
func (*DataCloneError) Error() string
type DeferredStreamFunc ¶
type DevinOAuthOptions ¶
type DevinOAuthOptions struct {
Client *http.Client
PKCE func() (PKCE, error)
CallbackHost func() string
StartCallback func(*OAuthCallbackServerOptions) (*OAuthCallbackServer, error)
Now func() float64
// contains filtered or unexported fields
}
type DevinProvider ¶
type DevinProvider struct {
// BaseURL overrides the Cascade origin; empty uses DevinDefaultBaseURL.
BaseURL string
// HTTP serves listing and unary calls.
HTTP *http.Client
// StreamHTTP serves the chat stream; nil falls back to HTTP.
StreamHTTP *http.Client
// contains filtered or unexported fields
}
DevinProvider is the wire client for the Devin Cascade surface. It is stateless per request apart from the account credential the pooled wrapper installs via applyOAuthToken.
func NewDevinProvider ¶
func NewDevinProvider() *DevinProvider
func (*DevinProvider) Chat ¶
func (p *DevinProvider) Chat(ctx context.Context, req protocol.ChatRequest) (protocol.ChatResponse, error)
Chat consumes the stream to completion (Cascade has no non-streaming mode).
func (*DevinProvider) ModelCapabilities ¶
func (p *DevinProvider) ModelCapabilities(model string) protocol.ModelCapabilities
func (*DevinProvider) Name ¶
func (p *DevinProvider) Name() string
func (*DevinProvider) Stream ¶
func (p *DevinProvider) Stream(ctx context.Context, req protocol.ChatRequest) (<-chan protocol.ChatDelta, error)
Stream adapts the native Devin event stream onto the neutral ChatDelta contract.
func (*DevinProvider) SupportsTools ¶
func (p *DevinProvider) SupportsTools() bool
type EventStream ¶
type EventStream[T, R any] struct { // contains filtered or unexported fields }
EventStream ports Pi's unbounded FIFO stream and independent final-result promise. Push never waits for a reader. Next consumes one event; Result does not drain events. Multiple readers share the queue, rather than receiving broadcast copies. Values retain their identity, including Partial pointers.
Go callers may cancel a Next or Result wait through context. Cancellation of a reader does not cancel the provider or discard a queued event.
func NewEventStream ¶
func NewEventStream[T, R any](isComplete func(T) bool, extractResult func(T) R) *EventStream[T, R]
func (*EventStream[T, R]) End ¶
func (s *EventStream[T, R]) End(result ...R)
End stops production and wakes waiting consumers. Queued events remain readable. With no result argument, Result remains pending, exactly as Pi's end(undefined). A later End(result) may settle it; the first result wins.
func (*EventStream[T, R]) Next ¶
func (s *EventStream[T, R]) Next(ctx context.Context) (T, bool, error)
func (*EventStream[T, R]) Push ¶
func (s *EventStream[T, R]) Push(event T)
type FallbackAttempt ¶
type FallbackAttempt struct {
Number int
Outcome AttemptOutcome
Failure error
FailureKind string
}
FallbackAttempt is delivered after producer settlement, including failed attempts whose terminal event is withheld from the logical request. Observers can account for usage without inserting failed messages into the transcript.
type FallbackCandidate ¶
type FallbackCandidate struct {
SafeLabels telemetry.SafeLabels
Model json.RawMessage
Stream StreamFn
}
FallbackCandidate binds a model to its native provider transport.
type FallbackProvider ¶
type FallbackProvider struct {
Candidates []FallbackCandidate
Retry protocol.RetryPolicy
// Wait optionally supplies a clock for backoff. Nil uses a cancellable timer.
Wait func(context.Context, time.Duration) error
}
FallbackProvider treats retry and ordered provider fallback as one native stream. Each candidate receives its own retry budget. Once content is visible, a request can never replay on another attempt or candidate. Fields are immutable while requests are running; request selection belongs to the host and is supplied separately to Run.
func (FallbackProvider) Run ¶
func (p FallbackProvider) Run(ctx context.Context, out *AssistantMessageEventStream, request FallbackRequest) (AttemptOutcome, error)
Run publishes to a caller-owned stream without ending it. Hosts with durable selection/request admission use this entry point; ordinary consumers use Stream.
func (FallbackProvider) Stream ¶
func (p FallbackProvider) Stream(ctx context.Context, model json.RawMessage, transcript TranscriptContext, options map[string]any) (*AssistantMessageEventStream, error)
Stream implements StreamFn. Candidate models override the supplied model; an omitted candidate model uses it. It keeps no mutable selection between calls.
type FallbackRequest ¶
type FallbackRequest struct {
// SafeLabels is an optional immutable host-approved binding per candidate.
// An omitted entry uses the corresponding candidate's approved labels.
SafeLabels []telemetry.SafeLabels
Start int
Candidates int
Open func(context.Context, int, int) (*AssistantMessageEventStream, error)
Commit func(context.Context, int) error
Observe func(context.Context, int, FallbackAttempt) error
Validate func() error
BeforePublish func(AttemptOutcome) error
// Recover may repair one confirmed pre-output provider rejection per rung.
// AI still owns replay safety and the hard one-recovery bound.
Recover func(context.Context, int, FallbackAttempt) (bool, error)
}
FallbackRequest binds host policy to one logical request. Open prepares a fresh candidate on every attempt; Commit durably selects it before visible output. Observe sees all settled attempts, including discarded failures. Validate fences host state after settlement, before escalation/publication. BeforePublish runs once before the chosen terminal event is exposed.
type FallbackSelection ¶
type FallbackSelection struct {
Version int `json:"version"`
Generation uint64 `json:"generation"`
Rung int `json:"rung"`
ProviderID string `json:"providerId"`
Model json.RawMessage `json:"model"`
}
FallbackSelection records the selected candidate without credentials or pool handles. Provider-row identity disambiguates candidates with equal models.
func ParseFallbackSelection ¶
func ParseFallbackSelection(raw string, previous *FallbackSelection) (FallbackSelection, error)
ParseFallbackSelection validates one committed transition after the previous selection. The host must additionally check the selected provider/model against its bound candidates before opening a request.
type GitHubCopilotOAuthOptions ¶
type GitHubCopilotOAuthOptions struct {
KnownModels *Object
Client *http.Client
Now func() float64
Sleep func(float64, context.Context) error
PollSleep func(float64, context.Context, string) error
RequestTimeout time.Duration
RetryBudget time.Duration
}
GitHubCopilotOAuthOptions supplies native dependencies. KnownModels is the source GITHUB_COPILOT_MODELS object: own keys determine which unconfigured account models may be enabled. The pinned source omits its generated data, so callers must supply this catalog rather than silently using today's list.
type GrammarConstrainedSampling ¶
type GrammarConstrainedSampling struct {
Format string `json:"format"`
Definition string `json:"definition"`
InputProperty string `json:"inputProperty"`
}
func ResolveGrammarConstrainedSampling ¶
func ResolveGrammarConstrainedSampling(tool Tool, supported bool) (*GrammarConstrainedSampling, error)
type HTTPDoer ¶
HTTPDoer is the injected HTTP seam so list-models and Chat share a client and tests can point at httptest servers.
type InMemoryCredentialStore ¶
type InMemoryCredentialStore struct {
// contains filtered or unexported fields
}
InMemoryCredentialStore follows Pi's live-reference credential storage. Read and List inspect committed state immediately, without waiting for queued writes. Modify/Delete serialize per provider; other providers are independent. Returned credentials are not cloned. Callers synchronize mutations of shared credential objects, including mutations from a Modify callback.
func NewInMemoryCredentialStore ¶
func NewInMemoryCredentialStore() *InMemoryCredentialStore
func (*InMemoryCredentialStore) Delete ¶
func (s *InMemoryCredentialStore) Delete(ctx context.Context, providerID string) error
func (*InMemoryCredentialStore) List ¶
func (s *InMemoryCredentialStore) List(ctx context.Context) (*Array, error)
type InMemoryModelsStore ¶
type InMemoryModelsStore struct {
// contains filtered or unexported fields
}
InMemoryModelsStore keeps provider-scoped catalog snapshots. Entries use Object/Array values so optional fields, unknown metadata, shared references, sparse arrays and cycles survive cloning. Read and Write detach the entire graph; callers own their inputs and returned values. It is safe for concurrent store operations, but callers must not mutate an input during Write. This is Pi's catalog store, separate from the workspace's live model list.
func NewInMemoryModelsStore ¶
func NewInMemoryModelsStore() *InMemoryModelsStore
func (*InMemoryModelsStore) Delete ¶
func (s *InMemoryModelsStore) Delete(ctx context.Context, providerID string) error
type JSONMethod ¶
type JSONMethod = jsonjs.JSONMethod
type KimiCodingOAuthOptions ¶
type KimiCodingOAuthOptions struct {
Client *http.Client
Env func(string) string
Now func() float64
PollSleep func(float64, context.Context, string) error
RetrySleep func(float64, context.Context) error
RequestTimeout time.Duration
}
KimiCodingOAuthOptions supplies concrete HTTP, environment and clock hooks. RetrySleep retains the cancellation cause; PollSleep uses the poller's fixed cancellation message. Env is consulted anew for each login/refresh operation.
type LazyAPICapabilities ¶
type LazyAPICapabilities struct{ FetchDeferred, CancelDeferred bool }
type LazyOAuthOptions ¶
type LazyOAuthOptions struct {
Name string
IsSubscription *bool
LoginLabel *string
Load func() (*OAuthAuth, error)
}
LazyOAuthOptions advertises auth metadata without loading the flow. Load is read on first use. Its returned error is a rejected load and remains cached; a panic models a synchronous throw before Pi has assigned the load promise. Input/implementation callback mutations require caller synchronization.
type ListError ¶
ListError is a per-provider list-models failure. The rest of the collection still contributes its models.
type ListResult ¶
type ListResult struct {
Models []Model `json:"models"`
Errors []ListError `json:"errors,omitempty"`
}
ListResult is the union of listed models plus any per-provider errors.
type Message ¶
type Message struct {
Role string `json:"role"`
Content MessageContent `json:"content"`
Timestamp int64 `json:"timestamp"`
Sections SystemSections `json:"sections,omitempty"`
ToolsAdded []Tool `json:"toolsAdded,omitempty"`
ToolsRemoved []ToolReference `json:"toolsRemoved,omitempty"`
API string `json:"api,omitempty"`
Provider string `json:"provider,omitempty"`
Model string `json:"model,omitempty"`
ResponseModel *string `json:"responseModel,omitempty"`
ResponseID *string `json:"responseId,omitempty"`
ProviderThinkingLevel *string `json:"providerThinkingLevel,omitempty"`
ThinkingLevel *string `json:"thinkingLevel,omitempty"`
Diagnostics json.RawMessage `json:"diagnostics,omitempty"`
Usage *Usage `json:"usage,omitempty"`
StopReason string `json:"stopReason,omitempty"`
Deferred json.RawMessage `json:"deferred,omitempty"`
ErrorMessage *string `json:"errorMessage,omitempty"`
RawStopReason *string `json:"rawStopReason,omitempty"`
EndTurn *bool `json:"endTurn,omitempty"`
ToolCallID string `json:"toolCallId,omitempty"`
ToolName string `json:"toolName,omitempty"`
// Details is a live JSON-shaped value. nil/Undefined mean absent; Null is explicit null.
Details any `json:"-"`
NestedCalls json.RawMessage `json:"nestedCalls,omitempty"`
IsError bool `json:"isError,omitempty"`
Extra map[string]json.RawMessage `json:"-"`
// contains filtered or unexported fields
}
Message is Pi's transcript union. Role selects the applicable fields. Unlike agentcore.Message, content order, signatures, tool results, and system deltas remain in their original representation. Extra retains application metadata.
func GetCurrentSystemMessage ¶
func GetInitialSystemMessage ¶
func TransformMessageReferences ¶
func TransformMessageReferences(messages []*Message, model *Model, normalize ToolCallIDNormalizer) []*Message
TransformMessageReferences retains Pi's message/model object identity and shallow-copy boundaries, including edits made by the ID normalizer. Callers must synchronize concurrent access to these shared objects.
func TransformMessages ¶
func TransformMessages(messages []Message, model Model, normalize ToolCallIDNormalizer) []Message
TransformMessages ports Pi's api/transform-messages.ts: replay signatures only on the originating model, repair unanswered tool calls, and downgrade images for models that accept only text. Transformation itself does not mutate input messages; the normalizer receives their shared objects and may edit them. This value adapter snapshots the model and returned message structs. Use TransformMessageReferences to retain message/model identity throughout.
func WithToolChanges ¶
func WithToolChanges(message Message, changes ToolStateChanges) Message
WithToolChanges copies a system message with replacement declarations, matching Pi's agent-loop helper. Empty lists remove the fields, including retained null provenance. Other wire fields and the original message remain untouched; unchanged declaration slices may share storage as in Pi.
func (Message) HasContent ¶
HasContent reports own-field presence without serializing the message. Native zero content is explicit null; decoded messages can retain an absent field. Assigning text or a non-nil block list makes the field present again.
func (Message) MarshalJSON ¶
func (*Message) UnmarshalJSON ¶
type MessageContent ¶
type MessageContent struct {
Text *string
Blocks *BlockList
// contains filtered or unexported fields
}
MessageContent preserves the distinction between string and array content. Its zero value is JSON null, as accepted by transformMessages for old logs.
func BlockContent ¶
func BlockContent(blocks ...ContentBlock) MessageContent
func BlockReferences ¶
func BlockReferences(blocks ...*ContentBlock) MessageContent
BlockReferences constructs a new list retaining the supplied block objects. Copy MessageContent to share the list itself. BlockContent also constructs new block objects when the caller does not need to retain references.
func TextContent ¶
func TextContent(text string) MessageContent
func (MessageContent) MarshalJSON ¶
func (c MessageContent) MarshalJSON() ([]byte, error)
func (*MessageContent) UnmarshalJSON ¶
func (c *MessageContent) UnmarshalJSON(data []byte) error
type MetaOAuthOptions ¶
type MetaOAuthOptions struct {
Client *http.Client
Now func() float64
PollSleep func(float64, context.Context, string) error
RequestTimeout time.Duration
}
MetaOAuthOptions supplies concrete HTTP and clock hooks. Each request has its own timeout; the derived signal stays live after the request returns.
type Model ¶
type Model struct {
// API and Provider identify Pi's wire implementation and provider separately.
// Input declares supported modalities for transcript replay. The workspace
// catalog fields below remain available to existing Go callers.
API string `json:"api,omitempty"`
Provider string `json:"provider,omitempty"`
Input []string `json:"input,omitempty"`
ProviderID string `json:"provider_id"`
ProviderVendor string `json:"provider_vendor"`
ProviderName string `json:"provider_name"`
ID string `json:"id"`
// ContextWindow is the model's input window in tokens: the vendor's own
// figure when its list-models response carried one, otherwise this package's
// fallback, otherwise 0 for "unknown". It is what the compaction budget is
// capped against, so 0 must stay distinguishable from a real number.
ContextWindow int `json:"context_window,omitempty"`
// Capabilities is sparse tri-state metadata. Missing fields mean unknown,
// not unsupported; explicit live discovery values refine adapter defaults.
Capabilities protocol.ModelCapabilities `json:"capabilities,omitempty"`
}
Model is one entry from a provider's live list-models response. IDs are whatever the vendor returned — never a hardcoded catalog.
type ModelPrice ¶
type ModelPrice struct {
InputPerM float64 `json:"input_per_m"` // USD per 1M input tokens
OutputPerM float64 `json:"output_per_m"` // USD per 1M output tokens
// CacheReadPerM / CacheWritePerM are optional explicit cache rates; zero means
// "derive from InputPerM" (see Cost).
CacheReadPerM float64 `json:"cache_read_per_m,omitempty"`
CacheWritePerM float64 `json:"cache_write_per_m,omitempty"`
}
ModelPrice is the per-million-token price of a model, in USD. Providers quote input and output tokens separately, so we keep them separate and let the caller weigh them by the run's actual usage. Cache-read/write rates are optional: when left zero they derive from InputPerM by the standard provider multipliers (a cache hit is far cheaper than fresh input; an Anthropic cache write carries a small premium), so the default table stays terse while cache cost is still priced honestly.
type ModelProvider ¶
type ModelProvider struct {
ID string
Name string
BaseURL any
Headers any
Auth *ProviderAuth
GetModels func() (*Array, error)
GetAllModels func() (*Array, error)
RefreshModels func(ModelRefreshContext) error
FilterModels func(*Array, any) (*Array, error)
FilterAllModels func(*Array, any) (*Array, error)
Stream ModelStreamFunc
StreamSimple ModelStreamFunc
FetchDeferred DeferredStreamFunc
CancelDeferred func(context.Context, any, any, *Object) error
GenerateImages func(context.Context, any, any, *Object) (any, error)
Classify func(context.Context, any, any, *Object) (any, error)
}
ModelProvider carries Pi's concrete provider identity and catalog callbacks. GetAllModels and RefreshModels are optional. Callbacks retain their returned model/list identities and run without the registry lock. Provider mutation outside SetProvider is caller-synchronized, as are live model objects.
func NewModelProvider ¶
func NewModelProvider(input *ProviderFactoryOptions) (*ModelProvider, error)
NewModelProvider ports Pi's createProvider. Optional capabilities are fixed at construction, while dispatch observes replacements in the captured maps. Metadata does not implicitly override model or request options.
func RadiusProvider ¶
func RadiusProvider(input RadiusProviderOptions, dependencies RadiusProviderDependencies) (*ModelProvider, error)
type ModelRefreshContext ¶
type ModelStreamFunc ¶
type ModelStreamFunc func(context.Context, any, TranscriptContext, *Object) (*ProviderEventSource, error)
type Models ¶
type Models struct {
// contains filtered or unexported fields
}
Models is Pi's native provider registry. It is distinct from the workspace Collection, whose live discovery and owner-by-model-ID rules differ from Pi. Native auth lifecycle, availability, catalog refresh and request dispatch share this registry. Concrete provider factories are wired separately.
func NewModels ¶
func NewModels(options ...ModelsOptions) *Models
func (*Models) CancelDeferred ¶
func (*Models) CheckAuth ¶
CheckAuth checks configuration without refreshing stored OAuth credentials. API-key providers may supply a side-effect-free Check; otherwise resolution performs its own second credential read, as in Pi.
func (*Models) ClearProviders ¶
func (m *Models) ClearProviders()
func (*Models) Complete ¶
func (m *Models) Complete(ctx context.Context, model any, input Context, options ...*ModelsRequestOptions) (*Message, error)
Completion waits for the stream result; request cancellation is handled by auth/provider callbacks and normally becomes an error message, not a Go wait error. Callers needing independent reader cancellation can use Stream.Result.
func (*Models) CompleteSimple ¶
func (*Models) DeleteProvider ¶
func (*Models) FetchDeferred ¶
func (*Models) GenerateImages ¶
func (*Models) GetAllAvailable ¶
func (*Models) GetAllModels ¶
func (*Models) GetAuth ¶
func (m *Models) GetAuth(ctx context.Context, providerOrModel any, options ...AuthResolutionOverrides) (any, error)
GetAuth accepts a provider ID or a model Object. Model headers override auth headers case-insensitively without changing the credential or auth result. An unknown provider returns Undefined even when the caller is already aborted.
func (*Models) GetAvailable ¶
Availability reads propagate provider errors. Passing "" selects every provider, unlike the synchronous catalog accessors. Returned unions skip sparse holes and preserve model identity.
func (*Models) GetAvailableOfType ¶
func (*Models) GetModelOfType ¶
func (*Models) GetModels ¶
GetModels/GetAllModels return the provider's original list when scoped and a fresh union when unscoped. Passing no argument differs from passing "".
func (*Models) GetModelsOfType ¶
func (*Models) GetProvider ¶
func (m *Models) GetProvider(id string) *ModelProvider
func (*Models) GetProviders ¶
func (m *Models) GetProviders() []*ModelProvider
func (*Models) Login ¶
func (m *Models) Login(providerID, authType string, interaction ProviderAuthInteraction, options ...*OAuthLoginOptions) (any, error)
Login preserves the provider's credential object and waits for its write to settle. Cancellation stops waiting for a queued mutation. Once the store invokes the mutation callback, it owns the commit/abort boundary.
func (*Models) Logout ¶
Logout also accepts an unknown provider ID: stored credentials outlive provider registration. An admitted delete follows the persistence callback's settlement, including stores that finish a deletion after caller cancellation.
func (*Models) Refresh ¶
func (m *Models) Refresh(ctx context.Context, options ...ModelsRefreshOptions) ModelsRefreshResult
Refresh restores cached catalogs before resolving auth or fetching models. Providers run independently. Cancellation stops waiting without releasing an unfinished persistence transaction or admitting a superseded publication.
func (*Models) SetProvider ¶
func (m *Models) SetProvider(provider *ModelProvider)
func (*Models) Stream ¶
func (m *Models) Stream(ctx context.Context, model any, input Context, options ...*ModelsRequestOptions) *AssistantMessageEventStream
func (*Models) StreamDeferred ¶
func (m *Models) StreamDeferred(ctx context.Context, model, handle any, options ...*ModelsRequestOptions) *AssistantMessageEventStream
func (*Models) StreamSimple ¶
func (m *Models) StreamSimple(ctx context.Context, model any, input Context, options ...*ModelsRequestOptions) *AssistantMessageEventStream
type ModelsError ¶
ModelsError retains Pi's public error category and the underlying Go cause.
func NewModelsError ¶
func NewModelsError(code, message string, cause error) *ModelsError
func (*ModelsError) Error ¶
func (e *ModelsError) Error() string
func (*ModelsError) Unwrap ¶
func (e *ModelsError) Unwrap() error
type ModelsOptions ¶
type ModelsOptions struct {
Persistence *ModelsPersistence
Credentials *CredentialPersistence
AuthContext *AuthContext
}
type ModelsPersistence ¶
type ModelsPersistence struct {
Read func(context.Context, string) (any, error)
Write func(context.Context, string, any) error
Delete func(context.Context, string) error
}
ModelsPersistence is the concrete callback boundary for provider snapshots. Blocking callbacks must honor ctx; a canceled publication returns promptly, while its per-provider queue remains held until the callback actually settles.
type ModelsPublication ¶
ModelsPublication applies persistence before the provider's in-memory update. nil/Undefined leave persistence unchanged; Null deletes the stored entry.
type ModelsPublicationScope ¶
type ModelsPublicationScope struct {
Context context.Context
// contains filtered or unexported fields
}
func (*ModelsPublicationScope) Finish ¶
func (s *ModelsPublicationScope) Finish()
Finish retires only this scope's active controller. Publications retained by a provider may still run until its generation is superseded. Finishing an old refresh must not retire a replacement refresh's controller.
func (*ModelsPublicationScope) Publish ¶
func (s *ModelsPublicationScope) Publish(publication ModelsPublication) (bool, error)
Publish queues the entire transaction, including update, before waiting. Even when cancellation wins the caller's wait, the queue cannot be released early: a store that ignores cancellation could otherwise overwrite a newer snapshot.
type ModelsPublisher ¶
type ModelsPublisher struct {
// contains filtered or unexported fields
}
ModelsPublisher implements the generation check and publication chain used by Pi's Models collection. Provider catalog updates run only after successful persistence and a second generation check. Independent providers do not block each other. Construct it with the same persistence used for refresh reads.
func NewModelsPublisher ¶
func NewModelsPublisher(persistence ModelsPersistence) *ModelsPublisher
func (*ModelsPublisher) Begin ¶
func (p *ModelsPublisher) Begin(ctx context.Context, providerID string) *ModelsPublicationScope
Begin supersedes any earlier refresh for the provider, including cancellation of its network/persistence context. Other provider scopes are unaffected.
func (*ModelsPublisher) Clear ¶
func (p *ModelsPublisher) Clear()
func (*ModelsPublisher) Invalidate ¶
func (p *ModelsPublisher) Invalidate(providerID string)
Invalidate is the publication boundary for deleting/replacing a provider.
type ModelsRefreshFailure ¶
type ModelsRefreshOptions ¶
type ModelsRefreshResult ¶
type ModelsRefreshResult struct {
Aborted bool
// Errors retain completion order, like Pi's Map, and original error identity.
Errors []ModelsRefreshFailure
}
type ModelsRequestOptions ¶
type ModelsRequestOptions struct {
Values any
TransformHeaders func(any) (any, error)
Now func() int64
}
Values retains provider-specific options and their field presence. The Models-only header callback runs after auth/model/request headers are merged. Now supplies the clock for failure results, never a provider request option.
type NativeClient ¶
type NativeClient struct {
// contains filtered or unexported fields
}
NativeClient binds host credentials and endpoint defaults to the lossless stream API. Credentials stay in the client, never in checkpointed models.
func NewNativeClient ¶
func NewNativeClient(spec ClientSpec) (*NativeClient, error)
func (*NativeClient) Candidate ¶
func (c *NativeClient) Candidate(id string, contextWindow int) (FallbackCandidate, error)
Candidate resolves a model for the AI fallback provider. ContextWindow is a host limit; the model metadata and opaque native messages survive checkpoints.
type NativeGenerationControls ¶
type NativeGenerationControls struct {
ToolChoice protocol.ToolChoice
ParallelToolCalls *bool
OutputSchema *protocol.OutputSchema
}
NativeGenerationControls are host-requested generation policies shared by native clients. They do not re-encode provider history or opaque signatures.
type NativeProvider ¶
type NativeProvider struct{ Tokens TokenSource }
NativeProvider selects a native Go transport from a model's declared API. Tokens optionally binds account rotation for Codex, Antigravity or Claude Code. It implements the same StreamFn contract as FallbackProvider.
func (NativeProvider) Stream ¶
func (p NativeProvider) Stream(ctx context.Context, model json.RawMessage, transcript TranscriptContext, options map[string]any) (*AssistantMessageEventStream, error)
type NativeProviderFailure ¶
type NativeProviderFailure struct {
// contains filtered or unexported fields
}
NativeProviderFailure is a request-scoped, passive host observation. It never changes Pi's event JSON or wraps the error passed to Pi's display formatter.
func WithNativeProviderFailure ¶
func WithNativeProviderFailure(ctx context.Context) (context.Context, *NativeProviderFailure)
func (*NativeProviderFailure) Failure ¶
func (c *NativeProviderFailure) Failure() error
func (*NativeProviderFailure) HostFailure ¶
func (c *NativeProviderFailure) HostFailure() bool
func (*NativeProviderFailure) Kind ¶
func (c *NativeProviderFailure) Kind() string
type OAuthAuth ¶
type OAuthAuth struct {
Name string
IsSubscription *bool
LoginLabel *string
Login func(ProviderAuthInteraction, *OAuthLoginOptions) (any, error)
Refresh func(context.Context, any) (any, error)
ToAuth func(any) (any, error)
}
func AnthropicOAuth ¶
func AnthropicOAuth(settings ...AnthropicOAuthOptions) *OAuthAuth
func AntigravityOAuth ¶
func AntigravityOAuth(settings ...AntigravityOAuthOptions) *OAuthAuth
AntigravityOAuth builds the interactive login for the google-antigravity vendor: Google code grant → userinfo email → Cloud Code project discovery (+free-tier onboarding). The returned credential carries access/refresh/ expires plus `project` and `account` for the vault.
func DevinOAuth ¶
func DevinOAuth(settings ...DevinOAuthOptions) *OAuthAuth
DevinOAuth builds the interactive login for the Devin vendor. The credential it returns is the same shape the other OAuth vendors store: a JS object with type/access/refresh/expires.
func GitHubCopilotOAuth ¶
func GitHubCopilotOAuth(options GitHubCopilotOAuthOptions) (*OAuthAuth, error)
func KimiCodingOAuth ¶
func KimiCodingOAuth(settings ...KimiCodingOAuthOptions) *OAuthAuth
func LazyOAuth ¶
func LazyOAuth(input *LazyOAuthOptions) *OAuthAuth
LazyOAuth shares one load across login, refresh and toAuth, including concurrent callers. Waiting for that load does not inspect the operation's cancellation signal; the loaded flow owns cancellation and its result.
func MetaOAuth ¶
func MetaOAuth(settings ...MetaOAuthOptions) *OAuthAuth
func OpenAIChatGPTOAuth ¶
func OpenAIChatGPTOAuth(settings ...OpenAIChatGPTOAuthOptions) *OAuthAuth
func OpenAICodexOAuth ¶
func OpenAICodexOAuth(settings ...OpenAICodexOAuthOptions) *OAuthAuth
func OpenRouterOAuth ¶
func OpenRouterOAuth(settings ...OpenRouterOAuthOptions) *OAuthAuth
func RadiusOAuth ¶
func RadiusOAuth(input *RadiusOAuthOptions, settings ...RadiusOAuthRuntimeOptions) *OAuthAuth
func XaiOAuth ¶
func XaiOAuth(settings ...XaiOAuthOptions) *OAuthAuth
type OAuthCallbackServer ¶
type OAuthCallbackServer struct {
RedirectURI string
Wait func() (any, error)
Cancel func()
Close func()
}
OAuthCallbackServer is a concrete callback boundary. Wait shares one result, retaining its identity; cancellation without a claimed code returns Undefined. Close stops accepting connections and lets active completion callbacks finish.
func StartOAuthCallbackServer ¶
func StartOAuthCallbackServer(options *OAuthCallbackServerOptions) (*OAuthCallbackServer, error)
StartOAuthCallbackServer binds the requested TCP port without fallback. Path, state and Complete remain live; callers synchronize any option mutations.
type OAuthDeviceCodePollOptions ¶
type OAuthDeviceCodePollOptions struct {
IntervalSeconds any
ExpiresInSeconds any
WaitBeforeFirstPoll bool
Poll func() (OAuthDeviceCodePollResult, error)
Context context.Context
Now func() float64
Sleep func(float64, context.Context, string) error
}
OAuthDeviceCodePollOptions keeps polling callbacks and cancellation live. Numeric controls accept JSON values to preserve Pi's coercion and number-only expiry/server-interval checks. Nil/Undefined mean omitted for initial timing. Now and Sleep are concrete native clock dependencies; defaults use real time.
type OAuthDiagnosticError ¶
type OAuthDiagnosticError struct {
Name *string
Message string
Code any
Errno any
Cause any
Stack string
// contains filtered or unexported fields
}
OAuthDiagnosticError carries the Error fields read by Pi's OAuth diagnostic formatter. Nil metadata means absent; Null represents an explicit null cause or errno. Stack is host-provided: a Go stack is not a JavaScript stack.
func (*OAuthDiagnosticError) Error ¶
func (e *OAuthDiagnosticError) Error() string
func (*OAuthDiagnosticError) Unwrap ¶
func (e *OAuthDiagnosticError) Unwrap() error
type OAuthLoginOptions ¶
type OAuthLoginOptions struct{ GetDeviceID func() string }
type OAuthManualPrompt ¶
type OAuthManualPrompt struct{ Message, Placeholder string }
type OAuthToken ¶
type OAuthToken struct {
// AccountID is the workspace_provider_accounts row id — the handle Report
// uses to block/disable the account that served a failed request.
AccountID string
// AccessToken is the OAuth access token sent as the Bearer credential.
AccessToken string
// ProviderAccountID is the vendor-side account id (Codex chatgpt_account_id).
ProviderAccountID string
// ProjectID is the Cloud Code Assist project (Antigravity only).
ProjectID string
// Email is the account's login email, for diagnostics.
Email string
}
OAuthToken is one acquired account credential, handed to a wire client for a single request. The extra identity fields ride along because the wires need them: Codex sends AccountID as the chatgpt-account-id header, Antigravity puts ProjectID in the request envelope.
type Object ¶
Object and Array carry mutable metadata without serializing nested values.
func FlattenChatModelCatalog ¶
FlattenChatModelCatalog projects Pi's grouped catalog onto model IDs. The provider argument is type metadata in Pi; it never rewrites model fields. Use Object/Array containers to retain enumeration order and model identity. Missing type is excluded here, unlike the model collection's legacy-chat rule.
func GetRadiusCredentialConfig ¶
A nil result corresponds to the source's undefined configuration.
func LoadRadiusGatewayConfig ¶
func WaitForCallbackOrManualInput ¶
func WaitForCallbackOrManualInput(interaction ProviderAuthInteraction, callback *OAuthCallbackServer, prompt OAuthManualPrompt) (*Object, error)
WaitForCallbackOrManualInput cancels the manual prompt on every exit. A claimed browser callback keeps priority over successful manual input; manual failure is checked after a successful callback wait, as in Pi.
type OpenAICodexOAuthOptions ¶
type OpenAICompletionsCompat ¶
type OpenAICompletionsCompat struct {
SupportsStore bool `json:"supportsStore"`
SupportsDeveloperRole bool `json:"supportsDeveloperRole"`
SupportsReasoningEffort bool `json:"supportsReasoningEffort"`
SupportsUsageInStreaming bool `json:"supportsUsageInStreaming"`
SupportsFinishReason bool `json:"supportsFinishReason"`
MaxTokensField string `json:"maxTokensField"`
RequiresToolResultName bool `json:"requiresToolResultName"`
RequiresAssistantAfterToolResult bool `json:"requiresAssistantAfterToolResult"`
RequiresThinkingAsText bool `json:"requiresThinkingAsText"`
RequiresReasoningContentOnAssistantMessages bool `json:"requiresReasoningContentOnAssistantMessages"`
ThinkingFormat string `json:"thinkingFormat"`
OpenRouterRouting json.RawMessage `json:"openRouterRouting"`
VercelGatewayRouting json.RawMessage `json:"vercelGatewayRouting"`
ChatTemplateKwargs json.RawMessage `json:"chatTemplateKwargs"`
ChatTemplateArgs json.RawMessage `json:"chatTemplateArgs"`
ZaiToolStream bool `json:"zaiToolStream"`
SupportsThinkingTokenBudget bool `json:"supportsThinkingTokenBudget"`
ThinkingTokenBudgetField string `json:"thinkingTokenBudgetField,omitempty"`
SupportsStrictMode bool `json:"supportsStrictMode"`
SupportsOpenAIGrammarTools bool `json:"supportsOpenAIGrammarTools"`
SupportsMidConvoSystemMessages bool `json:"supportsMidConvoSystemMessages"`
SupportsMidConvoToolAdditions bool `json:"supportsMidConvoToolAdditions"`
CacheControlFormat string `json:"cacheControlFormat,omitempty"`
SendSessionAffinityHeaders bool `json:"sendSessionAffinityHeaders"`
SessionAffinityFormat string `json:"sessionAffinityFormat"`
SupportsLongCacheRetention bool `json:"supportsLongCacheRetention"`
VLLMPriority json.RawMessage `json:"vllmPriority,omitempty"`
}
OpenAICompletionsCompat is the resolved Pi contract. Explicit model.compat overrides are applied by ResolveOpenAICompletionsCompat; this is independent of the legacy provider's smaller Compat table.
func ResolveOpenAICompletionsCompat ¶
func ResolveOpenAICompletionsCompat(raw json.RawMessage) (OpenAICompletionsCompat, error)
type OpenAICompletionsStreamOptions ¶
type OpenAICompletionsStreamOptions struct {
Options json.RawMessage
Client *http.Client
OnPayload func(context.Context, json.RawMessage, json.RawMessage) (json.RawMessage, error)
OnResponse func(context.Context, CompletionsResponse, json.RawMessage) error
OnProviderStreamEvent func(context.Context, *json.RawMessage, json.RawMessage) error
Now func() int64
}
OpenAICompletionsStreamOptions carries serializable Pi provider controls in Options and concrete Go dependencies/callbacks separately. A nil OnPayload result preserves the payload. OnProviderStreamEvent may replace the decoded JSON through its pointer before accumulation, just as Pi permits mutation. Callbacks run serially on the producer and must honor their context.
func BindNativeStreamOptions ¶
func BindNativeStreamOptions(raw json.RawMessage, options map[string]any, defaults OpenAICompletionsStreamOptions) (OpenAICompletionsStreamOptions, error)
BindNativeStreamOptions binds the host's concrete callbacks and HTTP client separately from serializable controls. Nil callback entries preserve defaults. Provider-specific validation and request cancellation belong to the stream.
type OpenAIEmbedder ¶
OpenAIEmbedder calls the OpenAI /embeddings endpoint, reusing the project's BYO key and base_url. It backs the protocol.Embedder seam so semantic recall works for the OpenAI-wire provider family (OpenAI + compatible vendors) with no extra credentials.
func NewOpenAIEmbedder ¶
func NewOpenAIEmbedder(apiKey, baseURL, model string) *OpenAIEmbedder
NewOpenAIEmbedder builds an embedder. Empty baseURL/model fall back to the OpenAI defaults.
type OpenAIProvider ¶
type OpenAIProvider struct {
APIKey string
BaseURL string
Compat Compat
// HTTP serves the non-streamed Chat path (absolute cap, no header deadline —
// a buffered completion's headers arrive only once generation is done).
HTTP *http.Client
// StreamHTTP serves the SSE path, where a header deadline is meaningful. Nil
// falls back to HTTP, so a caller that overrides only HTTP still works.
StreamHTTP *http.Client
// Vendor, when set, overrides Name() for OpenAI-compatible vendors served
// through this wire (the pi-ai pattern: most providers are this API at a
// different base URL). Empty keeps the stock "openai" identity.
Vendor string
// contains filtered or unexported fields
}
OpenAIProvider speaks the OpenAI chat-completions wire format. base_url resolution (per-config -> OPENAI_BASE_URL env -> vendor default) is performed by the caller; this struct receives the resolved BaseURL.
func NewGeminiProvider ¶
func NewGeminiProvider(apiKey string) *OpenAIProvider
NewGeminiProvider builds a Google Gemini provider on the OpenAI-compatible wire — provider breadth as config, not new code (Compat table pattern). Use Gemini model ids ("gemini-2.5-pro", "gemini-2.5-flash"); tools, streaming, and usage accounting ride the shared implementation, and Name() reports "google" so traces, tiers, and RefreshKey attribute to the right vendor.
func NewOpenAIProvider ¶
func NewOpenAIProvider(apiKey, baseURL string, compat Compat) *OpenAIProvider
NewOpenAIProvider builds a provider. An empty baseURL falls back to the vendor default; an empty Compat falls back to DefaultCompat.
func (*OpenAIProvider) Chat ¶
func (p *OpenAIProvider) Chat(ctx context.Context, req protocol.ChatRequest) (protocol.ChatResponse, error)
Chat performs one non-streaming completion. OpenAI-compatible endpoints are inconsistent about optional request hints; when one explicitly rejects a hint, the adapter removes only that hint and retains the endpoint lesson in the logical provider session.
func (*OpenAIProvider) ModelCapabilities ¶
func (p *OpenAIProvider) ModelCapabilities(model string) protocol.ModelCapabilities
func (*OpenAIProvider) Name ¶
func (p *OpenAIProvider) Name() string
func (*OpenAIProvider) Stream ¶
func (p *OpenAIProvider) Stream(ctx context.Context, req protocol.ChatRequest) (<-chan protocol.ChatDelta, error)
Stream performs a real token-streaming completion: it sets stream=true and emits one protocol.ChatDelta per server-sent content chunk, accumulating fragmented tool-call deltas (which arrive piecewise, keyed by index) into whole ToolCalls flushed before the terminal Done delta. The loop forwards content deltas to a live SSE sink; tool execution is unchanged.
func (*OpenAIProvider) SupportsTools ¶
func (p *OpenAIProvider) SupportsTools() bool
func (*OpenAIProvider) UpdateAPIKey ¶
func (p *OpenAIProvider) UpdateAPIKey(key string)
UpdateAPIKey swaps the key used for subsequent requests (protocol.KeyUpdater), letting the loop refresh an expiring BYO token between turns.
type OpenAIResponsesCompat ¶
type OpenAIResponsesCompat struct {
SupportsDeveloperRole bool `json:"supportsDeveloperRole"`
SupportsMidConvoSystemMessages bool `json:"supportsMidConvoSystemMessages"`
SessionAffinityFormat string `json:"sessionAffinityFormat"`
SupportsLongCacheRetention bool `json:"supportsLongCacheRetention"`
SupportsStrictMode bool `json:"supportsStrictMode"`
SupportsOpenAIGrammarTools bool `json:"supportsOpenAIGrammarTools"`
SupportsAdditionalTools bool `json:"supportsAdditionalTools"`
SupportsToolSearch bool `json:"supportsToolSearch"`
SupportsExplicitPromptCacheMode bool `json:"supportsExplicitPromptCacheMode"`
SupportsMaxOutputTokens bool `json:"supportsMaxOutputTokens"`
}
func ResolveOpenAIResponsesCompat ¶
func ResolveOpenAIResponsesCompat(raw json.RawMessage) (OpenAIResponsesCompat, error)
type OpenAIResponsesProvider ¶
type OpenAIResponsesProvider struct {
APIKey string
BaseURL string
HTTP *http.Client
StreamHTTP *http.Client
// Vendor overrides Name() when this Responses wire is selected for a model
// belonging to an ordinary OpenAI provider row. Empty preserves the explicit
// openai-responses identity.
Vendor string
// contains filtered or unexported fields
}
OpenAIResponsesProvider speaks the public OpenAI Responses API. It is an explicit provider rather than a silent replacement for OpenAI Chat Completions: Responses can retain server-side conversation state, and that privacy/compatibility choice belongs to the workspace operator.
When agentcore supplies a logical provider session, successful turns are chained with previous_response_id after an exact wire-prefix check. Without that state (or after an eviction) the adapter sends the complete transcript, so correctness never depends on process affinity.
func NewOpenAIResponsesProvider ¶
func NewOpenAIResponsesProvider(apiKey, baseURL string) *OpenAIResponsesProvider
NewOpenAIResponsesProvider builds the public Responses wire client. baseURL is the API root (normally https://api.openai.com/v1); /responses is appended unless the caller already supplied the endpoint itself.
func (*OpenAIResponsesProvider) Chat ¶
func (p *OpenAIResponsesProvider) Chat(ctx context.Context, req protocol.ChatRequest) (protocol.ChatResponse, error)
Chat drains the public API's streaming transport into one neutral response. Keeping one decoder prevents Chat and Stream from learning different chain baselines or tool-call semantics.
func (*OpenAIResponsesProvider) ModelCapabilities ¶
func (p *OpenAIResponsesProvider) ModelCapabilities(model string) protocol.ModelCapabilities
func (*OpenAIResponsesProvider) Name ¶
func (p *OpenAIResponsesProvider) Name() string
func (*OpenAIResponsesProvider) Stream ¶
func (p *OpenAIResponsesProvider) Stream(ctx context.Context, req protocol.ChatRequest) (<-chan protocol.ChatDelta, error)
func (*OpenAIResponsesProvider) SupportsTools ¶
func (p *OpenAIResponsesProvider) SupportsTools() bool
func (*OpenAIResponsesProvider) UpdateAPIKey ¶
func (p *OpenAIResponsesProvider) UpdateAPIKey(key string)
type OpenAIResponsesStreamOptions ¶
type OpenAIResponsesStreamOptions = OpenAICompletionsStreamOptions
OpenAIResponsesStreamOptions uses the same concrete HTTP client and callback contracts as Completions; Options carries Responses-specific request controls.
type OpenAIWire ¶
type OpenAIWire string
OpenAIWire is the API dialect used behind an OpenAI provider identity. Provider identity owns credentials and routing; the selected model owns the wire. Explicit openai-responses vendors remain supported for operators who want the wire fixed for the whole provider row.
const ( OpenAIWireChat OpenAIWire = "chat-completions" OpenAIWireResponses OpenAIWire = "responses" )
type OpenRouterOAuthOptions ¶
type OpenRouterOAuthOptions struct {
Client *http.Client
PKCE func() (PKCE, error)
RandomUUID func() (string, error)
StartCallback func(*OAuthCallbackServerOptions) (*OAuthCallbackServer, error)
CallbackHost func() string
LoginTimeoutMS *float64
ExchangeTimeout time.Duration
}
OpenRouterOAuthOptions supplies native transport, entropy and listener dependencies. Defaults use crypto/rand, a real loopback server, a five-minute login deadline and a thirty-second exchange deadline. It does not change the provider's authorization or token URLs.
type PKCE ¶
func GeneratePKCE ¶
GeneratePKCE uses 32 cryptographically random bytes and an S256 challenge, matching Pi's Web Crypto flow. The challenge hashes the encoded verifier.
type PayloadYield ¶
PayloadYield releases payload access while work waits, then reacquires it before returning (or propagating a panic). It is the explicit Go counterpart of yielding a JavaScript callback at await. Do not access live payload fields inside work; use them before yielding or after it returns.
type PiMessagesResponseError ¶
type PiMessagesResponseError struct {
Message string
Code *string
DiagnosticDetails *Object
Status int
Headers http.Header
}
func (*PiMessagesResponseError) Error ¶
func (e *PiMessagesResponseError) Error() string
type PiMessagesStreamOptions ¶
type PiMessagesStreamOptions struct {
Values *Object
Client *http.Client
OnPayload func(context.Context, any, *Object) (any, error)
OnResponse func(context.Context, CompletionsResponse, *Object) error
OnProviderStreamEvent func(context.Context, any, *Object) error
Now func() float64
ErrorStack func(error) string
// contains filtered or unexported fields
}
PiMessagesStreamOptions retains JSON-shaped values and serial callbacks. Undefined from OnPayload preserves its input; nil replaces it with JSON null. Model, Values and callback values are live. Callers synchronize edits while a request is active. ErrorStack supplies host stack metadata; Go cannot produce the JavaScript source locations of the reference runtime.
type PreparationError ¶
type PreparationError struct{ Cause error }
PreparationError marks host setup failures, which must never retry or fall back even when their cause is a typed provider error.
func (*PreparationError) Error ¶
func (e *PreparationError) Error() string
func (*PreparationError) Unwrap ¶
func (e *PreparationError) Unwrap() error
type Pricing ¶
type Pricing map[string]ModelPrice
Pricing maps a model name to its price. Lookup is exact first, then by longest matching prefix, so a table keyed on a family ("gpt-4o", "claude-3-5-sonnet") still prices a dated or suffixed variant ("gpt-4o-2024-08-06"). A model with no entry prices at zero — the cost field stays an honest 0 rather than a guess.
func DefaultPricing ¶
func DefaultPricing() Pricing
DefaultPricing is the built-in price table (USD per 1M tokens), keyed by model family so dated variants resolve by prefix. These are list prices and drift; they are a code-defined default, overridable per Runner via WithPricing. Keep entries sorted by family for easy auditing.
func (Pricing) Cost ¶
Cost weighs a usage record against the model's price, and reports whether the model was actually found in the table. The second return is load-bearing: an unknown model must never be mistaken for a genuinely free one. A caller that discards it and treats the float as gospel reproduces the bug this type exists to prevent — a model with no entry silently billing as "$0.00 spent" instead of "we don't know what this cost." Cache-read and cache-write tokens are priced as their own categories (so a long run whose stable prefix is served from cache is billed at the discounted rate, not the full input rate); absent explicit cache rates they derive from InputPerM by the standard multipliers.
type Provider ¶
type Provider interface {
protocol.LLMProvider
ID() string
Vendor() string
DisplayName() string
BaseURL() string
APIKey() string
ListModels(ctx context.Context) ([]Model, error)
}
Provider is the runtime unit: identity, auth, live model list, Chat/Stream.
type ProviderAuth ¶
type ProviderAuth struct {
APIKey *APIKeyAuth
OAuth *OAuthAuth
}
type ProviderAuthInteraction ¶
type ProviderClassifier ¶
type ProviderEventSource ¶
type ProviderEventSource struct {
Next func(context.Context) (AssistantMessageEvent, bool, error)
Result func(context.Context) (*Message, bool, error)
}
ProviderEventSource is a concrete async-iteration boundary. Result is optional; its bool distinguishes undefined (no result) from explicit null. Iteration and result callbacks remain live until forwarding finishes.
func SourceFromAssistantStream ¶
func SourceFromAssistantStream(stream *AssistantMessageEventStream) *ProviderEventSource
type ProviderFactoryOptions ¶
type ProviderFactoryOptions struct {
ID string
Name *string
BaseURL any
Headers any
Auth *ProviderAuth
Models *Array
FetchModels func(ModelRefreshContext) (*Array, error)
FilterModels func(*Array, any) (*Array, error)
FilterAllModels func(*Array, any) (*Array, error)
API *ProviderStreams
APIs map[string]*ProviderStreams
Images map[string]*ProviderImages
Classifiers map[string]*ProviderClassifier
Now func() int64
}
ProviderFactoryOptions composes a provider from concrete callbacks. API is a single implementation for every chat model; otherwise APIs dispatches by model.api. Maps, auth, baseline models and their entries retain identity. Callers synchronize mutations, as for ModelProvider itself.
type ProviderImages ¶
type ProviderModelCatalog ¶
type ProviderModelCatalog struct {
RefreshModels func(ModelRefreshContext) error
// contains filtered or unexported fields
}
ProviderModelCatalog is the catalog state owned by Pi's createProvider. The baseline and model objects retain identity; callers synchronize mutations of those objects. Refresh publication and catalog list construction are safe to call concurrently. Static catalogs have no RefreshModels callback.
func NewProviderModelCatalog ¶
func NewProviderModelCatalog(providerID string, baseline *Array, fetch func(ModelRefreshContext) (*Array, error), now func() int64) *ProviderModelCatalog
func (*ProviderModelCatalog) GetAllModels ¶
func (c *ProviderModelCatalog) GetAllModels() (*Array, error)
func (*ProviderModelCatalog) GetModels ¶
func (c *ProviderModelCatalog) GetModels() (*Array, error)
type ProviderOption ¶
type ProviderOption struct {
ID string `json:"id"`
Label string `json:"label"`
Auth string `json:"auth"`
Purpose string `json:"purpose"`
BaseURL string `json:"base_url,omitempty"`
}
ProviderOption describes supported host setup; credentials are never metadata.
func ProviderOptions ¶
func ProviderOptions() []ProviderOption
type ProviderStreams ¶
type ProviderStreams struct {
Stream ModelStreamFunc
StreamSimple ModelStreamFunc
FetchDeferred DeferredStreamFunc
CancelDeferred func(context.Context, any, any, *Object) error
}
func AnthropicMessagesAPI ¶
func AnthropicMessagesAPI(defaults OpenAICompletionsStreamOptions) *ProviderStreams
func AzureOpenAIResponsesAPI ¶
func AzureOpenAIResponsesAPI(defaults OpenAICompletionsStreamOptions) *ProviderStreams
func LazyAPI ¶
func LazyAPI(load func(context.Context) (*ProviderStreams, error), capabilities LazyAPICapabilities, now func() int64) *ProviderStreams
LazyAPI invokes load per request. Import/module caching belongs to the caller's loader, as in Pi; this adapter does not cache load failures.
func OpenAICodexResponsesAPI ¶
func OpenAICodexResponsesAPI(defaults OpenAICompletionsStreamOptions) *ProviderStreams
func OpenAICompletionsAPI ¶
func OpenAICompletionsAPI(defaults OpenAICompletionsStreamOptions) *ProviderStreams
Native API factories compose the existing stream implementations with Pi's lazy provider boundary. Defaults carry concrete HTTP/clock/callback hooks; request controls and callbacks override them for that request.
func OpenAIResponsesAPI ¶
func OpenAIResponsesAPI(defaults OpenAICompletionsStreamOptions) *ProviderStreams
func PiMessagesAPI ¶
func PiMessagesAPI(defaults PiMessagesStreamOptions) *ProviderStreams
PiMessagesAPI is the concrete provider factory used by Radius and custom providers. Request callbacks use the same signatures as the option fields.
type RadiusOAuthOptions ¶
type RadiusOAuthOptions struct{ Name, Gateway string }
Name remains live for prompts; the gateway and advertised name are captured at construction. Callers synchronize any later options mutations.
type RadiusOAuthResponseError ¶
RadiusOAuthResponseError preserves the status and optional OAuth error read by the device poller. An empty OAuth error is distinct from an absent one.
func (*RadiusOAuthResponseError) Error ¶
func (e *RadiusOAuthResponseError) Error() string
type RadiusProviderDependencies ¶
type RadiusProviderDependencies struct {
BaselineModels *Object
Streams *ProviderStreams
Client *http.Client
Now func() float64
LoadOAuth func(*RadiusOAuthOptions) (*OAuthAuth, error)
}
The pinned tree omits generated Radius model data. Supply that catalog explicitly; no current external catalog is silently substituted. Streams defaults to the native Go pi-messages implementation.
type RadiusProviderOptions ¶
type RadiusProviderOptions struct{ ID, Name, Gateway *string }
type ResponsesMessagesOptions ¶
type ResponsesMessagesOptions struct {
IncludeSystemPrompt *bool `json:"includeSystemPrompt,omitempty"`
GrammarToolInputProperties map[string]string `json:"grammarToolInputProperties,omitempty"`
SupportsMidConvoSystemMessages bool `json:"supportsMidConvoSystemMessages,omitempty"`
SupportsAdditionalTools bool `json:"supportsAdditionalTools,omitempty"`
SupportsToolSearch bool `json:"supportsToolSearch,omitempty"`
ToolOptions ResponsesToolsOptions `json:"toolOptions,omitempty"`
}
type ResponsesToolsOptions ¶
type ResponsesToolsOptions struct {
Strict json.RawMessage `json:"strict,omitempty"`
SupportsStrictMode *bool `json:"supportsStrictMode,omitempty"`
SupportsOpenAIGrammarTools bool `json:"supportsOpenAIGrammarTools,omitempty"`
ToolSearchResult bool `json:"toolSearchResult,omitempty"`
}
ResponsesToolsOptions preserves the distinction between absent and null strict. SupportsStrictMode defaults to true, as in Pi's shared Responses converter.
type Spec ¶
type Spec struct {
ID string
Vendor string
Name string
APIKey string
BaseURL string
HTTP HTTPDoer
// Compat selects OpenAI-compatible wire details such as max_tokens versus
// max_completion_tokens. Zero uses the vendor default.
Compat Compat
// TokenSource is the OAuth account pool a subscription vendor draws
// per-request credentials from. Required for OAuth vendors, ignored by
// the rest.
TokenSource TokenSource
}
Spec constructs a provider. Vendor is openai | openai-responses | anthropic | google | gemini, an OAuth subscription vendor (claude-code | openai-codex | google-antigravity), or any OpenAI-compatible name (which requires BaseURL). ID is the caller's stable handle (a workspace provider row id); empty ID falls back to Vendor.
type SpeechSpec ¶
type SpeechSpec struct{ Vendor, BaseURL, APIKey, Model, Voice, Text string }
type StreamFn ¶
type StreamFn func(context.Context, json.RawMessage, TranscriptContext, map[string]any) (*AssistantMessageEventStream, error)
StreamFn is the native provider contract shared by providers and the engine.
func ScriptedStream ¶
ScriptedStream emits a sequence of native assistant messages for offline demos and tests. Each invocation owns its script; concurrent runs must create separate streams. Exhaustion returns an empty completed answer.
type SystemSection ¶
SystemSections is an ordered JSON object, not a Go map: section order changes the actual model prompt. A nil Value removes a section. As in JavaScript, array-index names enumerate before other names, in numeric order.
type SystemSections ¶
type SystemSections []SystemSection
func (SystemSections) MarshalJSON ¶
func (s SystemSections) MarshalJSON() ([]byte, error)
func (*SystemSections) UnmarshalJSON ¶
func (s *SystemSections) UnmarshalJSON(data []byte) error
type ThinkingTokenBudget ¶
type ThinkingTokenBudget struct {
MaxTokens float64 `json:"maxTokens"`
ThinkingBudget float64 `json:"thinkingBudget"`
}
func AdjustMaxTokensForThinking ¶
func AdjustMaxTokensForThinking(base *float64, modelMax float64, level string, custom map[string]json.RawMessage) ThinkingTokenBudget
type TokenSource ¶
type TokenSource interface {
// Acquire returns the next usable account's token. It MUST return an error
// (never an empty AccessToken) when no account can serve — the caller
// surfaces that as the provider error.
Acquire(ctx context.Context) (OAuthToken, error)
// Report records the outcome of a request made with tok. err nil means the
// account served fine (clears nothing, updates last_used_at); a rate-limit
// or auth failure marks the account so the next Acquire skips it.
Report(ctx context.Context, tok OAuthToken, err error)
}
TokenSource is the seam between a pooled (multi-account OAuth) provider and the account store. Acquire picks the next usable account — refreshing its access token when near expiry — and Report feeds the request's outcome back so a rate-limited or dead account is rotated out of the pool.
type Tool ¶
type Tool struct {
Name string `json:"name"`
Description string `json:"description"`
Parameters json.RawMessage `json:"parameters"`
ConstrainedSampling json.RawMessage `json:"constrainedSampling,omitempty"`
Extra map[string]json.RawMessage `json:"-"`
// contains filtered or unexported fields
}
func GetCurrentTools ¶
func GetDeclaredTools ¶
func ToToolDeclaration ¶
func (Tool) MarshalJSON ¶
func (*Tool) UnmarshalJSON ¶
type ToolCallIDNormalizer ¶
type ToolReference ¶
type ToolReference struct {
Name string `json:"name"`
Extra map[string]json.RawMessage `json:"-"`
// contains filtered or unexported fields
}
func (ToolReference) MarshalJSON ¶
func (t ToolReference) MarshalJSON() ([]byte, error)
func (*ToolReference) UnmarshalJSON ¶
func (t *ToolReference) UnmarshalJSON(data []byte) error
type ToolStateChanges ¶
type ToolStateChanges struct {
ToolsAdded []Tool `json:"toolsAdded"`
ToolsRemoved []ToolReference `json:"toolsRemoved"`
}
func GetToolStateChanges ¶
func GetToolStateChanges(previous, current []Tool) ToolStateChanges
type TranscriptContext ¶
type TranscriptContext struct {
// contains filtered or unexported fields
}
TranscriptContext can only be constructed through NormalizeContext. Provider implementations receive the prompt and tools through system messages.
func CollapseSystemMessages ¶
func CollapseSystemMessages(context TranscriptContext) TranscriptContext
func NormalizeContext ¶
func NormalizeContext(context Context) TranscriptContext
func ResolveTranscript ¶
func ResolveTranscript(context TranscriptContext, supportsMidConvoSystemMessages bool) TranscriptContext
func (TranscriptContext) MarshalJSON ¶
func (c TranscriptContext) MarshalJSON() ([]byte, error)
func (TranscriptContext) Messages ¶
func (c TranscriptContext) Messages() []Message
type TranscriptTools ¶
type TranscriptTools struct {
RequestTools []Tool `json:"requestTools"`
AnchorsAdditions bool `json:"anchorsAdditions"`
}
func ResolveTranscriptTools ¶
func ResolveTranscriptTools(messages []Message, supportsToolAdditions bool) TranscriptTools
type UnsupportedStrictJSONSchemaError ¶
type UnsupportedStrictJSONSchemaError struct{ Reason string }
func (*UnsupportedStrictJSONSchemaError) Error ¶
func (e *UnsupportedStrictJSONSchemaError) Error() string
type UnsupportedStrictSchemaKeywordCheck ¶
type UnsupportedStrictSchemaKeywordCheck func(string, json.RawMessage) bool
UnsupportedStrictSchemaKeywordCheck adds a provider-specific strict-schema restriction. The value retains its original JSON, including opaque numbers.
type Usage ¶
type Usage struct {
// Observation is producer evidence for safe settlement only. It is never
// serialized into Pi transcripts/checkpoints and is not inferred on decode.
Observation UsageObservation `json:"-"`
Input float64 `json:"input"`
Output float64 `json:"output"`
CacheRead float64 `json:"cacheRead"`
CacheWrite float64 `json:"cacheWrite"`
CacheWrite1h *float64 `json:"cacheWrite1h,omitempty"`
Reasoning *float64 `json:"reasoning,omitempty"`
TotalTokens float64 `json:"totalTokens"`
Cost UsageCost `json:"cost"`
// contains filtered or unexported fields
}
func (Usage) MarshalJSON ¶
func (*Usage) UnmarshalJSON ¶
type UsageCost ¶
type UsageCost struct {
Input float64 `json:"input"`
Output float64 `json:"output"`
CacheRead float64 `json:"cacheRead"`
CacheWrite float64 `json:"cacheWrite"`
Total float64 `json:"total"`
// contains filtered or unexported fields
}
func (UsageCost) MarshalJSON ¶
func (*UsageCost) UnmarshalJSON ¶
type UsageObservation ¶ added in v0.1.1
type UsageObservation struct {
UsageObserved bool
PricingObserved bool
UsageSource string
PricingSource string
}
UsageObservation distinguishes a wire/explicit observation from initialized zero values. Custom producers must supply explicit evidence; JSON decode and model cost-object presence cannot establish it.
Source Files
¶
- anthropic.go
- anthropic_native_client.go
- anthropic_native_federation.go
- anthropic_native_http.go
- anthropic_native_messages.go
- anthropic_native_params.go
- anthropic_native_pool.go
- anthropic_native_simple.go
- anthropic_native_sse.go
- anthropic_native_stream.go
- anthropic_native_token_cache.go
- antigravity.go
- antigravity_native.go
- auth_context.go
- auth_helpers.go
- auth_lazy.go
- auth_resolve.go
- azure_responses_http.go
- azure_responses_params.go
- azure_responses_simple.go
- capabilities.go
- client.go
- codex.go
- codex_native_auth.go
- codex_native_combined.go
- codex_native_continuation.go
- codex_native_debug.go
- codex_native_errors.go
- codex_native_http.go
- codex_native_params.go
- codex_native_pool.go
- codex_native_pool_error.go
- codex_native_prepare.go
- codex_native_simple.go
- codex_native_socket_cache.go
- codex_native_sse.go
- codex_native_stream.go
- codex_native_transport_policy.go
- codex_native_websocket_conn.go
- codex_native_websocket_request.go
- codex_native_websocket_stream.go
- collection.go
- constrained_sampling.go
- content.go
- content_blocks.go
- context_overflow.go
- credential_store.go
- devin.go
- devin_native.go
- devin_proto.go
- doc.go
- estimate.go
- event_stream.go
- factory.go
- failure_kind.go
- fallback.go
- fallback_attempt.go
- fallback_retry.go
- fallback_selection.go
- httpclient.go
- js_number.go
- json_parse.go
- json_stringify.go
- lazy_stream.go
- list.go
- model_catalog.go
- model_pricing.go
- models_auth.go
- models_availability.go
- models_error.go
- models_login.go
- models_publication.go
- models_refresh.go
- models_registry.go
- models_request.go
- models_store.go
- models_stream.go
- native_api_factory.go
- native_client.go
- native_controls.go
- native_failure.go
- native_oauth_pool.go
- native_provider.go
- native_stream_options.go
- native_usage_observation.go
- oauth.go
- oauth_anthropic.go
- oauth_antigravity.go
- oauth_callback.go
- oauth_chatgpt.go
- oauth_chatgpt_callback.go
- oauth_codex.go
- oauth_copilot.go
- oauth_device_code.go
- oauth_devin.go
- oauth_diagnostics.go
- oauth_form.go
- oauth_http.go
- oauth_input.go
- oauth_kimi.go
- oauth_meta.go
- oauth_openrouter.go
- oauth_page.go
- oauth_pkce.go
- oauth_radius.go
- oauth_sleep.go
- oauth_xai.go
- openai.go
- openai_native_compat.go
- openai_native_errors.go
- openai_native_http.go
- openai_native_messages.go
- openai_native_params.go
- openai_native_platform_unix.go
- openai_native_sse.go
- openai_native_stream.go
- openai_pool.go
- openai_prompt_cache.go
- openai_responses.go
- openai_responses_native_http.go
- openai_responses_native_messages.go
- openai_responses_native_params.go
- openai_responses_native_stream.go
- pi_messages_http.go
- pi_messages_json.go
- pi_messages_protocol.go
- pi_messages_typed.go
- pooled.go
- provider.go
- provider_error.go
- provider_factory.go
- provider_model_catalog.go
- provider_operation_queue.go
- provider_radius.go
- provider_session.go
- radius_config.go
- safe_accounting.go
- safe_failure.go
- safe_physical.go
- sanitize_unicode.go
- scripted.go
- settings_catalog.go
- simple_options.go
- speech.go
- stream_payload.go
- stream_snapshot.go
- telemetry.go
- tool_choice.go
- tool_schema.go
- tool_strict.go
- transcript.go
- transcript_types.go
- transform_messages.go
- url_query.go
- usage_json.go
- values.go
- window.go