Expand description
Put a queue, deduplication and a circuit breaker between your app and an LLM provider, so repeated prompts cost one call and an outage fails fast instead of piling up.

§Quick start
[dependencies]
tokio-prompt-orchestrator = "1.4"
tokio = { version = "1", features = ["rt-multi-thread", "macros"] }Send one prompt through the pipeline with the offline EchoWorker (no API
key), read the answer, and check the dead-letter queue for anything dropped:
use std::{collections::HashMap, sync::Arc};
use tokio_prompt_orchestrator::{spawn_pipeline, EchoWorker, ModelWorker, PromptRequest, SessionId};
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
// Swap EchoWorker for AnthropicWorker, OpenAiWorker, LlamaCppWorker or VllmWorker.
let worker: Arc<dyn ModelWorker> = Arc::new(EchoWorker::new());
let handles = spawn_pipeline(worker);
let mut output = handles.take_output_rx().await.ok_or("output already taken")?;
handles.input_tx.send(PromptRequest {
session: SessionId::new("demo"),
request_id: "req-1".into(),
input: "Hello, pipeline!".into(),
meta: HashMap::new(),
deadline: None,
}).await?;
let answer = output.recv().await.ok_or("pipeline closed")?;
assert!(answer.text.contains("Hello, pipeline!"));
println!("{}", answer.text);
for dropped in handles.dlq.drain() {
println!("dropped {}: {}", dropped.request_id, dropped.reason);
}
Ok(())
}For a fuller example with deduplication and a simulated provider outage, run
cargo run --example llm_pipeline in the repository.
§Main types
spawn_pipelineandspawn_pipeline_with_configstart the five stages and returnPipelineHandles.PromptRequestgoes in throughPipelineHandles::input_tx;PostOutputcomes out ofPipelineHandles::take_output_rx.ModelWorkeris the trait for a backend:AnthropicWorker,OpenAiWorker,LlamaCppWorker,VllmWorker,EchoWorker.enhanced::Deduplicatorandenhanced::CircuitBreakerare the resilience building blocks to wrap around a worker.DeadLetterQueueholds every request that was dropped, with the reason.PipelineConfigloads the same pipeline from TOML.
§How it works
PromptRequest -> RAG(512) -> Assemble(512) -> Inference(1024) -> Post(512) -> Stream(256)Each stage runs as its own Tokio task, joined by bounded channels. When a
downstream channel is full, send_with_shed drops the item into the
DeadLetterQueue instead of blocking the stage above it. Stage 3 checks the
request deadline, calls the worker through a shared circuit breaker and
enforces a timeout.
§Binaries and features
The orchestrator binary (prebuilt on the
releases page,
or cargo binstall tokio-prompt-orchestrator) serves the pipeline over a
terminal prompt and an HTTP API: orchestrator --provider echo runs with no key.
No features are on by default. web-api adds the REST, SSE and WebSocket
server, full adds metrics, Redis caching and rate limiting, tui a terminal
dashboard, mcp an MCP server, and self-improving a self-tuning control loop.
§Modules
| Module | Description |
|---|---|
stages | Five pipeline stage implementations and channel wiring |
worker | ModelWorker trait and the provider implementations |
enhanced | Resilience primitives: circuit breaker, dedup, semantic dedup (SimHash), retry, cache, rate limiter, smart batching |
config | TOML-deserialisable PipelineConfig with hot-reload support |
metrics | Prometheus metrics initialisation and helper functions |
routing | ModelRouter for complexity-scored routing; ArbitrageEngine for SLA-aware provider selection; PoolSizer for worker scaling |
security | PromptGuard: prompt injection and jailbreak detection (no external I/O) |
session | Multi-turn conversation context: injects history per session |
cascade | Multi-turn cascading inference: tool call loops with pluggable executors |
multi_pipeline | Named pipeline fleet with prompt classification and per-class routing |
ab_test | Prompt A/B testing: consistent hashing assignment, Welch’s t-test, Cohen’s d |
coordination | Agent fleet management and task claiming |
adaptive_pool | Kalman-filter worker pool controller: predicts queue depth and recommends scale events |
web_api | REST/SSE/WebSocket server (feature: web-api) |
distributed | Redis dedup and NATS coordination (feature: distributed) |
tui | Ratatui terminal dashboard (feature: tui) |
self_tune, self_modify, intelligence, evolution, self_improve | Self-improving control loop (features of the same names) |
Re-exports§
pub use cache::CacheConfig;pub use cache::CacheEntry;pub use cache::CacheStats;pub use cache::PromptCache;pub use rate_limiter::RateLimitError;pub use rate_limiter::RateLimiterConfig;pub use rate_limiter::RateLimiterRegistry;pub use rate_limiter::ModelRateLimiter;pub use ab_test::AbTestConfig;pub use ab_test::AbTestResult;pub use ab_test::AbTestRunner;pub use ab_test::SuccessMetric;pub use ab_test::Variant;pub use conversation::ConversationConfig;pub use conversation::ConversationManager;pub use conversation::PromptFormat;pub use conversation::Role;pub use conversation::Turn;pub use stages::spawn_pipeline;pub use stages::spawn_pipeline_with_config;pub use stages::LogSink;pub use stages::OutputSink;pub use stages::PipelineHandles;pub use stages::SinkError;pub use templates::AbExperiment;pub use templates::ExperimentReport;pub use templates::ExperimentVariant;pub use templates::PromptTemplate;pub use templates::TemplateError;pub use templates::TemplateRegistry;pub use worker::stream_worker;pub use worker::AnthropicWorker;pub use worker::EchoWorker;pub use worker::LlamaCppWorker;pub use worker::LoadBalancedWorker;pub use worker::ModelWorker;pub use worker::OpenAiWorker;pub use worker::VllmWorker;pub use failover::FailoverChain;pub use provider_health::ProviderHealth;pub use provider_health::ProviderHealthMonitor;pub use smart_router::ModelPricing;pub use smart_router::RoutingDecision;pub use smart_router::RoutingRequirements;pub use smart_router::SmartRouter;pub use load_balancer::BalancerConfig;pub use load_balancer::EndpointStats;pub use load_balancer::LoadBalancer;pub use load_balancer::LoadBalancerStats;pub use load_balancer::ModelEndpoint;pub use template::TemplateContext;pub use template::TemplateLibrary;pub use template::TemplateValue;pub use pipeline::AppendStage;pub use pipeline::LanguageDetectStage;pub use pipeline::Pipeline;pub use pipeline::PipelineBuilder;pub use pipeline::PipelineError;pub use pipeline::PipelineResult;pub use pipeline::PipelineStats;pub use pipeline::PrependStage;pub use pipeline::RegexReplaceStage;pub use pipeline::TrimStage;pub use pipeline::TruncateStage;pub use pipeline::PipelineStage as PromptPipelineStage;pub use audit::AuditEntry;pub use audit::AuditFilter;pub use audit::AuditLog;pub use audit::AuditQueryResponse;pub use audit::AuditStats;pub use audit::AuditStatsResponse;
Modules§
- ab_test
- A/B Testing Framework
- ab_
testing - A/B Testing Framework
- adaptive_
pool - Adaptive worker pool with Kalman-filter latency prediction.
- adaptive_
timeout - Adaptive timeout management for LLM requests.
- admission_
control - AIMD Adaptive Admission Control
- audit
- Request Replay / Audit Log
- cache
- Prompt Cache
- cache_
warmer - Cache Pre-Warming
- cascade
- Cascading multi-turn inference engine.
- chain_
of_ thought - Chain-of-thought prompting with structured scratchpad and step decomposition.
- circuit_
breaker - Circuit Breaker State Machine
- compression
- Prompt compression — reduce token count before sending to the model.
- config
- Stage: Declarative Pipeline Configuration
- context_
compression - Context window compression and summarization utilities.
- context_
mgr - Context Manager
- conversation
- Conversation History Manager
- conversation_
analyzer - Conversation quality and engagement analysis.
- conversation_
graph - DAG-based conversation branching.
- conversation_
state - Conversation State Machine
- coordination
- Coordination — programmatic agent fleet management
- cost_
estimator - Pre-flight cost estimation for LLM requests.
- cost_
optimizer - Module: Cost Optimizer
- enhanced
- Enhanced resilience and performance features for the pipeline.
- eval_
harness - Evaluation Harness
- experiment_
runner - A/B test runner with statistical significance testing.
- failover
- Provider Failover Chain
- feedback_
loop - Reinforcement-learning-style feedback collector and reward modeler.
- hot_
config - Hot-Reloadable Configuration
- intent_
classifier - Intent Classifier
- job_
scheduler - Job Scheduler
- load_
balancer - Multi-Model Load Balancer
- metrics
- Prometheus metrics for the orchestrator pipeline.
- model_
fallback - Model fallback chains for resilience — automatic failover across model tiers.
- model_
registry - Model Registry
- model_
selector - Intelligent model selection based on cost / quality / latency tradeoffs.
- multi_
modal - Multi-modal content handling for prompts and responses.
- multi_
pipeline - Multi-pipeline routing with prompt classification.
- observability
- Observability — OpenTelemetry-Compatible Tracing and Metrics
- output_
cache - Semantic output deduplication cache with TTL and pluggable eviction policies.
- persona_
manager - Persona management for AI assistants.
- pipeline
- Prompt Pipeline
- pipeline_
builder - Fluent builder for prompt processing pipelines.
- plugin
- Custom plugin stage system for the LLM inference pipeline.
- priority_
queue - Multi-Level Priority Queue with Anti-Starvation Aging
- prompt_
optimizer - Prompt Optimizer
- prompt_
router - Route prompts to different handlers based on content and intent.
- prompt_
safety - Content moderation, toxicity detection, and PII detection.
- prompt_
template - Jinja-lite template engine for prompt construction.
- prompt_
validator - Prompt Validator
- prompt_
versioning - Prompt version control with diff and rollback.
- provider_
health - Provider Health Monitor
- provider_
manager - Multi-provider LLM manager with failover, load-balancing, and rate-limit tracking.
- rate_
limiter - Token-bucket and sliding-window rate limiter per model.
- request_
dedup - Request Deduplication
- response_
classifier - LLM response classification and quality scoring.
- response_
validator - Response Validator
- retry_
budget - Retry budgeting with exponential backoff tracking.
- retry_
policy - Retry Policy
- routing
- Stage: Model Routing Intelligence
- scheduler
- Cron-style scheduler for periodic LLM prompt submissions.
- security
- Prompt Security — Injection and Jailbreak Detection
- semantic_
cache - Semantic Cache
- session
- Session Context Manager
- session_
manager - Session tracking and context-window management.
- session_
mgr - Conversation Session Manager
- smart_
router - Cost-Aware Smart Router
- stages
- Pipeline stage implementations with structured tracing.
- stream_
agg - Streaming Response Aggregator
- streaming_
processor - Streaming token processor for SSE/streaming LLM responses.
- template
- Prompt Template Engine
- templates
- Prompt Template Engine
- token_
budget - Token Budget Middleware
- token_
counter - Multi-Model Token Counting with BPE Approximation
- tool_
call_ parser - Tool call parsing and schema validation for LLM output.
- trace_
export - Module: Distributed Tracing Export
- trace_
ui - Module: Request Lifecycle Tracer TUI Panel
- webhooks
- Module: Webhook Notifications
- worker
- Model worker abstraction and implementations
- worker_
pool - Dynamic worker pool with auto-scaling.
Structs§
- Assemble
Output - Output from the prompt assembly stage.
- Dead
Letter Queue - In-memory dead-letter queue for shed pipeline requests.
- Dropped
Request - A request that was dropped (shed) by the pipeline due to backpressure or
failure. Stored in the
DeadLetterQueuefor inspection and replay. - Inference
Output - Output from the inference stage.
- Post
Output - Output from the post-processing stage.
- Prompt
Request - Initial prompt request from client
- RagOutput
- Output from the RAG (retrieval-augmented generation) stage.
- Session
Id - Unique session identifier for request tracking and affinity
Enums§
- Orchestrator
Error - Orchestrator-specific errors.
- Pipeline
Stage - Pipeline stage identifier for metrics and logging.
- Send
Outcome - Outcome of a
send_with_shedcall.
Functions§
- init_
tracing - Initialise tracing with env-filter support. Call once at binary startup.
- send_
with_ shed - Send with graceful shedding on backpressure.
- shard_
session - Session affinity sharding helper.
- try_
build_ otel_ layer - Attempt to build an OpenTelemetry tracing layer, returning
Noneon error.
Type Aliases§
- Otel
Layer - Type alias for the optional OpenTelemetry tracing layer used in main.rs.