availability_evidence.rs (16660B)
1 //! Exact verified public listing evidence under the sole governed SQLite host. 2 3 use harvestcircle_domain::{ 4 AvailabilityEventVersion, AvailabilityObservation, AvailabilityUnsupportedReason, 5 AvailabilityVersionView, SafeError, SafeErrorCode, SafeMessage, UnixTimestamp, 6 }; 7 use radroots_event::envelope::event_head::EventHeadCoordinate; 8 use radroots_event::listing::classified::ClassifiedListingPartition; 9 use radroots_event::wire::{DEFAULT_RAW_JSON_MAX_BYTES, Nip01EventWire}; 10 use radroots_event_codec::verify::verify_nip01_event; 11 use radroots_service_sqlite::ServiceSqliteTransaction; 12 use sqlx::Row; 13 use sqlx::sqlite::SqliteRow; 14 15 use crate::Database; 16 use crate::db::{corrupt_storage, map_transaction_error, storage_unavailable}; 17 18 const VERSION_CAPACITY: i64 = 4096; 19 const TOTAL_PAYLOAD_BYTES: i64 = 134_217_728; 20 const ORDINARY_PAYLOAD_BYTES: i64 = 125_829_120; 21 const FIXED_VERSION_BYTES: usize = 92; 22 const RAW_IDENTIFIER_BYTES: usize = 4096; 23 const SOURCE_BYTES: usize = 2048; 24 const ADMISSION_BYTES: usize = 64; 25 const REJECTION_BYTES: usize = 64; 26 27 // Every byte-bearing column is projected before any Rust hydration or decoding. 28 // The retained full length distinguishes a bounded prefix from complete evidence. 29 pub(crate) const SELECT_AVAILABILITY_VERSION_SQL: &str = "SELECT \ 30 substr(event_id, 1, 33) AS event_id, length(event_id) AS event_id_bytes, \ 31 substr(author, 1, 33) AS author, length(author) AS author_bytes, kind, \ 32 substr(CAST(raw_d AS BLOB), 1, 4097) AS raw_d, \ 33 length(CAST(raw_d AS BLOB)) AS raw_d_bytes, \ 34 substr(signed_at, 1, 9) AS signed_at, length(signed_at) AS signed_at_bytes, \ 35 substr(published_at, 1, 9) AS published_at, length(published_at) AS published_at_bytes, \ 36 substr(CAST(original_json AS BLOB), 1, 262145) AS original_json, \ 37 length(CAST(original_json AS BLOB)) AS original_json_bytes, \ 38 substr(CAST(admission_label AS BLOB), 1, 65) AS admission_label, \ 39 length(CAST(admission_label AS BLOB)) AS admission_label_bytes, \ 40 substr(CAST(rejection_code AS BLOB), 1, 65) AS rejection_code, \ 41 length(CAST(rejection_code AS BLOB)) AS rejection_code_bytes, \ 42 substr(CAST(source AS BLOB), 1, 2049) AS source, \ 43 length(CAST(source AS BLOB)) AS source_bytes, observed_at_unix_s, payload_bytes \ 44 FROM availability_versions WHERE event_id = ? LIMIT 2"; 45 46 impl Database { 47 /// Retains the first exact signed wire and named provenance for a verified version. 48 /// 49 /// Public evidence does not install an account or authorize a local signer. 50 /// Duplicate IDs preserve their original evidence without consuming more quota. 51 pub async fn retain_availability_version( 52 &self, 53 view: AvailabilityVersionView, 54 ) -> Result<(), SafeError> { 55 self.host() 56 .transaction(|transaction| { 57 Box::pin(async move { retain_availability_version_on(transaction, view).await }) 58 }) 59 .await 60 .map_err(map_transaction_error) 61 } 62 63 /// Loads bounded original wire and re-verifies it with the selected shared codec. 64 /// 65 /// Missing evidence returns `None`. Corrupt evidence or global accounting 66 /// fails closed without resetting state or constructing a verified substitute. 67 pub async fn load_availability_version( 68 &self, 69 version: AvailabilityEventVersion, 70 ) -> Result<Option<AvailabilityVersionView>, SafeError> { 71 self.host() 72 .transaction(|transaction| { 73 Box::pin(async move { 74 validate_public_usage(transaction).await?; 75 let row = sqlx::query(SELECT_AVAILABILITY_VERSION_SQL) 76 .bind(version.event_id().as_bytes().as_slice()) 77 .fetch_optional(&mut *transaction) 78 .await 79 .map_err(|_| corrupt_storage())?; 80 row.map(decode_availability_row).transpose() 81 }) 82 }) 83 .await 84 .map_err(map_transaction_error) 85 } 86 } 87 88 /// The real insertion path, also exercised inside governed rollback fixtures. 89 pub(crate) async fn retain_availability_version_on( 90 transaction: &mut ServiceSqliteTransaction<'_>, 91 view: AvailabilityVersionView, 92 ) -> Result<(), SafeError> { 93 let (count, bytes) = validate_public_usage(transaction).await?; 94 let existing = sqlx::query(SELECT_AVAILABILITY_VERSION_SQL) 95 .bind(view.version().event_id().as_bytes().as_slice()) 96 .fetch_optional(&mut *transaction) 97 .await 98 .map_err(|_| corrupt_storage())?; 99 if let Some(row) = existing { 100 let existing = decode_availability_row(row)?; 101 // Signature-verified ID correlation binds signed data. Original JSON 102 // formatting, signature variants, extras and later provenance can differ. 103 if existing.version() != view.version() 104 || existing.publisher() != view.publisher() 105 || existing.created_at() != view.created_at() 106 || existing.raw_coordinate() != view.raw_coordinate() 107 || existing.focused() != view.focused() 108 || existing.unsupported_reason() != view.unsupported_reason() 109 { 110 return Err(corrupt_storage()); 111 } 112 return Ok(()); 113 } 114 115 let (label, reason) = admission(&view)?; 116 let charge = version_charge(&view, label, reason)?; 117 let next_bytes = bytes.checked_add(charge).ok_or_else(corrupt_storage)?; 118 if count >= VERSION_CAPACITY || next_bytes > ORDINARY_PAYLOAD_BYTES { 119 return Err(capacity()); 120 } 121 let reserved = sqlx::query( 122 "UPDATE public_payload_usage SET version_count = version_count + 1, \ 123 payload_bytes = payload_bytes + ? \ 124 WHERE singleton = 1 AND version_count = ? AND payload_bytes = ?", 125 ) 126 .bind(charge) 127 .bind(count) 128 .bind(bytes) 129 .execute(&mut *transaction) 130 .await 131 .map_err(|_| storage_unavailable())?; 132 if reserved.rows_affected() != 1 { 133 return Err(corrupt_storage()); 134 } 135 136 let signed_at = view.created_at().as_u64().to_be_bytes(); 137 let published_at = view 138 .focused() 139 .map(|projection| projection.published_at().as_u64().to_be_bytes()); 140 // Borrow the owned verified view after quota reservation: no full-wire clone. 141 let inserted = sqlx::query( 142 "INSERT INTO availability_versions \ 143 (event_id, author, kind, raw_d, signed_at, published_at, original_json, \ 144 admission_label, rejection_code, source, observed_at_unix_s, payload_bytes) \ 145 VALUES (?, ?, 30402, ?, ?, ?, ?, ?, ?, ?, ?, ?)", 146 ) 147 .bind(view.version().event_id().as_bytes().as_slice()) 148 .bind(view.publisher().public_key().as_bytes().as_slice()) 149 .bind(identifier(&view)?) 150 .bind(signed_at.as_slice()) 151 .bind(published_at.as_ref().map(|bytes| bytes.as_slice())) 152 .bind(view.original_json()) 153 .bind(label) 154 .bind(reason) 155 .bind(view.observation().source().as_str()) 156 .bind(view.observation().observed_at().as_seconds()) 157 .bind(charge) 158 .execute(&mut *transaction) 159 .await 160 .map_err(|_| storage_unavailable())?; 161 if inserted.rows_affected() != 1 { 162 return Err(corrupt_storage()); 163 } 164 Ok(()) 165 } 166 167 async fn validate_public_usage( 168 transaction: &mut ServiceSqliteTransaction<'_>, 169 ) -> Result<(i64, i64), SafeError> { 170 let rows = sqlx::query( 171 "SELECT singleton, version_count, payload_bytes FROM public_payload_usage LIMIT 2", 172 ) 173 .fetch_all(&mut *transaction) 174 .await 175 .map_err(|_| corrupt_storage())?; 176 let [row] = rows.as_slice() else { 177 return Err(corrupt_storage()); 178 }; 179 let singleton: i64 = row.try_get("singleton").map_err(|_| corrupt_storage())?; 180 let count: i64 = row 181 .try_get("version_count") 182 .map_err(|_| corrupt_storage())?; 183 let bytes: i64 = row 184 .try_get("payload_bytes") 185 .map_err(|_| corrupt_storage())?; 186 if singleton != 1 187 || !(0..=VERSION_CAPACITY).contains(&count) 188 || !(0..=TOTAL_PAYLOAD_BYTES).contains(&bytes) 189 { 190 return Err(corrupt_storage()); 191 } 192 // LIMIT bounds corruption inspection as well as normal state. The sums 193 // operate on byte lengths and integers without hydrating retained wire. 194 let actual = sqlx::query( 195 "SELECT count(*) AS actual_count, \ 196 coalesce(sum(payload_bytes), 0) AS recorded_bytes, \ 197 coalesce(sum(actual_charge), 0) AS actual_bytes, \ 198 coalesce(max(CASE WHEN payload_bytes = actual_charge THEN 0 ELSE 1 END), 0) AS mismatch \ 199 FROM (SELECT payload_bytes, 92 + length(CAST(original_json AS BLOB)) \ 200 + length(CAST(raw_d AS BLOB)) + length(CAST(source AS BLOB)) \ 201 + length(CAST(admission_label AS BLOB)) \ 202 + coalesce(length(CAST(rejection_code AS BLOB)), 0) AS actual_charge \ 203 FROM availability_versions LIMIT 4097)", 204 ) 205 .fetch_one(&mut *transaction) 206 .await 207 .map_err(|_| corrupt_storage())?; 208 let actual_count: i64 = actual 209 .try_get("actual_count") 210 .map_err(|_| corrupt_storage())?; 211 let recorded_bytes: i64 = actual 212 .try_get("recorded_bytes") 213 .map_err(|_| corrupt_storage())?; 214 let actual_bytes: i64 = actual 215 .try_get("actual_bytes") 216 .map_err(|_| corrupt_storage())?; 217 let mismatch: i64 = actual.try_get("mismatch").map_err(|_| corrupt_storage())?; 218 if actual_count != count || recorded_bytes != bytes || actual_bytes != bytes || mismatch != 0 { 219 return Err(corrupt_storage()); 220 } 221 Ok((count, bytes)) 222 } 223 224 pub(crate) fn decode_availability_row( 225 row: SqliteRow, 226 ) -> Result<AvailabilityVersionView, SafeError> { 227 let event_id = exact_bytes::<32>(&row, "event_id", "event_id_bytes")?; 228 let author = exact_bytes::<32>(&row, "author", "author_bytes")?; 229 let signed_at = u64::from_be_bytes(exact_bytes::<8>(&row, "signed_at", "signed_at_bytes")?); 230 let published_at = optional_bytes(&row, "published_at", "published_at_bytes", 8)? 231 .map(|bytes| { 232 bytes 233 .try_into() 234 .map(u64::from_be_bytes) 235 .map_err(|_| corrupt_storage()) 236 }) 237 .transpose()?; 238 let kind: i64 = row.try_get("kind").map_err(|_| corrupt_storage())?; 239 let observed_at: i64 = row 240 .try_get("observed_at_unix_s") 241 .map_err(|_| corrupt_storage())?; 242 let charge: i64 = row 243 .try_get("payload_bytes") 244 .map_err(|_| corrupt_storage())?; 245 if kind != 30402 || observed_at < 0 { 246 return Err(corrupt_storage()); 247 } 248 let raw_d = text(&row, "raw_d", "raw_d_bytes", RAW_IDENTIFIER_BYTES)?; 249 let source = text(&row, "source", "source_bytes", SOURCE_BYTES)?; 250 let label = text( 251 &row, 252 "admission_label", 253 "admission_label_bytes", 254 ADMISSION_BYTES, 255 )?; 256 let reason = optional_bytes( 257 &row, 258 "rejection_code", 259 "rejection_code_bytes", 260 REJECTION_BYTES, 261 )? 262 .map(|bytes| String::from_utf8(bytes).map_err(|_| corrupt_storage())) 263 .transpose()?; 264 if reason 265 .as_ref() 266 .is_some_and(|code| code.is_empty() || !code.is_ascii()) 267 { 268 return Err(corrupt_storage()); 269 } 270 let original_json = text( 271 &row, 272 "original_json", 273 "original_json_bytes", 274 DEFAULT_RAW_JSON_MAX_BYTES, 275 )?; 276 let verified = verify_nip01_event( 277 Nip01EventWire::parse_json_unverified(&original_json) 278 .map_err(|_| corrupt_storage())? 279 .into_unverified_envelope() 280 .map_err(|_| corrupt_storage())?, 281 ) 282 .map_err(|_| corrupt_storage())?; 283 let observation = AvailabilityObservation::parse( 284 &source, 285 UnixTimestamp::from_seconds(observed_at).ok_or_else(corrupt_storage)?, 286 ) 287 .map_err(|_| corrupt_storage())?; 288 let view = AvailabilityVersionView::from_verified(verified, &original_json, observation) 289 .map_err(|_| corrupt_storage())?; 290 let (expected_label, expected_reason) = admission(&view)?; 291 let expected_publication = view 292 .focused() 293 .map(|projection| projection.published_at().as_u64()); 294 if view.version().event_id().as_bytes() != &event_id 295 || view.publisher().public_key().as_bytes() != &author 296 || view.created_at().as_u64() != signed_at 297 || identifier(&view)? != raw_d 298 || view.observation().source().as_str() != source 299 || label != expected_label 300 || reason.as_deref() != expected_reason 301 || published_at != expected_publication 302 || charge != version_charge(&view, expected_label, expected_reason)? 303 { 304 return Err(corrupt_storage()); 305 } 306 Ok(view) 307 } 308 309 fn bounded_bytes( 310 row: &SqliteRow, 311 field: &str, 312 length_field: &str, 313 maximum: usize, 314 ) -> Result<Vec<u8>, SafeError> { 315 let length: i64 = row.try_get(length_field).map_err(|_| corrupt_storage())?; 316 let length = usize::try_from(length).map_err(|_| corrupt_storage())?; 317 if length > maximum { 318 return Err(corrupt_storage()); 319 } 320 let bytes: Vec<u8> = row.try_get(field).map_err(|_| corrupt_storage())?; 321 if bytes.len() != length { 322 return Err(corrupt_storage()); 323 } 324 Ok(bytes) 325 } 326 327 fn exact_bytes<const N: usize>( 328 row: &SqliteRow, 329 field: &str, 330 length_field: &str, 331 ) -> Result<[u8; N], SafeError> { 332 bounded_bytes(row, field, length_field, N)? 333 .try_into() 334 .map_err(|_| corrupt_storage()) 335 } 336 337 fn optional_bytes( 338 row: &SqliteRow, 339 field: &str, 340 length_field: &str, 341 maximum: usize, 342 ) -> Result<Option<Vec<u8>>, SafeError> { 343 let length: Option<i64> = row.try_get(length_field).map_err(|_| corrupt_storage())?; 344 if length.is_some() { 345 return bounded_bytes(row, field, length_field, maximum).map(Some); 346 } 347 let bytes: Option<Vec<u8>> = row.try_get(field).map_err(|_| corrupt_storage())?; 348 if bytes.is_some() { 349 return Err(corrupt_storage()); 350 } 351 Ok(None) 352 } 353 354 fn text( 355 row: &SqliteRow, 356 field: &str, 357 length_field: &str, 358 maximum: usize, 359 ) -> Result<String, SafeError> { 360 String::from_utf8(bounded_bytes(row, field, length_field, maximum)?) 361 .map_err(|_| corrupt_storage()) 362 } 363 364 fn identifier(view: &AvailabilityVersionView) -> Result<&str, SafeError> { 365 match view.raw_coordinate() { 366 EventHeadCoordinate::Addressable { 367 kind: 30402, d_tag, .. 368 } => Ok(d_tag), 369 _ => Err(corrupt_storage()), 370 } 371 } 372 373 fn admission( 374 view: &AvailabilityVersionView, 375 ) -> Result<(&'static str, Option<&'static str>), SafeError> { 376 match (view.focused(), view.unsupported_reason().copied()) { 377 (Some(_), None) => Ok(("focused", None)), 378 (None, Some(AvailabilityUnsupportedReason::Excluded(partition))) => Ok(( 379 match partition { 380 ClassifiedListingPartition::FocusedFoodAvailability => "excluded_focused", 381 ClassifiedListingPartition::OperationalListing => "excluded_operational", 382 ClassifiedListingPartition::GenericNip99 => "excluded_generic", 383 ClassifiedListingPartition::Ambiguous => "excluded_ambiguous", 384 }, 385 None, 386 )), 387 (None, Some(AvailabilityUnsupportedReason::ProjectionRejected(code))) 388 if !code.is_empty() && code.len() <= REJECTION_BYTES && code.is_ascii() => 389 { 390 Ok(("projection_rejected", Some(code))) 391 } 392 _ => Err(corrupt_storage()), 393 } 394 } 395 396 fn version_charge( 397 view: &AvailabilityVersionView, 398 label: &str, 399 reason: Option<&str>, 400 ) -> Result<i64, SafeError> { 401 let identifier = identifier(view)?; 402 let source = view.observation().source().as_str(); 403 if view.original_json().len() > DEFAULT_RAW_JSON_MAX_BYTES 404 || identifier.len() > RAW_IDENTIFIER_BYTES 405 || source.is_empty() 406 || source.len() > SOURCE_BYTES 407 || label.is_empty() 408 || label.len() > ADMISSION_BYTES 409 || reason 410 .is_some_and(|code| code.is_empty() || code.len() > REJECTION_BYTES || !code.is_ascii()) 411 { 412 return Err(corrupt_storage()); 413 } 414 let bytes = [ 415 view.original_json().len(), 416 identifier.len(), 417 source.len(), 418 label.len(), 419 reason.map_or(0, str::len), 420 ] 421 .into_iter() 422 .try_fold(FIXED_VERSION_BYTES, usize::checked_add) 423 .ok_or_else(corrupt_storage)?; 424 i64::try_from(bytes).map_err(|_| corrupt_storage()) 425 } 426 427 const fn capacity() -> SafeError { 428 SafeError::new( 429 SafeErrorCode::AvailabilityCapacity, 430 SafeMessage::new("The retained availability evidence is at capacity."), 431 ) 432 }