Skip to repository content487 lines · 17.3 KB · rust
tenant.openagents/omega
No repository description is available.
OpenAgents Git authority 2026-07-28T03:59:42.671Z 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
mercury.rs
1use crate::{
2 DebugEvent, EditPredictionFinishedDebugEvent, EditPredictionId, EditPredictionModelInput,
3 EditPredictionStartedDebugEvent, EditPredictionStore,
4 open_ai_response::text_from_response,
5 prediction::{EditPredictionInputs, EditPredictionResult},
6 zeta::compute_edits,
7};
8use anyhow::{Context as _, Result};
9use cloud_llm_client::EditPredictionRejectReason;
10use credentials_provider::CredentialsProvider;
11use futures::AsyncReadExt as _;
12use gpui::{
13 App, AppContext as _, Context, Entity, Global, SharedString, Task, TaskExt,
14 http_client::{self, AsyncBody, HttpClient, Method, StatusCode},
15};
16use language::{ToOffset, ToPoint as _};
17use language_model::{ApiKeyState, EnvVar, env_var};
18use release_channel::AppVersion;
19use serde::{Deserialize, Serialize};
20use std::{mem, ops::Range, path::Path, sync::Arc};
21use zeta_prompt::Zeta2PromptInput;
22
23const MERCURY_API_URL: &str = "https://api.inceptionlabs.ai/v1/edit/completions";
24
25pub struct Mercury {
26 pub api_token: Entity<ApiKeyState>,
27 payment_required_error: bool,
28}
29
30impl Mercury {
31 pub fn new(cx: &mut App) -> Self {
32 Mercury {
33 api_token: mercury_api_token(cx),
34 payment_required_error: false,
35 }
36 }
37
38 pub fn has_payment_required_error(&self) -> bool {
39 self.payment_required_error
40 }
41
42 pub fn set_payment_required_error(&mut self, payment_required_error: bool) {
43 self.payment_required_error = payment_required_error;
44 }
45
46 pub(crate) fn request_prediction(
47 &mut self,
48 EditPredictionModelInput {
49 buffer,
50 snapshot,
51 position,
52 events,
53 related_files,
54 debug_tx,
55 trigger,
56 ..
57 }: EditPredictionModelInput,
58 credentials_provider: Arc<dyn CredentialsProvider>,
59 cx: &mut Context<EditPredictionStore>,
60 ) -> Task<Result<Option<EditPredictionResult>>> {
61 self.api_token.update(cx, |key_state, cx| {
62 _ = key_state.load_if_needed(MERCURY_CREDENTIALS_URL, |s| s, credentials_provider, cx);
63 });
64 let Some(api_token) = self.api_token.read(cx).key(&MERCURY_CREDENTIALS_URL) else {
65 return Task::ready(Ok(None));
66 };
67 let full_path: Arc<Path> = snapshot
68 .file()
69 .map(|file| file.full_path(cx))
70 .unwrap_or_else(|| "untitled".into())
71 .into();
72
73 let http_client = cx.http_client();
74 let cursor_point = position.to_point(&snapshot);
75 let request_start = cx.background_executor().now();
76 let active_buffer = buffer.clone();
77
78 let result = cx.background_spawn(async move {
79 let cursor_offset = cursor_point.to_offset(&snapshot);
80 let (excerpt_point_range, excerpt_offset_range, cursor_offset_in_excerpt) =
81 crate::cursor_excerpt::compute_cursor_excerpt(&snapshot, cursor_offset);
82
83 let related_files = zeta_prompt::filter_redundant_excerpts(
84 related_files,
85 full_path.as_ref(),
86 excerpt_point_range.start.row..excerpt_point_range.end.row,
87 );
88
89 let cursor_excerpt: Arc<str> = snapshot
90 .text_for_range(excerpt_point_range.clone())
91 .collect::<String>()
92 .into();
93 let syntax_ranges = crate::cursor_excerpt::compute_syntax_ranges(
94 &snapshot,
95 cursor_offset,
96 &excerpt_offset_range,
97 );
98 let excerpt_ranges = zeta_prompt::compute_legacy_excerpt_ranges(
99 &cursor_excerpt,
100 cursor_offset_in_excerpt,
101 &syntax_ranges,
102 );
103
104 let editable_offset_range = (excerpt_offset_range.start
105 + excerpt_ranges.editable_350.start)
106 ..(excerpt_offset_range.start + excerpt_ranges.editable_350.end);
107
108 let inputs = zeta_prompt::Zeta2PromptInput {
109 events,
110 related_files: Some(related_files),
111 cursor_offset_in_excerpt: cursor_point.to_offset(&snapshot)
112 - excerpt_offset_range.start,
113 cursor_path: full_path.clone(),
114 cursor_excerpt,
115 excerpt_start_row: Some(excerpt_point_range.start.row),
116 excerpt_ranges,
117 syntax_ranges: Some(syntax_ranges),
118 active_buffer_diagnostics: vec![],
119 in_open_source_repo: false,
120 can_collect_data: false,
121 repo_url: None,
122 };
123
124 let prompt = build_prompt(&inputs);
125
126 if let Some(debug_tx) = &debug_tx {
127 debug_tx
128 .unbounded_send(DebugEvent::EditPredictionStarted(
129 EditPredictionStartedDebugEvent {
130 buffer: active_buffer.downgrade(),
131 prompt: Some(prompt.clone()),
132 position,
133 },
134 ))
135 .ok();
136 }
137
138 let request_body = open_ai::Request {
139 model: "mercury-coder".into(),
140 messages: vec![open_ai::RequestMessage::User {
141 content: open_ai::MessageContent::Plain(prompt),
142 }],
143 stream: false,
144 stream_options: None,
145 max_completion_tokens: None,
146 max_tokens: None,
147 stop: vec![],
148 temperature: None,
149 tool_choice: None,
150 parallel_tool_calls: None,
151 tools: vec![],
152 prompt_cache_key: None,
153 reasoning_effort: None,
154 service_tier: None,
155 };
156
157 let buf = serde_json::to_vec(&request_body)?;
158 let body: AsyncBody = buf.into();
159
160 let request = http_client::Request::builder()
161 .uri(MERCURY_API_URL)
162 .header("Content-Type", "application/json")
163 .header("Authorization", format!("Bearer {}", api_token))
164 .header("Connection", "keep-alive")
165 .method(Method::POST)
166 .body(body)
167 .context("Failed to create request")?;
168
169 let mut response = http_client
170 .send(request)
171 .await
172 .context("Failed to send request")?;
173
174 let mut body: Vec<u8> = Vec::new();
175 response
176 .body_mut()
177 .read_to_end(&mut body)
178 .await
179 .context("Failed to read response body")?;
180
181 if !response.status().is_success() {
182 if response.status() == StatusCode::PAYMENT_REQUIRED {
183 anyhow::bail!(MercuryPaymentRequiredError(
184 mercury_payment_required_message(&body),
185 ));
186 }
187
188 anyhow::bail!(
189 "Request failed with status: {:?}\nBody: {}",
190 response.status(),
191 String::from_utf8_lossy(&body),
192 );
193 };
194
195 let mut response: open_ai::Response =
196 serde_json::from_slice(&body).context("Failed to parse response")?;
197
198 let id = mem::take(&mut response.id);
199 let response_str = text_from_response(response).unwrap_or_default();
200
201 if let Some(debug_tx) = &debug_tx {
202 debug_tx
203 .unbounded_send(DebugEvent::EditPredictionFinished(
204 EditPredictionFinishedDebugEvent {
205 buffer: active_buffer.downgrade(),
206 model_output: Some(response_str.clone()),
207 position,
208 },
209 ))
210 .ok();
211 }
212
213 let response_str = response_str.strip_prefix("```\n").unwrap_or(&response_str);
214 let response_str = response_str.strip_suffix("\n```").unwrap_or(&response_str);
215
216 let mut edits = Vec::new();
217 const NO_PREDICTION_OUTPUT: &str = "None";
218
219 if response_str != NO_PREDICTION_OUTPUT {
220 let old_text = snapshot
221 .text_for_range(editable_offset_range.clone())
222 .collect::<String>();
223 edits = compute_edits(
224 old_text,
225 &response_str,
226 editable_offset_range.start,
227 &snapshot,
228 );
229 }
230
231 let editable_range = snapshot.anchor_range_inside(editable_offset_range);
232
233 anyhow::Ok((id, edits, snapshot, inputs, editable_range))
234 });
235
236 cx.spawn(async move |ep_store, cx| {
237 let result = result.await.context("Mercury edit prediction failed");
238
239 let has_payment_required_error = result
240 .as_ref()
241 .err()
242 .is_some_and(is_mercury_payment_required_error);
243
244 ep_store.update(cx, |store, cx| {
245 store
246 .mercury
247 .set_payment_required_error(has_payment_required_error);
248 cx.notify();
249 })?;
250
251 let (id, edits, old_snapshot, inputs, editable_range) = result?;
252 anyhow::Ok(Some(
253 EditPredictionResult::new(
254 EditPredictionId(id.into()),
255 &buffer,
256 &old_snapshot,
257 edits.into(),
258 None,
259 Some(editable_range),
260 EditPredictionInputs::V2(inputs),
261 None,
262 trigger,
263 cx.background_executor().now() - request_start,
264 cx,
265 )
266 .await,
267 ))
268 })
269 }
270}
271
272fn build_prompt(inputs: &Zeta2PromptInput) -> String {
273 const RECENTLY_VIEWED_SNIPPETS_START: &str = "<|recently_viewed_code_snippets|>\n";
274 const RECENTLY_VIEWED_SNIPPETS_END: &str = "<|/recently_viewed_code_snippets|>\n";
275 const RECENTLY_VIEWED_SNIPPET_START: &str = "<|recently_viewed_code_snippet|>\n";
276 const RECENTLY_VIEWED_SNIPPET_END: &str = "<|/recently_viewed_code_snippet|>\n";
277 const CURRENT_FILE_CONTENT_START: &str = "<|current_file_content|>\n";
278 const CURRENT_FILE_CONTENT_END: &str = "<|/current_file_content|>\n";
279 const CODE_TO_EDIT_START: &str = "<|code_to_edit|>\n";
280 const CODE_TO_EDIT_END: &str = "<|/code_to_edit|>\n";
281 const EDIT_DIFF_HISTORY_START: &str = "<|edit_diff_history|>\n";
282 const EDIT_DIFF_HISTORY_END: &str = "<|/edit_diff_history|>\n";
283 const CURSOR_TAG: &str = "<|cursor|>";
284 const CODE_SNIPPET_FILE_PATH_PREFIX: &str = "code_snippet_file_path: ";
285 const CURRENT_FILE_PATH_PREFIX: &str = "current_file_path: ";
286
287 let mut prompt = String::new();
288
289 push_delimited(
290 &mut prompt,
291 RECENTLY_VIEWED_SNIPPETS_START..RECENTLY_VIEWED_SNIPPETS_END,
292 |prompt| {
293 for related_file in inputs.related_files.as_deref().unwrap_or_default().iter() {
294 for related_excerpt in &related_file.excerpts {
295 push_delimited(
296 prompt,
297 RECENTLY_VIEWED_SNIPPET_START..RECENTLY_VIEWED_SNIPPET_END,
298 |prompt| {
299 prompt.push_str(CODE_SNIPPET_FILE_PATH_PREFIX);
300 prompt.push_str(related_file.path.to_string_lossy().as_ref());
301 prompt.push('\n');
302 prompt.push_str(related_excerpt.text.as_ref());
303 },
304 );
305 }
306 }
307 },
308 );
309
310 push_delimited(
311 &mut prompt,
312 CURRENT_FILE_CONTENT_START..CURRENT_FILE_CONTENT_END,
313 |prompt| {
314 prompt.push_str(CURRENT_FILE_PATH_PREFIX);
315 prompt.push_str(inputs.cursor_path.as_os_str().to_string_lossy().as_ref());
316 prompt.push('\n');
317
318 let editable_range = &inputs.excerpt_ranges.editable_350;
319 prompt.push_str(&inputs.cursor_excerpt[0..editable_range.start]);
320 push_delimited(prompt, CODE_TO_EDIT_START..CODE_TO_EDIT_END, |prompt| {
321 prompt.push_str(
322 &inputs.cursor_excerpt[editable_range.start..inputs.cursor_offset_in_excerpt],
323 );
324 prompt.push_str(CURSOR_TAG);
325 prompt.push_str(
326 &inputs.cursor_excerpt[inputs.cursor_offset_in_excerpt..editable_range.end],
327 );
328 });
329 prompt.push_str(&inputs.cursor_excerpt[editable_range.end..]);
330 },
331 );
332
333 push_delimited(
334 &mut prompt,
335 EDIT_DIFF_HISTORY_START..EDIT_DIFF_HISTORY_END,
336 |prompt| {
337 for event in inputs.events.iter() {
338 zeta_prompt::write_event(prompt, &event);
339 }
340 },
341 );
342
343 prompt
344}
345
346fn push_delimited(prompt: &mut String, delimiters: Range<&str>, cb: impl FnOnce(&mut String)) {
347 prompt.push_str(delimiters.start);
348 cb(prompt);
349 prompt.push('\n');
350 prompt.push_str(delimiters.end);
351}
352
353pub const MERCURY_CREDENTIALS_URL: SharedString =
354 SharedString::new_static("https://api.inceptionlabs.ai/v1/edit/completions");
355pub const MERCURY_CREDENTIALS_USERNAME: &str = "mercury-api-token";
356
357#[derive(Debug, thiserror::Error)]
358#[error("{0}")]
359struct MercuryPaymentRequiredError(SharedString);
360
361#[derive(Deserialize)]
362struct MercuryErrorResponse {
363 error: MercuryErrorMessage,
364}
365
366#[derive(Deserialize)]
367struct MercuryErrorMessage {
368 message: String,
369}
370
371fn is_mercury_payment_required_error(error: &anyhow::Error) -> bool {
372 error
373 .downcast_ref::<MercuryPaymentRequiredError>()
374 .is_some()
375}
376
377fn mercury_payment_required_message(body: &[u8]) -> SharedString {
378 serde_json::from_slice::<MercuryErrorResponse>(body)
379 .map(|response| response.error.message.into())
380 .unwrap_or_else(|_| String::from_utf8_lossy(body).trim().to_string().into())
381}
382
383pub static MERCURY_TOKEN_ENV_VAR: std::sync::LazyLock<EnvVar> = env_var!("MERCURY_AI_TOKEN");
384
385struct GlobalMercuryApiKey(Entity<ApiKeyState>);
386
387impl Global for GlobalMercuryApiKey {}
388
389pub fn mercury_api_token(cx: &mut App) -> Entity<ApiKeyState> {
390 if let Some(global) = cx.try_global::<GlobalMercuryApiKey>() {
391 return global.0.clone();
392 }
393 let entity =
394 cx.new(|_| ApiKeyState::new(MERCURY_CREDENTIALS_URL, MERCURY_TOKEN_ENV_VAR.clone()));
395 cx.set_global(GlobalMercuryApiKey(entity.clone()));
396 entity
397}
398
399pub fn load_mercury_api_token(cx: &mut App) -> Task<Result<(), language_model::AuthenticateError>> {
400 let credentials_provider = zed_credentials_provider::global(cx);
401 mercury_api_token(cx).update(cx, |key_state, cx| {
402 key_state.load_if_needed(MERCURY_CREDENTIALS_URL, |s| s, credentials_provider, cx)
403 })
404}
405
406const FEEDBACK_API_URL: &str = "https://api-feedback.inceptionlabs.ai/feedback";
407
408#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
409#[serde(rename_all = "snake_case")]
410enum MercuryUserAction {
411 Accept,
412 Reject,
413 Ignore,
414}
415
416#[derive(Serialize)]
417struct FeedbackRequest {
418 request_id: SharedString,
419 provider_name: &'static str,
420 user_action: MercuryUserAction,
421 provider_version: String,
422}
423
424pub(crate) fn edit_prediction_accepted(
425 prediction_id: EditPredictionId,
426 http_client: Arc<dyn HttpClient>,
427 cx: &App,
428) {
429 send_feedback(prediction_id, MercuryUserAction::Accept, http_client, cx);
430}
431
432pub(crate) fn edit_prediction_rejected(
433 prediction_id: EditPredictionId,
434 was_shown: bool,
435 reason: EditPredictionRejectReason,
436 http_client: Arc<dyn HttpClient>,
437 cx: &App,
438) {
439 if !was_shown {
440 return;
441 }
442 let action = match reason {
443 EditPredictionRejectReason::Rejected => MercuryUserAction::Reject,
444 EditPredictionRejectReason::Discarded => MercuryUserAction::Ignore,
445 _ => return,
446 };
447 send_feedback(prediction_id, action, http_client, cx);
448}
449
450fn send_feedback(
451 prediction_id: EditPredictionId,
452 action: MercuryUserAction,
453 http_client: Arc<dyn HttpClient>,
454 cx: &App,
455) {
456 let request_id = prediction_id.0;
457 let app_version = AppVersion::global(cx);
458 cx.background_spawn(async move {
459 let body = FeedbackRequest {
460 request_id,
461 provider_name: "zed",
462 user_action: action,
463 provider_version: app_version.to_string(),
464 };
465
466 let request = http_client::Request::builder()
467 .uri(FEEDBACK_API_URL)
468 .method(Method::POST)
469 .header("Content-Type", "application/json")
470 .body(AsyncBody::from(serde_json::to_vec(&body)?))?;
471
472 let response = http_client.send(request).await?;
473 if !response.status().is_success() {
474 anyhow::bail!("Feedback API returned status: {}", response.status());
475 }
476
477 log::debug!(
478 "Mercury feedback sent: request_id={}, action={:?}",
479 body.request_id,
480 body.user_action
481 );
482
483 anyhow::Ok(())
484 })
485 .detach_and_log_err(cx);
486}
487