1use 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#[derive(Debug)]
53pub enum VmIter {
54 Range {
56 next: i64,
57 end: i64,
58 inclusive: bool,
59 done: bool,
60 },
61 Vec {
63 items: Arc<Vec<VmValue>>,
64 idx: usize,
65 },
66 Dict {
68 entries: Arc<crate::value::DictMap>,
69 keys: Vec<crate::value::HarnStr>,
72 idx: usize,
73 },
74 Chars { s: arcstr::ArcStr, byte_idx: usize },
76 Gen { gen: Arc<VmGenerator> },
78 Stream { stream: Arc<VmStream> },
80 Chan { handle: Arc<VmChannelHandle> },
82 Map { inner: VmIterHandle, f: VmValue },
84 Tap { inner: VmIterHandle, f: VmValue },
86 Filter { inner: VmIterHandle, p: VmValue },
88 Scan {
90 inner: VmIterHandle,
91 acc: VmValue,
92 f: VmValue,
93 },
94 FlatMap {
96 inner: VmIterHandle,
97 f: VmValue,
98 cur: Option<VmIterHandle>,
99 },
100 Take {
102 inner: VmIterHandle,
103 remaining: usize,
104 },
105 Skip {
108 inner: VmIterHandle,
109 remaining: usize,
110 },
111 TakeWhile {
114 inner: VmIterHandle,
115 p: VmValue,
116 done: bool,
117 },
118 TakeUntil { inner: VmIterHandle, p: VmValue },
121 SkipWhile {
124 inner: VmIterHandle,
125 p: VmValue,
126 primed: bool,
127 },
128 Zip { a: VmIterHandle, b: VmIterHandle },
131 Enumerate { inner: VmIterHandle, i: i64 },
133 Chain {
135 a: VmIterHandle,
136 b: VmIterHandle,
137 on_a: bool,
138 },
139 Merge {
141 sources: Vec<Option<VmIterHandle>>,
142 cursor: usize,
143 },
144 Interleave {
146 sources: Vec<Option<VmIterHandle>>,
147 cursor: usize,
148 },
149 Race {
151 sources: Vec<Option<VmIterHandle>>,
152 winner: Option<VmIterHandle>,
153 },
154 Broadcast {
156 shared: Arc<Mutex<VmBroadcastState>>,
157 branch: usize,
158 index: usize,
159 },
160 Throttle {
162 inner: VmIterHandle,
163 interval_ms: u64,
164 next_ready: Option<tokio::time::Instant>,
165 },
166 Debounce { inner: VmIterHandle, window_ms: u64 },
169 Chunks { inner: VmIterHandle, n: usize },
172 Windows {
175 inner: VmIterHandle,
176 n: usize,
177 buf: VecDeque<VmValue>,
178 },
179 Exhausted,
181}
182
183impl VmIter {
184 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
802pub 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 *handle.lock() = state;
818 result
819}
820
821pub 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
834pub 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
1159pub 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}