reifydb_sub_flow/operator/window/
store.rs1use reifydb_codec::key::encoded::{EncodedKey, EncodedKeyRange};
5use reifydb_core::{
6 interface::catalog::flow::FlowNodeId,
7 key::{EncodableKey, flow_node_internal_state::FlowNodeInternalStateKey},
8 window::store::WindowStore,
9};
10use reifydb_sdk::state::{decode_payload, encode_payload};
11use reifydb_value::{Result, value::row_number::RowNumber};
12use serde::{Serialize, de::DeserializeOwned};
13
14use crate::{
15 operator::stateful::{
16 row::{RowNumberProvider, allocate_row_numbers},
17 utils::{internal_state_drop, state_drop},
18 },
19 transaction::FlowTransaction,
20};
21
22pub struct FlowWindowStore<'a> {
23 txn: &'a mut FlowTransaction,
24 node: FlowNodeId,
25 now_nanos: u64,
26}
27
28impl<'a> FlowWindowStore<'a> {
29 pub fn new(txn: &'a mut FlowTransaction, node: FlowNodeId) -> Self {
30 let now_nanos = txn.clock().now_nanos();
31 Self {
32 txn,
33 node,
34 now_nanos,
35 }
36 }
37}
38
39impl WindowStore for FlowWindowStore<'_> {
40 fn state_get<V: DeserializeOwned>(&mut self, key: &EncodedKey) -> Result<Option<V>> {
41 match self.txn.state_get(self.node, key)? {
42 Some(row) => Ok(Some(decode_payload::<V>(&row)?)),
43 None => Ok(None),
44 }
45 }
46
47 fn state_get_many_visit<V: DeserializeOwned>(
48 &mut self,
49 keys: &[EncodedKey],
50 visit: &mut dyn FnMut(EncodedKey, V) -> Result<()>,
51 ) -> Result<()> {
52 let batch = self.txn.state_get_many(self.node, keys)?;
53 for r in batch.items {
54 let value = decode_payload::<V>(&r.row)?;
55 visit(r.key, value)?;
56 }
57 Ok(())
58 }
59
60 fn state_set<V: Serialize>(&mut self, key: &EncodedKey, value: &V) -> Result<()> {
61 self.txn.state_set(self.node, key, encode_payload(value, self.now_nanos)?)
62 }
63
64 fn state_remove(&mut self, key: &EncodedKey) -> Result<()> {
65 self.txn.state_remove(self.node, key)
66 }
67
68 fn state_drop(&mut self, key: &EncodedKey) -> Result<()> {
69 state_drop(self.node, self.txn, key)
70 }
71
72 fn internal_get<V: DeserializeOwned>(&mut self, key: &EncodedKey) -> Result<Option<V>> {
73 match self.txn.internal_state_get(self.node, key)? {
74 Some(row) => Ok(Some(decode_payload::<V>(&row)?)),
75 None => Ok(None),
76 }
77 }
78
79 fn internal_get_many_visit<V: DeserializeOwned>(
80 &mut self,
81 keys: &[EncodedKey],
82 visit: &mut dyn FnMut(EncodedKey, V) -> Result<()>,
83 ) -> Result<()> {
84 let batch = self.txn.internal_state_get_many(self.node, keys)?;
85 for r in batch.items {
86 let value = decode_payload::<V>(&r.row)?;
87 visit(r.key, value)?;
88 }
89 Ok(())
90 }
91
92 fn internal_set<V: Serialize>(&mut self, key: &EncodedKey, value: &V) -> Result<()> {
93 self.txn.internal_state_set(self.node, key, encode_payload(value, self.now_nanos)?)
94 }
95
96 fn internal_remove(&mut self, key: &EncodedKey) -> Result<()> {
97 self.txn.internal_state_remove(self.node, key)
98 }
99
100 fn internal_drop(&mut self, key: &EncodedKey) -> Result<()> {
101 internal_state_drop(self.node, self.txn, key)
102 }
103
104 fn internal_range_visit<V: DeserializeOwned>(
105 &mut self,
106 range: EncodedKeyRange,
107 visit: &mut dyn FnMut(EncodedKey, V) -> Result<()>,
108 ) -> Result<()> {
109 let batch = self.txn.internal_state_range_all(self.node, range)?;
110 for r in batch.items {
111 if let Some(decoded) = FlowNodeInternalStateKey::decode(&r.key) {
112 let value = decode_payload::<V>(&r.row)?;
113 visit(EncodedKey::new(decoded.key), value)?;
114 }
115 }
116 Ok(())
117 }
118
119 fn get_or_create_row_number(&mut self, key: &EncodedKey) -> Result<(RowNumber, bool)> {
120 RowNumberProvider::new(self.node).get_or_create_row_number(self.txn, key)
121 }
122
123 fn get_or_create_row_numbers(&mut self, keys: &[EncodedKey]) -> Result<Vec<(RowNumber, bool)>> {
124 RowNumberProvider::new(self.node).get_or_create_row_numbers(self.txn, keys.iter())
125 }
126
127 fn allocate_row_numbers(&mut self, count: u64) -> Result<RowNumber> {
128 allocate_row_numbers(self.txn, self.node, count).map(RowNumber)
129 }
130
131 fn clock_now_nanos(&self) -> u64 {
132 self.now_nanos
133 }
134}