From d9a146b6df6d805464c8351b99d3347b519be64a Mon Sep 17 00:00:00 2001 From: dawn <90008@gaze.systems> Date: Sat, 21 Feb 2026 03:28:45 +0300 Subject: [PATCH] fjall: store serde_json::Value in db instead of storing RawValue --- src/mirror/fjall.rs | 39 ++++++++++++++++++++------------------- src/plc_fjall.rs | 24 +++++++++++++++++------- 2 files changed, 37 insertions(+), 26 deletions(-) diff --git a/src/mirror/fjall.rs b/src/mirror/fjall.rs index e1ca187..3ab3843 100644 --- a/src/mirror/fjall.rs +++ b/src/mirror/fjall.rs @@ -167,12 +167,12 @@ async fn fjall_resolve(req: &Request, Data(state): Data<&FjallState>) -> Result< match sub_path { "" => { - let parsed: Vec = ops - .iter() - .filter(|op| !op.nullified) - .filter_map(|op| serde_json::from_str(op.operation.get()).ok()) - .collect(); - let data = doc::apply_op_log(did_str, &parsed); + let data = doc::apply_op_log( + did_str, + ops.iter() + .filter(|op| !op.nullified) + .map(|op| &op.operation), + ); let Some(data) = data else { return Err(Error::from_string( format!("DID not available: {did_str}"), @@ -185,10 +185,10 @@ async fn fjall_resolve(req: &Request, Data(state): Data<&FjallState>) -> Result< .body(serde_json::to_string(&doc).unwrap())) } "/log" => { - let log: Vec = ops + let log: Vec<&serde_json::Value> = ops .iter() .filter(|op| !op.nullified) - .filter_map(|op| serde_json::from_str(op.operation.get()).ok()) + .map(|op| &op.operation) .collect(); Ok(Response::builder() .content_type("application/json") @@ -200,7 +200,7 @@ async fn fjall_resolve(req: &Request, Data(state): Data<&FjallState>) -> Result< .map(|op| { serde_json::json!({ "did": op.did, - "operation": serde_json::from_str::(op.operation.get()).unwrap_or_default(), + "operation": op.operation, "cid": op.cid, "nullified": op.nullified, "createdAt": op.created_at.to_rfc3339(), @@ -212,10 +212,11 @@ async fn fjall_resolve(req: &Request, Data(state): Data<&FjallState>) -> Result< .body(serde_json::to_string(&audit).unwrap())) } "/log/last" => { - let last = - ops.iter().filter(|op| !op.nullified).last().and_then(|op| { - serde_json::from_str::(op.operation.get()).ok() - }); + let last = ops + .iter() + .filter(|op| !op.nullified) + .last() + .map(|op| &op.operation); let Some(last) = last else { return Err(Error::from_string( format!("DID not available: {did_str}"), @@ -227,12 +228,12 @@ async fn fjall_resolve(req: &Request, Data(state): Data<&FjallState>) -> Result< .body(serde_json::to_string(&last).unwrap())) } "/data" => { - let parsed: Vec = ops - .iter() - .filter(|op| !op.nullified) - .filter_map(|op| serde_json::from_str(op.operation.get()).ok()) - .collect(); - let data = doc::apply_op_log(did_str, &parsed); + let data = doc::apply_op_log( + did_str, + ops.iter() + .filter(|op| !op.nullified) + .map(|op| &op.operation), + ); let Some(data) = data else { return Err(Error::from_string( format!("DID not available: {did_str}"), diff --git a/src/plc_fjall.rs b/src/plc_fjall.rs index 9afe499..990c4aa 100644 --- a/src/plc_fjall.rs +++ b/src/plc_fjall.rs @@ -1,4 +1,4 @@ -use crate::{Dt, ExportPage, Op, PageBoundaryState}; +use crate::{Dt, ExportPage, Op as CommonOp, PageBoundaryState}; use data_encoding::BASE32_NOPAD; use fjall::{Database, Keyspace, KeyspaceCreateOptions, OwnedWriteBatch, PersistMode}; use serde::{Deserialize, Serialize}; @@ -78,6 +78,16 @@ fn decode_timestamp(key: &[u8]) -> anyhow::Result
{ .ok_or_else(|| anyhow::anyhow!("invalid timestamp {micros}")) } +// we have our own Op struct for fjall since we dont want to have to convert Value back to RawValue +#[derive(Debug, Serialize)] +pub struct Op { + pub did: String, + pub cid: String, + pub created_at: Dt, + pub nullified: bool, + pub operation: serde_json::Value, +} + // this is basically Op, but without the cid and created_at fields // since we have them in the key already #[derive(Debug, Deserialize, Serialize)] @@ -86,7 +96,7 @@ struct DbOp { #[serde(with = "serde_bytes")] pub did: Vec, pub nullified: bool, - pub operation: Box, + pub operation: serde_json::Value, } #[derive(Clone)] @@ -142,7 +152,7 @@ impl FjallDb { .map(Some) } - pub fn insert_op(&self, batch: &mut OwnedWriteBatch, op: &Op) -> anyhow::Result { + pub fn insert_op(&self, batch: &mut OwnedWriteBatch, op: &CommonOp) -> anyhow::Result { let pk = by_did_key(&op.did, &op.created_at, &op.cid)?; if self.inner.by_did.get(&pk)?.is_some() { return Ok(0); @@ -155,9 +165,9 @@ impl FjallDb { let db_op = DbOp { did: encoded_did, nullified: op.nullified, - operation: op.operation.clone(), + operation: serde_json::to_value(&op.operation)?, }; - let value = rmp_serde::to_vec_named(&db_op)?; + let value = rmp_serde::to_vec(&db_op)?; batch.insert(&self.inner.ops, &ts_key, &value); batch.insert(&self.inner.by_did, &pk, &[]); Ok(1) @@ -180,10 +190,10 @@ impl FjallDb { let ts_bytes = key_rest .get(..8) - .ok_or_else(|| anyhow::anyhow!("invalid length"))?; + .ok_or_else(|| anyhow::anyhow!("invalid length: {key_rest:?}"))?; let cid_bytes = key_rest .get(9..) - .ok_or_else(|| anyhow::anyhow!("invalid length"))?; + .ok_or_else(|| anyhow::anyhow!("invalid length: {key_rest:?}"))?; let op_key = [ts_bytes, &[SEP][..], cid_bytes].concat(); let ts = decode_timestamp(ts_bytes)?; -- 2.51.2