mz_row_spine/
arc_batch.rs1use std::sync::Arc;
28
29use differential_dataflow::trace::{
30 Batch, BatchReader, Builder, Cursor, Description, Merger, Navigable,
31};
32use timely::progress::{Antichain, frontier::AntichainRef};
33
34pub struct ArcBatch<B>(pub Arc<B>);
39
40impl<B> Clone for ArcBatch<B> {
45 fn clone(&self) -> Self {
46 ArcBatch(Arc::clone(&self.0))
47 }
48}
49
50impl<B> ArcBatch<B> {
51 pub fn new(batch: B) -> Self {
53 ArcBatch(Arc::new(batch))
54 }
55}
56
57impl<B> std::ops::Deref for ArcBatch<B> {
58 type Target = B;
59 fn deref(&self) -> &B {
60 &self.0
61 }
62}
63
64impl<B: BatchReader + Navigable> Navigable for ArcBatch<B> {
65 type Cursor = ArcBatchCursor<B::Cursor>;
66 fn cursor(&self) -> Self::Cursor {
67 ArcBatchCursor::new(<B as Navigable>::cursor(&self.0))
70 }
71}
72
73impl<B: BatchReader> BatchReader for ArcBatch<B> {
74 type Time = B::Time;
75 fn len(&self) -> usize {
76 self.0.len()
77 }
78 fn description(&self) -> &Description<Self::Time> {
79 self.0.description()
80 }
81}
82
83pub struct ArcBatchCursor<C> {
85 cursor: C,
86}
87
88impl<C> ArcBatchCursor<C> {
89 fn new(cursor: C) -> Self {
90 ArcBatchCursor { cursor }
91 }
92}
93
94impl<C: Cursor> Cursor for ArcBatchCursor<C> {
95 type Storage = ArcBatch<C::Storage>;
96
97 type Key<'a> = C::Key<'a>;
98 type ValOwn = C::ValOwn;
99 type Val<'a> = C::Val<'a>;
100 type Time = C::Time;
101 type TimeGat<'a> = C::TimeGat<'a>;
102 type Diff = C::Diff;
103 type DiffGat<'a> = C::DiffGat<'a>;
104 type KeyContainer = C::KeyContainer;
105 type ValContainer = C::ValContainer;
106 type TimeContainer = C::TimeContainer;
107 type DiffContainer = C::DiffContainer;
108
109 #[inline]
110 fn key_valid(&self, storage: &Self::Storage) -> bool {
111 self.cursor.key_valid(&storage.0)
112 }
113 #[inline]
114 fn val_valid(&self, storage: &Self::Storage) -> bool {
115 self.cursor.val_valid(&storage.0)
116 }
117
118 #[inline]
119 fn key<'a>(&self, storage: &'a Self::Storage) -> Self::Key<'a> {
120 self.cursor.key(&storage.0)
121 }
122 #[inline]
123 fn val<'a>(&self, storage: &'a Self::Storage) -> Self::Val<'a> {
124 self.cursor.val(&storage.0)
125 }
126
127 #[inline]
128 fn get_key<'a>(&self, storage: &'a Self::Storage) -> Option<Self::Key<'a>> {
129 self.cursor.get_key(&storage.0)
130 }
131 #[inline]
132 fn get_val<'a>(&self, storage: &'a Self::Storage) -> Option<Self::Val<'a>> {
133 self.cursor.get_val(&storage.0)
134 }
135
136 #[inline]
137 fn map_times<L: FnMut(Self::TimeGat<'_>, Self::DiffGat<'_>)>(
138 &mut self,
139 storage: &Self::Storage,
140 logic: L,
141 ) {
142 self.cursor.map_times(&storage.0, logic)
143 }
144
145 #[inline]
146 fn step_key(&mut self, storage: &Self::Storage) {
147 self.cursor.step_key(&storage.0)
148 }
149 #[inline]
150 fn seek_key(&mut self, storage: &Self::Storage, key: Self::Key<'_>) {
151 self.cursor.seek_key(&storage.0, key)
152 }
153
154 #[inline]
155 fn step_val(&mut self, storage: &Self::Storage) {
156 self.cursor.step_val(&storage.0)
157 }
158 #[inline]
159 fn seek_val(&mut self, storage: &Self::Storage, val: Self::Val<'_>) {
160 self.cursor.seek_val(&storage.0, val)
161 }
162
163 #[inline]
164 fn rewind_keys(&mut self, storage: &Self::Storage) {
165 self.cursor.rewind_keys(&storage.0)
166 }
167 #[inline]
168 fn rewind_vals(&mut self, storage: &Self::Storage) {
169 self.cursor.rewind_vals(&storage.0)
170 }
171}
172
173impl<B: Batch> Batch for ArcBatch<B> {
174 type Merger = ArcMerger<B>;
175 fn empty(lower: Antichain<Self::Time>, upper: Antichain<Self::Time>) -> Self {
176 ArcBatch::new(B::empty(lower, upper))
177 }
178}
179
180pub struct ArcBuilder<B: Builder> {
182 builder: B,
183}
184
185impl<B: Builder> Builder for ArcBuilder<B> {
186 type Input = B::Input;
187 type Time = B::Time;
188 type Output = ArcBatch<B::Output>;
189 fn with_capacity(keys: usize, vals: usize, upds: usize) -> Self {
190 ArcBuilder {
191 builder: B::with_capacity(keys, vals, upds),
192 }
193 }
194 fn push(&mut self, input: &mut Self::Input) {
195 self.builder.push(input)
196 }
197 fn done(self, description: Description<Self::Time>) -> ArcBatch<B::Output> {
198 ArcBatch::new(self.builder.done(description))
199 }
200 fn seal(chain: &mut Vec<Self::Input>, description: Description<Self::Time>) -> Self::Output {
201 ArcBatch::new(B::seal(chain, description))
202 }
203}
204
205pub struct ArcMerger<B: Batch> {
207 merger: B::Merger,
208}
209
210impl<B: Batch> Merger<ArcBatch<B>> for ArcMerger<B> {
211 fn new(
212 source1: &ArcBatch<B>,
213 source2: &ArcBatch<B>,
214 compaction_frontier: AntichainRef<B::Time>,
215 ) -> Self {
216 ArcMerger {
217 merger: B::begin_merge(&source1.0, &source2.0, compaction_frontier),
218 }
219 }
220 fn work(&mut self, source1: &ArcBatch<B>, source2: &ArcBatch<B>, fuel: &mut isize) {
221 self.merger.work(&source1.0, &source2.0, fuel)
222 }
223 fn done(self) -> ArcBatch<B> {
224 ArcBatch::new(self.merger.done())
225 }
226}
227
228#[cfg(test)]
229mod tests {
230 use differential_dataflow::trace::cursor::Cursor;
231 use differential_dataflow::trace::implementations::ord_neu::OrdValBatcher;
232 use differential_dataflow::trace::{Batcher, Builder, Navigable};
233 use timely::container::PushInto;
234 use timely::progress::Antichain;
235
236 use crate::ArcOrdValBuilder;
237
238 #[mz_ore::test]
247 fn arc_batch_reads_from_other_thread() {
248 fn assert_send_sync<T: Send + Sync>(_: &T) {}
249
250 let mut batcher = OrdValBatcher::<u64, u64, usize, i64>::new(None, 0);
251 batcher.push_into(vec![((1, 2), 0, 1), ((2, 3), 1, 1)]);
252 let (mut chain, description) = batcher.seal(Antichain::from_elem(2));
253 let batch = ArcOrdValBuilder::<u64, u64, usize, i64>::seal(&mut chain, description);
254
255 assert_send_sync(&batch);
256
257 let read = std::thread::spawn(move || {
258 let mut cursor = batch.cursor();
259 cursor.to_vec(&batch, |k| *k, |v| *v)
260 })
261 .join()
262 .expect("reader thread panicked");
263
264 assert_eq!(read, vec![((1, 2), vec![(0, 1)]), ((2, 3), vec![(1, 1)])]);
265 }
266}