1use std::cell::RefCell;
16use std::collections::{BTreeMap, VecDeque};
17use std::rc::Rc;
18
19use dfir_rs::scheduled::SubgraphId;
20use dfir_rs::scheduled::graph::Dfir;
21use get_size2::GetSize;
22
23use crate::compute::types::ErrCollector;
24use crate::repr::{self, Timestamp};
25use crate::utils::{ArrangeHandler, Arrangement};
26
27#[derive(Debug, Default)]
30pub struct DataflowState {
31 schedule_subgraph: Rc<RefCell<BTreeMap<Timestamp, VecDeque<SubgraphId>>>>,
34 as_of: Rc<RefCell<Timestamp>>,
40 err_collector: ErrCollector,
43 arrange_used: Vec<ArrangeHandler>,
46 expire_after: Option<Timestamp>,
48 last_exec_time: Option<Timestamp>,
50 start_time: Option<Timestamp>,
52}
53
54impl DataflowState {
55 pub fn new_arrange(&mut self, name: Option<Vec<String>>) -> ArrangeHandler {
56 let arrange = name.map(Arrangement::new_with_name).unwrap_or_default();
57
58 let arr = ArrangeHandler::from(arrange);
59 self.arrange_used.push(
61 arr.clone_future_only()
62 .expect("No write happening at this point"),
63 );
64 arr
65 }
66
67 #[allow(clippy::swap_with_temporary)]
71 pub fn run_available_with_schedule(&mut self, df: &mut Dfir) -> bool {
72 let mut before = self
74 .schedule_subgraph
75 .borrow_mut()
76 .split_off(&(*self.as_of.borrow() + 1));
77 std::mem::swap(&mut before, &mut self.schedule_subgraph.borrow_mut());
78 for (_, v) in before {
79 for subgraph in v {
80 df.schedule_subgraph(subgraph);
81 }
82 }
83 df.run_available()
84 }
85 pub fn get_scheduler(&self) -> Scheduler {
86 Scheduler {
87 schedule_subgraph: self.schedule_subgraph.clone(),
88 cur_subgraph: Rc::new(RefCell::new(None)),
89 }
90 }
91
92 pub fn current_time_ref(&self) -> Rc<RefCell<Timestamp>> {
96 self.as_of.clone()
97 }
98
99 pub fn current_ts(&self) -> Timestamp {
100 *self.as_of.borrow()
101 }
102
103 pub fn set_current_ts(&mut self, ts: Timestamp) {
104 self.as_of.replace(ts);
105 }
106
107 pub fn get_err_collector(&self) -> ErrCollector {
108 self.err_collector.clone()
109 }
110
111 pub fn set_expire_after(&mut self, after: Option<repr::Duration>) {
112 self.expire_after = after;
113 }
114
115 pub fn expire_after(&self) -> Option<Timestamp> {
116 self.expire_after
117 }
118
119 pub fn get_state_size(&self) -> usize {
120 self.arrange_used.iter().map(|x| x.read().get_size()).sum()
121 }
122
123 pub fn set_last_exec_time(&mut self, time: Timestamp) {
124 self.last_exec_time = Some(time);
125 if self.start_time.is_none() {
126 self.start_time = Some(time);
129 }
130 }
131
132 pub fn last_exec_time(&self) -> Option<Timestamp> {
133 self.last_exec_time
134 }
135
136 pub fn start_time(&self) -> Option<Timestamp> {
138 self.start_time
139 }
140}
141
142#[derive(Debug, Clone)]
143pub struct Scheduler {
144 schedule_subgraph: Rc<RefCell<BTreeMap<Timestamp, VecDeque<SubgraphId>>>>,
146 cur_subgraph: Rc<RefCell<Option<SubgraphId>>>,
147}
148
149impl Scheduler {
150 pub fn schedule_at(&self, next_run_time: Timestamp) {
151 let mut schedule_subgraph = self.schedule_subgraph.borrow_mut();
152 let subgraph = self.cur_subgraph.borrow();
153 let subgraph = subgraph.as_ref().expect("Set SubgraphId before schedule");
154 let subgraph_queue = schedule_subgraph.entry(next_run_time).or_default();
155 subgraph_queue.push_back(*subgraph);
156 }
157
158 pub fn schedule_for_arrange(&self, arrange: &Arrangement, now: Timestamp) {
159 if let Some(i) = arrange.get_next_update_time(&now) {
160 self.schedule_at(i)
161 }
162 }
163
164 pub fn set_cur_subgraph(&self, subgraph: SubgraphId) {
165 self.cur_subgraph.replace(Some(subgraph));
166 }
167}