Skip to repository content872 lines · 29.2 KB · rust
tenant.openagents/omega
No repository description is available.
OpenAgents Git authority 2026-07-28T02:52:08.944Z 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
log_store.rs
1use std::{
2 collections::VecDeque,
3 sync::Arc,
4 time::{Duration, Instant},
5};
6
7use collections::HashMap;
8use futures::{StreamExt, channel::mpsc};
9use gpui::{
10 App, AppContext as _, Context, Entity, EventEmitter, Global, Subscription, TaskExt, WeakEntity,
11};
12use lsp::{
13 IoKind, LanguageServer, LanguageServerId, LanguageServerName, LanguageServerSelector,
14 MessageType, RequestId, TraceValue,
15};
16use rpc::proto;
17use serde::Deserialize;
18use settings::WorktreeId;
19
20use crate::{LanguageServerLogType, LspStore, Project, ProjectItem as _};
21
22const MAX_STORED_LOG_ENTRIES: usize = 2000;
23const MAX_PENDING_REQUESTS: usize = MAX_STORED_LOG_ENTRIES;
24
25pub fn init(on_headless_host: bool, cx: &mut App) -> Entity<LogStore> {
26 let log_store = cx.new(|cx| LogStore::new(on_headless_host, cx));
27 cx.set_global(GlobalLogStore(log_store.clone()));
28 log_store
29}
30
31pub struct GlobalLogStore(pub Entity<LogStore>);
32
33impl Global for GlobalLogStore {}
34
35#[derive(Debug)]
36pub enum Event {
37 NewServerLogEntry {
38 id: LanguageServerId,
39 kind: LanguageServerLogType,
40 text: String,
41 },
42}
43
44impl EventEmitter<Event> for LogStore {}
45
46pub struct LogStore {
47 on_headless_host: bool,
48 projects: HashMap<WeakEntity<Project>, ProjectState>,
49 pub language_servers: HashMap<LanguageServerId, LanguageServerState>,
50 io_tx: mpsc::UnboundedSender<(LanguageServerId, IoKind, String, Instant)>,
51}
52
53struct ProjectState {
54 _subscriptions: [Subscription; 2],
55 copilot_log_subscription: Option<lsp::Subscription>,
56}
57
58pub trait Message: AsRef<str> {
59 type Level: Copy + std::fmt::Debug;
60 fn should_include(&self, _: Self::Level) -> bool {
61 true
62 }
63}
64
65#[derive(Debug)]
66pub struct LogMessage {
67 message: String,
68 typ: MessageType,
69}
70
71impl AsRef<str> for LogMessage {
72 fn as_ref(&self) -> &str {
73 &self.message
74 }
75}
76
77impl Message for LogMessage {
78 type Level = MessageType;
79
80 fn should_include(&self, level: Self::Level) -> bool {
81 match (self.typ, level) {
82 (MessageType::ERROR, _) => true,
83 (_, MessageType::ERROR) => false,
84 (MessageType::WARNING, _) => true,
85 (_, MessageType::WARNING) => false,
86 (MessageType::INFO, _) => true,
87 (_, MessageType::INFO) => false,
88 _ => true,
89 }
90 }
91}
92
93#[derive(Debug)]
94pub struct TraceMessage {
95 message: String,
96 is_verbose: bool,
97}
98
99impl AsRef<str> for TraceMessage {
100 fn as_ref(&self) -> &str {
101 &self.message
102 }
103}
104
105impl Message for TraceMessage {
106 type Level = TraceValue;
107
108 fn should_include(&self, level: Self::Level) -> bool {
109 match level {
110 TraceValue::Off => false,
111 TraceValue::Messages => !self.is_verbose,
112 TraceValue::Verbose => true,
113 }
114 }
115}
116
117#[derive(Debug)]
118pub struct RpcMessage {
119 message: String,
120}
121
122impl AsRef<str> for RpcMessage {
123 fn as_ref(&self) -> &str {
124 &self.message
125 }
126}
127
128impl Message for RpcMessage {
129 type Level = ();
130}
131
132pub struct LanguageServerState {
133 pub name: Option<LanguageServerName>,
134 pub worktree_id: Option<WorktreeId>,
135 pub kind: LanguageServerKind,
136 log_messages: VecDeque<LogMessage>,
137 trace_messages: VecDeque<TraceMessage>,
138 pub rpc_state: Option<LanguageServerRpcState>,
139 pub trace_level: TraceValue,
140 pub log_level: MessageType,
141 io_logs_subscription: Option<lsp::Subscription>,
142 pub toggled_log_kind: Option<LogKind>,
143}
144
145impl std::fmt::Debug for LanguageServerState {
146 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
147 f.debug_struct("LanguageServerState")
148 .field("name", &self.name)
149 .field("worktree_id", &self.worktree_id)
150 .field("kind", &self.kind)
151 .field("log_messages", &self.log_messages)
152 .field("trace_messages", &self.trace_messages)
153 .field("rpc_state", &self.rpc_state)
154 .field("trace_level", &self.trace_level)
155 .field("log_level", &self.log_level)
156 .field("toggled_log_kind", &self.toggled_log_kind)
157 .finish_non_exhaustive()
158 }
159}
160
161#[derive(PartialEq, Clone)]
162pub enum LanguageServerKind {
163 Local { project: WeakEntity<Project> },
164 Remote { project: WeakEntity<Project> },
165 LocalSsh { lsp_store: WeakEntity<LspStore> },
166 Global,
167}
168
169impl std::fmt::Debug for LanguageServerKind {
170 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
171 match self {
172 LanguageServerKind::Local { .. } => write!(f, "LanguageServerKind::Local"),
173 LanguageServerKind::Remote { .. } => write!(f, "LanguageServerKind::Remote"),
174 LanguageServerKind::LocalSsh { .. } => write!(f, "LanguageServerKind::LocalSsh"),
175 LanguageServerKind::Global => write!(f, "LanguageServerKind::Global"),
176 }
177 }
178}
179
180impl LanguageServerKind {
181 pub fn project(&self) -> Option<&WeakEntity<Project>> {
182 match self {
183 Self::Local { project } => Some(project),
184 Self::Remote { project } => Some(project),
185 Self::LocalSsh { .. } => None,
186 Self::Global { .. } => None,
187 }
188 }
189}
190
191#[derive(Debug)]
192pub struct LanguageServerRpcState {
193 pub rpc_messages: VecDeque<RpcMessage>,
194 last_message_kind: Option<MessageKind>,
195 request_tracker: RpcRequestTracker,
196}
197
198#[derive(Debug, Default)]
199struct RpcRequestTracker {
200 pending_requests: HashMap<PendingRequestKey, Instant>,
201}
202
203#[derive(Debug, Clone, PartialEq, Eq, Hash)]
204struct PendingRequestKey {
205 kind: MessageKind,
206 id: RequestId,
207}
208
209#[derive(Deserialize)]
210struct RpcEnvelope<'a> {
211 id: Option<RequestId>,
212 method: Option<&'a str>,
213}
214
215#[derive(Debug, Copy, Clone, PartialEq, Eq, Hash)]
216enum MessageKind {
217 Send,
218 Receive,
219}
220
221impl MessageKind {
222 fn opposite(self) -> Self {
223 match self {
224 Self::Send => Self::Receive,
225 Self::Receive => Self::Send,
226 }
227 }
228}
229
230impl RpcRequestTracker {
231 fn observe(
232 &mut self,
233 kind: MessageKind,
234 message: &str,
235 observed_at: Instant,
236 ) -> Option<Duration> {
237 let envelope = serde_json::from_str::<RpcEnvelope>(message).ok()?;
238 let id = envelope.id?;
239 if envelope.method.is_some() {
240 self.insert(PendingRequestKey { kind, id }, observed_at);
241 None
242 } else {
243 self.pending_requests
244 .remove(&PendingRequestKey {
245 kind: kind.opposite(),
246 id,
247 })
248 .and_then(|started_at| observed_at.checked_duration_since(started_at))
249 }
250 }
251
252 fn insert(&mut self, key: PendingRequestKey, observed_at: Instant) {
253 if self.pending_requests.len() >= MAX_PENDING_REQUESTS
254 && !self.pending_requests.contains_key(&key)
255 && let Some(oldest_key) = self
256 .pending_requests
257 .iter()
258 .min_by_key(|(_, started_at)| **started_at)
259 .map(|(key, _)| key.clone())
260 {
261 self.pending_requests.remove(&oldest_key);
262 }
263 self.pending_requests.insert(key, observed_at);
264 }
265}
266
267#[cfg(feature = "test-support")]
268#[derive(Default)]
269pub struct TestRpcRequestTracker(RpcRequestTracker);
270
271#[cfg(feature = "test-support")]
272impl TestRpcRequestTracker {
273 pub fn new() -> Self {
274 Self::default()
275 }
276
277 pub fn observe(
278 &mut self,
279 received: bool,
280 message: &str,
281 observed_at: Instant,
282 ) -> Option<Duration> {
283 let kind = if received {
284 MessageKind::Receive
285 } else {
286 MessageKind::Send
287 };
288 self.0.observe(kind, message, observed_at)
289 }
290
291 pub fn pending_request_count(&self) -> usize {
292 self.0.pending_requests.len()
293 }
294
295 pub fn max_pending_requests() -> usize {
296 MAX_PENDING_REQUESTS
297 }
298}
299
300enum RpcTiming {
301 ObservedAt(Instant),
302 Forwarded(Option<Duration>),
303}
304
305fn format_duration(duration: Duration) -> String {
306 let seconds = duration.as_secs_f64();
307 if seconds < 0.001 {
308 format!("{:.0}µs", seconds * 1_000_000.0)
309 } else if seconds < 1.0 {
310 format!("{:.1}ms", seconds * 1_000.0)
311 } else {
312 format!("{seconds:.2}s")
313 }
314}
315
316#[derive(Clone, Copy, Debug, Default, PartialEq)]
317pub enum LogKind {
318 Rpc,
319 Trace,
320 #[default]
321 Logs,
322 ServerInfo,
323}
324
325impl LogKind {
326 pub fn from_server_log_type(log_type: &LanguageServerLogType) -> Self {
327 match log_type {
328 LanguageServerLogType::Log(_) => Self::Logs,
329 LanguageServerLogType::Trace { .. } => Self::Trace,
330 LanguageServerLogType::Rpc { .. } => Self::Rpc,
331 }
332 }
333}
334
335impl LogStore {
336 pub fn new(on_headless_host: bool, cx: &mut Context<Self>) -> Self {
337 let (io_tx, mut io_rx) = mpsc::unbounded();
338
339 let log_store = Self {
340 projects: HashMap::default(),
341 language_servers: HashMap::default(),
342
343 on_headless_host,
344 io_tx,
345 };
346 cx.spawn(async move |log_store, cx| {
347 while let Some((server_id, io_kind, message, observed_at)) = io_rx.next().await {
348 if let Some(log_store) = log_store.upgrade() {
349 log_store.update(cx, |log_store, cx| {
350 log_store.on_io(server_id, io_kind, &message, observed_at, cx);
351 });
352 }
353 }
354 anyhow::Ok(())
355 })
356 .detach_and_log_err(cx);
357
358 log_store
359 }
360
361 pub fn add_project(&mut self, project: &Entity<Project>, cx: &mut Context<Self>) {
362 let weak_project = project.downgrade();
363 self.projects.insert(
364 project.downgrade(),
365 ProjectState {
366 _subscriptions: [
367 cx.observe_release(project, move |this, _, _| {
368 this.projects.remove(&weak_project);
369 this.language_servers
370 .retain(|_, state| state.kind.project() != Some(&weak_project));
371 }),
372 cx.subscribe(project, move |log_store, project, event, cx| {
373 let server_kind = if project.read(cx).is_local() {
374 LanguageServerKind::Local {
375 project: project.downgrade(),
376 }
377 } else {
378 LanguageServerKind::Remote {
379 project: project.downgrade(),
380 }
381 };
382 match event {
383 crate::Event::LanguageServerAdded(id, name, worktree_id) => {
384 log_store.add_language_server(
385 server_kind,
386 *id,
387 Some(name.clone()),
388 *worktree_id,
389 project
390 .read(cx)
391 .lsp_store()
392 .read(cx)
393 .language_server_for_id(*id),
394 cx,
395 );
396 }
397 crate::Event::LanguageServerBufferRegistered {
398 server_id,
399 buffer_id,
400 name,
401 ..
402 } => {
403 let worktree_id = project
404 .read(cx)
405 .buffer_for_id(*buffer_id, cx)
406 .and_then(|buffer| {
407 Some(buffer.read(cx).project_path(cx)?.worktree_id)
408 });
409 let name = name.clone().or_else(|| {
410 project
411 .read(cx)
412 .lsp_store()
413 .read(cx)
414 .language_server_statuses
415 .get(server_id)
416 .map(|status| status.name.clone())
417 });
418 log_store.add_language_server(
419 server_kind,
420 *server_id,
421 name,
422 worktree_id,
423 None,
424 cx,
425 );
426 }
427 crate::Event::LanguageServerRemoved(id) => {
428 log_store.remove_language_server(*id, cx);
429 }
430 crate::Event::LanguageServerLog(id, typ, message) => {
431 log_store.add_language_server(
432 server_kind,
433 *id,
434 None,
435 None,
436 None,
437 cx,
438 );
439 match typ {
440 crate::LanguageServerLogType::Log(typ) => {
441 log_store.add_language_server_log(*id, *typ, message, cx);
442 }
443 crate::LanguageServerLogType::Trace { verbose_info } => {
444 log_store.add_language_server_trace(
445 *id,
446 message,
447 verbose_info.clone(),
448 cx,
449 );
450 }
451 crate::LanguageServerLogType::Rpc { received, elapsed } => {
452 let kind = if *received {
453 MessageKind::Receive
454 } else {
455 MessageKind::Send
456 };
457 log_store.add_language_server_rpc(
458 *id,
459 kind,
460 message,
461 RpcTiming::Forwarded(*elapsed),
462 cx,
463 );
464 }
465 }
466 }
467 crate::Event::ToggleLspLogs {
468 server_id,
469 enabled,
470 toggled_log_kind,
471 } => {
472 log_store.toggle_lsp_logs(*server_id, *enabled, *toggled_log_kind);
473 }
474 _ => {}
475 }
476 }),
477 ],
478 copilot_log_subscription: None,
479 },
480 );
481 }
482
483 pub fn get_language_server_state(
484 &mut self,
485 id: LanguageServerId,
486 ) -> Option<&mut LanguageServerState> {
487 self.language_servers.get_mut(&id)
488 }
489
490 pub fn add_language_server(
491 &mut self,
492 kind: LanguageServerKind,
493 server_id: LanguageServerId,
494 name: Option<LanguageServerName>,
495 worktree_id: Option<WorktreeId>,
496 server: Option<Arc<LanguageServer>>,
497 cx: &mut Context<Self>,
498 ) -> Option<&mut LanguageServerState> {
499 let server_state = self.language_servers.entry(server_id).or_insert_with(|| {
500 cx.notify();
501 LanguageServerState {
502 name: None,
503 worktree_id: None,
504 kind,
505 rpc_state: None,
506 log_messages: VecDeque::with_capacity(MAX_STORED_LOG_ENTRIES),
507 trace_messages: VecDeque::with_capacity(MAX_STORED_LOG_ENTRIES),
508 trace_level: TraceValue::Off,
509 log_level: MessageType::LOG,
510 io_logs_subscription: None,
511 toggled_log_kind: None,
512 }
513 });
514
515 if let Some(name) = name {
516 server_state.name = Some(name);
517 }
518 if let Some(worktree_id) = worktree_id {
519 server_state.worktree_id = Some(worktree_id);
520 }
521
522 if let Some(server) = server.filter(|_| server_state.io_logs_subscription.is_none()) {
523 let io_tx = self.io_tx.clone();
524 let server_id = server.server_id();
525 server_state.io_logs_subscription = Some(server.on_io(move |io_kind, message| {
526 let observed_at = Instant::now();
527 io_tx
528 .unbounded_send((server_id, io_kind, message.to_string(), observed_at))
529 .ok();
530 }));
531 }
532
533 Some(server_state)
534 }
535
536 pub fn add_language_server_log(
537 &mut self,
538 id: LanguageServerId,
539 typ: MessageType,
540 message: &str,
541 cx: &mut Context<Self>,
542 ) -> Option<()> {
543 let store_logs = !self.on_headless_host;
544 let language_server_state = self.get_language_server_state(id)?;
545
546 let log_lines = &mut language_server_state.log_messages;
547 let message = message.trim_end().to_string();
548 if !store_logs {
549 // Send all messages regardless of the visibility in case of not storing, to notify the receiver anyway
550 self.emit_event(
551 Event::NewServerLogEntry {
552 id,
553 kind: LanguageServerLogType::Log(typ),
554 text: message,
555 },
556 cx,
557 );
558 } else if let Some(new_message) = Self::push_new_message(
559 log_lines,
560 LogMessage { message, typ },
561 language_server_state.log_level,
562 ) {
563 self.emit_event(
564 Event::NewServerLogEntry {
565 id,
566 kind: LanguageServerLogType::Log(typ),
567 text: new_message,
568 },
569 cx,
570 );
571 }
572 Some(())
573 }
574
575 fn add_language_server_trace(
576 &mut self,
577 id: LanguageServerId,
578 message: &str,
579 verbose_info: Option<String>,
580 cx: &mut Context<Self>,
581 ) -> Option<()> {
582 let store_logs = !self.on_headless_host;
583 let language_server_state = self.get_language_server_state(id)?;
584
585 let log_lines = &mut language_server_state.trace_messages;
586 if !store_logs {
587 // Send all messages regardless of the visibility in case of not storing, to notify the receiver anyway
588 self.emit_event(
589 Event::NewServerLogEntry {
590 id,
591 kind: LanguageServerLogType::Trace { verbose_info },
592 text: message.trim().to_string(),
593 },
594 cx,
595 );
596 } else if let Some(new_message) = Self::push_new_message(
597 log_lines,
598 TraceMessage {
599 message: message.trim().to_string(),
600 is_verbose: false,
601 },
602 TraceValue::Messages,
603 ) {
604 if let Some(verbose_message) = verbose_info.as_ref() {
605 Self::push_new_message(
606 log_lines,
607 TraceMessage {
608 message: verbose_message.clone(),
609 is_verbose: true,
610 },
611 TraceValue::Verbose,
612 );
613 }
614 self.emit_event(
615 Event::NewServerLogEntry {
616 id,
617 kind: LanguageServerLogType::Trace { verbose_info },
618 text: new_message,
619 },
620 cx,
621 );
622 }
623 Some(())
624 }
625
626 fn push_new_message<T: Message>(
627 log_lines: &mut VecDeque<T>,
628 message: T,
629 current_severity: <T as Message>::Level,
630 ) -> Option<String> {
631 while log_lines.len() + 1 >= MAX_STORED_LOG_ENTRIES {
632 log_lines.pop_front();
633 }
634 let visible = message.should_include(current_severity);
635
636 let visible_message = visible.then(|| message.as_ref().to_string());
637 log_lines.push_back(message);
638 visible_message
639 }
640
641 fn add_language_server_rpc(
642 &mut self,
643 language_server_id: LanguageServerId,
644 kind: MessageKind,
645 message: &str,
646 timing: RpcTiming,
647 cx: &mut Context<'_, Self>,
648 ) {
649 let store_logs = !self.on_headless_host;
650 let Some(state) = self
651 .get_language_server_state(language_server_id)
652 .and_then(|state| state.rpc_state.as_mut())
653 else {
654 return;
655 };
656
657 let elapsed = match timing {
658 RpcTiming::ObservedAt(observed_at) => {
659 state.request_tracker.observe(kind, message, observed_at)
660 }
661 RpcTiming::Forwarded(elapsed) => elapsed,
662 };
663
664 let received = kind == MessageKind::Receive;
665 let direction = if received { "Receive" } else { "Send" };
666 let mut header = None;
667 if state.last_message_kind != Some(kind) || elapsed.is_some() {
668 header = Some(match elapsed {
669 Some(elapsed) => format!("\n// {direction} (took {}):", format_duration(elapsed)),
670 None => format!("\n// {direction}:"),
671 });
672 }
673 state.last_message_kind = Some(kind);
674
675 if store_logs {
676 let rpc_log_lines = &mut state.rpc_messages;
677 while rpc_log_lines.len() + 1 >= MAX_STORED_LOG_ENTRIES {
678 rpc_log_lines.pop_front();
679 }
680 let message = message.trim();
681 rpc_log_lines.push_back(RpcMessage {
682 message: match &header {
683 Some(header) => format!("{header}\n{message}"),
684 None => message.to_owned(),
685 },
686 });
687 }
688
689 if let Some(header) = header {
690 // Do not send a synthetic message over the wire, it will be derived from the actual RPC message
691 cx.emit(Event::NewServerLogEntry {
692 id: language_server_id,
693 kind: LanguageServerLogType::Rpc {
694 received,
695 elapsed: None,
696 },
697 text: header,
698 });
699 }
700
701 self.emit_event(
702 Event::NewServerLogEntry {
703 id: language_server_id,
704 kind: LanguageServerLogType::Rpc { received, elapsed },
705 text: message.to_owned(),
706 },
707 cx,
708 );
709 }
710
711 pub fn remove_language_server(&mut self, id: LanguageServerId, cx: &mut Context<Self>) {
712 self.language_servers.remove(&id);
713 cx.notify();
714 }
715
716 pub fn server_logs(&self, server_id: LanguageServerId) -> Option<&VecDeque<LogMessage>> {
717 Some(&self.language_servers.get(&server_id)?.log_messages)
718 }
719
720 pub fn server_trace(&self, server_id: LanguageServerId) -> Option<&VecDeque<TraceMessage>> {
721 Some(&self.language_servers.get(&server_id)?.trace_messages)
722 }
723
724 pub fn server_ids_for_project<'a>(
725 &'a self,
726 lookup_project: &'a WeakEntity<Project>,
727 ) -> impl Iterator<Item = LanguageServerId> + 'a {
728 self.language_servers
729 .iter()
730 .filter_map(move |(id, state)| match &state.kind {
731 LanguageServerKind::Local { project } | LanguageServerKind::Remote { project } => {
732 if project == lookup_project {
733 Some(*id)
734 } else {
735 None
736 }
737 }
738 LanguageServerKind::Global | LanguageServerKind::LocalSsh { .. } => Some(*id),
739 })
740 }
741
742 pub fn enable_rpc_trace_for_language_server(
743 &mut self,
744 server_id: LanguageServerId,
745 ) -> Option<&mut LanguageServerRpcState> {
746 let rpc_state = self
747 .language_servers
748 .get_mut(&server_id)?
749 .rpc_state
750 .get_or_insert_with(|| LanguageServerRpcState {
751 rpc_messages: VecDeque::with_capacity(MAX_STORED_LOG_ENTRIES),
752 last_message_kind: None,
753 request_tracker: RpcRequestTracker::default(),
754 });
755 Some(rpc_state)
756 }
757
758 pub fn disable_rpc_trace_for_language_server(
759 &mut self,
760 server_id: LanguageServerId,
761 ) -> Option<()> {
762 self.language_servers.get_mut(&server_id)?.rpc_state.take();
763 Some(())
764 }
765
766 pub fn has_server_logs(&self, server: &LanguageServerSelector) -> bool {
767 match server {
768 LanguageServerSelector::Id(id) => self.language_servers.contains_key(id),
769 LanguageServerSelector::Name(name) => self
770 .language_servers
771 .iter()
772 .any(|(_, state)| state.name.as_ref() == Some(name)),
773 }
774 }
775
776 fn on_io(
777 &mut self,
778 language_server_id: LanguageServerId,
779 io_kind: IoKind,
780 message: &str,
781 observed_at: Instant,
782 cx: &mut Context<Self>,
783 ) -> Option<()> {
784 let is_received = match io_kind {
785 IoKind::StdOut => true,
786 IoKind::StdIn => false,
787 IoKind::StdErr => {
788 self.add_language_server_log(language_server_id, MessageType::LOG, message, cx);
789 return Some(());
790 }
791 };
792
793 let kind = if is_received {
794 MessageKind::Receive
795 } else {
796 MessageKind::Send
797 };
798
799 self.add_language_server_rpc(
800 language_server_id,
801 kind,
802 message,
803 RpcTiming::ObservedAt(observed_at),
804 cx,
805 );
806 cx.notify();
807 Some(())
808 }
809
810 fn emit_event(&mut self, e: Event, cx: &mut Context<Self>) {
811 match &e {
812 Event::NewServerLogEntry { id, kind, text } => {
813 if let Some(state) = self.get_language_server_state(*id) {
814 let downstream_client = match &state.kind {
815 LanguageServerKind::Remote { project }
816 | LanguageServerKind::Local { project } => project
817 .upgrade()
818 .map(|project| project.read(cx).lsp_store()),
819 LanguageServerKind::LocalSsh { lsp_store } => lsp_store.upgrade(),
820 LanguageServerKind::Global => None,
821 }
822 .and_then(|lsp_store| lsp_store.read(cx).downstream_client());
823 if let Some((client, project_id)) = downstream_client {
824 if Some(LogKind::from_server_log_type(kind)) == state.toggled_log_kind {
825 client
826 .send(proto::LanguageServerLog {
827 project_id,
828 language_server_id: id.to_proto(),
829 message: text.clone(),
830 log_type: Some(kind.to_proto()),
831 })
832 .ok();
833 }
834 }
835 }
836 }
837 }
838
839 cx.emit(e);
840 }
841
842 pub fn toggle_lsp_logs(
843 &mut self,
844 server_id: LanguageServerId,
845 enabled: bool,
846 toggled_log_kind: LogKind,
847 ) {
848 if let Some(server_state) = self.get_language_server_state(server_id) {
849 if enabled {
850 server_state.toggled_log_kind = Some(toggled_log_kind);
851 } else {
852 server_state.toggled_log_kind = None;
853 }
854 }
855 if LogKind::Rpc == toggled_log_kind {
856 if enabled {
857 self.enable_rpc_trace_for_language_server(server_id);
858 } else {
859 self.disable_rpc_trace_for_language_server(server_id);
860 }
861 }
862 }
863 pub fn copilot_state_for_project(
864 &mut self,
865 project: &WeakEntity<Project>,
866 ) -> Option<&mut Option<lsp::Subscription>> {
867 self.projects
868 .get_mut(project)
869 .map(|project| &mut project.copilot_log_subscription)
870 }
871}
872