Developer-facing guide to building, testing, and extending busbar. For the operator/runtime view see operations.md and configuration.md; for the public request-lifecycle overview see architecture.md; for the design deep-dive see internals.md and the ADRs.
Contribution mechanics (PR checklist, formatting, the exhaustive-match invariant) live in CONTRIBUTING.md, this doc covers the codebase map and the two common extension tasks.
| Module | Owns |
|---|---|
main.rs |
Startup. Loads providers.yaml + config.yaml (with ${ENV} interpolation), resolve()s them, validates, builds lanes/pools/App, wires governance + observability + plugins, spawns health probers, and builds both listeners (the data-plane axum router and the separate admin router). |
config/mod.rs |
The deploy/provider/pool schema (DeployCfg, ProviderDef, ProviderDeploy, ModelCfg, PoolCfg, PoolMember, FailoverCfg, AffinityCfg, BreakerCfg, HealthCfg, GovernanceCfg, ObservabilityCfg, PluginsCfg, OnExhausted), ${ENV} interpolation, and resolve() (merge catalog def + deployment override). config/overlay.rs handles the live config-apply overlay. |
config_validate/ |
Post-resolve config validation (fail-loud diagnostics before lanes are built), including the SSRF host guard shared with the webhook path. |
state.rs |
Runtime types: Lane, WeightedLane, PoolRuntime, and the App shared state. Re-exports StateStore from store/. |
ingress/mod.rs |
named (POST /{name}/v1/messages) and adhoc (POST /{provider}/{model}/v1/messages) are the only two protocol-specific axum handlers still registered directly; every other protocol (openai chat, cohere, responses, gemini, bedrock converse) is served by the single fallback route ingress::protocol_dispatch, which identifies the protocol (proto::detect::protocol_id, mostly path-based with a small header-first exception, see docs/protocols.md) and hands off to ingress::dispatch::operation_ingress for the actual body/path-model resolution, governance pre-checks, and affinity-header resolution. ingress/dispatch.rs carries both protocol_dispatch and operation_ingress. |
auth/mod.rs |
AuthMiddleware and the auth_middleware layer: it runs the data-plane auth.chain (an ordered list of AuthModules; an empty chain is the open front door), opens /healthz, gates /metrics like any other route, resolves virtual keys, and threads the caller token. UpstreamCreds (Own / Passthrough) is the separate egress-credential mode. The old AuthMode enum is gone. Constant-time token compare lives in the busbar-api auth contract. |
proxy/engine/mod.rs |
The forwarding engine: forward / forward_with_pool (selection → translate → sign → POST → classify → stream/failover), RequestCtx (deadline + exclusions + visited-pools), the before-first-byte failover boundary + cross-protocol stream wiring, lane_auth_headers (the api-key auth-adapter seam), and the on_exhausted handlers (Status503/FallbackPool/LeastBad/Queue). proxy/select.rs, proxy/egress.rs, proxy/hooks.rs, and proxy/usage.rs split out selection, egress, the hook seam, and usage metering. |
breaker.rs |
The protocol-agnostic Stage 1b/2 classifier: StatusClass, Disposition, RawUpstreamError, CanonicalSignal, normalize_raw_error, classify (exhaustive). |
store/mod.rs |
The breaker FSM + lane state: StateStore trait, LaneState, BreakerCell / BreakerCellAccess, OutcomeWindow, SWRR select_weighted, the lane-default vs _in(pool, …) method split, BreakerCfg/TripConfig, test time injection. The concrete InMemoryStore is in store/in_memory.rs. (This is the runtime breaker store, distinct from the governance Store trait below.) |
ir/mod.rs |
The superset IR (ADR-0005): IrRequest, IrResponse, IrMessage, IrBlock, IrTool, IrUsage, IrStreamEvent, IrDelta, StreamDecodeState. Modality-specific IR (audio, image, embeddings, moderation, rerank) sits in sibling files under ir/. |
proto/mod.rs |
The protocol seam: ProtocolReader / ProtocolWriter traits, Protocol, ProtocolRegistry, SigningContext, probe_body default. proto/detect.rs sniffs the ingress protocol; proto/openai_family.rs holds the shared OpenAI-family bits; proto/stream.rs is the cross-protocol stream translator and SSE reframing. |
proto/{anthropic,openai_chat,openai_responses,gemini,bedrock,cohere}/ |
One folder-module per protocol: each holds the Reader (wire→IR + error extraction) and Writer (IR→wire + auth + paths). Bedrock's writer overrides sign_request for SigV4. |
sigv4.rs |
Hand-rolled AWS SigV4 (RustCrypto sha2 + hmac, no AWS SDK): sign_v4, signing_key, uri_encode_path, format_amz_time, sha256_hex. |
governance/mod.rs |
Signed virtual keys + the generic group limit engine: GovState, VirtualKey, try_admit over the per-(group, window) buckets, the revocation denylist, and the token-ledger cost model. The governance Store trait itself lives in the busbar-api crate (crates/api/src/store.rs); concrete backends are separate crates (busbar-store-memory compiled in by default, busbar-store-sqlite / -postgres / -redis as static or dynamically-loaded plugins chosen by store.module). |
admin/ |
The admin API: admin/mod.rs mounts the /api/v1/admin/* handlers (keys, usage, config, hooks, plugins) on the separate admin listener, admin/v1/ is the frozen JSON contract, and admin/rate.rs / admin/audit.rs carry admin rate-limiting and the hash-chained audit log. |
health.rs |
Active health probing (spawn_probers, probe_lane using each protocol's probe_body) and the /stats + /healthz handlers. |
metrics.rs |
Prometheus recorder init + the busbar_* metric name constants. |
observability.rs |
Optional OTLP tracer init + the fire-and-forget request-log webhook (with its own SSRF guard). |
eventstream.rs |
Codec for Bedrock's binary application/vnd.amazon.eventstream frames: drain_frames decodes ConverseStream responses; encode_frame/encode_exception_frame re-encode CRC32-valid frames for Bedrock-ingress streaming. |
test_support/ |
#[cfg(test)] in-crate mock-upstream harness (MockServer, MockServerState, MockResponse). Each module also carries its own #[cfg(test)] mod tests. See testing.md. |
Single Rust binary, stable toolchain, edition 2021.
cargo build # debug build
cargo build --release # release binary -> target/release/busbar
cargo test # full in-crate suite
cargo clippy --all-targets -- -D warnings # lints must be clean (treat warnings as errors)
cargo fmt --all # format (rustfmt.toml in repo)The test suite is in-crate: a shared
#[cfg(test)] mod test_support provides the MockServer harness, and each module
carries its own #[cfg(test)] mod tests. There are no tests/ integration
binaries: everything runs under cargo test. See testing.md.
Busbar reads two YAML files, located via env vars:
| Env var | Default | Purpose |
|---|---|---|
BUSBAR_PROVIDERS |
/etc/busbar/providers.yaml |
The verified provider catalog (shipped). |
BUSBAR_CONFIG |
/etc/busbar/config.yaml |
Your deployment. |
Both files support ${VAR} interpolation expanded at load time; an unset
referenced variable is a hard startup failure. Provider keys are supplied via the
env vars (or files/secret plugins) named by each provider's api_key secret reference, never written into the files.
export BUSBAR_CLIENT_TOKEN=dev-token
export ANTHROPIC_KEY=sk-ant-...
BUSBAR_PROVIDERS=./providers.yaml BUSBAR_CONFIG=./config.yaml cargo run
curl -s localhost:8080/healthz
curl -s -H "Authorization: Bearer $BUSBAR_CLIENT_TOKEN" localhost:8080/stats | jqFull field reference: configuration.md.
A protocol is the unit of Busbar's scope (the count to grow is 6, not the provider count). To add one:
- Implement
ProtocolReader(crates/busbar/src/proto/mod.rsdefines the trait):read_request(body) -> IrRequest: wire JSON → IR (ADR-0005 contract: model every field you can; stash adjacent fields inIrRequest.extra; holdtemperatureas the f64 it already is).read_response(body) -> IrResponseandread_response_event(s), wire → IR. For a flat stream, use the&mut StreamDecodeStateto synthesize the IR's block boundaries (one chunk →0..nevents); for a 1:1 stream, ignore it.extract_error(status, body) -> RawUpstreamError, Stage 1a: pull out the HTTP status and any in-bodyprovider_code.classify, the simple two-stage convenience wrapper.clone_box.
- Implement
ProtocolWriter:write_request(ir) -> Value,write_response(ir),write_response_event(ir): IR → wire.rewrite_model(body, model), set the selected lane's model on the body.upstream_path(+ optionallyupstream_path_for/upstream_path_for_streamif the path embeds the model or differs for streaming, as Gemini's does).auth_headers(key)for static headers; overridesign_request(key, ctx)only if the protocol signs the whole request (as Bedrock does for SigV4).- You get
probe_bodyfor free from the default impl: it serializes a one-token IR request through your ownwrite_request, so active health probing works with no extra code. clone_box.
- Register it in
crates/busbar/src/proto/mod.rs: add aProtocol::<name>()constructor, aprotocol_forarm, and an entry inProtocolRegistry::with_builtins. Add theStreamTranslate::newflags if it has a non-SSE wire (like Bedrock's binary eventstream) or a special terminator. - IR contract: the IR is a superset. If your protocol introduces a content
kind the IR can't represent, extend the
IrBlock/ event enums: and then every other writer must handle the new variant (the exhaustive matches will tell you). - Test it through the
MockServerharness and the cross-protocol round-trip tests incrates/busbar/src/proto/tests/tests.rs(test_probe_body_valid_for_all_protocolsalready asserts every protocol produces a valid probe body).
The Reader/Writer files (crates/busbar/src/proto/<name>/) are the only per-protocol code;
the registry + IR + forward path are protocol-agnostic.
A provider is just a catalog entry: no code. Add it to providers.yaml:
my-provider:
protocol: openai # one of the 6 implemented protocols
base_url: https://api.example.com
error_map: # optional: map vendor codes -> StatusClass (Stage 1b)
"insufficient_quota": billing
path: /chat/completions # optional: override the protocol's default path
auth: api-key # optional: 'bearer' (default) | 'api-key'
health: # optional: active probing
mode: dead # none | dead | active
interval_secs: 30
timeout_secs: 5Then reference it from config.yaml (supplying only the env var that holds the
key) and point a model at it:
providers:
my-provider:
api_key: { env: MY_PROVIDER_KEY }
models:
my-model:
provider: my-provider
max_concurrent: 20Notes on the seams:
error_mapis the data-driven Stage 1b override (see internals.md). Keys are the provider's in-body codes; values areStatusClassstrings (billing,rate_limit,auth,server_error,timeout,network,overloaded,context_length,client_error). The deployment'serror_mapinconfig.yamlmerges over the catalog's.pathoverrides the protocol's default upstream path verbatim: used by OpenAI-compatible providers that embed the API version inbase_urland serve/chat/completions(no/v1), and by Azure (which carries?api-version=and the deployment in the path).auth: api-keyis the auth-adapter seam (lane_auth_headersinproxy/engine/mod.rs): it sends anapi-key: <key>header instead of the protocol's native auth (used by Azure OpenAI). For genuinely new auth shapes (e.g. an OAuth2 token mint), the seam to extend isProtocolWriter::sign_request, the same hook Bedrock uses for SigV4: see the roadmap in roadmap.md.
resolve() (crates/busbar/src/config/mod.rs) merges the deployment over the catalog def; a
config.yaml provider name not present in providers.yaml is a fail-loud startup
error.
These are conventions visible in the code; treat the CONTRIBUTING.md checklist as authoritative.
-
SPDX header. Every
src/**/*.rsfile (including eachproto/<name>/module) starts with// SPDX-License-Identifier: Apache-2.0+// Copyright (C) 2026 Busbar Inc and contributors. -
No
_ =>catch-all in the disposition/breaker matches. The exhaustive match onStatusClass/Dispositionis how the compiler enforces that every failure mode is handled; the arms even useunreachable!()for classes that cannot reach a given arm. This is a stated project invariant (CONTRIBUTING.md, "Fixing a defect: the remediation contract"). -
error_mapis data, not code. Provider quirks belong in YAML, not in a match arm. -
Test time is injectable, not real. Breaker/FSM logic reads time via
store::now()(the public crate function), whichInMemoryStoreinternally wraps in a privatenow_secs()that, under#[cfg(test)], is shadowed to delegate tonow_for_test(); tests inject time viastore::set_now_for_test. Don't callSystemTime::now()directly in breaker-adjacent code. -
#[cfg_attr(not(test), allow(dead_code))]marks the lane-default breaker methods that release code reaches only via the_invariants but tests exercise directly: keep that pattern when adding parallel default/_inmethods. -
No
memchrdependency. Byte scanning (e.g. the SSE frame splitting and translation-body boundary scans) is done with plain slice iteration, not thememchrcrate. Keep it that way, don't addmemchr(or pull it in transitively for scanning) when a small hand-rolled scan will do.