Skip to repository content659 lines · 23.2 KB · rust
tenant.openagents/omega
No repository description is available.
OpenAgents Git authority 2026-07-28T02:36:55.241Z 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
deepseek.rs
1use anyhow::{Result, anyhow};
2use collections::{HashMap, IndexMap};
3use credentials_provider::CredentialsProvider;
4use deepseek::DEEPSEEK_API_URL;
5
6use futures::Stream;
7use futures::{FutureExt, StreamExt, future::BoxFuture, stream::BoxStream};
8use gpui::{App, AppContext, AsyncApp, Context, Entity, SharedString, Task};
9use http_client::{CustomHeaders, HttpClient};
10use language_model::{
11 ApiKeyConfiguration, ApiKeyState, AuthenticateError, EnvVar, IconOrSvg, LanguageModel,
12 LanguageModelCompletionError, LanguageModelCompletionEvent, LanguageModelEffortLevel,
13 LanguageModelId, LanguageModelName, LanguageModelProvider, LanguageModelProviderId,
14 LanguageModelProviderName, LanguageModelProviderState, LanguageModelRequest,
15 LanguageModelToolChoice, LanguageModelToolResultContent, LanguageModelToolUse, MessageContent,
16 ProviderSettingsView, RateLimiter, Role, StopReason, TokenUsage, env_var,
17};
18pub use settings::DeepseekAvailableModel as AvailableModel;
19use settings::{Settings, SettingsStore};
20use std::pin::Pin;
21use std::sync::{Arc, LazyLock};
22
23use ui::IconName;
24
25use language_model::util::{fix_streamed_json, parse_tool_arguments};
26
27const PROVIDER_ID: LanguageModelProviderId = LanguageModelProviderId::new("deepseek");
28const PROVIDER_NAME: LanguageModelProviderName = LanguageModelProviderName::new("DeepSeek");
29
30const API_KEY_ENV_VAR_NAME: &str = "DEEPSEEK_API_KEY";
31static API_KEY_ENV_VAR: LazyLock<EnvVar> = env_var!(API_KEY_ENV_VAR_NAME);
32
33#[derive(Default)]
34struct RawToolCall {
35 id: String,
36 name: String,
37 arguments: String,
38}
39
40#[derive(Default, Clone, Debug, PartialEq)]
41pub struct DeepSeekSettings {
42 pub api_url: String,
43 pub available_models: Vec<AvailableModel>,
44 pub custom_headers: CustomHeaders,
45}
46pub struct DeepSeekLanguageModelProvider {
47 http_client: Arc<dyn HttpClient>,
48 state: Entity<State>,
49}
50
51pub struct State {
52 api_key_state: ApiKeyState,
53 credentials_provider: Arc<dyn CredentialsProvider>,
54}
55
56impl State {
57 fn is_authenticated(&self) -> bool {
58 self.api_key_state.has_key()
59 }
60
61 fn set_api_key(&mut self, api_key: Option<String>, cx: &mut Context<Self>) -> Task<Result<()>> {
62 let credentials_provider = self.credentials_provider.clone();
63 let api_url = DeepSeekLanguageModelProvider::api_url(cx);
64 self.api_key_state.store(
65 api_url,
66 api_key,
67 |this| &mut this.api_key_state,
68 credentials_provider,
69 cx,
70 )
71 }
72
73 fn authenticate(&mut self, cx: &mut Context<Self>) -> Task<Result<(), AuthenticateError>> {
74 let credentials_provider = self.credentials_provider.clone();
75 let api_url = DeepSeekLanguageModelProvider::api_url(cx);
76 self.api_key_state.load_if_needed(
77 api_url,
78 |this| &mut this.api_key_state,
79 credentials_provider,
80 cx,
81 )
82 }
83}
84
85impl DeepSeekLanguageModelProvider {
86 pub fn new(
87 http_client: Arc<dyn HttpClient>,
88 credentials_provider: Arc<dyn CredentialsProvider>,
89 cx: &mut App,
90 ) -> Self {
91 let state = cx.new(|cx| {
92 cx.observe_global::<SettingsStore>(|this: &mut State, cx| {
93 let credentials_provider = this.credentials_provider.clone();
94 let api_url = Self::api_url(cx);
95 this.api_key_state.handle_url_change(
96 api_url,
97 |this| &mut this.api_key_state,
98 credentials_provider,
99 cx,
100 );
101 cx.notify();
102 })
103 .detach();
104 State {
105 api_key_state: ApiKeyState::new(Self::api_url(cx), (*API_KEY_ENV_VAR).clone()),
106 credentials_provider,
107 }
108 });
109
110 Self { http_client, state }
111 }
112
113 fn create_language_model(&self, model: deepseek::Model) -> Arc<dyn LanguageModel> {
114 Arc::new(DeepSeekLanguageModel {
115 id: LanguageModelId::from(model.id().to_string()),
116 model,
117 state: self.state.clone(),
118 http_client: self.http_client.clone(),
119 request_limiter: RateLimiter::new(4),
120 })
121 }
122
123 fn settings(cx: &App) -> &DeepSeekSettings {
124 &crate::AllLanguageModelSettings::get_global(cx).deepseek
125 }
126
127 fn api_url(cx: &App) -> SharedString {
128 let api_url = &Self::settings(cx).api_url;
129 if api_url.is_empty() {
130 DEEPSEEK_API_URL.into()
131 } else {
132 SharedString::new(api_url.as_str())
133 }
134 }
135}
136
137impl LanguageModelProviderState for DeepSeekLanguageModelProvider {
138 type ObservableEntity = State;
139
140 fn observable_entity(&self) -> Option<Entity<Self::ObservableEntity>> {
141 Some(self.state.clone())
142 }
143}
144
145impl LanguageModelProvider for DeepSeekLanguageModelProvider {
146 fn id(&self) -> LanguageModelProviderId {
147 PROVIDER_ID
148 }
149
150 fn name(&self) -> LanguageModelProviderName {
151 PROVIDER_NAME
152 }
153
154 fn icon(&self) -> IconOrSvg {
155 IconOrSvg::Icon(IconName::AiDeepSeek)
156 }
157
158 fn default_model(&self, _cx: &App) -> Option<Arc<dyn LanguageModel>> {
159 Some(self.create_language_model(deepseek::Model::default()))
160 }
161
162 fn default_fast_model(&self, _cx: &App) -> Option<Arc<dyn LanguageModel>> {
163 Some(self.create_language_model(deepseek::Model::default_fast()))
164 }
165
166 fn provided_models(&self, cx: &App) -> Vec<Arc<dyn LanguageModel>> {
167 let mut models = IndexMap::default();
168
169 models.insert("deepseek-v4-flash", deepseek::Model::V4Flash);
170 models.insert("deepseek-v4-pro", deepseek::Model::V4Pro);
171
172 for available_model in &Self::settings(cx).available_models {
173 models.insert(
174 &available_model.name,
175 deepseek::Model::Custom {
176 name: available_model.name.clone(),
177 display_name: available_model.display_name.clone(),
178 max_tokens: available_model.max_tokens,
179 max_output_tokens: available_model.max_output_tokens,
180 },
181 );
182 }
183
184 models
185 .into_values()
186 .map(|model| self.create_language_model(model))
187 .collect()
188 }
189
190 fn is_authenticated(&self, cx: &App) -> bool {
191 self.state.read(cx).is_authenticated()
192 }
193
194 fn authenticate(&self, cx: &mut App) -> Task<Result<(), AuthenticateError>> {
195 self.state.update(cx, |state, cx| state.authenticate(cx))
196 }
197
198 fn settings_view(&self, cx: &mut App) -> Option<ProviderSettingsView> {
199 let state = self.state.read(cx);
200 Some(ProviderSettingsView::ApiKey(ApiKeyConfiguration::new(
201 state.api_key_state.has_key(),
202 state.api_key_state.is_from_env_var(),
203 state.api_key_state.env_var_name().clone(),
204 "https://platform.deepseek.com/api_keys".into(),
205 )))
206 }
207
208 fn set_api_key(&self, api_key: Option<String>, cx: &mut App) -> Task<Result<()>> {
209 self.state
210 .update(cx, |state, cx| state.set_api_key(api_key, cx))
211 }
212}
213
214pub struct DeepSeekLanguageModel {
215 id: LanguageModelId,
216 model: deepseek::Model,
217 state: Entity<State>,
218 http_client: Arc<dyn HttpClient>,
219 request_limiter: RateLimiter,
220}
221
222impl DeepSeekLanguageModel {
223 fn stream_completion(
224 &self,
225 request: deepseek::Request,
226 cx: &AsyncApp,
227 ) -> BoxFuture<'static, Result<BoxStream<'static, Result<deepseek::StreamResponse>>>> {
228 let http_client = self.http_client.clone();
229
230 let (api_key, api_url, extra_headers) = self.state.read_with(cx, |state, cx| {
231 let api_url = DeepSeekLanguageModelProvider::api_url(cx);
232 let extra_headers = DeepSeekLanguageModelProvider::settings(cx)
233 .custom_headers
234 .clone();
235 (state.api_key_state.key(&api_url), api_url, extra_headers)
236 });
237
238 let future = self.request_limiter.stream(async move {
239 let Some(api_key) = api_key else {
240 return Err(LanguageModelCompletionError::NoApiKey {
241 provider: PROVIDER_NAME,
242 });
243 };
244 let request = deepseek::stream_completion(
245 http_client.as_ref(),
246 &api_url,
247 &api_key,
248 request,
249 &extra_headers,
250 );
251 let response = request.await?;
252 Ok(response)
253 });
254
255 async move { Ok(future.await?.boxed()) }.boxed()
256 }
257}
258
259impl LanguageModel for DeepSeekLanguageModel {
260 fn id(&self) -> LanguageModelId {
261 self.id.clone()
262 }
263
264 fn name(&self) -> LanguageModelName {
265 LanguageModelName::from(self.model.display_name().to_string())
266 }
267
268 fn provider_id(&self) -> LanguageModelProviderId {
269 PROVIDER_ID
270 }
271
272 fn provider_name(&self) -> LanguageModelProviderName {
273 PROVIDER_NAME
274 }
275
276 fn supports_tools(&self) -> bool {
277 true
278 }
279
280 fn supports_streaming_tools(&self) -> bool {
281 true
282 }
283
284 fn supports_thinking(&self) -> bool {
285 matches!(
286 self.model,
287 deepseek::Model::V4Flash | deepseek::Model::V4Pro
288 )
289 }
290
291 fn supported_effort_levels(&self) -> Vec<LanguageModelEffortLevel> {
292 if !self.supports_thinking() {
293 return Vec::new();
294 }
295
296 vec![
297 LanguageModelEffortLevel {
298 name: "High".into(),
299 value: "high".into(),
300 is_default: true,
301 },
302 LanguageModelEffortLevel {
303 name: "Max".into(),
304 value: "max".into(),
305 is_default: false,
306 },
307 ]
308 }
309
310 fn supports_tool_choice(&self, _choice: LanguageModelToolChoice) -> bool {
311 true
312 }
313
314 fn supports_images(&self) -> bool {
315 false
316 }
317
318 fn telemetry_id(&self) -> String {
319 format!("deepseek/{}", self.model.id())
320 }
321
322 fn max_token_count(&self) -> u64 {
323 self.model.max_token_count()
324 }
325
326 fn max_output_tokens(&self) -> Option<u64> {
327 self.model.max_output_tokens()
328 }
329
330 fn stream_completion(
331 &self,
332 request: LanguageModelRequest,
333 cx: &AsyncApp,
334 ) -> BoxFuture<
335 'static,
336 Result<
337 BoxStream<'static, Result<LanguageModelCompletionEvent, LanguageModelCompletionError>>,
338 LanguageModelCompletionError,
339 >,
340 > {
341 let request = match into_deepseek(request, &self.model, self.max_output_tokens()) {
342 Ok(request) => request,
343 Err(error) => return async move { Err(error.into()) }.boxed(),
344 };
345 let stream = self.stream_completion(request, cx);
346
347 async move {
348 let mapper = DeepSeekEventMapper::new();
349 Ok(mapper.map_stream(stream.await?).boxed())
350 }
351 .boxed()
352 }
353}
354
355pub fn into_deepseek(
356 request: LanguageModelRequest,
357 model: &deepseek::Model,
358 max_output_tokens: Option<u64>,
359) -> Result<deepseek::Request> {
360 if request.contains_custom_tool_input() {
361 anyhow::bail!("DeepSeek does not support custom tools");
362 }
363
364 let thinking = deepseek_thinking(model, request.thinking_allowed);
365 let thinking_enabled = thinking
366 .as_ref()
367 .is_some_and(|thinking| thinking.kind == deepseek::ThinkingType::Enabled);
368
369 let mut messages = Vec::new();
370 let mut current_reasoning: Option<String> = None;
371
372 for message in request.messages {
373 for content in message.content {
374 match content {
375 MessageContent::Text(text) => {
376 let should_add = if message.role == Role::User {
377 !text.trim().is_empty()
378 } else {
379 !text.is_empty()
380 };
381
382 if should_add {
383 messages.push(match message.role {
384 Role::User => deepseek::RequestMessage::User { content: text },
385 Role::Assistant => deepseek::RequestMessage::Assistant {
386 content: Some(text),
387 tool_calls: Vec::new(),
388 reasoning_content: current_reasoning.take(),
389 },
390 Role::System => deepseek::RequestMessage::System { content: text },
391 });
392 }
393 }
394 MessageContent::Thinking { text, .. } => {
395 // Accumulate reasoning content for next assistant message
396 current_reasoning.get_or_insert_default().push_str(&text);
397 }
398 MessageContent::RedactedThinking(_) => {}
399 MessageContent::Image(_) => {}
400 MessageContent::Compaction(_) => {}
401 MessageContent::ToolUse(tool_use) => {
402 let input = tool_use
403 .input
404 .as_json()
405 .ok_or_else(|| anyhow!("DeepSeek does not support custom tool calls"))?;
406 let tool_call = deepseek::ToolCall {
407 id: tool_use.id.to_string(),
408 content: deepseek::ToolCallContent::Function {
409 function: deepseek::FunctionContent {
410 name: tool_use.name.to_string(),
411 arguments: serde_json::to_string(input).unwrap_or_default(),
412 },
413 },
414 };
415
416 if let Some(deepseek::RequestMessage::Assistant { tool_calls, .. }) =
417 messages.last_mut()
418 {
419 tool_calls.push(tool_call);
420 } else {
421 messages.push(deepseek::RequestMessage::Assistant {
422 content: None,
423 tool_calls: vec![tool_call],
424 reasoning_content: current_reasoning.take(),
425 });
426 }
427 }
428 MessageContent::ToolResult(tool_result) => {
429 let mut text_parts: Vec<String> = Vec::new();
430 for part in &tool_result.content {
431 match part {
432 LanguageModelToolResultContent::Text(text) => {
433 text_parts.push(text.to_string());
434 }
435 LanguageModelToolResultContent::Image(_) => {
436 text_parts.push("[Tool responded with an image]".to_string());
437 }
438 }
439 }
440 let content = if text_parts.is_empty() {
441 "<Tool returned an empty string>".to_string()
442 } else {
443 text_parts.join("\n")
444 };
445 messages.push(deepseek::RequestMessage::Tool {
446 content,
447 tool_call_id: tool_result.tool_use_id.to_string(),
448 });
449 }
450 }
451 }
452 }
453
454 Ok(deepseek::Request {
455 model: model.id().to_string(),
456 messages,
457 stream: true,
458 max_tokens: max_output_tokens,
459 temperature: if thinking_enabled {
460 None
461 } else {
462 request.temperature
463 },
464 thinking,
465 reasoning_effort: if thinking_enabled {
466 into_deepseek_reasoning_effort(request.thinking_effort.as_deref())
467 } else {
468 None
469 },
470 response_format: None,
471 tool_choice: request.tool_choice.map(|choice| match choice {
472 LanguageModelToolChoice::Auto => deepseek::ToolChoice::Auto,
473 LanguageModelToolChoice::Any => deepseek::ToolChoice::Required,
474 LanguageModelToolChoice::None => deepseek::ToolChoice::None,
475 }),
476 tools: request
477 .tools
478 .into_iter()
479 .map(|tool| {
480 let input_schema = match tool.input {
481 language_model::LanguageModelRequestToolInput::Function {
482 input_schema,
483 ..
484 } => input_schema,
485 language_model::LanguageModelRequestToolInput::Custom { .. } => {
486 return Err(anyhow::anyhow!("DeepSeek does not support custom tools"));
487 }
488 };
489 Ok(deepseek::ToolDefinition::Function {
490 function: deepseek::FunctionDefinition {
491 name: tool.name,
492 description: Some(tool.description),
493 parameters: Some(input_schema),
494 },
495 })
496 })
497 .collect::<Result<_>>()?,
498 })
499}
500
501fn deepseek_thinking(
502 model: &deepseek::Model,
503 thinking_allowed: bool,
504) -> Option<deepseek::Thinking> {
505 let kind = match model {
506 deepseek::Model::V4Flash | deepseek::Model::V4Pro => {
507 if thinking_allowed {
508 deepseek::ThinkingType::Enabled
509 } else {
510 deepseek::ThinkingType::Disabled
511 }
512 }
513 deepseek::Model::Custom { .. } => return None,
514 };
515
516 Some(deepseek::Thinking { kind })
517}
518
519fn into_deepseek_reasoning_effort(effort: Option<&str>) -> Option<deepseek::ReasoningEffort> {
520 match effort {
521 Some("high") => Some(deepseek::ReasoningEffort::High),
522 Some("max") => Some(deepseek::ReasoningEffort::Max),
523 _ => None,
524 }
525}
526
527pub struct DeepSeekEventMapper {
528 tool_calls_by_index: HashMap<usize, RawToolCall>,
529}
530
531impl DeepSeekEventMapper {
532 pub fn new() -> Self {
533 Self {
534 tool_calls_by_index: HashMap::default(),
535 }
536 }
537
538 pub fn map_stream(
539 mut self,
540 events: Pin<Box<dyn Send + Stream<Item = Result<deepseek::StreamResponse>>>>,
541 ) -> impl Stream<Item = Result<LanguageModelCompletionEvent, LanguageModelCompletionError>>
542 {
543 events.flat_map(move |event| {
544 futures::stream::iter(match event {
545 Ok(event) => self.map_event(event),
546 Err(error) => vec![Err(LanguageModelCompletionError::from(error))],
547 })
548 })
549 }
550
551 pub fn map_event(
552 &mut self,
553 event: deepseek::StreamResponse,
554 ) -> Vec<Result<LanguageModelCompletionEvent, LanguageModelCompletionError>> {
555 let Some(choice) = event.choices.first() else {
556 return vec![Err(LanguageModelCompletionError::from(anyhow!(
557 "Response contained no choices"
558 )))];
559 };
560
561 let mut events = Vec::new();
562 if let Some(content) = choice.delta.content.clone()
563 && !content.is_empty()
564 {
565 events.push(Ok(LanguageModelCompletionEvent::Text(content)));
566 }
567
568 if let Some(reasoning_content) = choice.delta.reasoning_content.clone() {
569 events.push(Ok(LanguageModelCompletionEvent::Thinking {
570 text: reasoning_content,
571 signature: None,
572 }));
573 }
574
575 if let Some(tool_calls) = choice.delta.tool_calls.as_ref() {
576 for tool_call in tool_calls {
577 let entry = self.tool_calls_by_index.entry(tool_call.index).or_default();
578
579 if let Some(tool_id) = tool_call.id.clone() {
580 entry.id = tool_id;
581 }
582
583 if let Some(function) = tool_call.function.as_ref() {
584 if let Some(name) = function.name.clone() {
585 entry.name = name;
586 }
587
588 if let Some(arguments) = function.arguments.clone() {
589 entry.arguments.push_str(&arguments);
590 }
591 }
592
593 if !entry.id.is_empty() && !entry.name.is_empty() {
594 if let Ok(input) = serde_json::from_str::<serde_json::Value>(
595 &fix_streamed_json(&entry.arguments),
596 ) {
597 events.push(Ok(LanguageModelCompletionEvent::ToolUse(
598 LanguageModelToolUse {
599 id: entry.id.clone().into(),
600 name: entry.name.as_str().into(),
601 is_input_complete: false,
602 input: language_model::LanguageModelToolUseInput::Json(input),
603 raw_input: entry.arguments.clone(),
604 thought_signature: None,
605 },
606 )));
607 }
608 }
609 }
610 }
611
612 if let Some(usage) = event.usage {
613 events.push(Ok(LanguageModelCompletionEvent::UsageUpdate(TokenUsage {
614 input_tokens: usage.prompt_tokens,
615 output_tokens: usage.completion_tokens,
616 cache_creation_input_tokens: 0,
617 cache_read_input_tokens: 0,
618 })));
619 }
620
621 match choice.finish_reason.as_deref() {
622 Some("stop") => {
623 events.push(Ok(LanguageModelCompletionEvent::Stop(StopReason::EndTurn)));
624 }
625 Some("tool_calls") => {
626 events.extend(self.tool_calls_by_index.drain().map(|(_, tool_call)| {
627 match parse_tool_arguments(&tool_call.arguments) {
628 Ok(input) => Ok(LanguageModelCompletionEvent::ToolUse(
629 LanguageModelToolUse {
630 id: tool_call.id.clone().into(),
631 name: tool_call.name.as_str().into(),
632 is_input_complete: true,
633 input: language_model::LanguageModelToolUseInput::Json(input),
634 raw_input: tool_call.arguments.clone(),
635 thought_signature: None,
636 },
637 )),
638 Err(error) => Ok(LanguageModelCompletionEvent::ToolUseJsonParseError {
639 id: tool_call.id.clone().into(),
640 tool_name: tool_call.name.as_str().into(),
641 raw_input: tool_call.arguments.into(),
642 json_parse_error: error.to_string(),
643 }),
644 }
645 }));
646
647 events.push(Ok(LanguageModelCompletionEvent::Stop(StopReason::ToolUse)));
648 }
649 Some(stop_reason) => {
650 log::error!("Unexpected DeepSeek stop_reason: {stop_reason:?}",);
651 events.push(Ok(LanguageModelCompletionEvent::Stop(StopReason::EndTurn)));
652 }
653 None => {}
654 }
655
656 events
657 }
658}
659