mirror of
https://github.com/Bluemangoo/sekai-unpacker.git
synced 2026-09-20 00:06:52 +08:00
fix ram
This commit is contained in:
@@ -95,6 +95,30 @@ pub enum TunnelEndpoint {
|
||||
Server(Arc<ServerManager>),
|
||||
}
|
||||
|
||||
#[derive(Clone)]
|
||||
pub enum WeakTunnelEndpoint {
|
||||
Client(std::sync::Weak<ClientManager>),
|
||||
Server(std::sync::Weak<ServerManager>),
|
||||
}
|
||||
|
||||
impl TunnelEndpoint {
|
||||
pub fn downgrade(&self) -> WeakTunnelEndpoint {
|
||||
match self {
|
||||
Self::Client(c) => WeakTunnelEndpoint::Client(Arc::downgrade(c)),
|
||||
Self::Server(s) => WeakTunnelEndpoint::Server(Arc::downgrade(s)),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl WeakTunnelEndpoint {
|
||||
pub fn upgrade(&self) -> Option<TunnelEndpoint> {
|
||||
match self {
|
||||
Self::Client(c) => c.upgrade().map(TunnelEndpoint::Client),
|
||||
Self::Server(s) => s.upgrade().map(TunnelEndpoint::Server),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
pub struct ClientManager {
|
||||
pub session_id: AtomicU64,
|
||||
pub current_client: Mutex<Option<client::SendRequest<Bytes>>>,
|
||||
@@ -231,9 +255,9 @@ enum ResumeResult {
|
||||
pub struct TunnelListener {
|
||||
listener: TcpListener,
|
||||
config: ServerTunnelConfig,
|
||||
pending_plain_sessions: Mutex<HashSet<u64>>,
|
||||
pending_plain_sessions: Mutex<HashMap<u64, std::time::Instant>>,
|
||||
next_session_id: AtomicU64,
|
||||
active_sessions: Mutex<HashMap<u64, TunnelEndpoint>>,
|
||||
active_sessions: Mutex<HashMap<u64, WeakTunnelEndpoint>>,
|
||||
}
|
||||
|
||||
impl TunnelListener {
|
||||
@@ -243,7 +267,7 @@ impl TunnelListener {
|
||||
Ok(Self {
|
||||
listener,
|
||||
config,
|
||||
pending_plain_sessions: Mutex::new(HashSet::new()),
|
||||
pending_plain_sessions: Mutex::new(HashMap::new()),
|
||||
next_session_id: AtomicU64::new(1),
|
||||
active_sessions: Mutex::new(HashMap::new()),
|
||||
})
|
||||
@@ -264,7 +288,9 @@ impl TunnelListener {
|
||||
let mut stream = match self.try_resume_plain_session(stream, peer_addr).await? {
|
||||
ResumeResult::NewSession(ep_raw, sid) => {
|
||||
let ep = wrap_raw_endpoint(sid, ep_raw, None);
|
||||
self.active_sessions.lock().await.insert(sid, ep.clone());
|
||||
let mut sessions = self.active_sessions.lock().await;
|
||||
sessions.retain(|_, weak_ep| weak_ep.upgrade().is_some());
|
||||
sessions.insert(sid, ep.downgrade());
|
||||
return Ok(ep);
|
||||
}
|
||||
ResumeResult::ResumedExisting => continue,
|
||||
@@ -290,7 +316,9 @@ impl TunnelListener {
|
||||
|
||||
let sid = self.next_session_id.fetch_add(1, Ordering::Relaxed);
|
||||
let ep = wrap_raw_endpoint(sid, ep_raw, None);
|
||||
self.active_sessions.lock().await.insert(sid, ep.clone());
|
||||
let mut sessions = self.active_sessions.lock().await;
|
||||
sessions.retain(|_, weak_ep| weak_ep.upgrade().is_some());
|
||||
sessions.insert(sid, ep.downgrade());
|
||||
return Ok(ep);
|
||||
}
|
||||
}
|
||||
@@ -320,7 +348,8 @@ impl TunnelListener {
|
||||
let session_id = self.next_session_id.fetch_add(1, Ordering::Relaxed);
|
||||
{
|
||||
let mut pending = self.pending_plain_sessions.lock().await;
|
||||
pending.insert(session_id);
|
||||
pending.retain(|_, time| time.elapsed() < Duration::from_secs(30));
|
||||
pending.insert(session_id, std::time::Instant::now());
|
||||
}
|
||||
|
||||
tls_stream.write_all(TLS_BOOTSTRAP_MAGIC).await?;
|
||||
@@ -351,7 +380,7 @@ impl TunnelListener {
|
||||
|
||||
let is_pending = {
|
||||
let mut pending = self.pending_plain_sessions.lock().await;
|
||||
pending.remove(&session_id)
|
||||
pending.remove(&session_id).is_some()
|
||||
};
|
||||
|
||||
if is_pending {
|
||||
@@ -366,7 +395,11 @@ impl TunnelListener {
|
||||
let ep_raw = upgrade_to_h2_raw(stream, self.config.identity)
|
||||
.await
|
||||
.map_err(|e| anyhow::anyhow!("{}", e))?;
|
||||
update_endpoint(&ep, ep_raw).await;
|
||||
update_endpoint(
|
||||
&ep.upgrade().ok_or(anyhow!("Connection is cleared"))?,
|
||||
ep_raw,
|
||||
)
|
||||
.await;
|
||||
info!(
|
||||
"[{}] Successfully resumed existing session {}",
|
||||
peer_addr, session_id
|
||||
|
||||
Reference in New Issue
Block a user