Skip to content

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:

  1. The provider speaks the OpenAI Chat Completions format: declare a dialect.
  2. The provider speaks the Anthropic Messages format: declare a Messages dialect.
  3. The provider has its own API: write a wire for it.

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 implements HttpClientExt is 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.

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.

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());

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:

PieceTraitJob
Config and wire structs—Plain data: base URL, key, model
Replay targetReplayTargetWhat the model can read, and how options map to the request
DecoderDecoderTurns each reply frame into text, tool calls and a finish
ReassemblerReassembleRebuilds the provider’s own reply document from a stream, for raw
The wireWireTies them together and encodes the HTTP request

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,
}

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.

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.

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 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.

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())
}
}

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());

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");

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.

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