11namespace SpacetimeDB . Internal ;
22
33using System . Buffers ;
4+ using System . Collections ;
45using SpacetimeDB . BSATN ;
56
6- internal abstract class RawTableIterBase < T >
7+ internal abstract class RawTableIterBase < T > : IEnumerable < T >
78 where T : IStructuralReadWrite , new ( )
89{
9- public sealed class Enumerator ( FFI . RowIter handle ) : IDisposable
10- {
11- private const int InitialBufferSize = 1024 ;
12- private byte [ ] ? buffer = ArrayPool < byte > . Shared . Rent ( InitialBufferSize ) ;
13- public ArraySegment < byte > Current { get ; private set ; } = ArraySegment < byte > . Empty ;
14-
15- public bool MoveNext ( )
16- {
17- if ( handle == FFI . RowIter . INVALID )
18- {
19- return false ;
20- }
10+ private const int InitialBufferSize = 1024 ;
2111
22- if ( buffer is null )
23- {
24- return false ;
25- }
12+ protected abstract void IterStart ( out FFI . RowIter handle ) ;
2613
27- uint buffer_len ;
28- while ( true )
14+ public IEnumerator < T > GetEnumerator ( )
15+ {
16+ IterStart ( out var handle ) ;
17+ var buffer = ArrayPool < byte > . Shared . Rent ( InitialBufferSize ) ;
18+ try
19+ {
20+ while ( handle != FFI . RowIter . INVALID )
2921 {
3022 var requested_len = ( uint ) buffer . Length ;
31- buffer_len = requested_len ;
23+ var buffer_len = requested_len ;
3224 var ret = FFI . row_iter_bsatn_advance ( handle , buffer , ref buffer_len ) ;
3325 if ( ret == Errno . EXHAUSTED )
3426 {
@@ -38,82 +30,50 @@ public bool MoveNext()
3830 buffer_len = 0 ;
3931 }
4032 }
33+
4134 // On success, the only way `buffer_len == 0` is for the iterator to be exhausted.
4235 // This happens when the host iterator was empty from the start.
4336 System . Diagnostics . Debug . Assert ( ! ( ret == Errno . OK && buffer_len == 0 ) ) ;
4437 switch ( ret )
4538 {
46- // Iterator advanced and may also be `EXHAUSTED`.
47- // When `OK`, we'll need to advance the iterator in the next call to `MoveNext`.
48- // In both cases, update `Current` to point at the valid range in the scratch `buffer`.
4939 case Errno . EXHAUSTED
5040 or Errno . OK :
51- Current = new ArraySegment < byte > ( buffer , 0 , ( int ) buffer_len ) ;
52- return buffer_len != 0 ;
53- // Couldn't find the iterator, error!
54- case Errno . NO_SUCH_ITER :
55- throw new NoSuchIterException ( ) ;
56- // The scratch `buffer` is too small to fit a row / chunk.
57- // Grow `buffer` and try again.
58- // The `buffer_len` will have been updated with the necessary size.
41+ {
42+ using var stream = new MemoryStream (
43+ buffer ,
44+ 0 ,
45+ ( int ) buffer_len ,
46+ writable : false ,
47+ publiclyVisible : true
48+ ) ;
49+ using var reader = new BinaryReader ( stream ) ;
50+ while ( stream . Position < stream . Length )
51+ {
52+ yield return IStructuralReadWrite . Read < T > ( reader ) ;
53+ }
54+ break ;
55+ }
5956 case Errno . BUFFER_TOO_SMALL :
6057 ArrayPool < byte > . Shared . Return ( buffer ) ;
6158 buffer = ArrayPool < byte > . Shared . Rent ( ( int ) buffer_len ) ;
62- continue ;
59+ break ;
6360 default :
64- throw new UnknownException ( ret ) ;
61+ ret . Check ( ) ;
62+ break ;
6563 }
6664 }
6765 }
68-
69- public void Dispose ( )
66+ finally
7067 {
7168 if ( handle != FFI . RowIter . INVALID )
7269 {
7370 FFI . row_iter_bsatn_close ( handle ) ;
74- handle = FFI . RowIter . INVALID ;
7571 }
76-
77- if ( buffer is not null )
78- {
79- ArrayPool < byte > . Shared . Return ( buffer ) ;
80- buffer = null ;
81- }
82- }
83-
84- public void Reset ( )
85- {
86- throw new NotImplementedException ( ) ;
72+ ArrayPool < byte > . Shared . Return ( buffer ) ;
8773 }
8874 }
8975
90- protected abstract void IterStart ( out FFI . RowIter handle ) ;
91-
92- // Note: using the GetEnumerator() duck-typing protocol instead of IEnumerable to avoid extra boxing.
93- public Enumerator GetEnumerator ( )
94- {
95- IterStart ( out var handle ) ;
96- return new ( handle ) ;
97- }
98-
99- public IEnumerable < T > Parse ( )
100- {
101- foreach ( var chunk in this )
102- {
103- using var stream = new MemoryStream (
104- chunk . Array ! ,
105- chunk . Offset ,
106- chunk . Count ,
107- writable : false ,
108- publiclyVisible : true
109- ) ;
110- using var reader = new BinaryReader ( stream ) ;
111- while ( stream . Position < stream . Length )
112- {
113- yield return IStructuralReadWrite . Read < T > ( reader ) ;
114- }
115- }
116- }
76+ IEnumerator IEnumerable . GetEnumerator ( ) => GetEnumerator ( ) ;
11777}
11878
11979public interface ITableView < View , T >
@@ -166,7 +126,7 @@ protected static ulong DoCount()
166126 return count ;
167127 }
168128
169- protected static IEnumerable < T > DoIter ( ) => new RawTableIter ( tableId ) . Parse ( ) ;
129+ protected static IEnumerable < T > DoIter ( ) => new RawTableIter ( tableId ) ;
170130
171131 protected static T DoInsert ( T row )
172132 {
0 commit comments