use sea_orm::{ColumnTrait, EntityTrait, QueryFilter, QuerySelect};
use signer_core::{SignerKeys, SignerUser};
use signer_crdt::{SignerMeta, entity::crdt_event, view::CrdtEventVO};
use std::collections::HashMap;
use crate::{
error::{RemoteError, RemoteResult},
remote::{CrdtCryptedEventVO, HttpClient, HttpClientConfig},
};
#[derive(Debug, Clone, Default)]
pub struct CrdtFrontier(pub HashMap<String, i32>);
impl CrdtFrontier {
pub fn new() -> Self {
CrdtFrontier(HashMap::new())
}
pub async fn get_from_remote(addr: &str, keys: &SignerKeys, user: &SignerUser) -> RemoteResult<CrdtFrontier> {
let config = HttpClientConfig::new(keys.clone(), user.clone(), addr.to_string());
let client = HttpClient::new(config);
#[derive(serde::Serialize)]
struct QueryParams {
frontiers: Option<String>,
}
let query = QueryParams {
frontiers: Some("{}".to_string()), };
let response: serde_json::Value = client.get_with_query("/api/crdt-frontiers", &query).await?;
let frontiers_map: HashMap<String, i32> =
serde_json::from_value(response.clone()).map_err(|e| RemoteError::Json {
source: e,
reason: "Failed to deserialize frontiers from JSON Value".to_string(),
input: format!("{:?}", response),
})?;
Ok(CrdtFrontier(frontiers_map))
}
pub async fn get_from_signer(meta: &SignerMeta) -> RemoteResult<CrdtFrontier> {
let ce_vec: Vec<(i32, String)> = crdt_event::Entity::find()
.select_only()
.column_as(crdt_event::Column::Clock.max(), "clock")
.column(crdt_event::Column::Peer)
.group_by(crdt_event::Column::Peer)
.into_tuple()
.all(&meta.conn)
.await
.map_err(|e| RemoteError::Internal(format!("查询 CRDT 事件前沿信息失败: {}", e)))?;
let mut frontiers = HashMap::new();
for ce in ce_vec {
frontiers.insert(ce.1, ce.0);
}
Ok(CrdtFrontier(frontiers))
}
pub async fn get_events_from_signer(
&self,
meta: &SignerMeta,
) -> RemoteResult<Vec<CrdtEventVO>> {
let query = crdt_event::Entity::find();
let mut filter = crdt_event::Column::Peer.is_null();
for (key, value) in self.0.iter() {
filter = filter.or(crdt_event::Column::Clock
.gt(*value)
.and(crdt_event::Column::Peer.eq(key)));
}
filter = filter.or(crdt_event::Column::Peer.is_not_in(self.0.keys()));
let events = query
.filter(filter)
.all(&meta.conn)
.await
.map_err(|e| RemoteError::Internal(format!("查询 CRDT 事件失败: {}", e)))?;
Ok(events.into_iter().map(|i| i.into()).collect())
}
pub async fn get_events_from_remote(
&self,
addr: &str,
keys: &SignerKeys,
user: &SignerUser,
) -> RemoteResult<Vec<CrdtEventVO>> {
let frontiers_json = serde_json::to_string(&self.0).map_err(|e| RemoteError::Json {
source: e,
reason: "Failed to serialize frontiers to JSON".to_string(),
input: format!("{:?}", self.0),
})?;
let crypted_events = CrdtCryptedEventVO::pull(addr, keys, user, &frontiers_json).await?;
let events = crypted_events
.into_iter()
.map(|e| {
e.decrypt(keys)
.map_err(|err| RemoteError::Internal(format!("解密 CRDT 加密事件失败: {}", err)))
})
.collect::<Result<Vec<CrdtEventVO>, RemoteError>>()?;
Ok(events)
}
pub fn inner(&self) -> &HashMap<String, i32> {
&self.0
}
pub fn inner_mut(&mut self) -> &mut HashMap<String, i32> {
&mut self.0
}
}