reifydb_core/actors/
flow.rs1use std::{collections::BTreeSet, sync::Arc};
5
6use reifydb_runtime::actor::system::ActorHandle;
7use reifydb_value::Result;
8
9use crate::{
10 common::CommitVersion,
11 interface::{
12 catalog::{flow::FlowId, shape::ShapeId},
13 cdc::Cdc,
14 },
15};
16
17pub type FlowActorHandle = ActorHandle<FlowActorMessage>;
18
19pub enum FlowActorMessage {
20 Drain,
21
22 Wake,
23
24 Ingest {
25 cdcs: Arc<Vec<Cdc>>,
26 covers_from: CommitVersion,
27 up_to: CommitVersion,
28 },
29
30 Tick,
31
32 UpdateSources {
33 source_shapes: Arc<BTreeSet<ShapeId>>,
34 },
35
36 CommitDone {
37 advance_to: CommitVersion,
38 more: bool,
39 result: Result<()>,
40 },
41
42 Stop {
43 delete_checkpoint: bool,
44 reply: Box<dyn FnOnce() + Send>,
45 },
46}
47
48pub type FlowSupervisorHandle = ActorHandle<FlowSupervisorMessage>;
49
50pub enum FlowSupervisorMessage {
51 Bootstrap {
52 flows: Vec<(FlowId, bool)>,
53 },
54
55 Consume {
56 cdcs: Vec<Cdc>,
57 current_version: CommitVersion,
58 reply: Box<dyn FnOnce(Result<()>) + Send>,
59 },
60}