Skip to main content

isb_core/client/
console.rs

1//! Each instance's console output, kept in one SQLite database per host so
2//! that every reader sees all of it.
3//!
4//! incus hands a running container's console output out once: each
5//! `GET .../console` drains what it returns, so a second reader sees only
6//! what was written since the first. Every isb process on the host (the
7//! daemon, the TUI, `isb logs`, `isb up`), whichever user runs it, therefore
8//! records what it drains in [`dir`]`/console.db` and reads from there. The
9//! drain happens inside the database's write transaction, so two readers
10//! never record their chunks out of order.
11//!
12//! The directory is `root:incus-admin` 2770 (`isb host setup` makes it):
13//! whoever may use the incus socket may use the log, and the database and
14//! its WAL files are 0660 in the directory's group. Where it does not exist
15//! or cannot be written, the user's own `$XDG_STATE_HOME/isb/console.db`
16//! stands in, shared only by that user's processes.
17//!
18//! A log is chunks, each at its position in bytes since the log began, so
19//! a follower's position survives trimming the oldest.
20
21use 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
29/// How much of each instance's console output is kept (whole chunks: up to
30/// one chunk more).
31pub const KEEP: u64 = 2 << 20;
32
33/// Where the host's console database lives, unless `ISB_CONSOLE_DIR` says
34/// otherwise.
35pub const DEFAULT_DIR: &str = "/var/lib/isb/console";
36
37/// A log nobody has read for this long belongs to an instance that is gone.
38const STALE: Duration = Duration::from_secs(30 * 24 * 3600);
39
40/// The host's console directory.
41pub 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    // SQLite opens a file it may not write read-only, without saying so.
66    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    // A drain holds the write lock for one incus request.
73    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
82/// The host's database, else the user's own.
83fn open() -> Result<Connection> {
84    let shared = dir().join("console.db");
85    if let Ok(c) = open_at(&shared) {
86        // The group shares it; SQLite gives the WAL files the same mode.
87        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
101/// Read an instance's console through the log: `read` drains incus (and
102/// says whether its answer is a stopped instance's undrained rest, which
103/// it gives on every read), the log records it, and the answer is what was
104/// logged from position `pos` on, with the position after it. Without a
105/// usable database the read is still a read: what `read` returned, at 0.
106pub(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
125/// Forget an instance's log: a new instance under its name starts afresh.
126pub(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
135/// Where the log ends, and its last chunk.
136fn 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        // Reading is what keeps a log from going stale.
166        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
190/// What was logged from `pos` on (all that is kept, if `pos` was trimmed
191/// away), and the position after it.
192fn 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
216/// Remove the logs nobody has read for [`STALE`]: instances deleted by
217/// something other than isb.
218fn 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        // Two readers each drained a chunk; both see both.
252        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        // A running instance may print the same line again.
266        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        // A follower's position still means the same bytes.
283        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        // Reading a log keeps it.
293        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}