Skip to repository content

tenant.openagents/omega

No repository description is available.

OpenAgents Git authority 2026-07-28T03:38:12.010Z Public web read
NIP-34 coordinate30617:7649603503856e5148d571eac2766b288a8ff1e9e35d380337a1d2b0015b4f92:omega
MaintainersHidden in public view
References2 branches · 1 tag
Read-only clonegit clone https://openagents.com/git/tenant.openagents/omega.git
Browse files

open_ai_compatible.rs

735 lines · 23.8 KB · rust
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
Served at tenant.openagents/omega Member data and write actions are omitted.