1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
use eventsource_client::{Client, SSE};
use futures::StreamExt;
use serde::{Deserialize, Serialize};
use signer_auth::{SignerJWT, SignerJWTClaims, SignerJWTHeader};
use signer_crdt::SignerMeta;
use crate::{
error::RemoteError,
remote::{SignerRemote, Envelope},
};
/// 事件数据
#[derive(Debug, Clone, Serialize, Deserialize)]
pub enum EventData {
NewCrdtEvents,
NewEnvelopes,
}
pub struct SignerRemoteEventSource {
addr: String,
meta: SignerMeta,
rx: tokio::sync::mpsc::Receiver<()>,
}
impl SignerRemoteEventSource {
pub fn new(
addr: &str,
meta: SignerMeta,
) -> (tokio::sync::mpsc::Sender<()>, Self) {
let (tx, rx) = tokio::sync::mpsc::channel(1);
(
tx,
Self {
addr: addr.to_string(),
meta,
rx,
},
)
}
pub async fn open_eventsource(&mut self) -> crate::error::RemoteResult<()> {
let user = self.meta.get_current_user().await
.map_err(|e| RemoteError::Internal(format!("获取当前用户失败: {}", e)))?;
let keys = &self.meta.keys;
let remote = SignerRemote::new(&self.addr);
// 指数退避参数
let mut retry_delay = std::time::Duration::from_secs(1); // 初始延迟1秒
const MAX_RETRY_DELAY: std::time::Duration = std::time::Duration::from_secs(3 * 60); // 最大延迟3分钟
const BACKOFF_FACTOR: u32 = 2; // 指数因子
'outer: loop {
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 client =
eventsource_client::ClientBuilder::for_url(&format!("{}/api/events", self.addr))
.map_err(|e| {
RemoteError::Internal(format!(
"创建事件源客户端失败: {}",
e
))
})?
.header(
"Authorization",
&format!("Bearer {}", &jwt.encode(keys).unwrap()),
)
.map_err(|e| {
RemoteError::Internal(format!(
"添加请求头失败: {}",
e
))
})?
.build_http();
let mut stream = Box::pin(client.stream());
// 为连接建立设置超时
let initial_sync_result =
tokio::time::timeout(std::time::Duration::from_secs(15), async {
// 执行初始同步
remote.sync_crdt_event(&self.meta).await?;
// 拉取并处理信封(消息队列模式)
if let Err(e) = Envelope::pull_and_process(&self.addr, keys, &user, &self.meta).await {
tracing::warn!("拉取并处理信封失败: {}", e);
}
Ok::<(), RemoteError>(())
})
.await;
match initial_sync_result {
Ok(Ok(())) => {
// 初始同步成功,重置重试延迟
retry_delay = std::time::Duration::from_secs(1);
}
Ok(Err(e)) => {
tracing::warn!(
"事件源初始同步失败: {}. 等待 {:?} 后重试...",
e,
retry_delay
);
tokio::time::sleep(retry_delay).await;
// 增加重试延迟(指数退避)
retry_delay = std::cmp::min(retry_delay * BACKOFF_FACTOR, MAX_RETRY_DELAY);
continue;
}
Err(_) => {
tracing::warn!("事件源初始同步超时,等待 {:?} 后重试...", retry_delay);
tokio::time::sleep(retry_delay).await;
// 增加重试延迟(指数退避)
retry_delay = std::cmp::min(retry_delay * BACKOFF_FACTOR, MAX_RETRY_DELAY);
continue;
}
}
'inner: loop {
let e = tokio::select! {
event = stream.next() => event,
_ = self.rx.recv() => {
break 'outer;
},
_ = tokio::time::sleep(std::time::Duration::from_secs(120)) => {
tracing::warn!("事件源连接超时(120秒无活动),重新连接");
break 'inner;
}
}
.ok_or(RemoteError::Internal("事件流结束".to_string()))?;
let e = match e {
Ok(SSE::Event(e)) => e,
Ok(SSE::Comment(_)) => continue,
Ok(SSE::Connected(_)) => {
// 连接建立成功,立即同步 CRDT 和 Envelope 以保持状态一致
tracing::info!("事件源连接建立成功,开始同步数据");
// 同步 CRDT 事件
if let Err(e) = remote.sync_crdt_event(&self.meta).await {
tracing::warn!("连接后同步 CRDT 事件失败: {}", e);
} else {
tracing::debug!("连接后 CRDT 事件同步成功");
}
// 拉取并处理信封(消息队列模式)
if let Err(e) = Envelope::pull_and_process(&self.addr, keys, &user, &self.meta).await {
tracing::warn!("连接后拉取并处理信封失败: {}", e);
} else {
tracing::debug!("连接后信封同步成功");
}
// 同步成功,重置重试延迟
retry_delay = std::time::Duration::from_secs(1);
continue;
}
Err(e) => {
tracing::warn!("事件源错误: {}. 等待 {:?} 后重试...", e, retry_delay);
tokio::time::sleep(retry_delay).await;
// 增加重试延迟(指数退避)
retry_delay = std::cmp::min(retry_delay * BACKOFF_FACTOR, MAX_RETRY_DELAY);
break 'inner;
}
};
// 处理事件数据,支持直接字符串和JSON序列化的字符串
let event_data = &e.data;
tracing::debug!("接收到事件数据: {:?}", event_data);
// 智能解析事件数据:处理可能的JSON序列化
let parsed_data = event_data.trim();
let final_data = if parsed_data.starts_with('"') && parsed_data.ends_with('"') {
// 如果数据被双引号包围,说明被JSON序列化了,需要反序列化
match serde_json::from_str::<String>(parsed_data) {
Ok(parsed) => {
tracing::debug!("JSON反序列化成功: {:?} -> {:?}", parsed_data, parsed);
parsed
},
Err(e) => {
tracing::warn!("JSON反序列化失败: {:?}, 错误: {}, 使用原始数据", parsed_data, e);
parsed_data.to_string()
}
}
} else {
// 直接使用原始字符串
parsed_data.to_string()
};
tracing::debug!("最终解析的事件数据: {:?}", final_data);
// 匹配事件类型
let ed: EventData = match final_data.as_str() {
"NewCrdtEvents" => EventData::NewCrdtEvents,
"NewEnvelopes" => EventData::NewEnvelopes,
_ => {
tracing::warn!("未知的事件类型: {:?} (原始数据: {:?})", final_data, event_data);
continue;
}
};
match ed {
EventData::NewCrdtEvents => {
if let Err(e) = remote.sync_crdt_event(&self.meta).await {
tracing::error!(
"同步 CRDT 事件失败: {}. 等待 {:?} 后重试...",
e,
retry_delay
);
tokio::time::sleep(retry_delay).await;
// 增加重试延迟(指数退避)
retry_delay = std::cmp::min(retry_delay * BACKOFF_FACTOR, MAX_RETRY_DELAY);
break 'inner;
}
// 同步成功,重置重试延迟
retry_delay = std::time::Duration::from_secs(1);
}
EventData::NewEnvelopes => {
// 拉取并处理信封(消息队列模式)
if let Err(e) = Envelope::pull_and_process(&self.addr, keys, &user, &self.meta).await {
tracing::warn!("处理新信封失败: {}", e);
}
// 同步成功,重置重试延迟
retry_delay = std::time::Duration::from_secs(1);
}
}
}
}
Ok(())
}
}