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}