Skip to main content

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}