signer-remote 0.4.1

Signer remote communication package.
Documentation
//! 远程客户端实现

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(())
    }

    /// 根据前沿信息从远程服务器获取 CRDT 事件
    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
    }

    /// 使用 frontiers 同步 CRDT 信息
    pub async fn sync_crdt_event(&self, meta: &SignerMeta) -> RemoteResult<()> {
        // 从 SignerMeta 获取当前用户
        let user = meta.get_current_user().await
            .map_err(|e| RemoteError::Internal(format!("获取当前用户失败: {}", e)))?;
        let keys = &meta.keys;

        // 从远端拉取 CRDT 信息进行同步
        // 从本地获取前沿信息
        let frontiers = CrdtFrontier::get_from_signer(meta).await?;

        // 将前沿信息转换为 JSON 字符串
        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>>()?;

        // 使用 CrdtEventVO 的 insert_many 方法批量插入事件
        if !events.is_empty() {
            signer_crdt::view::CrdtEventVO::insert_many(events, meta)
                .await
                .map_err(|e| RemoteError::Internal(format!("CRDT 插入事件失败: {}", e)))?;
        }

        // 将本地 CRDT 信息推送到远端
        // 从远程服务器获取前沿信息
        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 方法推送事件
        CrdtCryptedEventVO::push(events, &self.addr, keys, &user).await?;

        // 使用我们新建的 reconcile 函数替换原来的 apply_all
        reconcile(meta)
            .await
            .map_err(|e| RemoteError::Internal(format!("CRDT 协调失败: {}", e)))?;

        // 注意:这里我们移除了 daemon.sync_user_public().await? 调用,
        // 因为在新的方法签名中我们没有 SignerDaemon 对象。
        // 如果需要同步用户信息,应该在其他地方处理。

        Ok(())
    }

    /// 从远程服务器拉取用户信息
    pub async fn pull_user(&self, meta: &SignerMeta, user_key: &str) -> RemoteResult<()> {
        // 从 SignerMeta 获取当前用户
        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);

            // 生成 JWT token 以获取当前用户的资源使用量
            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()))?
    }

    /// 从远程服务器获取 CRDT 前沿信息
    pub async fn get_crdt_frontiers(&self, keys: &SignerKeys, user: &SignerUser) -> RemoteResult<CrdtFrontier> {
        CrdtFrontier::get_from_remote(&self.addr, keys, user).await
    }
}