signer_remote/remote/
client.rs1use 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 pub fn new(addr: &str) -> Self {
26 Self {
27 addr: addr.to_string(),
28 event_source_terminator: None,
29 }
30 }
31
32 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 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 pub async fn sync_crdt_event(&self, meta: &SignerMeta) -> RemoteResult<()> {
54 let user = meta.get_current_user().await
56 .map_err(|e| RemoteError::Internal(format!("获取当前用户失败: {}", e)))?;
57 let keys = &meta.keys;
58
59 let frontiers = CrdtFrontier::get_from_signer(meta).await?;
62
63 let frontiers_json = serde_json::to_string(frontiers.inner())
65 .map_err(|e| RemoteError::Internal(format!("序列化前沿信息失败: {}", e)))?;
66
67 let crypted_events = CrdtCryptedEventVO::pull(&self.addr, keys, &user, &frontiers_json).await?;
69
70 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 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 let frontiers = CrdtFrontier::get_from_remote(&self.addr, keys, &user).await?;
89
90 let local_events = frontiers.get_events_from_signer(meta).await?;
92
93 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(events, &self.addr, keys, &user).await?;
102
103 reconcile(meta)
105 .await
106 .map_err(|e| RemoteError::Internal(format!("CRDT 协调失败: {}", e)))?;
107
108 Ok(())
113 }
114
115 pub async fn pull_user(&self, meta: &SignerMeta, user_key: &str) -> RemoteResult<()> {
117 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 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 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 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 let config = HttpClientConfig::new_no_auth(self.addr.clone());
179 let client = HttpClient::new(config);
180
181 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 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}