diff --git a/src/mirror/fjall.rs b/src/mirror/fjall.rs index 420ee5b..abbd044 100644 --- a/src/mirror/fjall.rs +++ b/src/mirror/fjall.rs @@ -285,7 +285,7 @@ async fn export( .await?; let stream = futures::stream::iter(ops).map(|op| { - let mut json = serde_json::to_string(&op).unwrap(); + let mut json = serde_json::to_string(&op.to_sequenced_json()).unwrap(); json.push('\n'); Ok::<_, std::io::Error>(json) }); @@ -340,7 +340,9 @@ async fn export_stream( // check that the provided cursor is not stale if (chrono::Utc::now() - created_at).num_days() > 1 { return Err(Error::from_string( - format!("cursor {cursor} is stale, catch up using /export first"), + format!( + "cursor {cursor} is stale (older than a day), catch up using /export first" + ), StatusCode::BAD_REQUEST, )); } @@ -382,12 +384,12 @@ async fn export_stream( }); while let Some(op) = op_rx.recv().await { - cursor = op.seq; - let json = serde_json::to_string(&op).unwrap(); + let json = serde_json::to_string(&op.to_sequenced_json()).unwrap(); if let Err(e) = socket.send(Message::Text(json)).await { log::warn!("closing export stream: {e}"); return; } + cursor = op.seq; } tokio::select! { diff --git a/src/plc_fjall.rs b/src/plc_fjall.rs index 21a56d2..588a107 100644 --- a/src/plc_fjall.rs +++ b/src/plc_fjall.rs @@ -834,11 +834,23 @@ pub struct Op { pub seq: u64, pub did: String, pub cid: String, + #[serde(rename = "createdAt")] pub created_at: Dt, pub nullified: bool, pub operation: serde_json::Value, } +impl Op { + /// adds the `type` field to the op + pub fn to_sequenced_json(&self) -> serde_json::Value { + let mut val = serde_json::to_value(self).expect("Op is serializable"); + if let serde_json::Value::Object(ref mut map) = val { + map.insert("type".to_string(), "sequenced_op".into()); + } + val + } +} + #[derive(Clone)] pub struct FjallDb { inner: Arc, @@ -1313,8 +1325,10 @@ impl FjallDb { // we can start two threads, one for forward iteration and one for reverse iteration // this way we have two scans in parallel which should be faster! - let f_handle = spawn_scan_thread!(next, 0, false, ops / 2); - let b_handle = spawn_scan_thread!(next_back, workers / 2, true, ops - (ops / 2)); + let f_count = ops / 2; + let f_handle = spawn_scan_thread!(next, 0, false, f_count); + let b_count = ops - f_count; + let b_handle = spawn_scan_thread!(next_back, workers / 2, true, b_count); f_handle.join().unwrap()?; b_handle.join().unwrap()?;