Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions .github/workflows/rust.yml
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,8 @@ jobs:
run: cargo fmt --all --check
- name: cargo clippy
run: cargo clippy --workspace --all-targets -- -D warnings
- name: No hand-rolled subscription handshakes
run: etc/ci/check-accept-uni.sh

test:
name: Test
Expand Down
39 changes: 2 additions & 37 deletions crates/ctl/src/logs.rs
Original file line number Diff line number Diff line change
Expand Up @@ -31,45 +31,10 @@ async fn run_log_session(
json_mode: bool,
follow: bool,
) -> Result<(), String> {
let req_bytes = serde_json::to_vec(&serde_json::json!({
"method": "/logs/stream",
"params": params,
}))
.expect("serialisation");

let (mut send, mut recv) = client
.open_bi()
.await
.map_err(|e| format!("open_bi: {e}"))?;

send.write_all(&req_bytes)
.await
.map_err(|e| format!("write: {e}"))?;
let _ = send.finish();

let resp = recv
.read_to_end(64 * 1024)
.await
.map_err(|e| format!("read response: {e}"))?;

if let Ok(v) = serde_json::from_slice::<serde_json::Value>(&resp)
&& let Some(err) = v.get("error")
{
let code = err
.get("code")
.and_then(|c| c.as_str())
.unwrap_or("unknown");
let msg = err
.get("message")
.and_then(|m| m.as_str())
.unwrap_or("unknown error");
return Err(format!("[{code}] {msg}"));
}

let mut log_stream = client
.accept_uni()
.open_subscription("/logs/stream", params)
.await
.map_err(|e| format!("accept_uni: {e}"))?;
.map_err(|e| e.to_string())?;

let mut buf = Vec::new();
let mut tmp = [0u8; 4096];
Expand Down
49 changes: 15 additions & 34 deletions crates/ctl/src/subscribe.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@ use std::{net::SocketAddr, time::Duration};

use seedling_protocol::{
actor::Actor,
client::{ClientAuth, OiClient},
client::{ClientAuth, ClientError, OiClient},
keys::ClientIdentity,
};

Expand Down Expand Up @@ -45,7 +45,13 @@ pub async fn subscribe(
backoff = Duration::from_secs(1);

match run_subscribe_session(&client).await {
SessionOutcome::GracefulClose => return,
// The server refused the subscription — retrying will not change
// its mind, and exiting 0 would tell a script the feed was
// consumed to its end.
SessionOutcome::Rejected(e) => {
eprintln!("subscription rejected: {e}");
std::process::exit(1);
}
SessionOutcome::Error(e) => {
// Reset the deadline when we start reconnecting, not when the
// session began — otherwise a long-lived session causes the
Expand All @@ -61,45 +67,20 @@ pub async fn subscribe(
}

enum SessionOutcome {
GracefulClose,
/// The server answered the subscription request with an error.
Rejected(String),
Error(String),
Interrupted,
}

// i[impl ctl.graceful-shutdown]
async fn run_subscribe_session(client: &OiClient) -> SessionOutcome {
let req_bytes = serde_json::to_vec(&serde_json::json!({
"method": "/events/subscribe",
"actor": client.actor(),
"params": {},
}))
.expect("serialisation");

let (mut send, mut recv) = match client.open_bi().await {
Ok(s) => s,
Err(e) => return SessionOutcome::Error(format!("open_bi: {e}")),
};

if let Err(e) = send.write_all(&req_bytes).await {
return SessionOutcome::Error(format!("write: {e}"));
}
let _ = send.finish();

let resp = match recv.read_to_end(64 * 1024).await {
Ok(r) => r,
Err(e) => return SessionOutcome::Error(format!("read response: {e}")),
};

if let Ok(v) = serde_json::from_slice::<serde_json::Value>(&resp)
&& v.get("error").is_some()
{
eprintln!("{}", serde_json::to_string_pretty(&v).unwrap_or_default());
return SessionOutcome::GracefulClose;
}

let mut event_stream = match client.accept_uni().await {
let mut event_stream = match client.subscribe_events().await {
Ok(s) => s,
Err(e) => return SessionOutcome::Error(format!("accept_uni: {e}")),
Err(ClientError::Api { code, message }) => {
return SessionOutcome::Rejected(format!("[{code}] {message}"));
}
Err(e) => return SessionOutcome::Error(e.to_string()),
};

let mut buf = Vec::new();
Expand Down
Loading