char/plc-mirror
git clone https://git.t4t.associates/char/plc-mirror
c2f2166
main
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 { 15pool : Pool , 16strings : 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 21let connection =Connection :: open_with_flags ( path, OpenFlags :: SQLITE_OPEN_READ_ONLY ) 22. with_context ( ||format! ( "opening {}" , path. display ())) ?; 23let version: i64 = connection. pragma_query_value ( None , "user_version" , |row| row. get ( 0 )) ?; 24ensure! ( 25 version ==1 , 26"unsupported plox database version {version}; expected compact format version 1" 27); 28let compacting: bool = connection. query_row ( 29"SELECT EXISTS(SELECT 1 FROM sqlite_schema WHERE name = 'compaction')" , 30[], 31 |row| row. get ( 0 ), 32) ?; 33if compacting{ 34let complete: bool = connection. query_row ( 35"SELECT complete FROM compaction WHERE singleton = 1" , 36[], 37 |row| row. get ( 0 ), 38) ?; 39ensure! ( 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 51let pool =DatabaseConfig :: new ( path) 52. builder ( Runtime :: Tokio1 ) ? 53. max_size ( max_size) 54. post_create ( Hook :: async_fn ( |connection, _|{ 55Box :: pin ( async move { 56match 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{ 63Ok ( Ok (())) =>Ok (()), 64Ok ( Err ( error)) =>Err ( HookError :: Backend ( error)), 65Err ( error) =>Err ( HookError :: message ( error. to_string ())), 66} 67}) 68})) 69. build () 70. context ( "creating database connection pool" ) ?; 71Ok ( Arc :: new ( Db { 72 pool, 73strings : RwLock :: new ( Vec :: new ()), 74})) 75} 76 77pub async fn resolve ( db : & Arc < Db >, did : & str ) ->Result < Option < Resolution >> { 78let connection = db 79. pool 80. get () 81. await 82. context ( "checking out database connection" ) ?; 83let query_did = opcodec:: encode_did ( did); 84let db =Arc :: clone ( db); 85let operation = connection 86. interact ( move |connection| ->Result < _ > { 87let tx = connection. transaction () ?; 88let Some ( program) = tx 89. prepare_cached ( LATEST_OPERATION ) ? 90. query_row ( params! [ query_did], |row| row. get ::< _ , Vec < u8 >>( 0 )) 91. optional () ? 92else { 93return Ok ( None ); 94}; 95let strings = db 96. strings 97. read () 98. map_err ( |_|anyhow! ( "string dictionary lock poisoned" )) ?; 99if let Ok ( operation) = 100 opcodec:: decode ( & program, |id| strings. get ( id). map ( String :: as_str)) 101{ 102return Ok ( Some ( operation)); 103} 104drop ( strings); 105 106// plox's dictionary is append-only; retry with additions from the operation's snapshot. 107let mut strings = db 108. strings 109. write () 110. map_err ( |_|anyhow! ( "string dictionary lock poisoned" )) ?; 111let mut stmt = 112 tx. prepare_cached ( "SELECT id, value FROM strings WHERE id >= ?1 ORDER BY id" ) ?; 113let mut rows = stmt. query ([ strings. len () as i64 ]) ?; 114while let Some ( row) = rows. next () ?{ 115let id: i64 = row. get ( 0 ) ?; 116ensure! ( 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}