Skip to repository content735 lines · 23.8 KB · rust
tenant.openagents/omega
No repository description is available.
OpenAgents Git authority 2026-07-28T03:38:12.010Z Public web read
NIP-34 coordinate
30617:7649603503856e5148d571eac2766b288a8ff1e9e35d380337a1d2b0015b4f92:omegaMaintainersHidden in public view
References2 branches · 1 tag
Read-only clone
git clone https://openagents.com/git/tenant.openagents/omega.gitBrowse files
open_ai_compatible.rs
1use anyhow::Result;
2use credentials_provider::CredentialsProvider;
3use futures::{FutureExt, StreamExt, future::BoxFuture};
4use gpui::{App, AppContext, AsyncApp, Entity, Task};
5use http_client::{CustomHeaders, HttpClient};
6use language_model::{
7 AuthenticateError, IconOrSvg, LanguageModel, LanguageModelCompletionError,
8 LanguageModelCompletionEvent, LanguageModelEffortLevel, LanguageModelId, LanguageModelName,
9 LanguageModelProvider, LanguageModelProviderId, LanguageModelProviderName,
10 LanguageModelProviderState, LanguageModelRequest, LanguageModelToolChoice,
11 LanguageModelToolSchemaFormat, ProviderSettingsView, RateLimiter, SubPageProviderSettings,
12};
13use open_ai::{
14 ResponseStreamEvent,
15 responses::{Request as ResponseRequest, StreamEvent as ResponsesStreamEvent, stream_response},
16 stream_completion,
17};
18use settings::Settings;
19use std::sync::Arc;
20use ui::IconName;
21
22use crate::provider::api_compatible::{
23 ApiCompatibleProviderConfigurationView, ApiCompatibleProviderSettings,
24 ApiCompatibleProviderState,
25};
26use crate::provider::open_ai::{
27 OpenAiEventMapper, OpenAiResponseEventMapper, into_open_ai, into_open_ai_response,
28};
29pub use settings::OpenAiCompatibleAvailableModel as AvailableModel;
30pub use settings::OpenAiCompatibleModelCapabilities as ModelCapabilities;
31
32const API_KEY_PLACEHOLDER: &str = "000000000000000000000000000000000000000000000000000";
33
34#[derive(Default, Clone, Debug, PartialEq)]
35pub struct OpenAiCompatibleSettings {
36 pub api_url: String,
37 pub available_models: Vec<AvailableModel>,
38 pub custom_headers: CustomHeaders,
39}
40
41impl ApiCompatibleProviderSettings for OpenAiCompatibleSettings {
42 fn api_url(&self) -> &str {
43 &self.api_url
44 }
45}
46
47pub type State = ApiCompatibleProviderState<OpenAiCompatibleSettings>;
48
49pub struct OpenAiCompatibleLanguageModelProvider {
50 id: LanguageModelProviderId,
51 name: LanguageModelProviderName,
52 http_client: Arc<dyn HttpClient>,
53 state: Entity<State>,
54}
55
56impl OpenAiCompatibleLanguageModelProvider {
57 pub fn new(
58 id: Arc<str>,
59 http_client: Arc<dyn HttpClient>,
60 credentials_provider: Arc<dyn CredentialsProvider>,
61 cx: &mut App,
62 ) -> Self {
63 let state = State::new(
64 id.clone(),
65 credentials_provider,
66 |id, cx| {
67 crate::AllLanguageModelSettings::get_global(cx)
68 .openai_compatible
69 .get(id)
70 },
71 cx,
72 );
73
74 Self {
75 id: id.clone().into(),
76 name: id.into(),
77 http_client,
78 state,
79 }
80 }
81
82 fn create_language_model(&self, model: AvailableModel) -> Arc<dyn LanguageModel> {
83 Arc::new(OpenAiCompatibleLanguageModel {
84 id: LanguageModelId::from(model.name.clone()),
85 provider_id: self.id.clone(),
86 provider_name: self.name.clone(),
87 model,
88 state: self.state.clone(),
89 http_client: self.http_client.clone(),
90 request_limiter: RateLimiter::new(4),
91 })
92 }
93}
94
95impl LanguageModelProviderState for OpenAiCompatibleLanguageModelProvider {
96 type ObservableEntity = State;
97
98 fn observable_entity(&self) -> Option<Entity<Self::ObservableEntity>> {
99 Some(self.state.clone())
100 }
101}
102
103impl LanguageModelProvider for OpenAiCompatibleLanguageModelProvider {
104 fn id(&self) -> LanguageModelProviderId {
105 self.id.clone()
106 }
107
108 fn name(&self) -> LanguageModelProviderName {
109 self.name.clone()
110 }
111
112 fn icon(&self) -> IconOrSvg {
113 IconOrSvg::Icon(IconName::AiOpenAiCompat)
114 }
115
116 fn default_model(&self, cx: &App) -> Option<Arc<dyn LanguageModel>> {
117 self.state
118 .read(cx)
119 .settings
120 .available_models
121 .first()
122 .map(|model| self.create_language_model(model.clone()))
123 }
124
125 fn default_fast_model(&self, _cx: &App) -> Option<Arc<dyn LanguageModel>> {
126 None
127 }
128
129 fn provided_models(&self, cx: &App) -> Vec<Arc<dyn LanguageModel>> {
130 self.state
131 .read(cx)
132 .settings
133 .available_models
134 .iter()
135 .map(|model| self.create_language_model(model.clone()))
136 .collect()
137 }
138
139 fn is_authenticated(&self, cx: &App) -> bool {
140 self.state.read(cx).is_authenticated()
141 }
142
143 fn authenticate(&self, cx: &mut App) -> Task<Result<(), AuthenticateError>> {
144 self.state.update(cx, |state, cx| state.authenticate(cx))
145 }
146
147 fn settings_view(&self, _cx: &mut App) -> Option<ProviderSettingsView> {
148 let state = self.state.clone();
149 Some(ProviderSettingsView::SubPage(SubPageProviderSettings::new(
150 move |window, cx| {
151 cx.new(|cx| {
152 ApiCompatibleProviderConfigurationView::new(
153 state.clone(),
154 "OpenAI",
155 API_KEY_PLACEHOLDER,
156 window,
157 cx,
158 )
159 })
160 .into()
161 },
162 )))
163 }
164
165 fn set_api_key(&self, api_key: Option<String>, cx: &mut App) -> Task<Result<()>> {
166 self.state
167 .update(cx, |state, cx| state.set_api_key(api_key, cx))
168 }
169}
170
171pub struct OpenAiCompatibleLanguageModel {
172 id: LanguageModelId,
173 provider_id: LanguageModelProviderId,
174 provider_name: LanguageModelProviderName,
175 model: AvailableModel,
176 state: Entity<State>,
177 http_client: Arc<dyn HttpClient>,
178 request_limiter: RateLimiter,
179}
180
181impl OpenAiCompatibleLanguageModel {
182 fn stream_completion(
183 &self,
184 request: open_ai::Request,
185 cx: &AsyncApp,
186 ) -> BoxFuture<
187 'static,
188 Result<
189 futures::stream::BoxStream<'static, Result<ResponseStreamEvent>>,
190 LanguageModelCompletionError,
191 >,
192 > {
193 let http_client = self.http_client.clone();
194
195 let (api_key, api_url, extra_headers) = self.state.read_with(cx, |state, _cx| {
196 let api_url = &state.settings.api_url;
197 (
198 state.api_key_state.key(api_url),
199 state.settings.api_url.clone(),
200 state.settings.custom_headers.clone(),
201 )
202 });
203
204 let provider = self.provider_name.clone();
205 let future = self.request_limiter.stream(async move {
206 let Some(api_key) = api_key else {
207 return Err(LanguageModelCompletionError::NoApiKey { provider });
208 };
209 let request = stream_completion(
210 http_client.as_ref(),
211 provider.0.as_str(),
212 &api_url,
213 &api_key,
214 request,
215 &extra_headers,
216 );
217 let response = request.await?;
218 Ok(response)
219 });
220
221 async move { Ok(future.await?.boxed()) }.boxed()
222 }
223
224 fn stream_response(
225 &self,
226 request: ResponseRequest,
227 cx: &AsyncApp,
228 ) -> BoxFuture<'static, Result<futures::stream::BoxStream<'static, Result<ResponsesStreamEvent>>>>
229 {
230 let http_client = self.http_client.clone();
231
232 let (api_key, api_url, extra_headers) = self.state.read_with(cx, |state, _cx| {
233 let api_url = &state.settings.api_url;
234 (
235 state.api_key_state.key(api_url),
236 state.settings.api_url.clone(),
237 state.settings.custom_headers.clone(),
238 )
239 });
240
241 let provider = self.provider_name.clone();
242 let future = self.request_limiter.stream(async move {
243 let Some(api_key) = api_key else {
244 return Err(LanguageModelCompletionError::NoApiKey { provider });
245 };
246 let request = stream_response(
247 http_client.as_ref(),
248 provider.0.as_str(),
249 &api_url,
250 &api_key,
251 request,
252 &extra_headers,
253 );
254 let response = request.await?;
255 Ok(response)
256 });
257
258 async move { Ok(future.await?.boxed()) }.boxed()
259 }
260}
261
262fn default_thinking_reasoning_effort(model: &AvailableModel) -> Option<open_ai::ReasoningEffort> {
263 model
264 .reasoning_effort
265 .filter(|effort| *effort != open_ai::ReasoningEffort::None)
266}
267
268fn supported_thinking_effort_levels(model: &AvailableModel) -> Vec<LanguageModelEffortLevel> {
269 let Some(default_effort) = default_thinking_reasoning_effort(model) else {
270 return Vec::new();
271 };
272
273 open_ai::ReasoningEffort::OPENAI_COMPATIBLE_SELECTABLE
274 .into_iter()
275 .map(|effort| LanguageModelEffortLevel {
276 name: effort.label().into(),
277 value: effort.value().into(),
278 is_default: effort == default_effort,
279 })
280 .collect()
281}
282
283fn selected_thinking_reasoning_effort(
284 request: &LanguageModelRequest,
285) -> Option<open_ai::ReasoningEffort> {
286 request
287 .thinking_effort
288 .as_deref()
289 .and_then(|effort| effort.parse::<open_ai::ReasoningEffort>().ok())
290 .filter(|effort| *effort != open_ai::ReasoningEffort::None)
291}
292
293fn chat_completion_max_tokens_parameter(
294 model: &AvailableModel,
295) -> crate::provider::open_ai::ChatCompletionMaxTokensParameter {
296 if model.capabilities.max_tokens_parameter {
297 crate::provider::open_ai::ChatCompletionMaxTokensParameter::MaxTokens
298 } else {
299 crate::provider::open_ai::ChatCompletionMaxTokensParameter::MaxCompletionTokens
300 }
301}
302
303fn supports_none_reasoning_effort(model: &AvailableModel) -> bool {
304 model.reasoning_effort.is_some()
305}
306
307fn chat_completion_reasoning_effort(
308 request: &LanguageModelRequest,
309 model: &AvailableModel,
310) -> Option<open_ai::ReasoningEffort> {
311 if model.reasoning_effort == Some(open_ai::ReasoningEffort::None) {
312 return Some(open_ai::ReasoningEffort::None);
313 }
314
315 if request.thinking_allowed {
316 selected_thinking_reasoning_effort(request)
317 .or_else(|| default_thinking_reasoning_effort(model))
318 } else if supports_none_reasoning_effort(model) {
319 Some(open_ai::ReasoningEffort::None)
320 } else {
321 None
322 }
323}
324
325fn disable_response_thinking_for_none_effort(
326 request: &mut LanguageModelRequest,
327 model: &AvailableModel,
328) {
329 if model.reasoning_effort == Some(open_ai::ReasoningEffort::None) {
330 request.thinking_allowed = false;
331 request.thinking_effort = None;
332 }
333}
334
335impl LanguageModel for OpenAiCompatibleLanguageModel {
336 fn id(&self) -> LanguageModelId {
337 self.id.clone()
338 }
339
340 fn name(&self) -> LanguageModelName {
341 LanguageModelName::from(
342 self.model
343 .display_name
344 .clone()
345 .unwrap_or_else(|| self.model.name.clone()),
346 )
347 }
348
349 fn provider_id(&self) -> LanguageModelProviderId {
350 self.provider_id.clone()
351 }
352
353 fn provider_name(&self) -> LanguageModelProviderName {
354 self.provider_name.clone()
355 }
356
357 fn supports_tools(&self) -> bool {
358 self.model.capabilities.tools
359 }
360
361 fn tool_input_format(&self) -> LanguageModelToolSchemaFormat {
362 LanguageModelToolSchemaFormat::JsonSchemaSubset
363 }
364
365 fn supports_images(&self) -> bool {
366 self.model.capabilities.images
367 }
368
369 fn supports_tool_choice(&self, choice: LanguageModelToolChoice) -> bool {
370 match choice {
371 LanguageModelToolChoice::Auto => self.model.capabilities.tools,
372 LanguageModelToolChoice::Any => self.model.capabilities.tools,
373 LanguageModelToolChoice::None => true,
374 }
375 }
376
377 fn supports_streaming_tools(&self) -> bool {
378 true
379 }
380
381 fn supports_thinking(&self) -> bool {
382 default_thinking_reasoning_effort(&self.model).is_some()
383 }
384
385 fn supported_effort_levels(&self) -> Vec<LanguageModelEffortLevel> {
386 supported_thinking_effort_levels(&self.model)
387 }
388
389 fn supports_split_token_display(&self) -> bool {
390 true
391 }
392
393 fn telemetry_id(&self) -> String {
394 format!("openai/{}", self.model.name)
395 }
396
397 fn max_token_count(&self) -> u64 {
398 self.model.max_tokens
399 }
400
401 fn max_output_tokens(&self) -> Option<u64> {
402 self.model.max_output_tokens
403 }
404
405 fn stream_completion(
406 &self,
407 mut request: LanguageModelRequest,
408 cx: &AsyncApp,
409 ) -> BoxFuture<
410 'static,
411 Result<
412 futures::stream::BoxStream<
413 'static,
414 Result<LanguageModelCompletionEvent, LanguageModelCompletionError>,
415 >,
416 LanguageModelCompletionError,
417 >,
418 > {
419 // `speed` can leak in from a parent thread's model; this provider never
420 // supports fast mode, and arbitrary compatible endpoints reject `service_tier`.
421 if !self.supports_fast_mode() {
422 request.speed = None;
423 }
424
425 if self.model.capabilities.chat_completions {
426 let reasoning_effort = chat_completion_reasoning_effort(&request, &self.model);
427 let request = match into_open_ai(
428 request,
429 &self.model.name,
430 self.model.capabilities.parallel_tool_calls,
431 self.model.capabilities.prompt_cache_key,
432 self.max_output_tokens(),
433 chat_completion_max_tokens_parameter(&self.model),
434 reasoning_effort,
435 self.model.capabilities.interleaved_reasoning,
436 ) {
437 Ok(request) => request,
438 Err(error) => return async move { Err(error.into()) }.boxed(),
439 };
440 let completions = self.stream_completion(request, cx);
441 async move {
442 let mapper = OpenAiEventMapper::new();
443 Ok(mapper.map_stream(completions.await?).boxed())
444 }
445 .boxed()
446 } else {
447 disable_response_thinking_for_none_effort(&mut request, &self.model);
448 let request = match into_open_ai_response(
449 request,
450 &self.model.name,
451 self.model.capabilities.parallel_tool_calls,
452 self.model.capabilities.prompt_cache_key,
453 self.max_output_tokens(),
454 default_thinking_reasoning_effort(&self.model),
455 supports_none_reasoning_effort(&self.model),
456 &self.provider_id,
457 ) {
458 Ok(request) => request,
459 Err(error) => return async move { Err(error.into()) }.boxed(),
460 };
461 let completions = self.stream_response(request, cx);
462 let compaction_state_owner = self.provider_id.clone();
463 async move {
464 let mapper = OpenAiResponseEventMapper::new(compaction_state_owner);
465 Ok(mapper.map_stream(completions.await?).boxed())
466 }
467 .boxed()
468 }
469 }
470}
471
472#[cfg(test)]
473mod tests {
474 use super::*;
475
476 use serde_json::json;
477
478 fn available_model(reasoning_effort: Option<open_ai::ReasoningEffort>) -> AvailableModel {
479 AvailableModel {
480 name: "custom-model".to_string(),
481 display_name: None,
482 max_tokens: 128_000,
483 max_output_tokens: None,
484 max_completion_tokens: None,
485 reasoning_effort,
486 capabilities: ModelCapabilities {
487 chat_completions: false,
488 ..Default::default()
489 },
490 }
491 }
492
493 #[test]
494 fn configured_reasoning_effort_supports_thinking() {
495 assert_eq!(
496 default_thinking_reasoning_effort(&available_model(Some(
497 open_ai::ReasoningEffort::High
498 ))),
499 Some(open_ai::ReasoningEffort::High)
500 );
501 }
502
503 #[test]
504 fn missing_or_none_reasoning_effort_does_not_support_thinking() {
505 assert_eq!(
506 default_thinking_reasoning_effort(&available_model(None)),
507 None
508 );
509 assert_eq!(
510 default_thinking_reasoning_effort(&available_model(Some(
511 open_ai::ReasoningEffort::None
512 ))),
513 None
514 );
515 }
516
517 #[test]
518 fn supported_thinking_effort_levels_use_configured_effort_as_default() {
519 let effort_levels = supported_thinking_effort_levels(&available_model(Some(
520 open_ai::ReasoningEffort::High,
521 )));
522 let values = effort_levels
523 .iter()
524 .map(|level| level.value.as_ref())
525 .collect::<Vec<_>>();
526
527 assert_eq!(values, ["minimal", "low", "medium", "high", "xhigh", "max"]);
528 assert_eq!(
529 effort_levels
530 .iter()
531 .find(|level| level.is_default)
532 .map(|level| level.value.as_ref()),
533 Some("high")
534 );
535 }
536
537 #[test]
538 fn supported_thinking_effort_levels_hide_missing_or_none_effort() {
539 assert!(supported_thinking_effort_levels(&available_model(None)).is_empty());
540 assert!(
541 supported_thinking_effort_levels(&available_model(Some(
542 open_ai::ReasoningEffort::None
543 )))
544 .is_empty()
545 );
546 }
547
548 #[test]
549 fn chat_completion_reasoning_effort_honors_request_and_configured_effort() {
550 let model = available_model(Some(open_ai::ReasoningEffort::Medium));
551 let mut request = LanguageModelRequest {
552 thinking_allowed: true,
553 ..Default::default()
554 };
555
556 assert_eq!(
557 chat_completion_reasoning_effort(&request, &model),
558 Some(open_ai::ReasoningEffort::Medium)
559 );
560
561 request.thinking_effort = Some("high".to_string());
562 assert_eq!(
563 chat_completion_reasoning_effort(&request, &model),
564 Some(open_ai::ReasoningEffort::High)
565 );
566
567 request.thinking_effort = Some("not-supported".to_string());
568 assert_eq!(
569 chat_completion_reasoning_effort(&request, &model),
570 Some(open_ai::ReasoningEffort::Medium)
571 );
572
573 request.thinking_allowed = false;
574 assert_eq!(
575 chat_completion_reasoning_effort(&request, &model),
576 Some(open_ai::ReasoningEffort::None)
577 );
578 }
579
580 #[test]
581 fn chat_completion_reasoning_effort_omits_missing_effort() {
582 let model = available_model(None);
583 let request = LanguageModelRequest {
584 thinking_allowed: false,
585 ..Default::default()
586 };
587
588 assert_eq!(chat_completion_reasoning_effort(&request, &model), None);
589 }
590
591 #[test]
592 fn chat_completion_reasoning_effort_preserves_explicit_none() {
593 let model = available_model(Some(open_ai::ReasoningEffort::None));
594 let request = LanguageModelRequest {
595 thinking_allowed: true,
596 thinking_effort: Some("high".to_string()),
597 ..Default::default()
598 };
599
600 assert_eq!(
601 chat_completion_reasoning_effort(&request, &model),
602 Some(open_ai::ReasoningEffort::None)
603 );
604 }
605
606 #[test]
607 fn chat_completion_max_tokens_parameter_defaults_to_max_completion_tokens() {
608 let model = available_model(Some(open_ai::ReasoningEffort::Medium));
609
610 assert_eq!(
611 chat_completion_max_tokens_parameter(&model),
612 crate::provider::open_ai::ChatCompletionMaxTokensParameter::MaxCompletionTokens
613 );
614 }
615
616 #[test]
617 fn chat_completion_max_tokens_parameter_uses_max_tokens_when_configured() {
618 let mut model = available_model(Some(open_ai::ReasoningEffort::Medium));
619 model.capabilities.max_tokens_parameter = true;
620
621 assert_eq!(
622 chat_completion_max_tokens_parameter(&model),
623 crate::provider::open_ai::ChatCompletionMaxTokensParameter::MaxTokens
624 );
625 }
626
627 #[test]
628 fn response_request_includes_reasoning_when_effort_is_configured() {
629 let model = available_model(Some(open_ai::ReasoningEffort::High));
630 let request = LanguageModelRequest {
631 thinking_allowed: true,
632 ..Default::default()
633 };
634
635 let request = into_open_ai_response(
636 request,
637 &model.name,
638 model.capabilities.parallel_tool_calls,
639 model.capabilities.prompt_cache_key,
640 model.max_output_tokens,
641 default_thinking_reasoning_effort(&model),
642 supports_none_reasoning_effort(&model),
643 &LanguageModelProviderId::new("test-compatible-provider"),
644 )
645 .unwrap();
646 let serialized = serde_json::to_value(request).unwrap();
647
648 assert_eq!(
649 serialized["reasoning"],
650 json!({ "effort": "high", "summary": "auto" })
651 );
652 assert_eq!(
653 serialized["include"],
654 json!(["reasoning.encrypted_content"])
655 );
656 }
657
658 #[test]
659 fn response_request_omits_reasoning_when_effort_is_missing() {
660 let model = available_model(None);
661 let request = LanguageModelRequest {
662 thinking_allowed: true,
663 ..Default::default()
664 };
665
666 let request = into_open_ai_response(
667 request,
668 &model.name,
669 model.capabilities.parallel_tool_calls,
670 model.capabilities.prompt_cache_key,
671 model.max_output_tokens,
672 default_thinking_reasoning_effort(&model),
673 supports_none_reasoning_effort(&model),
674 &LanguageModelProviderId::new("test-compatible-provider"),
675 )
676 .unwrap();
677 let serialized = serde_json::to_value(request).unwrap();
678
679 assert_eq!(serialized.get("reasoning"), None);
680 assert_eq!(serialized.get("include"), None);
681 }
682
683 #[test]
684 fn chat_completion_request_includes_selected_reasoning_effort() {
685 let mut model = available_model(Some(open_ai::ReasoningEffort::Medium));
686 model.capabilities.chat_completions = true;
687 let request = LanguageModelRequest {
688 thinking_allowed: true,
689 thinking_effort: Some("high".to_string()),
690 ..Default::default()
691 };
692 let reasoning_effort = chat_completion_reasoning_effort(&request, &model);
693
694 let request = into_open_ai(
695 request,
696 &model.name,
697 model.capabilities.parallel_tool_calls,
698 model.capabilities.prompt_cache_key,
699 model.max_output_tokens,
700 chat_completion_max_tokens_parameter(&model),
701 reasoning_effort,
702 model.capabilities.interleaved_reasoning,
703 )
704 .unwrap();
705 let serialized = serde_json::to_value(request).unwrap();
706
707 assert_eq!(serialized["reasoning_effort"], json!("high"));
708 }
709
710 #[test]
711 fn configured_reasoning_effort_supports_none_reasoning_effort() {
712 assert!(supports_none_reasoning_effort(&available_model(Some(
713 open_ai::ReasoningEffort::Medium
714 ))));
715 assert!(supports_none_reasoning_effort(&available_model(Some(
716 open_ai::ReasoningEffort::None
717 ))));
718 assert!(!supports_none_reasoning_effort(&available_model(None)));
719 }
720
721 #[test]
722 fn response_thinking_effort_preserves_explicit_none() {
723 let model = available_model(Some(open_ai::ReasoningEffort::None));
724 let mut request = LanguageModelRequest {
725 thinking_allowed: true,
726 thinking_effort: Some("high".to_string()),
727 ..Default::default()
728 };
729
730 disable_response_thinking_for_none_effort(&mut request, &model);
731 assert!(!request.thinking_allowed);
732 assert_eq!(request.thinking_effort, None);
733 }
734}
735