Skip to main content

locode_provider/openai/responses/
mod.rs

1//! The OpenAI Responses wire — the second live `Provider` (Task 18).
2//!
3//! `api_schema() = "openai-responses"`, `POST {base_url}/v1/responses`,
4//! non-streaming, always-Bearer, **stateless always** (`store:false`, no
5//! `previous_response_id`). Drives both OpenAI models (function + custom/
6//! grammar tools, encrypted-reasoning replay) and xAI grok models (function
7//! tools + encrypted reasoning; `custom_tools_supported=false` degrades
8//! freeform specs). Design: `tasks/plans/task-18-openai-responses-wire.md`.
9
10pub mod build;
11pub mod parse;
12mod stream;
13pub mod wire;
14
15use std::sync::Arc;
16
17use async_trait::async_trait;
18
19pub use build::{build_request, freeform_fallback_parameters, freeform_tool_names};
20pub use parse::response_to_completion;
21
22use crate::completion::{Completion, CompletionDelta};
23use crate::http::{self, HttpFailure, RetryPolicy};
24use crate::openai::{OpenAiModelConfig, classify};
25use crate::provider::{Provider, ProviderError};
26use crate::repair::repair_pairing;
27use crate::request::ConversationRequest;
28
29/// The live OpenAI Responses `Provider` (non-streaming, stateless).
30pub struct OpenAiResponsesProvider {
31    http: reqwest::Client,
32    config: OpenAiModelConfig,
33    retry: RetryPolicy,
34}
35
36impl OpenAiResponsesProvider {
37    /// Build a provider from a resolved [`OpenAiModelConfig`].
38    ///
39    /// # Errors
40    /// [`ProviderError::Transport`] when the HTTP client cannot be constructed.
41    pub fn new(config: OpenAiModelConfig) -> Result<Self, ProviderError> {
42        Ok(Self {
43            http: http::build_http_client()?,
44            config,
45            retry: RetryPolicy::default(),
46        })
47    }
48
49    /// Build from the environment (`LOCODE_API_KEY` / `LOCODE_BASE_URL`); the
50    /// model comes from `--model`/settings, not env (ADR-0024 §1.4).
51    ///
52    /// # Errors
53    /// [`ProviderError::Auth`] when the key is missing;
54    /// [`ProviderError::Transport`] when the client cannot be constructed.
55    pub fn from_env() -> Result<Self, ProviderError> {
56        Self::new(OpenAiModelConfig::from_env()?)
57    }
58
59    /// Override the transport retry policy.
60    #[must_use]
61    pub fn with_retry_policy(mut self, retry: RetryPolicy) -> Self {
62        self.retry = retry;
63        self
64    }
65
66    /// The active config (read-only view).
67    #[must_use]
68    pub fn config(&self) -> &OpenAiModelConfig {
69        &self.config
70    }
71
72    /// Mutable config access (the facade sets `prompt_cache_key` to the
73    /// session id, plan §A.5 Q4).
74    pub fn config_mut(&mut self) -> &mut OpenAiModelConfig {
75        &mut self.config
76    }
77
78    async fn send_once(
79        &self,
80        request: &wire::ResponsesRequest,
81        freeform_names: &std::collections::HashSet<String>,
82    ) -> Result<Completion, HttpFailure> {
83        let url = format!("{}/v1/responses", self.config.base_url);
84        let mut builder = self
85            .http
86            .post(&url)
87            .bearer_auth(&self.config.bearer)
88            .json(request);
89        for (name, value) in &self.config.extra_headers {
90            builder = builder.header(name, value);
91        }
92        let response = builder
93            .send()
94            .await
95            .map_err(|e| HttpFailure::transport(e.to_string()))?;
96
97        let status = response.status();
98        let retry_after = response
99            .headers()
100            .get(reqwest::header::RETRY_AFTER)
101            .and_then(|v| v.to_str().ok())
102            .and_then(http::parse_retry_after);
103
104        if status.is_success() {
105            let parsed: wire::ResponsesResponse = response
106                .json()
107                .await
108                .map_err(|e| HttpFailure::decode(format!("response body: {e}")))?;
109            return response_to_completion(parsed, freeform_names).map_err(|error| HttpFailure {
110                // A `failed` response's rate-limit/server-error codes ARE
111                // retryable — let the shared loop see them as such.
112                force_terminal: false,
113                retry_after,
114                error,
115            });
116        }
117
118        let text = response.text().await.unwrap_or_default();
119        let body: crate::openai::OpenAiErrorBody = match serde_json::from_str(&text) {
120            Ok(body) => body,
121            Err(_) => crate::openai::OpenAiErrorBody {
122                error: crate::openai::OpenAiErrorDetail {
123                    code: None,
124                    r#type: None,
125                    message: text,
126                },
127            },
128        };
129        Err(classify(status.as_u16(), retry_after, &body))
130    }
131}
132
133#[async_trait]
134impl Provider for OpenAiResponsesProvider {
135    #[allow(clippy::unnecessary_literal_bound)] // signature is the trait's
136    fn api_schema(&self) -> &str {
137        "openai-responses"
138    }
139
140    async fn complete(&self, request: &ConversationRequest) -> Result<Completion, ProviderError> {
141        // 1. Defensive transcript repair on a clone (ADR-0004 — same pass as
142        //    the anthropic wire; the engine runs the canonical one).
143        let mut repaired = request.clone();
144        let _ = repair_pairing(&mut repaired.messages);
145
146        // 2. Build (stateless, include, reasoning, tools + degradation).
147        let wire_request = build_request(&repaired, &self.config);
148        let freeform_names = freeform_tool_names(&repaired.tools);
149
150        // 3. Send with the shared transport retry.
151        let freeform_ref = &freeform_names;
152        http::run_with_retry(&self.retry, |_attempt| {
153            self.send_once(&wire_request, freeform_ref)
154        })
155        .await
156    }
157
158    async fn stream(
159        &self,
160        request: &ConversationRequest,
161        on_delta: &mut (dyn FnMut(CompletionDelta) + Send),
162    ) -> Result<Completion, ProviderError> {
163        // Same repair + build as `complete`, but streaming (the whole final
164        // response still rides the terminal SSE event → byte-identical result).
165        let mut repaired = request.clone();
166        let _ = repair_pairing(&mut repaired.messages);
167        let mut wire_request = build_request(&repaired, &self.config);
168        wire_request.stream = true;
169        let freeform_names = freeform_tool_names(&repaired.tools);
170
171        // Single attempt — a retryable failure surfaces to the engine's resample.
172        stream::send_once_streaming(
173            &self.http,
174            &self.config,
175            &wire_request,
176            &freeform_names,
177            on_delta,
178        )
179        .await
180        .map_err(|f| f.error)
181    }
182}
183
184/// The provider wrapped in [`Arc`] for the engine's `Arc<dyn Provider>` slot.
185#[must_use]
186pub fn into_provider(provider: OpenAiResponsesProvider) -> Arc<dyn Provider> {
187    Arc::new(provider)
188}