Skip to main content

Crate tokio_prompt_orchestrator

Crate tokio_prompt_orchestrator 

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

demo

§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

§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

ModuleDescription
stagesFive pipeline stage implementations and channel wiring
workerModelWorker trait and the provider implementations
enhancedResilience primitives: circuit breaker, dedup, semantic dedup (SimHash), retry, cache, rate limiter, smart batching
configTOML-deserialisable PipelineConfig with hot-reload support
metricsPrometheus metrics initialisation and helper functions
routingModelRouter for complexity-scored routing; ArbitrageEngine for SLA-aware provider selection; PoolSizer for worker scaling
securityPromptGuard: prompt injection and jailbreak detection (no external I/O)
sessionMulti-turn conversation context: injects history per session
cascadeMulti-turn cascading inference: tool call loops with pluggable executors
multi_pipelineNamed pipeline fleet with prompt classification and per-class routing
ab_testPrompt A/B testing: consistent hashing assignment, Welch’s t-test, Cohen’s d
coordinationAgent fleet management and task claiming
adaptive_poolKalman-filter worker pool controller: predicts queue depth and recommends scale events
web_apiREST/SSE/WebSocket server (feature: web-api)
distributedRedis dedup and NATS coordination (feature: distributed)
tuiRatatui terminal dashboard (feature: tui)
self_tune, self_modify, intelligence, evolution, self_improveSelf-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§

AssembleOutput
Output from the prompt assembly stage.
DeadLetterQueue
In-memory dead-letter queue for shed pipeline requests.
DroppedRequest
A request that was dropped (shed) by the pipeline due to backpressure or failure. Stored in the DeadLetterQueue for inspection and replay.
InferenceOutput
Output from the inference stage.
PostOutput
Output from the post-processing stage.
PromptRequest
Initial prompt request from client
RagOutput
Output from the RAG (retrieval-augmented generation) stage.
SessionId
Unique session identifier for request tracking and affinity

Enums§

OrchestratorError
Orchestrator-specific errors.
PipelineStage
Pipeline stage identifier for metrics and logging.
SendOutcome
Outcome of a send_with_shed call.

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 None on error.

Type Aliases§

OtelLayer
Type alias for the optional OpenTelemetry tracing layer used in main.rs.