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 actors::pending::Pending,
11 common::CommitVersion,
12 interface::{
13 catalog::{flow::FlowId, object::ObjectId},
14 cdc::Cdc,
15 },
16};
17
18pub type FlowActorHandle = ActorHandle<FlowActorMessage>;
19
20pub enum FlowActorMessage {
21 Drain,
22
23 Wake,
24
25 Loaded {
26 outcome: Result<(Vec<Arc<Cdc>>, CommitVersion)>,
27 },
28
29 Tick,
30
31 Sample,
32
33 PublishRestoredFrontiers,
34
35 UpdateSources {
36 source_objects: Arc<BTreeSet<ObjectId>>,
37 completeness_objects: Option<Arc<BTreeSet<u64>>>,
38 },
39
40 SliceCommitted {
41 advance_to: CommitVersion,
42 more: bool,
43 result: Result<()>,
44 committed: Option<(CommitVersion, Pending)>,
45 },
46
47 TickCommitted {
48 committed: Option<(CommitVersion, Pending)>,
49 },
50
51 Stop {
52 delete_checkpoint: bool,
53 reply: Box<dyn FnOnce() + Send>,
54 },
55}
56
57pub type FlowSupervisorHandle = ActorHandle<FlowSupervisorMessage>;
58
59pub enum FlowSupervisorMessage {
60 Bootstrap {
61 flows: Vec<FlowId>,
62 scan_from: Option<CommitVersion>,
63 },
64
65 Wake,
66
67 PersistFrontiers,
68}