Skip to main content

signer_remote/remote/
client.rs

1//! 远程客户端实现
2
3use signer_auth::{SignerJWT, SignerJWTClaims, SignerJWTHeader};
4use signer_core::{SignerKeys, SignerSigned, SignerUser};
5use signer_crdt::{SignerMeta, crdt::reconcile, view::CrdtEventVO};
6
7use crate::{
8    error::{RemoteError, RemoteResult},
9    remote::{
10        CrdtFrontier, HttpClient, HttpClientConfig,
11        eventsource::SignerRemoteEventSource, resource::ResourceUsage, CrdtCryptedEventVO,
12    },
13};
14
15#[derive(Clone)]
16pub struct SignerRemote {
17    addr: String,
18    event_source_terminator: Option<tokio::sync::mpsc::Sender<()>>,
19}
20
21
22
23impl SignerRemote {
24    /// 创建新的远程客户端
25    pub fn new(addr: &str) -> Self {
26        Self {
27            addr: addr.to_string(),
28            event_source_terminator: None,
29        }
30    }
31
32    /// 测试与远程服务器的连接
33    pub async fn ping(&self, keys: &SignerKeys, user: &SignerUser) -> RemoteResult<()> {
34        let config = HttpClientConfig::new(keys.clone(), user.clone(), self.addr.clone());
35        let client = HttpClient::new(config);
36
37        let _response: serde_json::Value = client.get("/api").await?;
38
39        Ok(())
40    }
41
42    /// 根据前沿信息从远程服务器获取 CRDT 事件
43    pub async fn get_crdt_events_from_remote(
44        &self,
45        frontiers: &CrdtFrontier,
46        keys: &SignerKeys,
47        user: &SignerUser,
48    ) -> RemoteResult<Vec<signer_crdt::view::CrdtEventVO>> {
49        frontiers.get_events_from_remote(&self.addr, keys, user).await
50    }
51
52    /// 使用 frontiers 同步 CRDT 信息
53    pub async fn sync_crdt_event(&self, meta: &SignerMeta) -> RemoteResult<()> {
54        // 从 SignerMeta 获取当前用户
55        let user = meta.get_current_user().await
56            .map_err(|e| RemoteError::Internal(format!("获取当前用户失败: {}", e)))?;
57        let keys = &meta.keys;
58
59        // 从远端拉取 CRDT 信息进行同步
60        // 从本地获取前沿信息
61        let frontiers = CrdtFrontier::get_from_signer(meta).await?;
62
63        // 将前沿信息转换为 JSON 字符串
64        let frontiers_json = serde_json::to_string(frontiers.inner())
65            .map_err(|e| RemoteError::Internal(format!("序列化前沿信息失败: {}", e)))?;
66
67        // 从远程服务器获取加密事件
68        let crypted_events = CrdtCryptedEventVO::pull(&self.addr, keys, &user, &frontiers_json).await?;
69
70        // 解密事件
71        let events = crypted_events
72            .into_iter()
73            .map(|e| {
74                e.decrypt(keys)
75                    .map_err(|err| RemoteError::Internal(format!("解密 CRDT 加密事件失败: {}", err)))
76            })
77            .collect::<Result<Vec<CrdtEventVO>, RemoteError>>()?;
78
79        // 使用 CrdtEventVO 的 insert_many 方法批量插入事件
80        if !events.is_empty() {
81            signer_crdt::view::CrdtEventVO::insert_many(events, meta)
82                .await
83                .map_err(|e| RemoteError::Internal(format!("CRDT 插入事件失败: {}", e)))?;
84        }
85
86        // 将本地 CRDT 信息推送到远端
87        // 从远程服务器获取前沿信息
88        let frontiers = CrdtFrontier::get_from_remote(&self.addr, keys, &user).await?;
89
90        // 根据前沿信息获取本地事件
91        let local_events = frontiers.get_events_from_signer(meta).await?;
92
93        // 加密事件
94        let events = local_events
95            .into_iter()
96            .map(|e| CrdtCryptedEventVO::encrypt(keys, &e))
97            .collect::<Result<Vec<CrdtCryptedEventVO>, _>>()
98            .map_err(|e| RemoteError::Internal(format!("创建 CRDT 加密事件失败: {}", e)))?;
99
100        // 使用 CrdtCryptedEventVO 的 push 方法推送事件
101        CrdtCryptedEventVO::push(events, &self.addr, keys, &user).await?;
102
103        // 使用我们新建的 reconcile 函数替换原来的 apply_all
104        reconcile(meta)
105            .await
106            .map_err(|e| RemoteError::Internal(format!("CRDT 协调失败: {}", e)))?;
107
108        // 注意:这里我们移除了 daemon.sync_user_public().await? 调用,
109        // 因为在新的方法签名中我们没有 SignerDaemon 对象。
110        // 如果需要同步用户信息,应该在其他地方处理。
111
112        Ok(())
113    }
114
115    /// 从远程服务器拉取用户信息
116    pub async fn pull_user(&self, meta: &SignerMeta, user_key: &str) -> RemoteResult<()> {
117        // 从 SignerMeta 获取当前用户
118        let user = meta.get_current_user().await
119            .map_err(|e| RemoteError::Internal(format!("获取当前用户失败: {}", e)))?;
120        let keys = &meta.keys;
121
122        let config = HttpClientConfig::new(keys.clone(), user.clone(), self.addr.clone());
123        let client = HttpClient::new(config);
124
125        let r: SignerSigned<SignerUser> =
126            client.get(&format!("/api/users/{}", user_key)).await?;
127
128        let signed_up = SignerSigned::<SignerUser> {
129            sig: r.sig,
130            msg: r.msg,
131            pubkey: r.pubkey,
132            _marker: Default::default(),
133        };
134        let up = signed_up
135            .verify_to_value()
136            .map_err(|e| RemoteError::Internal(format!("验证签名用户信息失败: {}", e)))?;
137
138        let user_vo = signer_crdt::view::UserVO::from_user_data(keys, &up).await
139            .map_err(|e| RemoteError::Internal(format!("创建用户视图对象失败: {}", e)))?;
140
141        user_vo.put(meta).await
142            .map_err(|e| RemoteError::Internal(format!("保存用户信息失败: {}", e)))?;
143
144        Ok(())
145    }
146
147    /// 打开事件源连接
148    pub async fn open_eventsource(
149        &mut self,
150        meta: SignerMeta,
151    ) -> RemoteResult<()> {
152        if self.event_source_terminator.is_some() {
153            return Err(RemoteError::Internal("事件源已存在".to_string()));
154        }
155
156        let (tx, mut es) = SignerRemoteEventSource::new(&self.addr, meta);
157        self.event_source_terminator = Some(tx);
158
159        tokio::spawn(async move {
160            es.open_eventsource().await.expect("打开事件源失败");
161        });
162
163        Ok(())
164    }
165
166    /// 关闭事件源连接
167    pub async fn close(&mut self) {
168        if let Some(event_source) = &mut self.event_source_terminator {
169            let _ = event_source.send(()).await;
170            self.event_source_terminator = None;
171        }
172    }
173
174    /// 查询用户在服务上的资源用量
175    pub async fn get_resource_usage(&self, keys: &SignerKeys, user: &SignerUser) -> RemoteResult<ResourceUsage> {
176        tokio::time::timeout(std::time::Duration::from_secs(10), async {
177            // 资源查询需要认证
178            let config = HttpClientConfig::new_no_auth(self.addr.clone());
179            let client = HttpClient::new(config);
180
181            // 生成 JWT token 以获取当前用户的资源使用量
182            let jwt = SignerJWT::new(
183                SignerJWTHeader::default(user),
184                SignerJWTClaims::default(keys, user, self.addr.clone(), uuid::Uuid::new_v4().to_string())
185                    .with_expired_duration(chrono::Duration::minutes(5)),
186            );
187
188            let token = jwt
189                .encode(keys)
190                .map_err(|e| RemoteError::Internal(format!("JWT 编码失败: {}", e)))?;
191
192            let usage: ResourceUsage = client
193                .get_with_header(
194                    "/api/resource-usage",
195                    "Authorization",
196                    &format!("Bearer {}", token),
197                )
198                .await?;
199
200            Ok::<ResourceUsage, RemoteError>(usage)
201        })
202        .await
203        .map_err(|_| RemoteError::Internal("资源用量查询超时".to_string()))?
204    }
205
206    /// 从远程服务器获取 CRDT 前沿信息
207    pub async fn get_crdt_frontiers(&self, keys: &SignerKeys, user: &SignerUser) -> RemoteResult<CrdtFrontier> {
208        CrdtFrontier::get_from_remote(&self.addr, keys, user).await
209    }
210}