Skip to main content

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;