use deno_core::op2;
use serde::Serialize;
use std::collections::HashMap;
use std::sync::{Arc, Mutex};
use tokio::sync::mpsc;
lazy_static::lazy_static! {
static ref SSE_CONNECTIONS: Arc<Mutex<SseStore>> = Arc::new(Mutex::new(SseStore::new()));
}
struct SseStore {
receivers: HashMap<i32, Arc<tokio::sync::Mutex<mpsc::UnboundedReceiver<SseEvent>>>>,
next_id: i32,
}
impl SseStore {
fn new() -> Self {
Self {
receivers: HashMap::new(),
next_id: 1,
}
}
}
pub struct SseState;
impl Default for SseState {
fn default() -> Self {
Self::new()
}
}
impl SseState {
pub fn new() -> Self {
Self
}
}
#[derive(Serialize, Clone, Debug)]
pub struct SseEvent {
pub event: String,
pub data: String,
pub id: String,
pub status: String,
}
#[derive(Serialize)]
pub struct SseConnectResult {
pub id: i32,
pub ok: bool,
pub error: String,
}
#[op2(async(lazy), fast)]
#[serde]
pub async fn op_sse_connect(
#[string] url: String,
) -> Result<SseConnectResult, deno_error::JsErrorBox> {
let (tx, rx) = mpsc::unbounded_channel::<SseEvent>();
let id = {
let mut store = SSE_CONNECTIONS.lock().unwrap_or_else(|e| e.into_inner());
let id = store.next_id;
store.next_id += 1;
store
.receivers
.insert(id, Arc::new(tokio::sync::Mutex::new(rx)));
id
};
let url_clone = url.clone();
tokio::spawn(async move {
if let Err(e) = sse_reader(&url_clone, &tx).await {
let _ = tx.send(SseEvent {
event: "error".to_string(),
data: e.to_string(),
id: String::new(),
status: "error".to_string(),
});
}
let _ = tx.send(SseEvent {
event: String::new(),
data: String::new(),
id: String::new(),
status: "closed".to_string(),
});
});
Ok(SseConnectResult {
id,
ok: true,
error: String::new(),
})
}
#[op2(async(lazy), fast)]
#[serde]
pub async fn op_sse_recv(#[smi] id: i32) -> Result<SseEvent, deno_error::JsErrorBox> {
let rx = {
let store = SSE_CONNECTIONS.lock().unwrap_or_else(|e| e.into_inner());
store.receivers.get(&id).cloned()
};
match rx {
Some(rx) => {
let mut rx = rx.lock().await;
match rx.recv().await {
Some(event) => Ok(event),
None => Ok(SseEvent {
event: String::new(),
data: String::new(),
id: String::new(),
status: "closed".to_string(),
}),
}
}
None => Ok(SseEvent {
event: String::new(),
data: String::new(),
id: String::new(),
status: "closed".to_string(),
}),
}
}
#[op2(fast)]
pub fn op_sse_close(#[smi] id: i32) {
let mut store = SSE_CONNECTIONS.lock().unwrap_or_else(|e| e.into_inner());
store.receivers.remove(&id);
}
async fn sse_reader(
url: &str,
tx: &mpsc::UnboundedSender<SseEvent>,
) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
let profile = crate::stealth::chrome_148_linux();
let client = crate::net::HttpClient::new(&profile)
.map_err(|e| format!("failed to create HTTP client: {e}"))?;
let resp = client
.get(url)
.await
.map_err(|e| format!("SSE fetch failed: {e}"))?;
let body = resp.text();
parse_sse_body(&body, tx);
Ok(())
}
fn parse_sse_body(body: &str, tx: &mpsc::UnboundedSender<SseEvent>) {
let mut event_type = String::new();
let mut data_buf = String::new();
let mut last_id = String::new();
for line in body.lines() {
if line.is_empty() {
if !data_buf.is_empty() {
if data_buf.ends_with('\n') {
data_buf.pop();
}
let event = SseEvent {
event: if event_type.is_empty() {
"message".to_string()
} else {
std::mem::take(&mut event_type)
},
data: std::mem::take(&mut data_buf),
id: last_id.clone(),
status: String::new(),
};
if tx.send(event).is_err() {
return;
}
}
event_type.clear();
continue;
}
if line.starts_with(':') {
continue; }
let (field, value) = if let Some(colon_pos) = line.find(':') {
let field = &line[..colon_pos];
let value = line[colon_pos + 1..]
.strip_prefix(' ')
.unwrap_or(&line[colon_pos + 1..]);
(field, value)
} else {
(line, "")
};
match field {
"event" => event_type = value.to_string(),
"data" => {
data_buf.push_str(value);
data_buf.push('\n');
}
"id" => last_id = value.to_string(),
"retry" => {}
_ => {}
}
}
}
deno_core::extension!(
sse_extension,
ops = [op_sse_connect, op_sse_recv, op_sse_close],
);