blockworx/storage/
writer.rs1use std::sync::mpsc::{Receiver, Sender, channel};
13
14use super::container::Container;
15use super::fs::FsStorage;
16use super::history;
17use crate::document::Document;
18
19struct Job {
21 document: Document,
22 seq: u64,
23 meta: history::Entry,
24}
25
26enum Message {
27 Write(Box<Job>),
28 Drained(Sender<()>),
30}
31
32pub struct Writer {
33 tx: Sender<Message>,
34 thread: Option<std::thread::JoinHandle<()>>,
35}
36
37impl Writer {
38 pub fn spawn(root: std::path::PathBuf) -> Self {
40 let (tx, rx) = channel();
41 let thread = std::thread::Builder::new()
42 .name("blockworx-writer".to_string())
43 .spawn(move || run(&Container::open(&root), &rx))
44 .ok();
45 Self { tx, thread }
46 }
47
48 pub fn write(&self, document: Document, seq: u64, meta: history::Entry) {
51 let _ = self.tx.send(Message::Write(Box::new(Job {
54 document,
55 seq,
56 meta,
57 })));
58 }
59
60 pub fn drain(&self) {
66 let (tx, rx) = channel();
67 if self.tx.send(Message::Drained(tx)).is_ok() {
68 let _ = rx.recv();
69 }
70 }
71}
72
73impl Drop for Writer {
74 fn drop(&mut self) {
75 self.drain();
76 let (dead_tx, _) = channel();
79 let _ = std::mem::replace(&mut self.tx, dead_tx);
80 if let Some(thread) = self.thread.take() {
81 let _ = thread.join();
82 }
83 }
84}
85
86fn run(container: &Container<FsStorage>, rx: &Receiver<Message>) {
87 while let Ok(message) = rx.recv() {
88 match message {
89 Message::Write(job) => {
90 if let Err(e) = container.save_and_record(&job.document, job.seq, &job.meta) {
91 tracing::error!("Autosave failed:\n{e:?}");
92 }
93 }
94 Message::Drained(reply) => {
95 let _ = reply.send(());
96 }
97 }
98 }
99}
100
101#[cfg(test)]
102mod tests {
103 use super::*;
104 use crate::storage::Storage as _;
105 use crate::storage::atomic::tests::TempDir;
106 use crate::storage::container::ROOT;
107
108 fn meta(changed: &[&str]) -> history::Entry {
109 history::Entry {
110 ts: history::now_millis(),
111 command: None,
112 changed: changed.iter().map(|s| (*s).to_string()).collect(),
113 }
114 }
115
116 #[test]
117 fn a_queued_document_reaches_the_disk() {
118 let dir = TempDir::new("writer-writes");
119 let root = dir.join("d.bwx");
120 let writer = Writer::spawn(root.clone());
121
122 writer.write(Document::default(), 0, meta(&["b1"]));
123 writer.drain();
124
125 assert!(root.join(ROOT).is_file());
126 let container = Container::open(&root);
127 assert_eq!(history::records(container.storage()).unwrap().len(), 1);
128 }
129
130 #[test]
133 fn a_burst_of_writes_all_land_in_order() {
134 let dir = TempDir::new("writer-burst");
135 let root = dir.join("d.bwx");
136 let writer = Writer::spawn(root.clone());
137
138 for seq in 0..8 {
139 writer.write(Document::default(), seq, meta(&[&format!("b{seq}")]));
140 }
141 writer.drain();
142
143 let container = Container::open(&root);
144 let records = history::records(container.storage()).unwrap();
145 let seqs: Vec<u64> = records.iter().map(|r| r.seq).collect();
146 assert_eq!(seqs, (0..8).collect::<Vec<_>>());
147 for (seq, record) in records.iter().enumerate() {
148 let meta = record.meta.as_ref().expect("a sidecar");
149 assert_eq!(meta.changed, vec![format!("b{seq}")]);
150 }
151 }
152
153 #[test]
155 fn dropping_the_writer_finishes_what_was_queued() {
156 let dir = TempDir::new("writer-drop");
157 let root = dir.join("d.bwx");
158 {
159 let writer = Writer::spawn(root.clone());
160 for seq in 0..4 {
161 writer.write(Document::default(), seq, meta(&[]));
162 }
163 }
165 let container = Container::open(&root);
166 assert_eq!(history::records(container.storage()).unwrap().len(), 4);
167 assert!(container.storage().exists(ROOT));
168 }
169}