use signer_auth::{SignerJWT, SignerJWTClaims, SignerJWTHeader};
use signer_core::{SignerKeys, SignerSigned, SignerUser};
use signer_crdt::{SignerMeta, crdt::reconcile, view::CrdtEventVO};
use crate::{
error::{RemoteError, RemoteResult},
remote::{
CrdtFrontier, HttpClient, HttpClientConfig,
eventsource::SignerRemoteEventSource, resource::ResourceUsage, CrdtCryptedEventVO,
},
};
#[derive(Clone)]
pub struct SignerRemote {
addr: String,
event_source_terminator: Option<tokio::sync::mpsc::Sender<()>>,
}
impl SignerRemote {
pub fn new(addr: &str) -> Self {
Self {
addr: addr.to_string(),
event_source_terminator: None,
}
}
pub async fn ping(&self, keys: &SignerKeys, user: &SignerUser) -> RemoteResult<()> {
let config = HttpClientConfig::new(keys.clone(), user.clone(), self.addr.clone());
let client = HttpClient::new(config);
let _response: serde_json::Value = client.get("/api").await?;
Ok(())
}
pub async fn get_crdt_events_from_remote(
&self,
frontiers: &CrdtFrontier,
keys: &SignerKeys,
user: &SignerUser,
) -> RemoteResult<Vec<signer_crdt::view::CrdtEventVO>> {
frontiers.get_events_from_remote(&self.addr, keys, user).await
}
pub async fn sync_crdt_event(&self, meta: &SignerMeta) -> RemoteResult<()> {
let user = meta.get_current_user().await
.map_err(|e| RemoteError::Internal(format!("获取当前用户失败: {}", e)))?;
let keys = &meta.keys;
let frontiers = CrdtFrontier::get_from_signer(meta).await?;
let frontiers_json = serde_json::to_string(frontiers.inner())
.map_err(|e| RemoteError::Internal(format!("序列化前沿信息失败: {}", e)))?;
let crypted_events = CrdtCryptedEventVO::pull(&self.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>>()?;
if !events.is_empty() {
signer_crdt::view::CrdtEventVO::insert_many(events, meta)
.await
.map_err(|e| RemoteError::Internal(format!("CRDT 插入事件失败: {}", e)))?;
}
let frontiers = CrdtFrontier::get_from_remote(&self.addr, keys, &user).await?;
let local_events = frontiers.get_events_from_signer(meta).await?;
let events = local_events
.into_iter()
.map(|e| CrdtCryptedEventVO::encrypt(keys, &e))
.collect::<Result<Vec<CrdtCryptedEventVO>, _>>()
.map_err(|e| RemoteError::Internal(format!("创建 CRDT 加密事件失败: {}", e)))?;
CrdtCryptedEventVO::push(events, &self.addr, keys, &user).await?;
reconcile(meta)
.await
.map_err(|e| RemoteError::Internal(format!("CRDT 协调失败: {}", e)))?;
Ok(())
}
pub async fn pull_user(&self, meta: &SignerMeta, user_key: &str) -> RemoteResult<()> {
let user = meta.get_current_user().await
.map_err(|e| RemoteError::Internal(format!("获取当前用户失败: {}", e)))?;
let keys = &meta.keys;
let config = HttpClientConfig::new(keys.clone(), user.clone(), self.addr.clone());
let client = HttpClient::new(config);
let r: SignerSigned<SignerUser> =
client.get(&format!("/api/users/{}", user_key)).await?;
let signed_up = SignerSigned::<SignerUser> {
sig: r.sig,
msg: r.msg,
pubkey: r.pubkey,
_marker: Default::default(),
};
let up = signed_up
.verify_to_value()
.map_err(|e| RemoteError::Internal(format!("验证签名用户信息失败: {}", e)))?;
let user_vo = signer_crdt::view::UserVO::from_user_data(keys, &up).await
.map_err(|e| RemoteError::Internal(format!("创建用户视图对象失败: {}", e)))?;
user_vo.put(meta).await
.map_err(|e| RemoteError::Internal(format!("保存用户信息失败: {}", e)))?;
Ok(())
}
pub async fn open_eventsource(
&mut self,
meta: SignerMeta,
) -> RemoteResult<()> {
if self.event_source_terminator.is_some() {
return Err(RemoteError::Internal("事件源已存在".to_string()));
}
let (tx, mut es) = SignerRemoteEventSource::new(&self.addr, meta);
self.event_source_terminator = Some(tx);
tokio::spawn(async move {
es.open_eventsource().await.expect("打开事件源失败");
});
Ok(())
}
pub async fn close(&mut self) {
if let Some(event_source) = &mut self.event_source_terminator {
let _ = event_source.send(()).await;
self.event_source_terminator = None;
}
}
pub async fn get_resource_usage(&self, keys: &SignerKeys, user: &SignerUser) -> RemoteResult<ResourceUsage> {
tokio::time::timeout(std::time::Duration::from_secs(10), async {
let config = HttpClientConfig::new_no_auth(self.addr.clone());
let client = HttpClient::new(config);
let jwt = SignerJWT::new(
SignerJWTHeader::default(user),
SignerJWTClaims::default(keys, user, self.addr.clone(), uuid::Uuid::new_v4().to_string())
.with_expired_duration(chrono::Duration::minutes(5)),
);
let token = jwt
.encode(keys)
.map_err(|e| RemoteError::Internal(format!("JWT 编码失败: {}", e)))?;
let usage: ResourceUsage = client
.get_with_header(
"/api/resource-usage",
"Authorization",
&format!("Bearer {}", token),
)
.await?;
Ok::<ResourceUsage, RemoteError>(usage)
})
.await
.map_err(|_| RemoteError::Internal("资源用量查询超时".to_string()))?
}
pub async fn get_crdt_frontiers(&self, keys: &SignerKeys, user: &SignerUser) -> RemoteResult<CrdtFrontier> {
CrdtFrontier::get_from_remote(&self.addr, keys, user).await
}
}