#![allow(dead_code)]
use std::io::{BufRead, BufReader, Write};
use std::path::Path;
use std::process::{Child, ChildStdin, Command, Stdio};
use std::sync::mpsc::{self, Receiver, RecvTimeoutError};
use std::time::Duration;
const ACK_TIMEOUT: Duration = Duration::from_secs(30);
pub struct HeldReader {
child: Child,
stdin: Option<ChildStdin>,
lines: Receiver<String>,
seq: u32,
}
impl HeldReader {
pub fn spawn(bin: &str, db: &Path) -> Self {
let mut child = Command::new(bin)
.arg(db)
.stdin(Stdio::piped())
.stdout(Stdio::piped())
.stderr(Stdio::null())
.spawn()
.unwrap();
let stdin = child.stdin.take().unwrap();
let stdout = child.stdout.take().unwrap();
let (tx, lines) = mpsc::channel();
std::thread::spawn(move || {
for line in BufReader::new(stdout).lines().map_while(Result::ok) {
if tx.send(line).is_err() {
break;
}
}
});
Self {
child,
stdin: Some(stdin),
lines,
seq: 0,
}
}
pub fn run(&mut self, sql: &str) {
self.seq += 1;
let token = format!("__ack_{}__", self.seq);
let stdin = self.stdin.as_mut().expect("reader stdin already released");
writeln!(stdin, "{sql}\nSELECT '{token}';").unwrap();
stdin.flush().unwrap();
loop {
match self.lines.recv_timeout(ACK_TIMEOUT) {
Ok(line) if line.trim() == token => return,
Ok(_) => {}
Err(RecvTimeoutError::Timeout) => {
panic!("sqlite3 did not acknowledge within {ACK_TIMEOUT:?}: {sql}")
}
Err(RecvTimeoutError::Disconnected) => {
panic!("sqlite3 reader exited before acknowledging: {sql}")
}
}
}
}
pub fn finish(mut self) {
if let Some(mut stdin) = self.stdin.take() {
let _ = writeln!(stdin, "COMMIT;\n.quit");
}
let _ = self.child.wait();
}
}
pub fn writer_sql(bin: &str, db: &Path, sql: &str) {
let out = Command::new(bin).arg(db).arg(sql).output().unwrap();
assert!(
out.status.success(),
"sqlite3 writer failed: {}",
String::from_utf8_lossy(&out.stderr)
);
}