harness_gateway_client/broker.rs
1//! The Gateway's inference broker: one model round per `Chat` effect on
2//! the Gateway client, and the Gateway's model list.
3//!
4//! The Harness takes only a round's finished reply, so the broker's
5//! `InferenceBroker` round returns it whole. A Host that shows the reply
6//! as it forms runs the round through [`GatewayBroker::chat_streaming`]
7//! instead, which hands each live piece to the Host's callback as it
8//! arrives, so no fragment ever reaches the Host's recorder. Either way
9//! the completion names the model the round's options sent it to, so a
10//! Host that routes a round to another model sees that model on its reply.
11
12use std::fmt;
13use std::sync::Arc;
14
15use harness::{BoxFuture, InferenceBroker};
16use promptforge::RunLimits;
17use promptforge::effect::Round;
18use promptforge::model::{
19 Completion, CompletionError, CompletionOptions, Message, ModelBinding, ModelCatalog, ToolSchema,
20};
21
22use crate::catalog::fetch_model_catalog;
23use crate::config::{GatewayEndpoint, SecretString};
24use crate::transport::GatewayChat;
25use crate::wire::delta::StreamDelta;
26
27/// A model broker that sends the Harness's model rounds to the PromptForge
28/// Gateway and lists the Gateway's models.
29///
30/// A Host passes this broker to the Harness as its [`InferenceBroker`].
31/// The broker performs each `Chat` effect from the Engine as one model
32/// round on a [`GatewayChat`]. It lists the Gateway's models through
33/// [`fetch_model_catalog`]. Both kinds of request go to the same API root
34/// and authenticate with the same bearer key.
35///
36/// Every round uses the Engine's default run limits from
37/// [`RunLimits::new`]: the wait limit for each part of the response and
38/// the cap on response size in bytes.
39///
40/// The broker's futures, [`models`](InferenceBroker::models) included, need
41/// a tokio runtime with its reactor and timer. A Host that uses this broker
42/// must therefore await `Harness::run` inside such a runtime. The Harness's
43/// own futures run on any executor.
44#[derive(Clone)]
45pub struct GatewayBroker {
46 client: GatewayChat,
47 endpoint: GatewayEndpoint,
48 key: SecretString,
49}
50
51impl fmt::Debug for GatewayBroker {
52 /// The bearer key is never written to logs or `Debug` output.
53 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
54 formatter
55 .debug_struct("GatewayBroker")
56 .field("client", &self.client)
57 .field("api_root", &self.endpoint.url())
58 .field("key", &"<redacted>")
59 .finish()
60 }
61}
62
63impl GatewayBroker {
64 /// Creates a broker that reaches the Gateway at `endpoint` with the
65 /// bearer key `key`.
66 ///
67 /// `endpoint` is the Gateway's OpenAI-compatible API root, the `/v1`
68 /// path. Model rounds and model-list requests both go to it.
69 #[must_use]
70 pub fn new(endpoint: GatewayEndpoint, key: SecretString) -> GatewayBroker {
71 let limits = RunLimits::new();
72 let client = GatewayChat::new(endpoint.clone(), key.clone())
73 .with_request_limits(limits.timeout(), limits.response_bytes());
74 GatewayBroker {
75 client,
76 endpoint,
77 key,
78 }
79 }
80
81 /// Runs one model round and passes each piece of the reply to
82 /// `on_piece` as it arrives.
83 ///
84 /// The round is the same one `InferenceBroker::chat` runs. `on_piece`
85 /// receives the pieces in the order the stream delivers them. The
86 /// returned completion holds the whole reply.
87 ///
88 /// `on_piece` is called inline as each piece is read, so it must not
89 /// block.
90 pub fn chat_streaming(
91 &self,
92 _binding: ModelBinding,
93 messages: Vec<Message>,
94 tools: Vec<ToolSchema>,
95 options: CompletionOptions,
96 on_piece: Arc<dyn Fn(StreamDelta) + Send + Sync>,
97 ) -> BoxFuture<Result<Box<Completion>, CompletionError>> {
98 self.round(messages, tools, options, move |piece| on_piece(piece))
99 }
100
101 /// One round on the client, handing each live piece to `on_piece`.
102 fn round(
103 &self,
104 messages: Vec<Message>,
105 tools: Vec<ToolSchema>,
106 options: CompletionOptions,
107 on_piece: impl Fn(StreamDelta) + Send + Sync + 'static,
108 ) -> BoxFuture<Result<Box<Completion>, CompletionError>> {
109 let client = self.client.clone();
110 Box::pin(async move {
111 // An empty advertisement sends no `tools` field at all, the
112 // plain chat-completions shape.
113 let tools = (!tools.is_empty()).then_some(tools.as_slice());
114 client
115 .complete(&messages, tools, &options, on_piece)
116 .await
117 .map(Box::new)
118 })
119 }
120}
121
122impl InferenceBroker for GatewayBroker {
123 fn models(&self) -> BoxFuture<Result<ModelCatalog, CompletionError>> {
124 let api_root = self.endpoint.url().to_owned();
125 let key = self.key.clone();
126 Box::pin(async move { fetch_model_catalog(&api_root, key.expose()).await })
127 }
128
129 fn chat(
130 &self,
131 _binding: ModelBinding,
132 messages: Vec<Message>,
133 tools: Vec<ToolSchema>,
134 options: CompletionOptions,
135 _round: Round,
136 ) -> BoxFuture<Result<Box<Completion>, CompletionError>> {
137 self.round(messages, tools, options, |_| {})
138 }
139}
140
141#[cfg(test)]
142#[path = "broker-tests.rs"]
143mod tests;