Skip to main content

lattice_runtime/
pending.rs

1//! `Pending<T>` -- the typed handle returned by every mutating
2//! actor call (DESIGN.md §5.2.1).
3//!
4//! A `Pending` wraps a `tokio::sync::oneshot::Receiver`. Three usage
5//! patterns:
6//!
7//! 1. **Async caller** (LSP client, plugin host, future async UI):
8//!    `pending.await` yields the typed result.
9//! 2. **Sync caller in a tokio context** (test fixtures running
10//!    `#[tokio::test]`): same as above.
11//! 3. **Sync caller outside tokio** (the TUI input loop, which is a
12//!    blocking `crossterm::event::read` loop on the main thread):
13//!    `pending.blocking_recv()` parks the current thread until the
14//!    actor responds. The TUI uses
15//!    [`crate::runtime::block_on`] which forwards to this.
16//!
17//! Errors are kept narrow: [`RuntimeError::ActorGone`] when the
18//! actor task has shut down before it could respond, and
19//! [`RuntimeError::Core`] for any inner [`lattice_core::CoreError`]
20//! (range out of bounds, etc.). The previous `Busy` variant was
21//! removed in audit slice 6 / H3 -- the document actor's mailbox
22//! is now unbounded, so backpressure surfaces as queue depth
23//! rather than per-call drops.
24
25use std::fmt;
26use std::sync::atomic::{AtomicU64, Ordering};
27
28use lattice_core::CoreError;
29use lattice_grammar::CommandError;
30use thiserror::Error;
31use tokio::sync::oneshot;
32
33/// Monotonic id assigned to every actor-bound invocation. Unique
34/// across the process -- not reused if an actor task dies and is
35/// respawned. Useful for telemetry, logging, and (post-Phase-7)
36/// for plugin-side correlation of request/response.
37#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
38pub struct InvocationId(pub u64);
39
40impl InvocationId {
41    /// Allocate the next id. Lock-free.
42    pub fn next() -> Self {
43        static SEQ: AtomicU64 = AtomicU64::new(1);
44        Self(SEQ.fetch_add(1, Ordering::Relaxed))
45    }
46}
47
48impl fmt::Display for InvocationId {
49    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
50        write!(f, "#{}", self.0)
51    }
52}
53
54/// Outcome of an actor-bound mutation. Wraps a oneshot receiver so
55/// the caller can await (or block on) the result.
56///
57/// `Pending` is neither `Clone` nor `Copy` -- the receiver is
58/// single-use, matching the "one response per request" contract.
59/// Dropping a `Pending` cancels the wait but does not interrupt the
60/// actor; the response is silently discarded.
61#[must_use = "the actor result is dropped if the Pending is not awaited or block_on'd"]
62pub struct Pending<T> {
63    pub id: InvocationId,
64    rx: oneshot::Receiver<Result<T, RuntimeError>>,
65    /// 2026-06-02: optional on-success transform applied at
66    /// poll / blocking_recv time. Lets a producer adapt the
67    /// inner actor's result coordinate space WITHOUT spawning
68    /// a second task — the transform runs on whichever thread
69    /// is driving the Pending (the consumer's thread).
70    /// `None` for vanilla actor-bound Pendings.
71    transform: Option<Box<dyn FnOnce(T) -> T + Send + 'static>>,
72}
73
74impl<T> Pending<T> {
75    pub(crate) fn new(id: InvocationId, rx: oneshot::Receiver<Result<T, RuntimeError>>) -> Self {
76        Self {
77            id,
78            rx,
79            transform: None,
80        }
81    }
82
83    /// 2026-06-02: attach an on-success transform that runs
84    /// when the inner result resolves. The transform runs on
85    /// the *consumer's* polling thread, NOT in a separately
86    /// spawned task — critical when the producer is inside a
87    /// `current_thread` runtime that's about to block in
88    /// `block_on`. Use this instead of `Pending::spawn` for
89    /// purely synchronous result-shape adaptation (e.g.
90    /// translating coordinate spaces).
91    ///
92    /// `map_ok` may be called only once per Pending. A second
93    /// call replaces the prior transform (last-write-wins).
94    pub fn map_ok<F>(mut self, f: F) -> Self
95    where
96        F: FnOnce(T) -> T + Send + 'static,
97    {
98        self.transform = Some(Box::new(f));
99        self
100    }
101
102    /// M.1 (2026-05-31): build a `Pending<T>` that resolves
103    /// immediately with `result`. Used by read-only impls of
104    /// [`crate::Document`] (`MultibufferDocumentHandle`) to
105    /// reject writes without ever spawning an actor — every
106    /// mutating method returns `Pending::ready(Err(
107    /// RuntimeError::ReadOnly))`.
108    pub fn ready(result: Result<T, RuntimeError>) -> Self {
109        let (tx, rx) = oneshot::channel();
110        let _ = tx.send(result);
111        Self {
112            id: InvocationId::next(),
113            rx,
114            transform: None,
115        }
116    }
117
118    /// OA.23b: wrap a channel that a task ALREADY RUNNING will
119    /// complete.
120    ///
121    /// The peer of [`Self::spawn`] for producers that have no task
122    /// to spawn — the work is already queued somewhere else and all
123    /// they hold is the reply end. `MultibufferDocumentHandle::
124    /// apply_to_source` is the case: the edit rides the view's
125    /// source-forwarder FIFO (spawned once, on the shared runtime,
126    /// at view construction) and the forwarder answers this channel.
127    ///
128    /// Not a stylistic preference over `spawn`. `spawn` calls
129    /// `tokio::spawn` and so needs a runtime in scope at the point
130    /// of construction; this path is reached from the editor actor —
131    /// a `current_thread` runtime about to `block_on` the result —
132    /// where spawning the awaiter onto the caller's own runtime is
133    /// the deadlock `map_ok` was added to avoid.
134    ///
135    /// A dropped sender resolves to `ActorGone`, as with any other
136    /// `Pending`.
137    pub fn from_channel(rx: oneshot::Receiver<Result<T, RuntimeError>>) -> Self {
138        Self::new(InvocationId::next(), rx)
139    }
140
141    /// M.3 (2026-06-01): build a `Pending<T>` that resolves when
142    /// the spawned future completes. Lets callers compose
143    /// multiple `Pending`s into one without blocking the runtime
144    /// (used by `MultibufferDocumentHandle::apply_edit_batch` to
145    /// fan a translated batch across N source handles, await each
146    /// source's Pending asynchronously, and combine the results).
147    ///
148    /// **Spawns on the SHARED runtime, never on the ambient one**, and that
149    /// is the whole of its correctness.
150    ///
151    /// It used to call bare `tokio::spawn`, which lands the task on whatever
152    /// runtime happens to be in scope at the point of construction. Every
153    /// consumer of a `Pending` is a document operation, and the host reaches
154    /// those from the editor actor — a `current_thread` runtime that then
155    /// **blocks itself** awaiting the result (`block_on(document.save())`).
156    /// The task is queued on the one thread that is now blocked, so nothing
157    /// drives it and the editor wedges permanently, with the operation never
158    /// performed.
159    ///
160    /// That is not a hypothetical. It shipped three times: the pre-M.11
161    /// `apply_edit` freeze, the `undo` freeze, and `MultibufferDocument::save`
162    /// — the user-visible form of the last being "`:w` on the agenda hangs
163    /// forever and the change is not on disk". The first two were fixed by
164    /// making those paths synchronous, which left this constructor still
165    /// carrying the trap for the one caller that genuinely needs async work.
166    ///
167    /// Naming the target runtime fixes the class rather than the instance:
168    /// document actors all live on the shared runtime, so it is where a
169    /// composition over them belongs. It also drops the old "requires a
170    /// current tokio runtime context" requirement — there is nothing left to
171    /// get wrong at a call site.
172    ///
173    /// [`Self::map_ok`] and [`Self::from_channel`] remain the right tools for
174    /// their own cases (a synchronous transform, and a task that is already
175    /// running); they are no longer *workarounds* for this hazard.
176    pub fn spawn<F>(future: F) -> Self
177    where
178        F: std::future::Future<Output = Result<T, RuntimeError>> + Send + 'static,
179        T: Send + 'static,
180    {
181        let (tx, rx) = oneshot::channel();
182        crate::runtime::shared_runtime().spawn(async move {
183            let result = future.await;
184            let _ = tx.send(result);
185        });
186        Self {
187            id: InvocationId::next(),
188            rx,
189            transform: None,
190        }
191    }
192
193    /// Block the current thread until the actor responds. Used by
194    /// the TUI input loop and by tests that don't drive a tokio
195    /// reactor explicitly. Panics only if the oneshot's internal
196    /// invariants are violated, which can't happen in safe code.
197    pub fn blocking_recv(mut self) -> Result<T, RuntimeError> {
198        match self.rx.blocking_recv() {
199            Ok(Ok(t)) => {
200                if let Some(f) = self.transform.take() {
201                    Ok(f(t))
202                } else {
203                    Ok(t)
204                }
205            }
206            Ok(Err(e)) => Err(e),
207            // Sender dropped without sending -- actor died mid-call.
208            Err(_) => Err(RuntimeError::ActorGone),
209        }
210    }
211}
212
213impl<T> std::future::Future for Pending<T> {
214    type Output = Result<T, RuntimeError>;
215
216    fn poll(
217        mut self: std::pin::Pin<&mut Self>,
218        cx: &mut std::task::Context<'_>,
219    ) -> std::task::Poll<Self::Output> {
220        use std::task::Poll;
221        match std::pin::Pin::new(&mut self.rx).poll(cx) {
222            Poll::Ready(Ok(Ok(t))) => {
223                let transformed = if let Some(f) = self.transform.take() {
224                    f(t)
225                } else {
226                    t
227                };
228                Poll::Ready(Ok(transformed))
229            }
230            Poll::Ready(Ok(Err(e))) => Poll::Ready(Err(e)),
231            Poll::Ready(Err(_)) => Poll::Ready(Err(RuntimeError::ActorGone)),
232            Poll::Pending => Poll::Pending,
233        }
234    }
235}
236
237impl<T> fmt::Debug for Pending<T> {
238    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
239        f.debug_struct("Pending").field("id", &self.id).finish()
240    }
241}
242
243/// Failure modes a runtime caller can observe. Kept narrow so the
244/// UI can branch on the failure shape rather than a generic message.
245#[derive(Debug, Error)]
246pub enum RuntimeError {
247    /// The actor task has terminated (panic, drop, or graceful
248    /// shutdown) before it could respond. Treated as a permanent
249    /// failure -- the caller should re-spawn the document.
250    #[error("document actor is no longer running")]
251    ActorGone,
252
253    /// An inner [`lattice_core::CoreError`] from
254    /// `Document::apply_edit` / `undo` / `redo`. The actor is healthy;
255    /// the operation itself was invalid.
256    #[error(transparent)]
257    Core(#[from] CoreError),
258
259    /// A [`lattice_grammar::CommandError`] from a
260    /// `lattice_grammar::execute` dispatch (unknown command, bad
261    /// args, motion out-of-bounds, ...). The actor is healthy; the
262    /// invocation was invalid.
263    #[error(transparent)]
264    Grammar(#[from] CommandError),
265
266    /// M.1 (2026-05-31): write attempted against a read-only
267    /// document. Returned by `MultibufferDocumentHandle`'s
268    /// mutating methods until M.3 lands edit propagation.
269    /// Distinguishes "this buffer doesn't accept writes by
270    /// design" from `ActorGone` (transient / recoverable) and
271    /// `Core` / `Grammar` (the write was tried but invalid).
272    #[error("document is read-only")]
273    ReadOnly,
274
275    /// SS.3 (2026-08-11): a multibuffer save refused one or more
276    /// sources because the file changed on disk after the view
277    /// snapshotted it.
278    ///
279    /// Typed rather than a message blob because the caller wants the
280    /// paths: the user has to go look at those files, and the recovery
281    /// (refresh the view — `gr` / `:copen` / `:search` — which re-reads
282    /// from disk) is per-view, not per-error-string.
283    ///
284    /// The other sources in the same save DID persist; this is a
285    /// partial-success report, not a failed write. See
286    /// `docs/dev/architecture/multibuffer-stale-sources.md` §2.1.
287    #[error(
288        "changed on disk since this view opened, not overwritten: {}. \
289         Refresh the view to pick up the new content.",
290        .paths.iter().map(|p| p.file_name()
291            .map(|n| n.to_string_lossy().into_owned())
292            .unwrap_or_else(|| p.display().to_string()))
293            .collect::<Vec<_>>().join(", ")
294    )]
295    SourcesChangedOnDisk { paths: Vec<std::path::PathBuf> },
296}
297
298#[cfg(test)]
299mod tests {
300    #![allow(clippy::unwrap_used, clippy::panic)]
301    use super::*;
302
303    #[test]
304    fn invocation_ids_are_monotonic() {
305        let a = InvocationId::next();
306        let b = InvocationId::next();
307        let c = InvocationId::next();
308        assert!(a.0 < b.0);
309        assert!(b.0 < c.0);
310    }
311
312    #[tokio::test]
313    async fn pending_resolves_to_sent_value() {
314        let (tx, rx) = oneshot::channel();
315        let p: Pending<i32> = Pending::new(InvocationId::next(), rx);
316        tx.send(Ok(42)).unwrap();
317        assert_eq!(p.await.unwrap(), 42);
318    }
319
320    #[tokio::test]
321    async fn pending_yields_actor_gone_when_sender_dropped() {
322        let (tx, rx) = oneshot::channel::<Result<i32, RuntimeError>>();
323        let p = Pending::new(InvocationId::next(), rx);
324        drop(tx);
325        match p.await {
326            Err(RuntimeError::ActorGone) => {}
327            other => panic!("expected ActorGone, got {other:?}"),
328        }
329    }
330
331    #[test]
332    fn pending_blocking_recv_returns_actor_gone_on_drop() {
333        let (tx, rx) = oneshot::channel::<Result<i32, RuntimeError>>();
334        let p = Pending::new(InvocationId::next(), rx);
335        drop(tx);
336        match p.blocking_recv() {
337            Err(RuntimeError::ActorGone) => {}
338            other => panic!("expected ActorGone, got {other:?}"),
339        }
340    }
341}