1use std::{marker::PhantomData, mem, ops::Bound};
5
6use reifydb_codec::{
7 encoded::{
8 row::EncodedRow,
9 shape::{RowShape, fingerprint::RowShapeFingerprint},
10 },
11 key::encoded::{EncodedKey, EncodedKeyRange},
12};
13use reifydb_core::{
14 common::CommitVersion,
15 interface::{
16 catalog::{
17 flow::FlowNodeId,
18 id::{NamespaceId, TableId},
19 namespace::Namespace,
20 table::Table,
21 },
22 change::Diff,
23 },
24 key::{EncodableKey, flow_node_internal_state::FlowNodeInternalStateKey, flow_node_state::FlowNodeStateKey},
25};
26use reifydb_sdk::{
27 error::{Result as SdkResult, SdkError},
28 operator::{
29 column::{row::Row, sink::native::NativeRowSink},
30 context::{
31 CatalogApi, DictionaryApi, InternalStateApi, OperatorContext, RowEmit, StateApi, StoreApi,
32 UpdateEmit,
33 },
34 },
35 state::{StateEntry, decode_payload, encode_payload, row::RowNumberProvider},
36};
37use reifydb_value::{
38 Result,
39 value::{
40 Value,
41 dictionary::{DictionaryEntryId, DictionaryId},
42 row_number::RowNumber,
43 },
44};
45use serde::{Serialize, de::DeserializeOwned};
46
47pub trait NativeBridge {
48 fn clock_now_nanos(&self) -> u64;
49
50 fn state_get(&mut self, key: &EncodedKey) -> Result<Option<EncodedRow>>;
51 fn state_get_many(&mut self, keys: &[EncodedKey]) -> Result<Vec<(EncodedKey, EncodedRow)>>;
52 fn state_set(&mut self, key: &EncodedKey, value: EncodedRow) -> Result<()>;
53 fn state_remove(&mut self, key: &EncodedKey) -> Result<()>;
54 fn state_drop(&mut self, key: &EncodedKey) -> Result<()>;
55 fn state_clear(&mut self) -> Result<()>;
56 fn state_range(&mut self, range: EncodedKeyRange) -> Result<Vec<(EncodedKey, EncodedRow)>>;
57
58 fn internal_state_get(&mut self, key: &EncodedKey) -> Result<Option<EncodedRow>>;
59 fn internal_state_get_many(&mut self, keys: &[EncodedKey]) -> Result<Vec<(EncodedKey, EncodedRow)>>;
60 fn internal_state_set(&mut self, key: &EncodedKey, value: EncodedRow) -> Result<()>;
61 fn internal_state_remove(&mut self, key: &EncodedKey) -> Result<()>;
62 fn internal_state_drop(&mut self, key: &EncodedKey) -> Result<()>;
63 fn internal_state_range(&mut self, range: EncodedKeyRange) -> Result<Vec<(EncodedKey, EncodedRow)>>;
64
65 fn allocate_row_numbers(&mut self, count: u64) -> Result<RowNumber>;
66
67 fn store_get(&mut self, key: &EncodedKey) -> Result<Option<EncodedRow>>;
68 fn store_contains(&mut self, key: &EncodedKey) -> Result<bool>;
69 fn store_prefix(&mut self, prefix: &EncodedKey) -> Result<Vec<(EncodedKey, EncodedRow)>>;
70 fn store_range(&mut self, range: EncodedKeyRange) -> Result<Vec<(EncodedKey, EncodedRow)>>;
71
72 fn catalog_find_namespace(
73 &mut self,
74 namespace: NamespaceId,
75 version: CommitVersion,
76 ) -> Result<Option<Namespace>>;
77 fn catalog_find_namespace_by_name(
78 &mut self,
79 namespace: &str,
80 version: CommitVersion,
81 ) -> Result<Option<Namespace>>;
82 fn catalog_find_table(&mut self, table: TableId, version: CommitVersion) -> Result<Option<Table>>;
83 fn catalog_find_table_by_name(
84 &mut self,
85 namespace: NamespaceId,
86 name: &str,
87 version: CommitVersion,
88 ) -> Result<Option<Table>>;
89 fn catalog_find_row_shape(&mut self, fingerprint: RowShapeFingerprint) -> Result<Option<RowShape>>;
90
91 fn dictionary_id_by_name(&mut self, name: &str) -> Result<Option<DictionaryId>>;
92 fn dictionary_find(&mut self, dictionary: DictionaryId, value: &Value) -> Result<Option<DictionaryEntryId>>;
93 fn dictionary_get(&mut self, dictionary: DictionaryId, id: DictionaryEntryId) -> Result<Option<Value>>;
94
95 fn state_get_many_visit(
96 &mut self,
97 keys: &[EncodedKey],
98 visit: &mut dyn FnMut(&EncodedKey, &EncodedRow) -> SdkResult<()>,
99 ) -> SdkResult<()>;
100 fn internal_state_get_many_visit(
101 &mut self,
102 keys: &[EncodedKey],
103 visit: &mut dyn FnMut(&EncodedKey, &EncodedRow) -> SdkResult<()>,
104 ) -> SdkResult<()>;
105 fn state_range_visit(
106 &mut self,
107 range: EncodedKeyRange,
108 visit: &mut dyn FnMut(&EncodedKey, &EncodedRow) -> SdkResult<()>,
109 ) -> SdkResult<()>;
110 fn store_range_visit(
111 &mut self,
112 range: EncodedKeyRange,
113 visit: &mut dyn FnMut(&EncodedKey, &EncodedRow) -> SdkResult<()>,
114 ) -> SdkResult<()>;
115 fn store_prefix_visit(
116 &mut self,
117 prefix: &EncodedKey,
118 visit: &mut dyn FnMut(&EncodedKey, &EncodedRow) -> SdkResult<()>,
119 ) -> SdkResult<()>;
120}
121
122fn to_sdk_err<E: ToString>(e: E) -> SdkError {
123 SdkError::Other(e.to_string())
124}
125
126fn decode<T: DeserializeOwned>(row: &EncodedRow) -> SdkResult<T> {
127 decode_payload(row)
128}
129
130fn strip_state_envelope(stored: &EncodedKey) -> EncodedKey {
131 FlowNodeStateKey::decode(stored).map(|k| EncodedKey::new(k.key)).unwrap_or_else(|| stored.clone())
132}
133
134fn strip_internal_envelope(stored: &EncodedKey) -> EncodedKey {
135 FlowNodeInternalStateKey::decode(stored).map(|k| EncodedKey::new(k.key)).unwrap_or_else(|| stored.clone())
136}
137
138fn encode<T: Serialize>(value: &T, now_nanos: u64) -> SdkResult<EncodedRow> {
139 encode_payload(value, now_nanos)
140}
141
142pub struct NativeOperatorContext<'a> {
143 bridge: *mut (dyn NativeBridge + 'a),
144 node: FlowNodeId,
145 now_nanos: u64,
146 diffs: Vec<Diff>,
147 _marker: PhantomData<&'a mut (dyn NativeBridge + 'a)>,
148}
149
150impl<'a> NativeOperatorContext<'a> {
151 pub fn new(bridge: &'a mut (dyn NativeBridge + 'a), node: FlowNodeId) -> Self {
152 let now_nanos = bridge.clock_now_nanos();
153 Self {
154 bridge: bridge as *mut (dyn NativeBridge + 'a),
155 node,
156 now_nanos,
157 diffs: Vec::new(),
158 _marker: PhantomData,
159 }
160 }
161
162 pub fn take_diffs(&mut self) -> Vec<Diff> {
163 mem::take(&mut self.diffs)
164 }
165}
166
167enum EmitKind {
168 Insert,
169 Remove,
170}
171
172pub struct NativeRowEmit<'a> {
173 sink: NativeRowSink,
174 diffs: &'a mut Vec<Diff>,
175 kind: EmitKind,
176 now_nanos: u64,
177}
178
179impl RowEmit for NativeRowEmit<'_> {
180 type Sink = NativeRowSink;
181 fn sink(&mut self) -> &mut NativeRowSink {
182 &mut self.sink
183 }
184 fn finish(self, row_numbers: &[RowNumber]) -> SdkResult<()> {
185 let columns = self.sink.finish(row_numbers.to_vec(), self.now_nanos)?;
186 match self.kind {
187 EmitKind::Insert => self.diffs.push(Diff::insert(columns)),
188 EmitKind::Remove => self.diffs.push(Diff::remove(columns)),
189 }
190 Ok(())
191 }
192}
193
194pub struct NativeUpdateEmit<'a> {
195 pre: NativeRowSink,
196 post: NativeRowSink,
197 diffs: &'a mut Vec<Diff>,
198 now_nanos: u64,
199}
200
201impl UpdateEmit for NativeUpdateEmit<'_> {
202 type Sink = NativeRowSink;
203 fn pre(&mut self) -> &mut NativeRowSink {
204 &mut self.pre
205 }
206 fn post(&mut self) -> &mut NativeRowSink {
207 &mut self.post
208 }
209 fn finish(self, row_numbers: &[RowNumber]) -> SdkResult<()> {
210 let pre_columns = self.pre.finish(row_numbers.to_vec(), self.now_nanos)?;
211 let post_columns = self.post.finish(row_numbers.to_vec(), self.now_nanos)?;
212 self.diffs.push(Diff::update(pre_columns, post_columns));
213 Ok(())
214 }
215}
216
217pub struct NativeState<'a> {
218 bridge: *mut (dyn NativeBridge + 'a),
219 now_nanos: u64,
220 _marker: PhantomData<&'a mut (dyn NativeBridge + 'a)>,
221}
222
223impl StateApi for NativeState<'_> {
224 fn get<T: DeserializeOwned>(&self, key: &EncodedKey) -> SdkResult<Option<T>> {
225 match unsafe { (*self.bridge).state_get(key) }.map_err(to_sdk_err)? {
226 Some(row) => Ok(Some(decode(&row)?)),
227 None => Ok(None),
228 }
229 }
230 fn set<T: Serialize>(&mut self, key: &EncodedKey, value: &T) -> SdkResult<()> {
231 let now = self.now_nanos;
232 unsafe { (*self.bridge).state_set(key, encode(value, now)?) }.map_err(to_sdk_err)
233 }
234 fn remove(&mut self, key: &EncodedKey) -> SdkResult<()> {
235 unsafe { (*self.bridge).state_remove(key) }.map_err(to_sdk_err)
236 }
237 fn drop(&mut self, key: &EncodedKey) -> SdkResult<()> {
238 unsafe { (*self.bridge).state_drop(key) }.map_err(to_sdk_err)
239 }
240 fn contains(&self, key: &EncodedKey) -> SdkResult<bool> {
241 Ok(unsafe { (*self.bridge).state_get(key) }.map_err(to_sdk_err)?.is_some())
242 }
243 fn clear(&mut self) -> SdkResult<()> {
244 unsafe { (*self.bridge).state_clear() }.map_err(to_sdk_err)
245 }
246 fn scan_prefix<T: DeserializeOwned>(&self, prefix: &EncodedKey) -> SdkResult<Vec<(EncodedKey, T)>> {
247 let rows = unsafe { (*self.bridge).state_range(EncodedKeyRange::prefix(prefix.as_ref())) }
248 .map_err(to_sdk_err)?;
249 rows.into_iter().map(|(k, r)| Ok((strip_state_envelope(&k), decode(&r)?))).collect()
250 }
251 fn get_many<T: DeserializeOwned>(&self, keys: &[EncodedKey]) -> SdkResult<Vec<(EncodedKey, T)>> {
252 let rows = unsafe { (*self.bridge).state_get_many(keys) }.map_err(to_sdk_err)?;
253 rows.into_iter().map(|(k, r)| Ok((strip_state_envelope(&k), decode(&r)?))).collect()
254 }
255 fn keys_with_prefix(&self, prefix: &EncodedKey) -> SdkResult<Vec<EncodedKey>> {
256 let rows = unsafe { (*self.bridge).state_range(EncodedKeyRange::prefix(prefix.as_ref())) }
257 .map_err(to_sdk_err)?;
258 Ok(rows.into_iter().map(|(k, _)| strip_state_envelope(&k)).collect())
259 }
260 fn range<T: DeserializeOwned>(
261 &self,
262 start: Bound<&EncodedKey>,
263 end: Bound<&EncodedKey>,
264 ) -> SdkResult<Vec<(EncodedKey, T)>> {
265 let range = EncodedKeyRange::new(start.map(|k| k.clone()), end.map(|k| k.clone()));
266 let rows = unsafe { (*self.bridge).state_range(range) }.map_err(to_sdk_err)?;
267 rows.into_iter().map(|(k, r)| Ok((strip_state_envelope(&k), decode(&r)?))).collect()
268 }
269 fn get_with_anchors<T: DeserializeOwned>(&self, key: &EncodedKey) -> SdkResult<Option<StateEntry<T>>> {
270 match unsafe { (*self.bridge).state_get(key) }.map_err(to_sdk_err)? {
271 Some(row) => Ok(Some(StateEntry {
272 created_at_nanos: row.created_at_nanos(),
273 updated_at_nanos: row.updated_at_nanos(),
274 value: decode(&row)?,
275 })),
276 None => Ok(None),
277 }
278 }
279 fn get_many_visit<T: DeserializeOwned>(
280 &self,
281 keys: &[EncodedKey],
282 visit: &mut dyn FnMut(EncodedKey, T) -> SdkResult<()>,
283 ) -> SdkResult<()> {
284 unsafe {
285 (*self.bridge).state_get_many_visit(keys, &mut |k, row| {
286 let value = decode::<T>(row)?;
287 visit(strip_state_envelope(k), value)
288 })
289 }
290 }
291 fn range_visit<T: DeserializeOwned>(
292 &self,
293 start: Bound<&EncodedKey>,
294 end: Bound<&EncodedKey>,
295 visit: &mut dyn FnMut(EncodedKey, T) -> SdkResult<()>,
296 ) -> SdkResult<()> {
297 let range = EncodedKeyRange::new(start.map(|k| k.clone()), end.map(|k| k.clone()));
298 unsafe {
299 (*self.bridge).state_range_visit(range, &mut |k, row| {
300 let value = decode::<T>(row)?;
301 visit(strip_state_envelope(k), value)
302 })
303 }
304 }
305 fn scan_prefix_visit<T: DeserializeOwned>(
306 &self,
307 prefix: &EncodedKey,
308 visit: &mut dyn FnMut(EncodedKey, T) -> SdkResult<()>,
309 ) -> SdkResult<()> {
310 unsafe {
311 (*self.bridge).state_range_visit(EncodedKeyRange::prefix(prefix.as_ref()), &mut |k, row| {
312 let value = decode::<T>(row)?;
313 visit(strip_state_envelope(k), value)
314 })
315 }
316 }
317}
318
319pub struct NativeInternalState<'a> {
320 bridge: *mut (dyn NativeBridge + 'a),
321 now_nanos: u64,
322 _marker: PhantomData<&'a mut (dyn NativeBridge + 'a)>,
323}
324
325impl InternalStateApi for NativeInternalState<'_> {
326 fn get<T: DeserializeOwned>(&self, key: &EncodedKey) -> SdkResult<Option<T>> {
327 match unsafe { (*self.bridge).internal_state_get(key) }.map_err(to_sdk_err)? {
328 Some(row) => Ok(Some(decode(&row)?)),
329 None => Ok(None),
330 }
331 }
332 fn get_many<T: DeserializeOwned>(&self, keys: &[EncodedKey]) -> SdkResult<Vec<(EncodedKey, T)>> {
333 let rows = unsafe { (*self.bridge).internal_state_get_many(keys) }.map_err(to_sdk_err)?;
334 rows.into_iter().map(|(k, r)| Ok((strip_internal_envelope(&k), decode(&r)?))).collect()
335 }
336 fn set<T: Serialize>(&mut self, key: &EncodedKey, value: &T) -> SdkResult<()> {
337 let now = self.now_nanos;
338 unsafe { (*self.bridge).internal_state_set(key, encode(value, now)?) }.map_err(to_sdk_err)
339 }
340 fn remove(&mut self, key: &EncodedKey) -> SdkResult<()> {
341 unsafe { (*self.bridge).internal_state_remove(key) }.map_err(to_sdk_err)
342 }
343 fn drop(&mut self, key: &EncodedKey) -> SdkResult<()> {
344 unsafe { (*self.bridge).internal_state_drop(key) }.map_err(to_sdk_err)
345 }
346 fn contains(&self, key: &EncodedKey) -> SdkResult<bool> {
347 Ok(unsafe { (*self.bridge).internal_state_get(key) }.map_err(to_sdk_err)?.is_some())
348 }
349 fn range<T: DeserializeOwned>(
350 &self,
351 start: Bound<&EncodedKey>,
352 end: Bound<&EncodedKey>,
353 ) -> SdkResult<Vec<(EncodedKey, T)>> {
354 let range = EncodedKeyRange::new(start.map(|k| k.clone()), end.map(|k| k.clone()));
355 let rows = unsafe { (*self.bridge).internal_state_range(range) }.map_err(to_sdk_err)?;
356 rows.into_iter().map(|(k, r)| Ok((strip_internal_envelope(&k), decode(&r)?))).collect()
357 }
358 fn get_many_visit<T: DeserializeOwned>(
359 &self,
360 keys: &[EncodedKey],
361 visit: &mut dyn FnMut(EncodedKey, T) -> SdkResult<()>,
362 ) -> SdkResult<()> {
363 unsafe {
364 (*self.bridge).internal_state_get_many_visit(keys, &mut |k, row| {
365 let value = decode::<T>(row)?;
366 visit(strip_internal_envelope(k), value)
367 })
368 }
369 }
370}
371
372pub struct NativeStore<'a> {
373 bridge: *mut (dyn NativeBridge + 'a),
374 _marker: PhantomData<&'a mut (dyn NativeBridge + 'a)>,
375}
376
377impl StoreApi for NativeStore<'_> {
378 fn get(&self, key: &EncodedKey) -> SdkResult<Option<EncodedRow>> {
379 unsafe { (*self.bridge).store_get(key) }.map_err(to_sdk_err)
380 }
381 fn contains(&self, key: &EncodedKey) -> SdkResult<bool> {
382 unsafe { (*self.bridge).store_contains(key) }.map_err(to_sdk_err)
383 }
384 fn prefix(&self, prefix: &EncodedKey) -> SdkResult<Vec<(EncodedKey, EncodedRow)>> {
385 unsafe { (*self.bridge).store_prefix(prefix) }.map_err(to_sdk_err)
386 }
387 fn range(
388 &self,
389 start: Bound<&EncodedKey>,
390 end: Bound<&EncodedKey>,
391 ) -> SdkResult<Vec<(EncodedKey, EncodedRow)>> {
392 let range = EncodedKeyRange::new(start.map(|k| k.clone()), end.map(|k| k.clone()));
393 unsafe { (*self.bridge).store_range(range) }.map_err(to_sdk_err)
394 }
395 fn range_visit(
396 &self,
397 start: Bound<&EncodedKey>,
398 end: Bound<&EncodedKey>,
399 visit: &mut dyn FnMut(EncodedKey, EncodedRow) -> SdkResult<()>,
400 ) -> SdkResult<()> {
401 let range = EncodedKeyRange::new(start.map(|k| k.clone()), end.map(|k| k.clone()));
402 unsafe { (*self.bridge).store_range_visit(range, &mut |k, row| visit(k.clone(), row.clone())) }
403 }
404 fn prefix_visit(
405 &self,
406 prefix: &EncodedKey,
407 visit: &mut dyn FnMut(EncodedKey, EncodedRow) -> SdkResult<()>,
408 ) -> SdkResult<()> {
409 unsafe { (*self.bridge).store_prefix_visit(prefix, &mut |k, row| visit(k.clone(), row.clone())) }
410 }
411}
412
413pub struct NativeCatalog<'a> {
414 bridge: *mut (dyn NativeBridge + 'a),
415 _marker: PhantomData<&'a mut (dyn NativeBridge + 'a)>,
416}
417
418impl CatalogApi for NativeCatalog<'_> {
419 fn find_namespace(&self, namespace: NamespaceId, version: CommitVersion) -> SdkResult<Option<Namespace>> {
420 unsafe { (*self.bridge).catalog_find_namespace(namespace, version) }.map_err(to_sdk_err)
421 }
422 fn find_namespace_by_name(&self, namespace: &str, version: CommitVersion) -> SdkResult<Option<Namespace>> {
423 unsafe { (*self.bridge).catalog_find_namespace_by_name(namespace, version) }.map_err(to_sdk_err)
424 }
425 fn find_table(&self, table: TableId, version: CommitVersion) -> SdkResult<Option<Table>> {
426 unsafe { (*self.bridge).catalog_find_table(table, version) }.map_err(to_sdk_err)
427 }
428 fn find_table_by_name(
429 &self,
430 namespace: NamespaceId,
431 name: &str,
432 version: CommitVersion,
433 ) -> SdkResult<Option<Table>> {
434 unsafe { (*self.bridge).catalog_find_table_by_name(namespace, name, version) }.map_err(to_sdk_err)
435 }
436 fn find_row_shape(&self, fingerprint: RowShapeFingerprint) -> SdkResult<Option<RowShape>> {
437 unsafe { (*self.bridge).catalog_find_row_shape(fingerprint) }.map_err(to_sdk_err)
438 }
439}
440
441pub struct NativeDictionary<'a> {
442 bridge: *mut (dyn NativeBridge + 'a),
443 _marker: PhantomData<&'a mut (dyn NativeBridge + 'a)>,
444}
445
446impl DictionaryApi for NativeDictionary<'_> {
447 fn id_by_name(&mut self, name: &str) -> SdkResult<Option<DictionaryId>> {
448 unsafe { (*self.bridge).dictionary_id_by_name(name) }.map_err(to_sdk_err)
449 }
450 fn find(&mut self, dictionary: DictionaryId, value: &Value) -> SdkResult<Option<DictionaryEntryId>> {
451 unsafe { (*self.bridge).dictionary_find(dictionary, value) }.map_err(to_sdk_err)
452 }
453 fn get(&mut self, dictionary: DictionaryId, id: DictionaryEntryId) -> SdkResult<Option<Value>> {
454 unsafe { (*self.bridge).dictionary_get(dictionary, id) }.map_err(to_sdk_err)
455 }
456}
457
458impl OperatorContext for NativeOperatorContext<'_> {
459 type InsertEmit<'a>
460 = NativeRowEmit<'a>
461 where
462 Self: 'a;
463 type UpdateEmit<'a>
464 = NativeUpdateEmit<'a>
465 where
466 Self: 'a;
467 type RemoveEmit<'a>
468 = NativeRowEmit<'a>
469 where
470 Self: 'a;
471
472 fn operator_id(&self) -> FlowNodeId {
473 self.node
474 }
475 fn clock_now_nanos(&self) -> u64 {
476 self.now_nanos
477 }
478 fn state(&mut self) -> impl StateApi + '_ {
479 NativeState {
480 bridge: self.bridge,
481 now_nanos: self.now_nanos,
482 _marker: PhantomData,
483 }
484 }
485 fn internal_state(&mut self) -> impl InternalStateApi + '_ {
486 NativeInternalState {
487 bridge: self.bridge,
488 now_nanos: self.now_nanos,
489 _marker: PhantomData,
490 }
491 }
492 fn store(&mut self) -> impl StoreApi + '_ {
493 NativeStore {
494 bridge: self.bridge,
495 _marker: PhantomData,
496 }
497 }
498 fn catalog(&mut self) -> impl CatalogApi + '_ {
499 NativeCatalog {
500 bridge: self.bridge,
501 _marker: PhantomData,
502 }
503 }
504 fn dictionary(&mut self) -> impl DictionaryApi + '_ {
505 NativeDictionary {
506 bridge: self.bridge,
507 _marker: PhantomData,
508 }
509 }
510 fn get_or_create_row_number(&mut self, key: &EncodedKey) -> SdkResult<(RowNumber, bool)> {
511 let provider = RowNumberProvider::new(self.node);
512 provider.get_or_create_row_number(self, key)
513 }
514 fn get_or_create_row_numbers(&mut self, keys: &[EncodedKey]) -> SdkResult<Vec<(RowNumber, bool)>> {
515 let provider = RowNumberProvider::new(self.node);
516 provider.get_or_create_row_numbers_batch(self, keys.iter())
517 }
518 fn allocate_row_numbers(&mut self, count: u64) -> SdkResult<RowNumber> {
519 unsafe { (*self.bridge).allocate_row_numbers(count) }.map_err(to_sdk_err)
520 }
521 fn shape_for_row(&mut self, row: &EncodedRow) -> SdkResult<RowShape> {
522 let fingerprint = row.fingerprint();
523 match self.catalog().find_row_shape(fingerprint)? {
524 Some(shape) => Ok(shape),
525 None => Err(SdkError::Other(format!(
526 "row shape with fingerprint {} not registered in catalog",
527 fingerprint.as_u64()
528 ))),
529 }
530 }
531 fn insert_emit<R: Row>(&mut self, _row_capacity: usize) -> SdkResult<NativeRowEmit<'_>> {
532 let now_nanos = self.now_nanos;
533 Ok(NativeRowEmit {
534 sink: NativeRowSink::new(R::COLUMNS)?,
535 diffs: &mut self.diffs,
536 kind: EmitKind::Insert,
537 now_nanos,
538 })
539 }
540 fn update_emit<R: Row>(&mut self, _row_capacity: usize) -> SdkResult<NativeUpdateEmit<'_>> {
541 let now_nanos = self.now_nanos;
542 Ok(NativeUpdateEmit {
543 pre: NativeRowSink::new(R::COLUMNS)?,
544 post: NativeRowSink::new(R::COLUMNS)?,
545 diffs: &mut self.diffs,
546 now_nanos,
547 })
548 }
549 fn remove_emit<R: Row>(&mut self, _row_capacity: usize) -> SdkResult<NativeRowEmit<'_>> {
550 let now_nanos = self.now_nanos;
551 Ok(NativeRowEmit {
552 sink: NativeRowSink::new(R::COLUMNS)?,
553 diffs: &mut self.diffs,
554 kind: EmitKind::Remove,
555 now_nanos,
556 })
557 }
558}