app

Local-first trade for farms and co-ops
git clone https://radroots.dev/git/app.git
Log | Files | Refs | README | LICENSE

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 }