diff --git a/src/tui/mod.rs b/src/tui/mod.rs index 04fcd19..cc5e22b 100644 --- a/src/tui/mod.rs +++ b/src/tui/mod.rs @@ -120,122 +120,149 @@ fn pump_core_events( tokio::spawn(async move { let client = reqwest::Client::new(); - let _ = tx - .send(Event::Core(crate::core::event::Event::Trace { - level: tracing::Level::INFO, - message: "Connecting to event stream...".into(), - })) - .await; - - let res = match client - .get(format!("{}/events", url)) - .header("Authorization", format!("Bearer {}", token)) - .send() - .await - { - Ok(r) => r, - Err(e) => { + loop { + let _ = tx + .send(Event::Core(crate::core::event::Event::Trace { + level: tracing::Level::INFO, + message: "Connecting to event stream...".into(), + })) + .await; + + let res = match client + .get(format!("{}/events", url)) + .header("Authorization", format!("Bearer {}", token)) + .send() + .await + { + Ok(r) => r, + Err(e) => { + let _ = tx + .send(Event::Core(crate::core::event::Event::Trace { + level: tracing::Level::ERROR, + message: format!("Failed to connect: {}. Retrying in 5s...", e), + })) + .await; + tokio::time::sleep(Duration::from_secs(5)).await; + continue; + } + }; + + if !res.status().is_success() { let _ = tx .send(Event::Core(crate::core::event::Event::Trace { - level: tracing::Level::INFO, - message: format!("Failed to connect: {}", e), + level: tracing::Level::ERROR, + message: format!("Server error: {}. Retrying in 5s...", res.status()), })) .await; - return; + tokio::time::sleep(Duration::from_secs(5)).await; + continue; } - }; - if !res.status().is_success() { let _ = tx .send(Event::Core(crate::core::event::Event::Trace { level: tracing::Level::INFO, - message: format!("Server error: {}", res.status()), + message: "Streaming events...".into(), })) .await; - return; - } - let _ = tx - .send(Event::Core(crate::core::event::Event::Trace { - level: tracing::Level::INFO, - message: "Streaming events...".into(), - })) - .await; - - let payload = serde_json::json!({ "type": "Bootstrap" }); - if let Err(e) = client - .post(format!("{}/messages", url)) - .header("Authorization", format!("Bearer {}", token)) - .json(&payload) - .send() - .await - { - let _ = tx - .send(Event::Core(crate::core::event::Event::Trace { - level: tracing::Level::INFO, - message: format!("Failed to send Bootstrap: {}", e), - })) - .await; - } else { - let _ = tx - .send(Event::Core(crate::core::event::Event::Trace { - level: tracing::Level::INFO, - message: "Sent Bootstrap message".into(), - })) - .await; - } + let payload = serde_json::json!({ "type": "Bootstrap" }); + if let Err(e) = client + .post(format!("{}/messages", url)) + .header("Authorization", format!("Bearer {}", token)) + .json(&payload) + .send() + .await + { + let _ = tx + .send(Event::Core(crate::core::event::Event::Trace { + level: tracing::Level::ERROR, + message: format!("Failed to send Bootstrap: {}. Retrying in 5s...", e), + })) + .await; + tokio::time::sleep(Duration::from_secs(5)).await; + continue; + } else { + let _ = tx + .send(Event::Core(crate::core::event::Event::Trace { + level: tracing::Level::INFO, + message: "Sent Bootstrap message".into(), + })) + .await; + } - let mut stream = res.bytes_stream(); - let mut buffer = Vec::new(); + let mut stream = res.bytes_stream(); + let mut buffer = Vec::new(); - loop { - tokio::select! { - res = stream.next() => { - match res { - Some(Ok(bytes)) => { - buffer.extend_from_slice(&bytes); - while let Some(pos) = buffer.windows(2).position(|w| w == b"\n\n") { - let event_bytes = buffer.drain(..pos + 2).collect::>(); - let text = String::from_utf8_lossy(&event_bytes); - for line in text.lines() { - if let Some(data) = line.strip_prefix("data: ") { - if let Ok(ev) = - serde_json::from_str::(data) - { - let _ = tx.send(Event::Core(ev)).await; - } else { - let _ = tx - .send(Event::Core(crate::core::event::Event::Trace { - level: tracing::Level::INFO, - message: data.to_string(), - })) - .await; + loop { + tokio::select! { + res = stream.next() => { + match res { + Some(Ok(bytes)) => { + buffer.extend_from_slice(&bytes); + while let Some(pos) = buffer.windows(2).position(|w| w == b"\n\n") { + let event_bytes = buffer.drain(..pos + 2).collect::>(); + let text = String::from_utf8_lossy(&event_bytes); + for line in text.lines() { + if let Some(data) = line.strip_prefix("data: ") { + if let Ok(ev) = + serde_json::from_str::(data) + { + let _ = tx.send(Event::Core(ev)).await; + } else { + let _ = tx + .send(Event::Core(crate::core::event::Event::Trace { + level: tracing::Level::INFO, + message: data.to_string(), + })) + .await; + } } } } } + Some(Err(e)) => { + let _ = tx + .send(Event::Core(crate::core::event::Event::Trace { + level: tracing::Level::ERROR, + message: format!("Stream error: {}. Retrying in 5s...", e), + })) + .await; + break; + } + None => { + let _ = tx + .send(Event::Core(crate::core::event::Event::Trace { + level: tracing::Level::INFO, + message: "Stream ended. Retrying in 5s...".into(), + })) + .await; + break; + } } - Some(Err(e)) => { - let _ = tx - .send(Event::Core(crate::core::event::Event::Trace { - level: tracing::Level::INFO, - message: format!("Stream error: {}", e), - })) - .await; - break; + } + res = payload_rx.recv() => { + match res { + Some(payload) => { + if let Err(e) = client + .post(format!("{}/messages", url)) + .header("Authorization", format!("Bearer {}", token)) + .json(&payload) + .send() + .await + { + let _ = tx.send(Event::Core(crate::core::event::Event::Trace { + level: tracing::Level::ERROR, + message: format!("Failed to send payload: {}", e), + })).await; + } + } + None => return, } - None => break, } } - Some(payload) = payload_rx.recv() => { - let _ = client - .post(format!("{}/messages", url)) - .header("Authorization", format!("Bearer {}", token)) - .json(&payload) - .send() - .await; - } } + + tokio::time::sleep(Duration::from_secs(5)).await; } }); }