1use std::os::unix::fs::PermissionsExt;
22use std::path::{Path, PathBuf};
23use std::time::{Duration, SystemTime, UNIX_EPOCH};
24
25use rusqlite::{Connection, OptionalExtension, Transaction, TransactionBehavior, params};
26
27use crate::error::{Error, Result};
28
29pub const KEEP: u64 = 2 << 20;
32
33pub const DEFAULT_DIR: &str = "/var/lib/isb/console";
36
37const STALE: Duration = Duration::from_secs(30 * 24 * 3600);
39
40pub fn dir() -> PathBuf {
42 std::env::var_os("ISB_CONSOLE_DIR")
43 .filter(|s| !s.is_empty())
44 .map(PathBuf::from)
45 .unwrap_or_else(|| PathBuf::from(DEFAULT_DIR))
46}
47
48const SCHEMA: &str = "
49CREATE TABLE IF NOT EXISTS chunks (
50 project TEXT NOT NULL,
51 instance TEXT NOT NULL,
52 pos INTEGER NOT NULL,
53 data BLOB NOT NULL,
54 at INTEGER NOT NULL,
55 PRIMARY KEY (project, instance, pos)
56);
57";
58
59fn db_err(what: &str, e: rusqlite::Error) -> Error {
60 Error::invalid(format!("console log database: {what}: {e}"))
61}
62
63fn open_at(path: &Path) -> Result<Connection> {
64 let conn = Connection::open(path).map_err(|e| db_err("open", e))?;
65 if conn
67 .is_readonly(rusqlite::MAIN_DB)
68 .map_err(|e| db_err("open", e))?
69 {
70 return Err(Error::invalid(format!("{} is read-only", path.display())));
71 }
72 conn.busy_timeout(Duration::from_secs(60))
74 .map_err(|e| db_err("busy timeout", e))?;
75 conn.execute_batch("PRAGMA journal_mode=WAL; PRAGMA synchronous=NORMAL;")
76 .map_err(|e| db_err("pragmas", e))?;
77 conn.execute_batch(SCHEMA)
78 .map_err(|e| db_err("schema", e))?;
79 Ok(conn)
80}
81
82fn open() -> Result<Connection> {
84 let shared = dir().join("console.db");
85 if let Ok(c) = open_at(&shared) {
86 let _ = std::fs::set_permissions(&shared, std::fs::Permissions::from_mode(0o660));
88 return Ok(c);
89 }
90 let own = crate::stack::Store::default_dir();
91 std::fs::create_dir_all(&own)?;
92 open_at(&own.join("console.db"))
93}
94
95fn now() -> i64 {
96 SystemTime::now()
97 .duration_since(UNIX_EPOCH)
98 .map_or(0, |d| d.as_secs() as i64)
99}
100
101pub(crate) fn read_through(
107 project: &str,
108 instance: &str,
109 pos: u64,
110 read: &mut dyn FnMut() -> Result<(Vec<u8>, bool)>,
111) -> Result<(Vec<u8>, u64)> {
112 let Ok(mut conn) = open() else {
113 return Ok((read()?.0, 0));
114 };
115 let Ok(tx) = conn.transaction_with_behavior(TransactionBehavior::Immediate) else {
116 return Ok((read()?.0, 0));
117 };
118 let (new, undrained) = read()?;
119 let kept = record(&tx, project, instance, &new, undrained, now())
120 .and_then(|()| since(&tx, project, instance, pos))
121 .and_then(|r| tx.commit().map(|()| r));
122 Ok(kept.unwrap_or((new, 0)))
123}
124
125pub(crate) fn forget(project: &str, instance: &str) {
127 if let Ok(c) = open() {
128 let _ = c.execute(
129 "DELETE FROM chunks WHERE project = ?1 AND instance = ?2",
130 params![project, instance],
131 );
132 }
133}
134
135fn tail(
137 tx: &Transaction,
138 project: &str,
139 instance: &str,
140) -> rusqlite::Result<(u64, Option<Vec<u8>>)> {
141 let last: Option<(i64, Vec<u8>)> = tx
142 .query_row(
143 "SELECT pos, data FROM chunks WHERE project = ?1 AND instance = ?2
144 ORDER BY pos DESC LIMIT 1",
145 params![project, instance],
146 |r| Ok((r.get(0)?, r.get(1)?)),
147 )
148 .optional()?;
149 Ok(match last {
150 Some((pos, data)) => (pos as u64 + data.len() as u64, Some(data)),
151 None => (0, None),
152 })
153}
154
155fn record(
156 tx: &Transaction,
157 project: &str,
158 instance: &str,
159 new: &[u8],
160 undrained: bool,
161 now: i64,
162) -> rusqlite::Result<()> {
163 let (end, last) = tail(tx, project, instance)?;
164 if new.is_empty() || (undrained && last.as_deref() == Some(new)) {
165 if end > 0 {
167 tx.execute(
168 "UPDATE chunks SET at = ?3 WHERE project = ?1 AND instance = ?2 AND pos = (
169 SELECT max(pos) FROM chunks WHERE project = ?1 AND instance = ?2)",
170 params![project, instance, now],
171 )?;
172 }
173 return Ok(());
174 }
175 if end == 0 {
176 prune(tx, now)?;
177 }
178 tx.execute(
179 "INSERT INTO chunks (project, instance, pos, data, at) VALUES (?1, ?2, ?3, ?4, ?5)",
180 params![project, instance, end as i64, new, now],
181 )?;
182 let end = end + new.len() as u64;
183 tx.execute(
184 "DELETE FROM chunks WHERE project = ?1 AND instance = ?2 AND pos + length(data) <= ?3",
185 params![project, instance, end.saturating_sub(KEEP) as i64],
186 )?;
187 Ok(())
188}
189
190fn since(
193 tx: &Transaction,
194 project: &str,
195 instance: &str,
196 pos: u64,
197) -> rusqlite::Result<(Vec<u8>, u64)> {
198 let mut st = tx.prepare(
199 "SELECT pos, data FROM chunks WHERE project = ?1 AND instance = ?2
200 AND pos + length(data) > ?3 ORDER BY pos",
201 )?;
202 let mut out = Vec::new();
203 let rows = st.query_map(
204 params![project, instance, pos.min(i64::MAX as u64) as i64],
205 |r| Ok((r.get::<_, i64>(0)? as u64, r.get::<_, Vec<u8>>(1)?)),
206 )?;
207 for row in rows {
208 let (at, data) = row?;
209 let skip = pos.saturating_sub(at).min(data.len() as u64) as usize;
210 out.extend_from_slice(&data[skip..]);
211 }
212 let (end, _) = tail(tx, project, instance)?;
213 Ok((out, end))
214}
215
216fn prune(tx: &Transaction, now: i64) -> rusqlite::Result<()> {
219 tx.execute(
220 "DELETE FROM chunks WHERE (project, instance) IN (
221 SELECT project, instance FROM chunks GROUP BY project, instance HAVING max(at) < ?1)",
222 params![now - STALE.as_secs() as i64],
223 )?;
224 Ok(())
225}
226
227#[cfg(test)]
228mod tests {
229 use super::*;
230
231 fn db() -> (tempfile::TempDir, Connection) {
232 let d = tempfile::tempdir().unwrap();
233 let c = open_at(&d.path().join("console.db")).unwrap();
234 (d, c)
235 }
236
237 fn add(c: &mut Connection, inst: &str, new: &[u8], undrained: bool, at: i64) {
238 let tx = c.transaction().unwrap();
239 record(&tx, "p", inst, new, undrained, at).unwrap();
240 tx.commit().unwrap();
241 }
242
243 fn read(c: &mut Connection, inst: &str, pos: u64) -> (Vec<u8>, u64) {
244 let tx = c.transaction().unwrap();
245 since(&tx, "p", inst, pos).unwrap()
246 }
247
248 #[test]
249 fn readers_share_the_log() {
250 let (_d, mut c) = db();
251 add(&mut c, "web", b"one\n", false, 1);
253 add(&mut c, "web", b"two\n", false, 1);
254 assert_eq!(read(&mut c, "web", 0), (b"one\ntwo\n".to_vec(), 8));
255 assert_eq!(read(&mut c, "web", 2), (b"e\ntwo\n".to_vec(), 8));
256 assert_eq!(read(&mut c, "web", 8), (Vec::new(), 8));
257 assert_eq!(read(&mut c, "web", u64::MAX), (Vec::new(), 8));
258 assert_eq!(read(&mut c, "db", 0), (Vec::new(), 0));
259 }
260
261 #[test]
262 fn a_stopped_instances_rest_is_recorded_once() {
263 let (_d, mut c) = db();
264 add(&mut c, "web", b"ping\n", false, 1);
265 add(&mut c, "web", b"ping\n", false, 1);
267 add(&mut c, "web", b"bye\n", true, 1);
268 add(&mut c, "web", b"bye\n", true, 1);
269 assert_eq!(read(&mut c, "web", 0).0, b"ping\nping\nbye\n".to_vec());
270 }
271
272 #[test]
273 fn trims_whole_chunks_and_keeps_positions() {
274 let (_d, mut c) = db();
275 add(&mut c, "web", b"old\n", false, 1);
276 let big = vec![b'x'; KEEP as usize];
277 add(&mut c, "web", &big, false, 1);
278 add(&mut c, "web", b"new\n", false, 1);
279 let (all, end) = read(&mut c, "web", 0);
280 assert_eq!(end, 4 + KEEP + 4);
281 assert_eq!(all.len() as u64, KEEP + 4, "the oldest chunk went");
282 assert_eq!(read(&mut c, "web", end - 4).0, b"new\n".to_vec());
284 }
285
286 #[test]
287 fn prunes_logs_nobody_reads() {
288 let (_d, mut c) = db();
289 add(&mut c, "gone", b"a\n", false, 1);
290 add(&mut c, "kept", b"a\n", false, 1);
291 let later = 1 + STALE.as_secs() as i64 + 1;
292 add(&mut c, "kept", b"", false, later);
294 add(&mut c, "new", b"b\n", false, later);
295 assert_eq!(read(&mut c, "gone", 0).0, Vec::<u8>::new());
296 assert_eq!(read(&mut c, "kept", 0).0, b"a\n".to_vec());
297 }
298}