signer-remote 0.4.1

Signer remote communication package.
Documentation
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(())
    }
}