Skip to main content

reifydb_sub_flow/operator/window/
store.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use 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}