Skip to main content

harn_vm/vm/
iter.rs

1//! Lazy iterator protocol for the Harn VM.
2//!
3//! `VmIter` is the backing enum for `VmValue::Iter`. It's a single-pass, fused
4//! iterator; once `next` returns `None` the variant is replaced with
5//! `Exhausted`.
6
7use std::collections::VecDeque;
8use std::sync::Arc;
9
10use parking_lot::Mutex;
11
12use crate::value::{VmChannelHandle, VmError, VmGenerator, VmStream, VmValue};
13
14pub type VmIterHandle = Arc<Mutex<VmIter>>;
15
16fn range_initial_done(start: i64, end: i64, inclusive: bool) -> bool {
17    if inclusive {
18        start > end
19    } else {
20        start >= end
21    }
22}
23
24fn range_next(next: &mut i64, end: i64, inclusive: bool, done: &mut bool) -> Option<i64> {
25    if *done {
26        return None;
27    }
28    let value = *next;
29    let at_end = if inclusive {
30        value >= end
31    } else {
32        value
33            .checked_add(1)
34            .is_none_or(|candidate| candidate >= end)
35    };
36    if at_end {
37        *done = true;
38    } else {
39        *next += 1;
40    }
41    Some(value)
42}
43
44#[derive(Debug)]
45pub struct VmBroadcastState {
46    source: VmIterHandle,
47    buffer: Vec<VmValue>,
48    exhausted: bool,
49}
50
51/// Backing enum for `VmValue::Iter`. See module docs.
52#[derive(Debug)]
53pub enum VmIter {
54    /// Step through a lazy integer range without materializing.
55    Range {
56        next: i64,
57        end: i64,
58        inclusive: bool,
59        done: bool,
60    },
61    /// Snapshot over a shared list / set backing store.
62    Vec {
63        items: Arc<Vec<VmValue>>,
64        idx: usize,
65    },
66    /// Snapshot over a dict; yields `Pair(key, value)` items.
67    Dict {
68        entries: Arc<crate::value::DictMap>,
69        /// `HarnStr` clones of the map's own keys — a refcount bump per key
70        /// to build, and yielding a pair re-shares the same allocation.
71        keys: Vec<crate::value::HarnStr>,
72        idx: usize,
73    },
74    /// Unicode scalar iteration over a string.
75    Chars { s: arcstr::ArcStr, byte_idx: usize },
76    /// Drains a generator's yield channel.
77    Gen { gen: Arc<VmGenerator> },
78    /// Drains a stream's emit channel.
79    Stream { stream: Arc<VmStream> },
80    /// Reads from a channel handle.
81    Chan { handle: Arc<VmChannelHandle> },
82    /// Maps each item through a closure.
83    Map { inner: VmIterHandle, f: VmValue },
84    /// Runs a callback for side effects, then yields the original item.
85    Tap { inner: VmIterHandle, f: VmValue },
86    /// Keeps only items for which the predicate is truthy.
87    Filter { inner: VmIterHandle, p: VmValue },
88    /// Running fold that yields each accumulator.
89    Scan {
90        inner: VmIterHandle,
91        acc: VmValue,
92        f: VmValue,
93    },
94    /// Maps each item to an iterable and flattens one level.
95    FlatMap {
96        inner: VmIterHandle,
97        f: VmValue,
98        cur: Option<VmIterHandle>,
99    },
100    /// Yields up to `remaining` items from `inner`, then becomes Exhausted.
101    Take {
102        inner: VmIterHandle,
103        remaining: usize,
104    },
105    /// Skips the first `remaining` items from `inner` on the first call, then
106    /// forwards. `remaining == 0` is the sentinel for "already primed".
107    Skip {
108        inner: VmIterHandle,
109        remaining: usize,
110    },
111    /// Yields items from `inner` while the predicate is truthy; after the
112    /// first falsy predicate or inner exhaustion, becomes Exhausted.
113    TakeWhile {
114        inner: VmIterHandle,
115        p: VmValue,
116        done: bool,
117    },
118    /// Yields items until the predicate is truthy. The matching sentinel item
119    /// is consumed but not yielded.
120    TakeUntil { inner: VmIterHandle, p: VmValue },
121    /// Discards items while the predicate is truthy; after the first falsy
122    /// item, forwards that item and all subsequent items from `inner`.
123    SkipWhile {
124        inner: VmIterHandle,
125        p: VmValue,
126        primed: bool,
127    },
128    /// Advances two inner iters in lockstep; yields `Pair(a, b)` until either
129    /// side is exhausted.
130    Zip { a: VmIterHandle, b: VmIterHandle },
131    /// Yields `Pair(i, item)` starting at `i = 0`.
132    Enumerate { inner: VmIterHandle, i: i64 },
133    /// Concatenates two iters: drains `a` first, then `b`.
134    Chain {
135        a: VmIterHandle,
136        b: VmIterHandle,
137        on_a: bool,
138    },
139    /// Drains any non-exhausted source in rotating order.
140    Merge {
141        sources: Vec<Option<VmIterHandle>>,
142        cursor: usize,
143    },
144    /// Strict round-robin over non-exhausted sources.
145    Interleave {
146        sources: Vec<Option<VmIterHandle>>,
147        cursor: usize,
148    },
149    /// First source to yield wins; subsequent pulls only read that source.
150    Race {
151        sources: Vec<Option<VmIterHandle>>,
152        winner: Option<VmIterHandle>,
153    },
154    /// One source fanned out into several single-pass branches.
155    Broadcast {
156        shared: Arc<Mutex<VmBroadcastState>>,
157        branch: usize,
158        index: usize,
159    },
160    /// Sleeps between emissions after the first item.
161    Throttle {
162        inner: VmIterHandle,
163        interval_ms: u64,
164        next_ready: Option<tokio::time::Instant>,
165    },
166    /// Coalesces immediately available bursts and emits the last item seen
167    /// after the quiet window.
168    Debounce { inner: VmIterHandle, window_ms: u64 },
169    /// Yields `VmValue::List` batches of up to `n` items from `inner`.
170    /// The final batch may be shorter; empty input yields no batches.
171    Chunks { inner: VmIterHandle, n: usize },
172    /// Yields sliding windows of exactly `n` items from `inner` as `VmValue::List`.
173    /// If the input has fewer than `n` items total, no windows are yielded.
174    Windows {
175        inner: VmIterHandle,
176        n: usize,
177        buf: VecDeque<VmValue>,
178    },
179    /// Terminal state: `next` always returns `None`.
180    Exhausted,
181}
182
183impl VmIter {
184    /// Produce the next value, or `None` when exhausted.
185    ///
186    /// Combinator variants (`Map`, `Filter`, `FlatMap`) invoke user-provided
187    /// closures through the `vm` parameter.
188    pub fn next<'a>(
189        &'a mut self,
190        vm: &'a mut crate::vm::Vm,
191    ) -> std::pin::Pin<
192        Box<dyn std::future::Future<Output = Result<Option<VmValue>, VmError>> + Send + 'a>,
193    > {
194        Box::pin(async move { self.next_impl(vm).await })
195    }
196
197    async fn next_impl(&mut self, vm: &mut crate::vm::Vm) -> Result<Option<VmValue>, VmError> {
198        match self {
199            VmIter::Exhausted => Ok(None),
200            VmIter::Range {
201                next,
202                end,
203                inclusive,
204                done,
205            } => {
206                if let Some(v) = range_next(next, *end, *inclusive, done) {
207                    Ok(Some(VmValue::Int(v)))
208                } else {
209                    *self = VmIter::Exhausted;
210                    Ok(None)
211                }
212            }
213            VmIter::Vec { items, idx } => {
214                if *idx < items.len() {
215                    let v = items[*idx].clone();
216                    *idx += 1;
217                    Ok(Some(v))
218                } else {
219                    *self = VmIter::Exhausted;
220                    Ok(None)
221                }
222            }
223            VmIter::Dict { entries, keys, idx } => {
224                if *idx < keys.len() {
225                    let k = keys[*idx].clone();
226                    let v = entries.get(&k).cloned().unwrap_or(VmValue::Nil);
227                    *idx += 1;
228                    Ok(Some(VmValue::Pair(std::sync::Arc::new((
229                        VmValue::String(k),
230                        v,
231                    )))))
232                } else {
233                    *self = VmIter::Exhausted;
234                    Ok(None)
235                }
236            }
237            VmIter::Chars { s, byte_idx } => {
238                if *byte_idx >= s.len() {
239                    *self = VmIter::Exhausted;
240                    return Ok(None);
241                }
242                let rest = &s[*byte_idx..];
243                if let Some(c) = rest.chars().next() {
244                    *byte_idx += c.len_utf8();
245                    let mut buf = [0u8; 4];
246                    let encoded = c.encode_utf8(&mut buf);
247                    Ok(Some(VmValue::String(arcstr::ArcStr::from(&*encoded))))
248                } else {
249                    *self = VmIter::Exhausted;
250                    Ok(None)
251                }
252            }
253            VmIter::Gen { gen } => {
254                if gen.is_done() {
255                    *self = VmIter::Exhausted;
256                    return Ok(None);
257                }
258                let rx = gen.receiver.clone();
259                let mut guard = rx.lock().await;
260                match guard.recv().await {
261                    Some(Ok(v)) => Ok(Some(v)),
262                    Some(Err(error)) => {
263                        gen.mark_done();
264                        drop(guard);
265                        *self = VmIter::Exhausted;
266                        Err(error)
267                    }
268                    None => {
269                        gen.mark_done();
270                        drop(guard);
271                        *self = VmIter::Exhausted;
272                        Ok(None)
273                    }
274                }
275            }
276            VmIter::Stream { stream } => {
277                if stream.is_done() {
278                    *self = VmIter::Exhausted;
279                    return Ok(None);
280                }
281                let rx = stream.receiver.clone();
282                let mut guard = rx.lock().await;
283                match guard.recv().await {
284                    Some(Ok(v)) => Ok(Some(v)),
285                    Some(Err(error)) => {
286                        stream.mark_done();
287                        drop(guard);
288                        *self = VmIter::Exhausted;
289                        Err(error)
290                    }
291                    None => {
292                        stream.mark_done();
293                        drop(guard);
294                        *self = VmIter::Exhausted;
295                        Ok(None)
296                    }
297                }
298            }
299            VmIter::Map { inner, f } => {
300                let f = f.clone();
301                let item = next_handle(inner, vm).await?;
302                match item {
303                    None => {
304                        *self = VmIter::Exhausted;
305                        Ok(None)
306                    }
307                    Some(v) => {
308                        let out = vm.call_callable_one(&f, &v).await?;
309                        Ok(Some(out))
310                    }
311                }
312            }
313            VmIter::Tap { inner, f } => {
314                let f = f.clone();
315                let item = next_handle(inner, vm).await?;
316                match item {
317                    None => {
318                        *self = VmIter::Exhausted;
319                        Ok(None)
320                    }
321                    Some(v) => {
322                        vm.call_callable_one(&f, &v).await?;
323                        Ok(Some(v))
324                    }
325                }
326            }
327            VmIter::Filter { inner, p } => {
328                let p = p.clone();
329                loop {
330                    let item = next_handle(inner, vm).await?;
331                    match item {
332                        None => {
333                            *self = VmIter::Exhausted;
334                            return Ok(None);
335                        }
336                        Some(v) => {
337                            let keep = vm.call_callable_one(&p, &v).await?;
338                            if keep.is_truthy() {
339                                return Ok(Some(v));
340                            }
341                        }
342                    }
343                }
344            }
345            VmIter::Scan { inner, acc, f } => {
346                let f = f.clone();
347                let item = next_handle(inner, vm).await?;
348                match item {
349                    None => {
350                        *self = VmIter::Exhausted;
351                        Ok(None)
352                    }
353                    Some(v) => {
354                        let next_acc = vm.call_callable_two(&f, acc, &v).await?;
355                        *acc = next_acc.clone();
356                        Ok(Some(next_acc))
357                    }
358                }
359            }
360            VmIter::FlatMap { inner, f, cur } => {
361                let f = f.clone();
362                loop {
363                    if let Some(cur_iter) = cur.clone() {
364                        let item = next_handle(&cur_iter, vm).await?;
365                        if let Some(v) = item {
366                            return Ok(Some(v));
367                        }
368                        *cur = None;
369                    }
370                    let item = next_handle(inner, vm).await?;
371                    match item {
372                        None => {
373                            *self = VmIter::Exhausted;
374                            return Ok(None);
375                        }
376                        Some(v) => {
377                            let result = vm.call_callable_one(&f, &v).await?;
378                            let lifted = iter_from_value(result)?;
379                            if let VmValue::Iter(h) = lifted {
380                                *cur = Some(h);
381                            } else {
382                                return Err(VmError::TypeError(
383                                    "flat_map: expected iterable result".to_string(),
384                                ));
385                            }
386                        }
387                    }
388                }
389            }
390            VmIter::Take { inner, remaining } => {
391                if *remaining == 0 {
392                    *self = VmIter::Exhausted;
393                    return Ok(None);
394                }
395                let item = next_handle(inner, vm).await?;
396                match item {
397                    None => {
398                        *self = VmIter::Exhausted;
399                        Ok(None)
400                    }
401                    Some(v) => {
402                        *remaining -= 1;
403                        if *remaining == 0 {
404                            *self = VmIter::Exhausted;
405                        }
406                        Ok(Some(v))
407                    }
408                }
409            }
410            VmIter::Skip { inner, remaining } => {
411                while *remaining > 0 {
412                    let item = next_handle(inner, vm).await?;
413                    match item {
414                        None => {
415                            *self = VmIter::Exhausted;
416                            return Ok(None);
417                        }
418                        Some(_) => {
419                            *remaining -= 1;
420                        }
421                    }
422                }
423                let item = next_handle(inner, vm).await?;
424                match item {
425                    None => {
426                        *self = VmIter::Exhausted;
427                        Ok(None)
428                    }
429                    Some(v) => Ok(Some(v)),
430                }
431            }
432            VmIter::TakeWhile { inner, p, done } => {
433                if *done {
434                    return Ok(None);
435                }
436                let p = p.clone();
437                let item = next_handle(inner, vm).await?;
438                match item {
439                    None => {
440                        *self = VmIter::Exhausted;
441                        Ok(None)
442                    }
443                    Some(v) => {
444                        let keep = vm.call_callable_one(&p, &v).await?;
445                        if keep.is_truthy() {
446                            Ok(Some(v))
447                        } else {
448                            *self = VmIter::Exhausted;
449                            Ok(None)
450                        }
451                    }
452                }
453            }
454            VmIter::TakeUntil { inner, p } => {
455                let p = p.clone();
456                let item = next_handle(inner, vm).await?;
457                match item {
458                    None => {
459                        *self = VmIter::Exhausted;
460                        Ok(None)
461                    }
462                    Some(v) => {
463                        let stop = vm.call_callable_one(&p, &v).await?;
464                        if stop.is_truthy() {
465                            *self = VmIter::Exhausted;
466                            Ok(None)
467                        } else {
468                            Ok(Some(v))
469                        }
470                    }
471                }
472            }
473            VmIter::SkipWhile { inner, p, primed } => {
474                if *primed {
475                    let item = next_handle(inner, vm).await?;
476                    return match item {
477                        None => {
478                            *self = VmIter::Exhausted;
479                            Ok(None)
480                        }
481                        Some(v) => Ok(Some(v)),
482                    };
483                }
484                let p = p.clone();
485                loop {
486                    let item = next_handle(inner, vm).await?;
487                    match item {
488                        None => {
489                            *self = VmIter::Exhausted;
490                            return Ok(None);
491                        }
492                        Some(v) => {
493                            let drop_it = vm.call_callable_one(&p, &v).await?;
494                            if !drop_it.is_truthy() {
495                                *primed = true;
496                                return Ok(Some(v));
497                            }
498                        }
499                    }
500                }
501            }
502            VmIter::Zip { a, b } => {
503                let ia = next_handle(a, vm).await?;
504                let x = match ia {
505                    None => {
506                        *self = VmIter::Exhausted;
507                        return Ok(None);
508                    }
509                    Some(v) => v,
510                };
511                let ib = next_handle(b, vm).await?;
512                let y = match ib {
513                    None => {
514                        *self = VmIter::Exhausted;
515                        return Ok(None);
516                    }
517                    Some(v) => v,
518                };
519                Ok(Some(VmValue::Pair(std::sync::Arc::new((x, y)))))
520            }
521            VmIter::Enumerate { inner, i } => {
522                let item = next_handle(inner, vm).await?;
523                match item {
524                    None => {
525                        *self = VmIter::Exhausted;
526                        Ok(None)
527                    }
528                    Some(v) => {
529                        let idx = *i;
530                        *i += 1;
531                        Ok(Some(VmValue::Pair(std::sync::Arc::new((
532                            VmValue::Int(idx),
533                            v,
534                        )))))
535                    }
536                }
537            }
538            VmIter::Chain { a, b, on_a } => {
539                if *on_a {
540                    let item = next_handle(a, vm).await?;
541                    if let Some(v) = item {
542                        return Ok(Some(v));
543                    }
544                    *on_a = false;
545                }
546                let item = next_handle(b, vm).await?;
547                match item {
548                    None => {
549                        *self = VmIter::Exhausted;
550                        Ok(None)
551                    }
552                    Some(v) => Ok(Some(v)),
553                }
554            }
555            VmIter::Merge { sources, cursor } => loop {
556                if sources.is_empty() || sources.iter().all(Option::is_none) {
557                    *self = VmIter::Exhausted;
558                    return Ok(None);
559                }
560                let len = sources.len();
561                let mut live = 0usize;
562                for offset in 0..len {
563                    let idx = (*cursor + offset) % len;
564                    let Some(handle) = sources[idx].clone() else {
565                        continue;
566                    };
567                    match try_next_ready(&handle, vm).await? {
568                        Some(v) => {
569                            *cursor = (idx + 1) % len;
570                            return Ok(Some(v));
571                        }
572                        None => {
573                            if is_exhausted_handle(&handle) {
574                                sources[idx] = None;
575                            } else {
576                                live += 1;
577                            }
578                        }
579                    }
580                }
581                if live == 0 {
582                    *self = VmIter::Exhausted;
583                    return Ok(None);
584                }
585                tokio::time::sleep(tokio::time::Duration::from_millis(1)).await;
586            },
587            VmIter::Interleave { sources, cursor } => {
588                if sources.is_empty() || sources.iter().all(Option::is_none) {
589                    *self = VmIter::Exhausted;
590                    return Ok(None);
591                }
592                let len = sources.len();
593                for offset in 0..len {
594                    let idx = (*cursor + offset) % len;
595                    let Some(handle) = sources[idx].clone() else {
596                        continue;
597                    };
598                    match next_handle(&handle, vm).await? {
599                        Some(v) => {
600                            *cursor = (idx + 1) % len;
601                            return Ok(Some(v));
602                        }
603                        None => {
604                            sources[idx] = None;
605                        }
606                    }
607                }
608                *self = VmIter::Exhausted;
609                Ok(None)
610            }
611            VmIter::Race { sources, winner } => {
612                if let Some(handle) = winner.clone() {
613                    let item = next_handle(&handle, vm).await?;
614                    return match item {
615                        Some(v) => Ok(Some(v)),
616                        None => {
617                            *self = VmIter::Exhausted;
618                            Ok(None)
619                        }
620                    };
621                }
622                loop {
623                    let mut live = 0usize;
624                    for source in sources.iter_mut() {
625                        let Some(handle) = source.clone() else {
626                            continue;
627                        };
628                        match try_next_ready(&handle, vm).await? {
629                            Some(v) => {
630                                *winner = Some(handle);
631                                sources.clear();
632                                return Ok(Some(v));
633                            }
634                            None => {
635                                if is_exhausted_handle(&handle) {
636                                    *source = None;
637                                } else {
638                                    live += 1;
639                                }
640                            }
641                        }
642                    }
643                    if live == 0 {
644                        *self = VmIter::Exhausted;
645                        return Ok(None);
646                    }
647                    tokio::time::sleep(tokio::time::Duration::from_millis(1)).await;
648                }
649            }
650            VmIter::Broadcast {
651                shared,
652                branch,
653                index,
654            } => {
655                let _ = branch;
656                loop {
657                    let mut state = std::mem::replace(&mut *shared.lock(), empty_broadcast_state());
658                    if *index < state.buffer.len() {
659                        let item = state.buffer[*index].clone();
660                        *index += 1;
661                        *shared.lock() = state;
662                        return Ok(Some(item));
663                    }
664                    if state.exhausted {
665                        *shared.lock() = state;
666                        *self = VmIter::Exhausted;
667                        return Ok(None);
668                    }
669                    let next = next_handle(&state.source, vm).await;
670                    match next {
671                        Err(err) => {
672                            *shared.lock() = state;
673                            return Err(err);
674                        }
675                        Ok(Some(v)) => {
676                            state.buffer.push(v);
677                            *shared.lock() = state;
678                        }
679                        Ok(None) => {
680                            state.exhausted = true;
681                            *shared.lock() = state;
682                        }
683                    }
684                }
685            }
686            VmIter::Throttle {
687                inner,
688                interval_ms,
689                next_ready,
690            } => {
691                if let Some(ready_at) = next_ready.take() {
692                    let now = tokio::time::Instant::now();
693                    if ready_at > now {
694                        tokio::time::sleep_until(ready_at).await;
695                    }
696                }
697                let item = next_handle(inner, vm).await?;
698                match item {
699                    None => {
700                        *self = VmIter::Exhausted;
701                        Ok(None)
702                    }
703                    Some(v) => {
704                        *next_ready = Some(
705                            tokio::time::Instant::now()
706                                + tokio::time::Duration::from_millis(*interval_ms),
707                        );
708                        Ok(Some(v))
709                    }
710                }
711            }
712            VmIter::Debounce { inner, window_ms } => {
713                let mut last = match next_handle(inner, vm).await? {
714                    Some(v) => v,
715                    None => {
716                        *self = VmIter::Exhausted;
717                        return Ok(None);
718                    }
719                };
720                if *window_ms > 0 {
721                    tokio::time::sleep(tokio::time::Duration::from_millis(*window_ms)).await;
722                }
723                while let Some(v) = try_next_ready(inner, vm).await? {
724                    last = v;
725                }
726                Ok(Some(last))
727            }
728            VmIter::Chunks { inner, n } => {
729                let n = *n;
730                let mut batch: Vec<VmValue> = Vec::with_capacity(n);
731                for _ in 0..n {
732                    let item = next_handle(inner, vm).await?;
733                    match item {
734                        Some(v) => {
735                            batch.push(v);
736                        }
737                        None => break,
738                    }
739                }
740                if batch.is_empty() {
741                    *self = VmIter::Exhausted;
742                    Ok(None)
743                } else {
744                    Ok(Some(VmValue::List(std::sync::Arc::new(batch))))
745                }
746            }
747            VmIter::Windows { inner, n, buf } => {
748                let n = *n;
749                if buf.is_empty() {
750                    while buf.len() < n {
751                        let item = next_handle(inner, vm).await?;
752                        match item {
753                            Some(v) => buf.push_back(v),
754                            None => {
755                                *self = VmIter::Exhausted;
756                                return Ok(None);
757                            }
758                        }
759                    }
760                } else {
761                    let item = next_handle(inner, vm).await?;
762                    match item {
763                        Some(v) => {
764                            buf.pop_front();
765                            buf.push_back(v);
766                        }
767                        None => {
768                            *self = VmIter::Exhausted;
769                            return Ok(None);
770                        }
771                    }
772                }
773                let snapshot: Vec<VmValue> = buf.iter().cloned().collect();
774                Ok(Some(VmValue::List(std::sync::Arc::new(snapshot))))
775            }
776            VmIter::Chan { handle } => {
777                let is_closed = handle.is_closed();
778                let rx = handle.receiver.clone();
779                let mut closed_rx = handle.subscribe_closed();
780                let mut guard = rx.lock().await;
781                let item = if is_closed {
782                    guard.try_recv().ok()
783                } else {
784                    tokio::select! {
785                        item = guard.recv() => item,
786                        _ = closed_rx.changed() => guard.try_recv().ok(),
787                    }
788                };
789                match item {
790                    Some(v) => Ok(Some(v)),
791                    None => {
792                        drop(guard);
793                        *self = VmIter::Exhausted;
794                        Ok(None)
795                    }
796                }
797            }
798        }
799    }
800}
801
802/// Advance a handle without holding the iterator-state lock across the await.
803///
804/// Swaps the iter state out into a local owned value (replacing it with
805/// `Exhausted`), runs `next` on the owned state, then swaps it back. This
806/// avoids blocking re-entrant iterator access while preserving single-pass
807/// semantics: a nested `next` call on the same handle during the await would
808/// see `Exhausted` (the iter protocol doesn't permit re-entrant stepping of
809/// the same handle anyway).
810pub async fn next_handle(
811    handle: &VmIterHandle,
812    vm: &mut crate::vm::Vm,
813) -> Result<Option<VmValue>, VmError> {
814    let mut state = std::mem::replace(&mut *handle.lock(), VmIter::Exhausted);
815    let result = state.next(vm).await;
816    // Restore state unless the inner call replaced it with Exhausted.
817    *handle.lock() = state;
818    result
819}
820
821/// Fully consume an iter handle into a Vec of values.
822pub async fn drain(handle: &VmIterHandle, vm: &mut crate::vm::Vm) -> Result<Vec<VmValue>, VmError> {
823    let mut out = Vec::new();
824    loop {
825        let v = next_handle(handle, vm).await?;
826        match v {
827            Some(v) => out.push(v),
828            None => break,
829        }
830    }
831    Ok(out)
832}
833
834/// Fully consume an iter handle into a Vec, failing before pushing item
835/// `max + 1`.
836pub async fn drain_capped(
837    handle: &VmIterHandle,
838    vm: &mut crate::vm::Vm,
839    max: usize,
840) -> Result<Vec<VmValue>, VmError> {
841    let mut out = Vec::new();
842    loop {
843        let v = next_handle(handle, vm).await?;
844        match v {
845            Some(v) => {
846                if out.len() >= max {
847                    return Err(VmError::Runtime(format!(
848                        "stream.collect: max cap {max} exceeded"
849                    )));
850                }
851                out.push(v);
852            }
853            None => break,
854        }
855    }
856    Ok(out)
857}
858
859pub fn iter_handle_from_value(v: VmValue) -> Result<VmIterHandle, VmError> {
860    match iter_from_value(v)? {
861        VmValue::Iter(handle) => Ok(handle),
862        _ => unreachable!("iter_from_value returns Iter"),
863    }
864}
865
866pub fn broadcast_branches(source: VmIterHandle, n: usize) -> Vec<VmValue> {
867    let shared = Arc::new(Mutex::new(VmBroadcastState {
868        source,
869        buffer: Vec::new(),
870        exhausted: false,
871    }));
872    (0..n)
873        .map(|branch| {
874            VmValue::Iter(Arc::new(Mutex::new(VmIter::Broadcast {
875                shared: Arc::clone(&shared),
876                branch,
877                index: 0,
878            })))
879        })
880        .collect()
881}
882
883fn empty_broadcast_state() -> VmBroadcastState {
884    VmBroadcastState {
885        source: Arc::new(Mutex::new(VmIter::Exhausted)),
886        buffer: Vec::new(),
887        exhausted: true,
888    }
889}
890
891fn try_next_ready<'a>(
892    handle: &'a VmIterHandle,
893    vm: &'a mut crate::vm::Vm,
894) -> std::pin::Pin<
895    Box<dyn std::future::Future<Output = Result<Option<VmValue>, VmError>> + Send + 'a>,
896> {
897    Box::pin(async move {
898        let mut state = std::mem::replace(&mut *handle.lock(), VmIter::Exhausted);
899        let result = state.try_next_ready_impl(vm).await;
900        *handle.lock() = state;
901        result
902    })
903}
904
905fn is_exhausted_handle(handle: &VmIterHandle) -> bool {
906    matches!(&*handle.lock(), VmIter::Exhausted)
907}
908
909impl VmIter {
910    async fn try_next_ready_impl(
911        &mut self,
912        vm: &mut crate::vm::Vm,
913    ) -> Result<Option<VmValue>, VmError> {
914        match self {
915            VmIter::Exhausted => Ok(None),
916            VmIter::Map { inner, f } => {
917                let f = f.clone();
918                match try_next_ready(inner, vm).await? {
919                    Some(v) => Ok(Some(vm.call_callable_one(&f, &v).await?)),
920                    None => {
921                        if is_exhausted_handle(inner) {
922                            *self = VmIter::Exhausted;
923                        }
924                        Ok(None)
925                    }
926                }
927            }
928            VmIter::Tap { inner, f } => {
929                let f = f.clone();
930                match try_next_ready(inner, vm).await? {
931                    Some(v) => {
932                        vm.call_callable_one(&f, &v).await?;
933                        Ok(Some(v))
934                    }
935                    None => {
936                        if is_exhausted_handle(inner) {
937                            *self = VmIter::Exhausted;
938                        }
939                        Ok(None)
940                    }
941                }
942            }
943            VmIter::Filter { inner, p } => {
944                let p = p.clone();
945                loop {
946                    match try_next_ready(inner, vm).await? {
947                        Some(v) => {
948                            let keep = vm.call_callable_one(&p, &v).await?;
949                            if keep.is_truthy() {
950                                return Ok(Some(v));
951                            }
952                        }
953                        None => {
954                            if is_exhausted_handle(inner) {
955                                *self = VmIter::Exhausted;
956                            }
957                            return Ok(None);
958                        }
959                    }
960                }
961            }
962            VmIter::Scan { inner, acc, f } => {
963                let f = f.clone();
964                match try_next_ready(inner, vm).await? {
965                    Some(v) => {
966                        let next_acc = vm.call_callable_two(&f, acc, &v).await?;
967                        *acc = next_acc.clone();
968                        Ok(Some(next_acc))
969                    }
970                    None => {
971                        if is_exhausted_handle(inner) {
972                            *self = VmIter::Exhausted;
973                        }
974                        Ok(None)
975                    }
976                }
977            }
978            VmIter::FlatMap { inner, f, cur } => {
979                let f = f.clone();
980                loop {
981                    if let Some(cur_iter) = cur.clone() {
982                        match try_next_ready(&cur_iter, vm).await? {
983                            Some(v) => return Ok(Some(v)),
984                            None => {
985                                if is_exhausted_handle(&cur_iter) {
986                                    *cur = None;
987                                }
988                                return Ok(None);
989                            }
990                        }
991                    }
992                    match try_next_ready(inner, vm).await? {
993                        Some(v) => {
994                            let result = vm.call_callable_one(&f, &v).await?;
995                            *cur = Some(iter_handle_from_value(result)?);
996                        }
997                        None => {
998                            if is_exhausted_handle(inner) {
999                                *self = VmIter::Exhausted;
1000                            }
1001                            return Ok(None);
1002                        }
1003                    }
1004                }
1005            }
1006            VmIter::Take { inner, remaining } => {
1007                if *remaining == 0 {
1008                    *self = VmIter::Exhausted;
1009                    return Ok(None);
1010                }
1011                match try_next_ready(inner, vm).await? {
1012                    Some(v) => {
1013                        *remaining -= 1;
1014                        if *remaining == 0 {
1015                            *self = VmIter::Exhausted;
1016                        }
1017                        Ok(Some(v))
1018                    }
1019                    None => {
1020                        if is_exhausted_handle(inner) {
1021                            *self = VmIter::Exhausted;
1022                        }
1023                        Ok(None)
1024                    }
1025                }
1026            }
1027            VmIter::TakeUntil { inner, p } => {
1028                let p = p.clone();
1029                match try_next_ready(inner, vm).await? {
1030                    Some(v) => {
1031                        let stop = vm.call_callable_one(&p, &v).await?;
1032                        if stop.is_truthy() {
1033                            *self = VmIter::Exhausted;
1034                            Ok(None)
1035                        } else {
1036                            Ok(Some(v))
1037                        }
1038                    }
1039                    None => {
1040                        if is_exhausted_handle(inner) {
1041                            *self = VmIter::Exhausted;
1042                        }
1043                        Ok(None)
1044                    }
1045                }
1046            }
1047            VmIter::Throttle {
1048                inner,
1049                interval_ms,
1050                next_ready,
1051            } => {
1052                if let Some(ready_at) = *next_ready {
1053                    if ready_at > tokio::time::Instant::now() {
1054                        return Ok(None);
1055                    }
1056                }
1057                match try_next_ready(inner, vm).await? {
1058                    Some(v) => {
1059                        *next_ready = Some(
1060                            tokio::time::Instant::now()
1061                                + tokio::time::Duration::from_millis(*interval_ms),
1062                        );
1063                        Ok(Some(v))
1064                    }
1065                    None => {
1066                        if is_exhausted_handle(inner) {
1067                            *self = VmIter::Exhausted;
1068                        }
1069                        Ok(None)
1070                    }
1071                }
1072            }
1073            VmIter::Range { .. }
1074            | VmIter::Vec { .. }
1075            | VmIter::Dict { .. }
1076            | VmIter::Chars { .. }
1077            | VmIter::Skip { .. }
1078            | VmIter::TakeWhile { .. }
1079            | VmIter::SkipWhile { .. }
1080            | VmIter::Zip { .. }
1081            | VmIter::Enumerate { .. }
1082            | VmIter::Chain { .. }
1083            | VmIter::Merge { .. }
1084            | VmIter::Interleave { .. }
1085            | VmIter::Race { .. }
1086            | VmIter::Broadcast { .. }
1087            | VmIter::Debounce { .. }
1088            | VmIter::Chunks { .. }
1089            | VmIter::Windows { .. } => self.next(vm).await,
1090            VmIter::Gen { gen } => {
1091                if gen.is_done() {
1092                    *self = VmIter::Exhausted;
1093                    return Ok(None);
1094                }
1095                let rx = gen.receiver.clone();
1096                let result = match rx.try_lock() {
1097                    Ok(mut guard) => match guard.try_recv() {
1098                        Ok(Ok(v)) => Ok(Some(v)),
1099                        Ok(Err(error)) => {
1100                            gen.mark_done();
1101                            *self = VmIter::Exhausted;
1102                            Err(error)
1103                        }
1104                        Err(tokio::sync::mpsc::error::TryRecvError::Empty) => Ok(None),
1105                        Err(tokio::sync::mpsc::error::TryRecvError::Disconnected) => {
1106                            gen.mark_done();
1107                            *self = VmIter::Exhausted;
1108                            Ok(None)
1109                        }
1110                    },
1111                    Err(_) => Ok(None),
1112                };
1113                result
1114            }
1115            VmIter::Stream { stream } => {
1116                if stream.is_done() {
1117                    *self = VmIter::Exhausted;
1118                    return Ok(None);
1119                }
1120                let rx = stream.receiver.clone();
1121                let result = match rx.try_lock() {
1122                    Ok(mut guard) => match guard.try_recv() {
1123                        Ok(Ok(v)) => Ok(Some(v)),
1124                        Ok(Err(error)) => {
1125                            stream.mark_done();
1126                            *self = VmIter::Exhausted;
1127                            Err(error)
1128                        }
1129                        Err(tokio::sync::mpsc::error::TryRecvError::Empty) => Ok(None),
1130                        Err(tokio::sync::mpsc::error::TryRecvError::Disconnected) => {
1131                            stream.mark_done();
1132                            *self = VmIter::Exhausted;
1133                            Ok(None)
1134                        }
1135                    },
1136                    Err(_) => Ok(None),
1137                };
1138                result
1139            }
1140            VmIter::Chan { handle } => {
1141                let rx = handle.receiver.clone();
1142                let result = match rx.try_lock() {
1143                    Ok(mut guard) => match guard.try_recv() {
1144                        Ok(v) => Ok(Some(v)),
1145                        Err(tokio::sync::mpsc::error::TryRecvError::Empty) => Ok(None),
1146                        Err(tokio::sync::mpsc::error::TryRecvError::Disconnected) => {
1147                            *self = VmIter::Exhausted;
1148                            Ok(None)
1149                        }
1150                    },
1151                    Err(_) => Ok(None),
1152                };
1153                result
1154            }
1155        }
1156    }
1157}
1158
1159/// Convenience: wrap a source value into a `VmValue::Iter`. Used by the
1160/// `iter()` builtin and by combinator/sink implementations in later steps.
1161pub fn iter_from_value(v: VmValue) -> Result<VmValue, VmError> {
1162    let inner = match v {
1163        VmValue::Iter(h) => return Ok(VmValue::Iter(h)),
1164        VmValue::Range(r) => VmIter::Range {
1165            next: r.start,
1166            end: r.end,
1167            inclusive: r.inclusive,
1168            done: range_initial_done(r.start, r.end, r.inclusive),
1169        },
1170        VmValue::List(items) => VmIter::Vec { items, idx: 0 },
1171        VmValue::Set(set) => VmIter::Vec {
1172            items: set.shared_items(),
1173            idx: 0,
1174        },
1175        VmValue::Dict(entries) => {
1176            let keys: Vec<crate::value::HarnStr> = entries.keys().cloned().collect();
1177            VmIter::Dict {
1178                entries,
1179                keys,
1180                idx: 0,
1181            }
1182        }
1183        VmValue::String(s) => VmIter::Chars { s, byte_idx: 0 },
1184        VmValue::Generator(gen) => VmIter::Gen { gen },
1185        VmValue::Stream(stream) => VmIter::Stream { stream },
1186        VmValue::Channel(handle) => VmIter::Chan { handle },
1187        other => {
1188            return Err(VmError::TypeError(format!(
1189                "iter: value of type {} is not iterable",
1190                other.type_name()
1191            )))
1192        }
1193    };
1194    Ok(VmValue::Iter(Arc::new(Mutex::new(inner))))
1195}
1196
1197#[cfg(test)]
1198mod tests {
1199    use super::*;
1200    use crate::value::VmRange;
1201
1202    fn run_iter_test(test: impl std::future::Future<Output = ()>) {
1203        let rt = tokio::runtime::Builder::new_current_thread()
1204            .enable_all()
1205            .build()
1206            .unwrap();
1207        rt.block_on(test);
1208    }
1209
1210    #[test]
1211    fn inclusive_range_at_i64_max_yields_the_endpoint() {
1212        run_iter_test(async {
1213            let mut vm = crate::vm::Vm::new();
1214            let VmValue::Iter(handle) = iter_from_value(VmValue::range(VmRange {
1215                start: i64::MAX,
1216                end: i64::MAX,
1217                inclusive: true,
1218            }))
1219            .unwrap() else {
1220                panic!("expected iter");
1221            };
1222
1223            let first = next_handle(&handle, &mut vm).await.unwrap();
1224            assert!(matches!(first, Some(VmValue::Int(value)) if value == i64::MAX));
1225            assert!(next_handle(&handle, &mut vm).await.unwrap().is_none());
1226        });
1227    }
1228
1229    #[test]
1230    fn exclusive_range_at_i64_max_is_empty() {
1231        run_iter_test(async {
1232            let mut vm = crate::vm::Vm::new();
1233            let VmValue::Iter(handle) = iter_from_value(VmValue::range(VmRange {
1234                start: i64::MAX,
1235                end: i64::MAX,
1236                inclusive: false,
1237            }))
1238            .unwrap() else {
1239                panic!("expected iter");
1240            };
1241
1242            assert!(next_handle(&handle, &mut vm).await.unwrap().is_none());
1243        });
1244    }
1245}