Skip to repository content

tenant.openagents/omega

No repository description is available.

OpenAgents Git authority 2026-07-28T02:52:08.944Z 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

log_store.rs

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