diff --git a/crates/buzz-cli/Cargo.toml b/crates/buzz-cli/Cargo.toml index 1476e60bfd4..c1f84cad566 100644 --- a/crates/buzz-cli/Cargo.toml +++ b/crates/buzz-cli/Cargo.toml @@ -23,7 +23,7 @@ clap = { version = "4", features = ["derive", "env"] } reqwest = { workspace = true, features = ["json"] } # Async runtime — tokio macros + multi-thread for reqwest -tokio = { workspace = true, features = ["macros", "rt-multi-thread"] } +tokio = { workspace = true, features = ["macros", "rt-multi-thread", "signal", "time"] } # Serialization — JSON body building and response passthrough serde = { workspace = true } diff --git a/crates/buzz-cli/README.md b/crates/buzz-cli/README.md index a2dcdce6d21..6c9f167ce42 100644 --- a/crates/buzz-cli/README.md +++ b/crates/buzz-cli/README.md @@ -92,8 +92,19 @@ buzz repos protect remove --id my-repo --ref refs/heads/main # Pipe to jq buzz channels list | jq '.[].name' + +# Activity feed +buzz feed get --limit 20 +buzz feed watch # stream NDJSON until Ctrl-C +buzz feed watch --types mentions | jq -r '.content' +buzz feed watch --since 1783497600 # replay backfill, then stay live ``` +`feed watch` writes one JSON event per line and flushes immediately, so it +composes with `jq --unbuffered` and other line-oriented consumers. Relay +notices and subscription closures go to stderr, leaving stdout pure NDJSON. +`feed get` keeps its JSON-array output unchanged. + `protect set` replaces every existing rule for the exact ref pattern. Any constraint omitted from the command is removed. `protect list` reports malformed stored rules in `validation_error` so an owner can remove and repair them. @@ -146,6 +157,7 @@ stored rules in `validation_error` so an owner can remove and repair them. | | `runs` | Get workflow run history | | | `approve` | Approve/deny a workflow step | | `feed` | `get` | Get your activity feed | +| | `watch` | Stream the activity feed as NDJSON until Ctrl-C | | `social` | `publish` | Publish a NIP-01 note | | | `set-contacts` | Set NIP-02 contact list | | | `event` | Get a Nostr event | diff --git a/crates/buzz-cli/src/client.rs b/crates/buzz-cli/src/client.rs index ee8868ad927..ca0706b513f 100644 --- a/crates/buzz-cli/src/client.rs +++ b/crates/buzz-cli/src/client.rs @@ -1067,6 +1067,20 @@ impl BuzzClient { /// Publish an ephemeral event via WebSocket with NIP-42 authentication. /// + /// WebSocket URL for this client's relay. + /// + /// Long-lived subscriptions (`buzz feed watch`) need to open their own + /// connection rather than going through the one-shot publish helper, so the + /// URL conversion is exposed instead of duplicated at the call site. + pub fn ws_url(&self) -> String { + to_ws_url(&self.relay_url) + } + + /// The NIP-OA authorization tag to include in AUTH events, when configured. + pub fn auth_tag(&self) -> Option<&Tag> { + self.auth_tag.as_ref() + } + /// The relay rejects ephemeral kinds (20000–29999) over HTTP. Delegates to /// `buzz_ws_client::publish_event` which handles connect, NIP-42 auth, /// EVENT send, OK wait, and graceful close. diff --git a/crates/buzz-cli/src/commands/feed.rs b/crates/buzz-cli/src/commands/feed.rs index d3d5c7f81a4..130933d15b8 100644 --- a/crates/buzz-cli/src/commands/feed.rs +++ b/crates/buzz-cli/src/commands/feed.rs @@ -1,10 +1,71 @@ use std::cmp::Reverse; +use std::io::Write; +use std::time::Duration; + +use tokio::time::Instant; + +use buzz_ws_client::{NostrWsConnection, RelayMessage, WsClientError}; use crate::client::{normalize_events, BuzzClient}; use crate::error::CliError; const VALID_FEED_TYPES: &[&str] = &["mentions", "needs_action", "activity", "agent_activity"]; +/// How long a single `next_event` call waits before returning control to the +/// watch loop. Short enough that Ctrl-C and idle-timeout accounting stay +/// responsive, long enough not to spin. +const WATCH_POLL_SECS: u64 = 5; + +/// Split and validate a `--types` value against [`VALID_FEED_TYPES`]. +/// +/// Shared by `feed get` and `feed watch` so both reject the same bad input +/// with the same message. +fn parse_feed_types(types_str: &str) -> Result, CliError> { + let type_list: Vec<&str> = types_str.split(',').map(str::trim).collect(); + for t in &type_list { + if !VALID_FEED_TYPES.contains(t) { + return Err(CliError::Usage(format!( + "invalid feed type {t:?} — must be one of: {}", + VALID_FEED_TYPES.join(", ") + ))); + } + } + Ok(type_list) +} + +/// Build the relay filter for the activity feed. +/// +/// Shared by `feed get` and `feed watch` so a stream and a one-shot read +/// select exactly the same events. `limit` is omitted for the streaming case, +/// where the relay should not cap the live subscription. +fn build_feed_filter( + my_pubkey: &str, + since: Option, + limit: Option, + types: Option<&[&str]>, + channel: Option<&str>, +) -> serde_json::Value { + let mut filter = serde_json::json!({ "#p": [my_pubkey] }); + if let Some(l) = limit { + filter["limit"] = serde_json::json!(l); + } + if let Some(s) = since { + filter["since"] = serde_json::json!(s); + } + if let Some(t) = types { + filter["feed_types"] = serde_json::json!(t); + } + if let Some(c) = channel { + filter["#h"] = serde_json::json!([c]); + } + filter +} + +/// Generate a NIP-01 subscription id (1–64 chars). +fn new_subscription_id() -> String { + format!("buzz-feed-watch-{:016x}", rand::random::()) +} + /// Get activity feed — query events mentioning our pubkey (via p-tag). pub async fn cmd_get_feed( client: &BuzzClient, @@ -16,27 +77,11 @@ pub async fn cmd_get_feed( let my_pk = client.keys().public_key().to_hex(); let limit = limit.unwrap_or(20).min(50); - let mut filter = serde_json::json!({ - "#p": [my_pk], - "limit": limit - }); - - if let Some(s) = since { - filter["since"] = serde_json::json!(s); - } - - if let Some(types_str) = types { - let type_list: Vec<&str> = types_str.split(',').map(str::trim).collect(); - for t in &type_list { - if !VALID_FEED_TYPES.contains(t) { - return Err(crate::error::CliError::Usage(format!( - "invalid feed type {t:?} — must be one of: {}", - VALID_FEED_TYPES.join(", ") - ))); - } - } - filter["feed_types"] = serde_json::json!(type_list); - } + let parsed_types = match types { + Some(t) => Some(parse_feed_types(t)?), + None => None, + }; + let filter = build_feed_filter(&my_pk, since, Some(limit), parsed_types.as_deref(), None); let resp = client.query(&filter).await?; let mut events: Vec = serde_json::from_str(&resp).unwrap_or_default(); @@ -64,6 +109,137 @@ pub async fn cmd_get_feed( Ok(()) } +/// Stream activity feed entries as NDJSON until Ctrl-C or idle timeout. +/// +/// Opens its own authenticated WebSocket connection and issues a NIP-01 `REQ` +/// with the same filter [`cmd_get_feed`] builds. One JSON object is written per +/// line to stdout and flushed immediately; relay notices and closures go to +/// stderr so a `| jq` consumer only ever sees events. +pub async fn cmd_watch_feed( + client: &BuzzClient, + types: Option<&str>, + since: Option, + channel: Option<&str>, + idle_timeout: u64, +) -> Result<(), CliError> { + let my_pk = client.keys().public_key().to_hex(); + + let parsed_types = match types { + Some(t) => Some(parse_feed_types(t)?), + None => None, + }; + if let Some(c) = channel { + crate::validate::validate_uuid(c)?; + } + + let filter = build_feed_filter(&my_pk, since, None, parsed_types.as_deref(), channel); + let sub_id = new_subscription_id(); + + let mut conn = NostrWsConnection::connect_authenticated( + &client.ws_url(), + client.keys(), + client.auth_tag(), + ) + .await + .map_err(|e| CliError::Other(e.to_string()))?; + + // Every exit path below — clean break, error, or Ctrl-C — must still close + // the subscription and drop the socket, so the loop result is captured and + // cleanup runs before it is returned. + let result = watch_loop(&mut conn, client, &sub_id, &filter, idle_timeout).await; + + let _ = conn.send_raw(&serde_json::json!(["CLOSE", sub_id])).await; + let _ = conn.disconnect().await; + result +} + +/// The `feed watch` receive loop, split out so [`cmd_watch_feed`] can run +/// subscription cleanup on every exit path including the error ones. +async fn watch_loop( + conn: &mut NostrWsConnection, + client: &BuzzClient, + sub_id: &str, + filter: &serde_json::Value, + idle_timeout: u64, +) -> Result<(), CliError> { + conn.send_raw(&serde_json::json!(["REQ", sub_id, filter])) + .await + .map_err(|e| CliError::Other(e.to_string()))?; + + let poll = Duration::from_secs(WATCH_POLL_SECS); + let idle_limit = (idle_timeout > 0).then(|| Duration::from_secs(idle_timeout)); + // Measured from the last event actually delivered, not from the last poll. + // The deadline is its own `select!` branch rather than a post-poll check: + // `next_event` answers a relay Ping internally and restarts its own timeout, + // so a chatty relay could otherwise keep it pending past the deadline. + let mut last_event = Instant::now(); + let mut stdout = std::io::stdout(); + + loop { + // Copied out before the future is built so the Event arm can still + // reset `last_event` without holding a borrow across the select. + let deadline = idle_limit.map(|limit| last_event + limit); + let idle_expired = async move { + match deadline { + Some(d) => tokio::time::sleep_until(d).await, + None => std::future::pending::<()>().await, + } + }; + + tokio::select! { + biased; + _ = tokio::signal::ctrl_c() => return Ok(()), + _ = idle_expired => return Ok(()), + msg = conn.next_event(poll) => { + match msg { + Ok(RelayMessage::Event { event, .. }) => { + last_event = Instant::now(); + let line = serde_json::to_string(&event) + .map_err(|e| CliError::Other(e.to_string()))?; + writeln!(stdout, "{line}").map_err(|e| CliError::Other(e.to_string()))?; + stdout.flush().map_err(|e| CliError::Other(e.to_string()))?; + } + Ok(RelayMessage::Eose { .. }) => {} + Ok(RelayMessage::Closed { message, .. }) => { + eprintln!("relay closed subscription: {message}"); + return Ok(()); + } + Ok(RelayMessage::Notice { message }) => eprintln!("relay notice: {message}"), + Ok(RelayMessage::Auth { .. }) => { + // Re-auth waits on the relay's OK, so it has to keep racing + // both Ctrl-C and the idle deadline. A relay that issues + // AUTH and withholds OK would otherwise swallow a SIGINT and + // stretch --idle-timeout by the auth budget. + match deadline { + Some(d) => tokio::select! { + biased; + _ = tokio::signal::ctrl_c() => return Ok(()), + _ = tokio::time::sleep_until(d) => return Ok(()), + r = conn.authenticate(client.keys(), client.auth_tag()) => { + r.map_err(|e| CliError::Other(e.to_string()))? + } + }, + None => tokio::select! { + biased; + _ = tokio::signal::ctrl_c() => return Ok(()), + r = conn.authenticate(client.keys(), client.auth_tag()) => { + r.map_err(|e| CliError::Other(e.to_string()))? + } + }, + } + } + Ok(RelayMessage::Ok(_)) | Ok(RelayMessage::Count { .. }) => {} + // A poll that expired with no frame is the normal idle path. + // Anything else — a closed socket, a transport failure — is + // terminal and must surface rather than spin. + Err(WsClientError::Timeout) => {} + Err(e) => return Err(CliError::Other(e.to_string())), + } + } + } + } +} + pub async fn dispatch( cmd: crate::FeedCmd, client: &BuzzClient, @@ -76,5 +252,96 @@ pub async fn dispatch( limit, types, } => cmd_get_feed(client, since, limit, types.as_deref(), format).await, + FeedCmd::Watch { + types, + since, + channel, + idle_timeout, + } => { + cmd_watch_feed( + client, + types.as_deref(), + since, + channel.as_deref(), + idle_timeout, + ) + .await + } + } +} + +#[cfg(test)] +mod tests { + use super::{build_feed_filter, new_subscription_id, parse_feed_types}; + use serde_json::json; + + const PUBKEY: &str = "cccccccccccccccccccccccccccccccccccccccccccccccccccccccccccccccc"; + const UUID: &str = "0b7f9c2e-3d4a-4b1c-8e5f-6a7b8c9d0e1f"; + + #[test] + fn every_valid_feed_type_is_accepted() { + for t in ["mentions", "needs_action", "activity", "agent_activity"] { + assert!(parse_feed_types(t).is_ok(), "{t} should be valid"); + } + assert_eq!( + parse_feed_types("mentions,activity").unwrap(), + vec!["mentions", "activity"] + ); + } + + #[test] + fn an_unknown_feed_type_is_rejected() { + let err = parse_feed_types("mentions,bogus").unwrap_err(); + assert!( + err.to_string().contains("invalid feed type"), + "unexpected error: {err}" + ); + } + + #[test] + fn surrounding_whitespace_is_trimmed() { + assert_eq!( + parse_feed_types(" mentions , activity ").unwrap(), + vec!["mentions", "activity"] + ); + } + + #[test] + fn the_no_flags_filter_matches_what_feed_get_sent_before() { + let f = build_feed_filter(PUBKEY, None, Some(20), None, None); + assert_eq!(f, json!({ "#p": [PUBKEY], "limit": 20 })); + } + + #[test] + fn optional_fields_appear_only_when_supplied() { + let bare = build_feed_filter(PUBKEY, None, None, None, None); + assert!(bare.get("since").is_none()); + assert!(bare.get("#h").is_none()); + assert!(bare.get("feed_types").is_none()); + assert!(bare.get("limit").is_none()); + + let full = build_feed_filter( + PUBKEY, + Some(1_783_497_600), + Some(10), + Some(&["mentions"]), + Some(UUID), + ); + assert_eq!(full["since"], json!(1_783_497_600)); + assert_eq!(full["limit"], json!(10)); + assert_eq!(full["feed_types"], json!(["mentions"])); + assert_eq!(full["#h"], json!([UUID])); + assert_eq!(full["#p"], json!([PUBKEY])); + } + + #[test] + fn subscription_ids_are_within_the_nip01_length_bound() { + let id = new_subscription_id(); + assert!( + (1..=64).contains(&id.len()), + "subscription id length {} out of NIP-01 range", + id.len() + ); + assert_ne!(id, new_subscription_id(), "ids should not repeat"); } } diff --git a/crates/buzz-cli/src/lib.rs b/crates/buzz-cli/src/lib.rs index 2b041da57b5..2502bffa692 100644 --- a/crates/buzz-cli/src/lib.rs +++ b/crates/buzz-cli/src/lib.rs @@ -981,6 +981,24 @@ pub enum FeedCmd { #[arg(long)] types: Option, }, + /// Stream activity feed entries as they arrive (NDJSON, one event per line) + #[command( + after_help = "Examples:\n buzz feed watch\n buzz feed watch --types mentions\n buzz feed watch --since 1783497600 | jq -r '.content'" + )] + Watch { + /// Comma-separated feed types to include: mentions, needs_action, activity, agent_activity + #[arg(long)] + types: Option, + /// Unix timestamp — replay entries after this time before going live + #[arg(long)] + since: Option, + /// Channel UUID — restrict the stream to one channel + #[arg(long)] + channel: Option, + /// Exit if no event arrives for this many seconds (0 = never) + #[arg(long, default_value_t = 0)] + idle_timeout: u64, + }, } #[derive(Subcommand)] @@ -2277,7 +2295,7 @@ mod tests { names(&cmd, "workflows"), vec!["approve", "create", "delete", "get", "list", "runs", "trigger", "update"] ); - assert_eq!(names(&cmd, "feed"), vec!["get"]); + assert_eq!(names(&cmd, "feed"), vec!["get", "watch"]); assert_eq!( names(&cmd, "social"), vec![ @@ -2359,7 +2377,7 @@ mod tests { ("channels", 16), ("dms", 4), ("emoji", 5), - ("feed", 1), + ("feed", 2), ("issues", 6), ("media", 1), ("messages", 8),