char/plc-mirror

git clone https://git.t4t.associates/char/plc-mirror

Charlotte Somuse new plox opcodec formatc2f2166

main
4.7 KiB131 linesraw
1use std::{
2    path::Path,
3    sync::{Arc, RwLock},
4};
5
6use anyhow::{Context, Result, anyhow, ensure};
7use deadpool_sqlite::{Config as DatabaseConfig, Hook, HookError, Pool, Runtime};
8use rusqlite::{Connection, OpenFlags, OptionalExtension, params};
9
10use crate::{Resolution, opcodec, resolution_from_operation};
11
12const LATEST_OPERATION: &str = "SELECT operation FROM ops WHERE did = ?1 ORDER BY seq DESC LIMIT 1";
13
14pub struct Db {
15    pool: Pool,
16    strings: RwLock<Vec<String>>,
17}
18
19pub fn open_database(path: &Path, max_size: usize) -> Result<Arc<Db>> {
20    // read-only open refuses to create an empty database where plox's should be
21    let connection = Connection::open_with_flags(path, OpenFlags::SQLITE_OPEN_READ_ONLY)
22        .with_context(|| format!("opening {}", path.display()))?;
23    let version: i64 = connection.pragma_query_value(None, "user_version", |row| row.get(0))?;
24    ensure!(
25        version == 1,
26        "unsupported plox database version {version}; expected compact format version 1"
27    );
28    let compacting: bool = connection.query_row(
29        "SELECT EXISTS(SELECT 1 FROM sqlite_schema WHERE name = 'compaction')",
30        [],
31        |row| row.get(0),
32    )?;
33    if compacting {
34        let complete: bool = connection.query_row(
35            "SELECT complete FROM compaction WHERE singleton = 1",
36            [],
37            |row| row.get(0),
38        )?;
39        ensure!(
40            complete,
41            "database compaction is unfinished; resume plox compact before serving it"
42        );
43    }
44    connection
45        .prepare(LATEST_OPERATION)
46        .context("reading plox ops schema")?;
47    connection
48        .prepare("SELECT id, value FROM strings")
49        .context("reading plox string dictionary schema")?;
50
51    let pool = DatabaseConfig::new(path)
52        .builder(Runtime::Tokio1)?
53        .max_size(max_size)
54        .post_create(Hook::async_fn(|connection, _| {
55            Box::pin(async move {
56                match connection
57                    .interact(|connection| {
58                        connection.busy_timeout(std::time::Duration::from_secs(5))?;
59                        connection.pragma_update(None, "query_only", true)
60                    })
61                    .await
62                {
63                    Ok(Ok(())) => Ok(()),
64                    Ok(Err(error)) => Err(HookError::Backend(error)),
65                    Err(error) => Err(HookError::message(error.to_string())),
66                }
67            })
68        }))
69        .build()
70        .context("creating database connection pool")?;
71    Ok(Arc::new(Db {
72        pool,
73        strings: RwLock::new(Vec::new()),
74    }))
75}
76
77pub async fn resolve(db: &Arc<Db>, did: &str) -> Result<Option<Resolution>> {
78    let connection = db
79        .pool
80        .get()
81        .await
82        .context("checking out database connection")?;
83    let query_did = opcodec::encode_did(did);
84    let db = Arc::clone(db);
85    let operation = connection
86        .interact(move |connection| -> Result<_> {
87            let tx = connection.transaction()?;
88            let Some(program) = tx
89                .prepare_cached(LATEST_OPERATION)?
90                .query_row(params![query_did], |row| row.get::<_, Vec<u8>>(0))
91                .optional()?
92            else {
93                return Ok(None);
94            };
95            let strings = db
96                .strings
97                .read()
98                .map_err(|_| anyhow!("string dictionary lock poisoned"))?;
99            if let Ok(operation) =
100                opcodec::decode(&program, |id| strings.get(id).map(String::as_str))
101            {
102                return Ok(Some(operation));
103            }
104            drop(strings);
105
106            // plox's dictionary is append-only; retry with additions from the operation's snapshot.
107            let mut strings = db
108                .strings
109                .write()
110                .map_err(|_| anyhow!("string dictionary lock poisoned"))?;
111            let mut stmt =
112                tx.prepare_cached("SELECT id, value FROM strings WHERE id >= ?1 ORDER BY id")?;
113            let mut rows = stmt.query([strings.len() as i64])?;
114            while let Some(row) = rows.next()? {
115                let id: i64 = row.get(0)?;
116                ensure!(
117                    id == strings.len() as i64,
118                    "non-contiguous string dictionary"
119                );
120                strings.push(row.get(1)?);
121            }
122            opcodec::decode(&program, |id| strings.get(id).map(String::as_str))
123                .context("decoding PLC operation")
124                .map(Some)
125        })
126        .await
127        .map_err(|error| anyhow!("querying database: {error}"))??;
128    operation
129        .map(|operation| resolution_from_operation(did, operation))
130        .transpose()
131}