faucet_source_singer/
assemble.rs1use faucet_core::{StreamPage, Value};
9
10pub struct PageAssembler {
13 target_stream: String,
14 batch_size: usize,
15 flush_on_state: bool,
16 buffer: Vec<Value>,
17 pending_state: Option<Value>,
19}
20
21impl PageAssembler {
22 pub fn new(target_stream: impl Into<String>, batch_size: usize, flush_on_state: bool) -> Self {
24 Self {
25 target_stream: target_stream.into(),
26 batch_size,
27 flush_on_state,
28 buffer: Vec::new(),
29 pending_state: None,
30 }
31 }
32
33 pub fn on_record(&mut self, stream: &str, record: Value) -> Option<StreamPage> {
36 if stream != self.target_stream {
37 return None;
38 }
39 self.buffer.push(record);
40 if self.batch_size != 0 && self.buffer.len() >= self.batch_size {
41 Some(self.flush())
42 } else {
43 None
44 }
45 }
46
47 pub fn on_state(&mut self, value: Value) -> Option<StreamPage> {
51 self.pending_state = Some(value);
52 if self.flush_on_state {
53 Some(self.flush())
54 } else {
55 None
56 }
57 }
58
59 pub fn on_eof(&mut self) -> Option<StreamPage> {
62 if self.buffer.is_empty() && self.pending_state.is_none() {
63 None
64 } else {
65 Some(self.flush())
66 }
67 }
68
69 fn flush(&mut self) -> StreamPage {
71 StreamPage {
72 records: std::mem::take(&mut self.buffer),
73 bookmark: self.pending_state.take(),
74 }
75 }
76}
77
78#[cfg(test)]
79mod tests {
80 use super::*;
81 use serde_json::json;
82
83 fn rec(id: i64) -> Value {
84 json!({ "id": id })
85 }
86
87 #[test]
88 fn flushes_at_batch_size_with_no_bookmark() {
89 let mut a = PageAssembler::new("s", 2, true);
90 assert!(a.on_record("s", rec(1)).is_none());
91 let page = a.on_record("s", rec(2)).expect("flush at size 2");
92 assert_eq!(page.records, vec![rec(1), rec(2)]);
93 assert!(
94 page.bookmark.is_none(),
95 "no STATE covers a size-based flush"
96 );
97 assert!(a.on_record("s", rec(3)).is_none());
99 let tail = a.on_eof().unwrap();
100 assert_eq!(tail.records, vec![rec(3)]);
101 }
102
103 #[test]
104 fn state_flush_attaches_bookmark() {
105 let mut a = PageAssembler::new("s", 1000, true);
106 assert!(a.on_record("s", rec(1)).is_none());
107 assert!(a.on_record("s", rec(2)).is_none());
108 let page = a.on_state(json!({"last_id": 2})).expect("flush on state");
109 assert_eq!(page.records, vec![rec(1), rec(2)]);
110 assert_eq!(page.bookmark, Some(json!({"last_id": 2})));
111 assert!(a.on_eof().is_none());
113 }
114
115 #[test]
116 fn empty_run_with_trailing_state_yields_empty_page_with_bookmark() {
117 let mut a = PageAssembler::new("s", 1000, true);
118 let page = a.on_state(json!({"last_id": 0})).unwrap();
119 assert!(page.records.is_empty());
120 assert_eq!(page.bookmark, Some(json!({"last_id": 0})));
121 }
122
123 #[test]
124 fn records_for_other_streams_are_ignored() {
125 let mut a = PageAssembler::new("wanted", 1, true);
126 assert!(a.on_record("other", rec(1)).is_none());
127 assert!(a.on_record("other", rec(2)).is_none());
128 let page = a.on_record("wanted", rec(3)).unwrap();
130 assert_eq!(page.records, vec![rec(3)]);
131 }
132
133 #[test]
134 fn no_flush_on_state_defers_bookmark_to_next_flush() {
135 let mut a = PageAssembler::new("s", 2, false);
136 assert!(a.on_state(json!({"last_id": 0})).is_none());
138 assert!(a.on_record("s", rec(1)).is_none());
139 let page = a.on_record("s", rec(2)).unwrap();
141 assert_eq!(page.records, vec![rec(1), rec(2)]);
142 assert_eq!(page.bookmark, Some(json!({"last_id": 0})));
143 }
144
145 #[test]
146 fn eof_with_empty_buffer_and_no_state_yields_nothing() {
147 let mut a = PageAssembler::new("s", 1000, true);
148 assert!(a.on_eof().is_none());
149 }
150}