Write your own provider
Rig ships many providers, and you can add your own outside the Rig repository. This guide covers the three ways to do it, from least to most work:
- The provider speaks the OpenAI Chat Completions format: declare a dialect.
- The provider speaks the Anthropic Messages format: declare a Messages dialect.
- The provider has its own API: write a wire for it.
How a Rig model is built
Section titled “How a Rig model is built”Every model in Rig is a rig::driver::Model<W, T>, which pairs two things:
- A wire (
W: rig::wire::Wire): plain data that describes one endpoint. It encodes a request into an HTTP request and gives a decoder for the reply. It never touches the network. - A transport (
T: rig::driver::Transport<W>): how the request travels. Any HTTP client that implementsHttpClientExtis a transport, so most wires never name one. The default is the shared reqwest client.
The driver sits between them and owns everything else: sending, splitting the reply into
frames (SSE, NDJSON, or one whole body), status checks, telemetry, and folding the
decoded events into a CompletionResponse. Single replies (model.call) and streamed replies
(model.stream) go through the same decoder, so a provider written once works both
ways. It also works as an agent model (AgentBuilder::new(model)), behind a DynModel,
and with every hook and tool the agent runtime has.
A provider client such as OpenAI is just a convenience: a config plus a transport, with
methods like completion(model) that return Models.
OpenAI-compatible providers
Section titled “OpenAI-compatible providers”If the provider accepts OpenAI’s Chat Completions requests, you don’t need a wire of your
own. Declare a Dialect with the provider’s name, default base URL and API-key variable,
and adjust the Quirks for anything it does differently:
use rig::providers::openai::OpenAIConfig;use rig::providers::openai::wire::{Dialect, Quirks};
/// Acme's gateway speaks Chat Completions but rejects `stream_options` and/// has no structured output.const ACME: Dialect = Dialect::gateway("acme", "https://api.acme.test/v1", "ACME_API_KEY") .with_quirks(Quirks::openai().without_stream_usage().without_response_format());
let acme = OpenAIConfig::from_env_with(&ACME)?.client();let model = acme.completion("acme-large");
let agent = AgentBuilder::new(model).preamble("You are terse.").build();let answer = agent.prompt("Say hello.").await?.output();Quirks has more switches, such as done_without_finish_reason() for streams that end with
[DONE] but no finish reason. For a one-off endpoint, OpenAIConfig::new(key).with_base_url(url)
also works.
Anthropic Messages-compatible providers
Section titled “Anthropic Messages-compatible providers”Providers that accept Anthropic’s Messages format get the same treatment through
anthropic::compatible:
use rig::providers::anthropic::{self, AnthropicConfig, Dialect};
const ACME_MESSAGES: Dialect = anthropic::compatible( "acme", "https://api.acme.test/anthropic", "ACME_API_KEY", Some("ACME_BASE_URL"),);
let model = AnthropicConfig::from_env_with(&ACME_MESSAGES)? .client() .completion("acme-large");let response = model.call("Say hello.").await?;println!("{}", response.text());A provider with its own API
Section titled “A provider with its own API”The rest of this guide builds a completion provider for a made-up API, Acme. It takes:
POST {base_url}/v1/generate{"model": "acme-large", "messages": [{"role": "user", "content": "Hi"}], "stream": false}A single reply is one JSON document. A streamed reply is newline-delimited JSON (NDJSON) in
the same shape, with each line carrying a piece of text and the last one saying done:
{"id": "gen_1", "model": "acme-large", "text": "Hello!", "done": true, "usage": {"input_tokens": 5, "output_tokens": 2}}You write five pieces:
| Piece | Trait | Job |
|---|---|---|
| Config and wire structs | — | Plain data: base URL, key, model |
| Replay target | ReplayTarget | What the model can read, and how options map to the request |
| Decoder | Decoder | Turns each reply frame into text, tool calls and a finish |
| Reassembler | Reassemble | Rebuilds the provider’s own reply document from a stream, for raw |
| The wire | Wire | Ties them together and encodes the HTTP request |
Config and wire
Section titled “Config and wire”A wire is plain data: Clone, and usually Debug, PartialEq, Serialize and
Deserialize, so it can be saved and loaded like any config. Keep the credential in a
rig::wire::Secret. It prints and serializes as [redacted], so a key never ends up in
logs.
use rig::wire::Secret;
/// Where Acme lives and the key it takes.#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]pub struct AcmeConfig { pub api_key: Secret, pub base_url: String,}
impl AcmeConfig { pub fn new(api_key: impl Into<Secret>) -> Self { Self { api_key: api_key.into(), base_url: "https://api.acme.test".to_owned(), } }
pub fn with_base_url(mut self, base_url: impl Into<String>) -> Self { self.base_url = base_url.into(); self }}
/// Acme's generate endpoint, for one model.#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]pub struct AcmeChat { pub config: AcmeConfig, pub model: String,}Replay target
Section titled “Replay target”Every completion wire implements ReplayTarget. Before your encoder runs, Rig uses it to
shape the conversation for this model. Content the model can’t read is downgraded: images
become placeholders, and tool calls and results become text when tools is off. That way,
a history written by any other provider can be sent to yours, and your encoder only sees
what you said you accept.
map_options says how each portable option on the request (top_p, seed, stop,
reasoning, …) is sent. Destructure fields without .., so that a new option added to
Rig fails to compile until you handle it. Mapping::Send merges JSON into the request
body. Mapping::unsupported refuses the option before anything is sent.
use rig::completion::options::{Mapping, OptionFields, OptionMap};use rig::completion::{Accepts, CompletionRequest, Media, ReplayTarget};use rig::message::Api;
impl ReplayTarget for AcmeChat { fn api(&self) -> Api { Api::from_static("acme.generate") }
fn provider(&self) -> &str { "acme" }
fn model(&self) -> &str { &self.model }
/// Text only: no images, and no tool calling. fn accepts(&self, _model: &str) -> Accepts { Accepts { tools: false, ..Accepts::TEXT } }
/// No media at all: text documents are sent as their text. fn encodes(&self, _model: &str, _media: Media<'_>) -> bool { false }
fn map_options(&self, _request: &CompletionRequest, fields: OptionFields<'_>) -> OptionMap { let OptionFields { reasoning, cache, service_tier, verbosity, parallel_tool_calls, top_p, seed, stop, } = fields; const NO_FIELD: &str = "Acme has no such field"; OptionMap { reasoning: Mapping::of(reasoning, |_| Mapping::unsupported(NO_FIELD)), cache: Mapping::of(cache, |_| Mapping::unsupported(NO_FIELD)), service_tier: Mapping::of(service_tier, |_| Mapping::unsupported(NO_FIELD)), verbosity: Mapping::of(verbosity, |_| Mapping::unsupported(NO_FIELD)), parallel_tool_calls: Mapping::of(parallel_tool_calls, |_| { Mapping::unsupported(NO_FIELD) }), top_p: Mapping::of(top_p, |top_p| Mapping::Send(json!({ "top_p": top_p }))), seed: Mapping::of(seed, |seed| Mapping::Send(json!({ "seed": seed }))), stop: Mapping::of_stop(stop, |stop| Mapping::Send(json!({ "stop": stop }))), } }}ReplayTarget has more methods with defaults, for providers with stricter rules: for
example, starts_with_user, alternates_roles and normalize_tool_call_id. Override the
ones your API needs.
Decoder
Section titled “Decoder”A decoder reads one reply. classify parses a frame into your typed event. For JSON frames,
use one of the shared classifiers in rig::providers::internal::wire, so that bad JSON is
reported as a corrupt reply and not skipped silently. decode writes the event into the
reply through Out:
out.run(Block::Text, text)appends text. Consecutive runs of the same kind join into one block.out.end(Finish { .. })ends the reply with usage, the finish reason, the response id and the model. A reply that runs out of frames without an end is reported as truncated.
use rig::completion::{FinishReason, Usage};use rig::operation::{Block, Completion, Finish};use rig::providers::internal::wire::classify_untyped_line;use rig::wire::{Decoder, Flow, Out, WireEvent, WireFrame};use rig::ProviderError;
/// One Acme reply document, or one line of a streamed reply.#[derive(Debug, Deserialize)]pub struct Chunk { #[serde(default)] pub id: Option<String>, #[serde(default)] pub model: Option<String>, #[serde(default)] pub text: String, #[serde(default)] pub done: bool, #[serde(default)] pub usage: Option<AcmeUsage>,}
#[derive(Debug, Deserialize)]pub struct AcmeUsage { pub input_tokens: u64, pub output_tokens: u64,}
#[derive(Debug, Default)]pub struct AcmeDecoder;
impl<'id> Decoder<'id, Completion> for AcmeDecoder { type Event = Chunk;
fn classify(&self, frame: WireFrame) -> WireEvent<Chunk> { classify_untyped_line(frame.as_str().as_bytes()) }
fn decode( &mut self, chunk: Chunk, mut out: Out<'id, Completion>, ) -> Result<Flow, ProviderError> { if !chunk.text.is_empty() { out.run(Block::Text, &chunk.text)?; } if !chunk.done { return Ok(Flow::More); } out.end_run()?; let usage = chunk.usage.map_or_else(Usage::new, |usage| { Usage::new() .input_tokens(usage.input_tokens) .output_tokens(usage.output_tokens) }); Ok(out.end(Finish { usage, reason: Some(FinishReason::Stop), response_id: chunk.id, model: chunk.model, ..Finish::default() })) }}The decoder never knows whether the reply was streamed. With a single reply, the whole body is one frame. With a stream, each NDJSON line is a frame.
Reassembler
Section titled “Reassembler”CompletionResponse::raw holds the provider’s own reply document. For a single reply,
Rig uses the body. For a stream, there is no body, so the wire’s reassembler sees every
frame and rebuilds the document. Here, that means joining the text and keeping the last
line’s other fields:
use rig::operation::Completion;use rig::wire::document::{Reassemble, Serves};use rig::wire::WireFrame;use serde_json::{Map, Value};
#[derive(Default)]pub struct AcmeDocument { text: String, last: Map<String, Value>,}
impl Reassemble<WireFrame> for AcmeDocument { fn absorb(&mut self, frame: &WireFrame) { let Ok(Value::Object(chunk)) = serde_json::from_str(&frame.as_str()) else { return; }; if let Some(text) = chunk.get("text").and_then(Value::as_str) { self.text.push_str(text); } self.last = chunk; }
fn finish(self) -> Value { if self.last.is_empty() { return Value::Null; } let mut document = self.last; document.insert("text".to_owned(), Value::String(self.text)); Value::Object(document) }}
impl Serves<Completion> for AcmeDocument {}The wire
Section titled “The wire”The Wire impl names the operation (Completion), the payload (Encoded, an HTTP request
plus how to split its reply), the frame type, the decoder and the reassembler. describe
gives the provider name and model for telemetry and points at the replay target.
encode must not do any I/O. Build your own body in a closure and pass it to
options::request_params. It adds the mapped options, any typed provider options and the
caller’s additional_params on top, in the same order every built-in provider uses.
use rig::completion::options::{request_params, RawAt};use rig::completion::CompletionRequest;use rig::error::EncodeError;use rig::http_client::{HeaderValue, Request};use rig::message::{AssistantContent, Message, UserContent};use rig::operation::Completion;use rig::wire::{Descriptor, Encoded, Framing, Mode, Wire, WireFrame};use serde_json::{Map, Value};
impl AcmeChat { /// Acme's own fields: model, messages, sampling and the stream flag. fn base(&self, request: &CompletionRequest, mode: Mode) -> Result<Map<String, Value>, EncodeError> { if !request.tools.is_empty() { return Err(EncodeError::request("Acme does not support tools")); } let mut messages = Vec::new(); for message in &request.chat_history { let (role, content) = match message { Message::System { content } => ("system", content.clone()), Message::User { content } => ("user", user_text(content)?), Message::Assistant(turn) => ("assistant", assistant_text(&turn.content)), }; messages.push(json!({ "role": role, "content": content })); } let mut body = Map::new(); // `prepare` has already resolved the model: the request's override, or ours. let model = request.model.clone().unwrap_or_else(|| self.model.clone()); body.insert("model".into(), json!(model)); body.insert("messages".into(), json!(messages)); body.insert("stream".into(), json!(mode == Mode::Streaming)); if let Some(temperature) = request.temperature { body.insert("temperature".into(), json!(temperature)); } if let Some(max_tokens) = request.max_tokens { body.insert("max_tokens".into(), json!(max_tokens)); } Ok(body) }}
/// The text of a user message. `accepts` and `encodes` mean the adapter/// has already replaced anything else.fn user_text(content: &[UserContent]) -> Result<String, EncodeError> { let mut texts = Vec::new(); for part in content { match part { UserContent::Text(text) => texts.push(text.text()), _ => return Err(EncodeError::request("Acme takes text only")), } } Ok(texts.join("\n"))}
fn assistant_text(content: &[AssistantContent]) -> String { content .iter() .filter_map(|block| match block { AssistantContent::Text(text) => Some(text.text()), _ => None, }) .collect()}
impl Wire for AcmeChat { type Op = Completion; type Payload = Encoded; type Frame = WireFrame; type Decoder<'id> = AcmeDecoder; type Reassembler = AcmeDocument;
fn describe(&self) -> Descriptor<'_> { Descriptor::new("acme").model(self.model.as_str()).replay(self) }
fn encode(&self, request: CompletionRequest, mode: Mode) -> Result<Encoded, EncodeError> { let body = request_params(self, &request, |_| self.base(&request, mode), RawAt::Top, &[])?; let mut authorization = HeaderValue::from_str(&format!("Bearer {}", self.config.api_key.expose())) .map_err(EncodeError::request)?; authorization.set_sensitive(true); let http = Request::post(format!("{}/v1/generate", self.config.base_url)) .header("content-type", "application/json") .header("authorization", authorization) .body(body.into_body())?; let framing = match mode { Mode::Unary => Framing::Whole, Mode::Streaming => Framing::Ndjson, }; Ok(Encoded::new(http, framing).with_route(Some("/v1/generate"))) }
fn decoder<'id>(&self) -> Self::Decoder<'id> { AcmeDecoder }}Use Framing::Sse for a server-sent-events stream. If the provider returns a request id in
a header, name it with Encoded::with_request_id_header, and Rig attaches it to responses
and errors.
A client
Section titled “A client”The client is ordinary Rust. Follow the shape of the built-in clients: a config, a
transport, and one method per kind of model. DynHttpClient holds any HttpClientExt,
and rig::rig_reqwest::shared() is the process-wide reqwest client every built-in
provider uses by default:
use rig::driver::Model;use rig::http_client::{DynHttpClient, HttpClientExt};use rig::wire::Secret;
#[derive(Clone)]pub struct Acme { config: AcmeConfig, http: DynHttpClient,}
impl Acme { pub fn new(api_key: impl Into<Secret>) -> Self { Self { config: AcmeConfig::new(api_key), http: rig::rig_reqwest::shared(), } }
pub fn from_env() -> Result<Self, std::env::VarError> { Ok(Self::new(std::env::var("ACME_API_KEY")?)) }
/// Send through another HTTP client: a configured reqwest client, /// middleware, or your own `HttpClientExt`. pub fn with_http(mut self, http: impl HttpClientExt + 'static) -> Self { self.http = DynHttpClient::new(http); self }
pub fn completion(&self, model: impl Into<String>) -> Model<AcmeChat> { let wire = AcmeChat { config: self.config.clone(), model: model.into(), }; Model::new(wire, self.http.clone()) }}Using it
Section titled “Using it”The model works anywhere a built-in one does: direct calls, streams, and agents.
use futures::StreamExt;use rig::streaming::{Item, StreamEvent};
let model = Acme::from_env()?.completion("acme-large");
// One request, no agent loop.let response = model.call("Say hello.").await?;println!("{} ({:?} output tokens)", response.text(), response.usage.output_tokens);
// The same request, streamed through the same decoder.let mut stream = model.stream("Count to five.")?;while let Some(item) = stream.next().await { if let Item::Event(StreamEvent::Text { text, .. }) = item? { print!("{text}"); }}
// As an agent model.let agent = AgentBuilder::new(model).preamble("You are terse.").build();println!("{}", agent.prompt("What is Rust?").await?.output());Testing without a network
Section titled “Testing without a network”Because the transport is a separate piece, you can test a wire by putting it on a
transport that replays a recorded body. Framing::split cuts the body into frames the same
way the HTTP transport does. This test runs the real encoder, decoder, fold and reassembler:
use futures::stream;use rig::driver::{Exchange, Model, Opened, Opening, Transport};use rig::wire::{Encoded, Wire, WireFrame};
/// Replays one recorded reply body for any HTTP wire.#[derive(Clone)]struct Replay(&'static str);
impl<W: Wire<Payload = Encoded, Frame = WireFrame>> Transport<W> for Replay { fn send(&self, payload: Encoded, _exchange: Exchange) -> Opening<WireFrame> { let frames = payload.framing.split(self.0.as_bytes()); Opening::ready(Opened::new(stream::iter(frames.into_iter().map(Ok)))) }}
let wire = AcmeChat { config: AcmeConfig::new("test-key"), model: "acme-large".to_owned(),};let streamed = concat!( r#"{"text":"Hel"}"#, "\n", r#"{"text":"lo"}"#, "\n", r#"{"id":"gen_1","done":true,"usage":{"input_tokens":3,"output_tokens":2}}"#, "\n",);let model = Model::new(wire, Replay(streamed));
let response = model.stream("hi")?.finish().await?;assert_eq!(response.text(), "Hello");assert_eq!(response.usage.output_tokens, Some(2));assert_eq!(response.raw["text"], "Hello");Tool calls
Section titled “Tool calls”To support tools, set tools: true in accepts. Then send request.tools (each has a
name, description and JSON-schema parameters) and encode AssistantContent::ToolCall
and UserContent::ToolResult in the history. On the decoding side, write each tool-call
delta with out.fragment. The writer gathers fragments by the provider’s index, accepts an
id or name that arrives late, and emits streaming events as the arguments grow. Before you
end the reply, close every open call:
use rig::operation::{CallFragment, Completion};use rig::wire::Out;use rig::ProviderError;
/// One streamed tool-call delta, as an OpenAI-style API sends it.fn write_call_delta( out: &mut Out<'_, Completion>, index: usize, id: Option<&str>, name: Option<&str>, arguments: &str,) -> Result<(), ProviderError> { out.fragment(Some(index), CallFragment { id, name, arguments: Some(arguments) })}
/// The provider's end of the reply: close the open calls, then `out.end(..)`.fn close_calls(out: &mut Out<'_, Completion>) -> Result<(), ProviderError> { out.finish_open()}If a call arrives whole, use out.whole(index, Block::Call { .. }, item, arguments), and
Block::Reasoning works the same way for reasoning text. The built-in providers are the best
reference for real APIs. Ollama’s native chat wire (rig::providers::ollama::Chat) is a
compact, complete example with tools, images and NDJSON streaming.
Other operations
Section titled “Other operations”Completion is one operation. Embeddings, reranking, transcription, image and audio
generation (rig::operation::{Embedding, Rerank, Transcription, ImageGeneration, AudioGeneration}) work the same way: a wire whose Op is that operation, with a decoder
for its reply. Look at the matching built-in wire, such as openai::wire::Embeddings.
For an endpoint that fits none of them, define your own rig::wire::Operation with its
request, response and fold. When the reply is one JSON document, the ready-made Json
decoder and Whole fold do the work. Rig’s
frozen_api examples
show this (p6_http_operation.rs), along with a custom OpenAI dialect
(p7_openai_compatible_vendor.rs) and a replay transport (p8_custom_transport.rs).
See also
Section titled “See also”- Providers & Clients: how clients, models and transports fit together.
- Model Providers: the built-in providers.
rig::wireon docs.rs: the fullWire,DecoderandOperationcontracts.
