lattice_lsp/fan_in.rs
1//! Per-actor `LspDocumentChanged` -> `ActorCmd::RecordEdit` fan-in.
2//!
3//! Subscribes one typed channel per actor to the editor's event
4//! bus (M.5.5: `EventBus::subscribe_typed::<LspDocumentChanged>`).
5//! On every event whose `path` resolves to a URI the actor cares
6//! about, it forwards one `RecordEdit` actor command per applied
7//! edit.
8//!
9//! The publisher (App's `publish_document_changed`) only fires
10//! `LspDocumentChanged` when `lsp-mode` is active for the edited
11//! buffer, so the gate happens at the publish site -- fan_in
12//! never sees edits the user gated off via `:lsp-mode`. The
13//! generic `Event::DocumentChanged` keeps firing for non-LSP
14//! subscribers regardless.
15//!
16//! ## Why per-actor and not one shared dispatcher
17//!
18//! Each actor owns its own DocSync mirror; the only writer to
19//! that mirror is the actor's own task. Routing edits straight
20//! into the actor's mailbox makes the edit path lock-free
21//! end-to-end (publish on the UI thread is a single mutex grab
22//! on the bus inner; the rest is fully async). A central
23//! dispatcher would re-introduce a shared lock on the
24//! `attachments` map and serialise every edit through it -- the
25//! exact contention pattern this refactor exists to remove.
26//!
27//! ## Lifecycle
28//!
29//! - The supervisor calls [`spawn`] right after a new actor is
30//! running. The returned [`SubscriptionId`] is stored next to
31//! the actor so it can be unsubscribed at shutdown.
32//! - The fan-in task exits when the bus drops the channel
33//! (supervisor called `unsubscribe`, dropping the sender) or
34//! when the actor's `record_edit` returns
35//! [`LspError::ActorGone`] (the supervisor dropped the
36//! handle).
37//!
38//! ## Filtering by attached URIs
39//!
40//! The fan-in does *not* know which URIs are attached to which
41//! actor. Instead it forwards every `LspDocumentChanged` whose
42//! `path` is `Some(_)` to its actor; the actor's DocSync
43//! warns + skips on URIs it doesn't track. This trades a small
44//! amount of per-event work (one `Uri` build + one mpsc send)
45//! against keeping the supervisor's attachment map out of the
46//! hot path.
47
48use std::sync::Arc;
49
50use lattice_protocol::edit::{Edit, EditKind};
51use lattice_runtime::{EventBus, SubscriptionId};
52
53use crate::actor::{ServerHandle, uri_from_path};
54use crate::error::LspError;
55use crate::events::LspDocumentChanged;
56use crate::logging::{LogLevel, LogSource};
57
58/// Subscribe `handle` to every `LspDocumentChanged` event on
59/// `bus` (M.5.5; previously `Event::DocumentChanged`) and spawn
60/// a tokio task that forwards them as `RecordEdit` actor
61/// commands. Returns the
62/// subscription id; the supervisor must hand this to
63/// [`EventBus::unsubscribe`] when the actor is dropped to keep
64/// the bus's bucket from accumulating dead entries.
65pub fn spawn(handle: ServerHandle, bus: Arc<EventBus>) -> SubscriptionId {
66 let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel::<LspDocumentChanged>();
67 let sub_id = bus.subscribe_typed(tx);
68
69 let instance = handle.instance();
70 let logger = handle.logger().clone();
71
72 tokio::spawn(async move {
73 while let Some(event) = rx.recv().await {
74 let Some(path) = event.path else {
75 // Scratch buffer / unsaved doc -- no URI to map.
76 continue;
77 };
78 let uri = uri_from_path(&path);
79 for ae in event.edits {
80 let edit = Edit {
81 range: ae.original_range,
82 kind: EditKind::Replace {
83 text: ae.inserted_text,
84 },
85 };
86 if let Err(LspError::ActorGone) = handle.record_edit(uri.clone(), edit) {
87 // The actor has shut down. Stop the fan-in;
88 // the supervisor will unsubscribe when it
89 // notices, but exiting promptly stops us
90 // accumulating events for a dead actor.
91 logger.log(
92 Some(&instance),
93 LogLevel::Debug,
94 LogSource::Client,
95 "fan_in: actor gone; exiting",
96 );
97 return;
98 }
99 }
100 }
101 // Sender dropped -> bus unsubscribed us. Nothing to do.
102 });
103
104 sub_id
105}