From b73bc3c9094c56ae09e44e82fbd18f8ff1bf8bc8 Mon Sep 17 00:00:00 2001 From: Relayflow Lead Date: Sat, 29 Aug 2026 02:51:25 -0400 Subject: [PATCH] drive: cloud run 8abf7774 Work produced by cloud run 8abf7774-4992-48ce-9996-7b96b527c0ff in a workflow sandbox and delivered from this host, because a sandbox has no remote and no GitHub token. Verification and adversarial review ran in-run; see ops/reviews/ in the diff. --- kernel/relayflowd/src/engine.rs | 2 + kernel/relayflowd/src/engine/hn_poller.rs | 79 +++++++++++++++++++++++ kernel/relayflowd/src/lib.rs | 2 +- kernel/relayflowd/tests/hn_poller.rs | 59 +++++++++++++++++ 4 files changed, 141 insertions(+), 1 deletion(-) create mode 100644 kernel/relayflowd/src/engine/hn_poller.rs create mode 100644 kernel/relayflowd/tests/hn_poller.rs diff --git a/kernel/relayflowd/src/engine.rs b/kernel/relayflowd/src/engine.rs index c1f7ec0ab..d6413fc41 100644 --- a/kernel/relayflowd/src/engine.rs +++ b/kernel/relayflowd/src/engine.rs @@ -17,9 +17,11 @@ use crate::worker::{JournalObserver, StepDispatcher}; mod drive; mod effects; +mod hn_poller; mod model; mod remote; mod wake; +pub use hn_poller::HnPoller; pub use model::{RunOutcome, RunSnapshot, RunStatus, StepSnapshot, StepStatus}; use model::{outcome_from_state, snapshot_from_state}; pub use remote::OutOfBandCompletion; diff --git a/kernel/relayflowd/src/engine/hn_poller.rs b/kernel/relayflowd/src/engine/hn_poller.rs new file mode 100644 index 000000000..ba5474dc8 --- /dev/null +++ b/kernel/relayflowd/src/engine/hn_poller.rs @@ -0,0 +1,79 @@ +use std::process::Command; + +use anyhow::{bail, Context, Result}; +use relayflowd_core::{Clock, Event, RunSpec}; +use serde_json::json; + +use super::{Engine, EventSubmitOutcome}; + +const TOP_STORIES_URL: &str = "https://hacker-news.firebaseio.com/v0/topstories.json"; +const DEFAULT_STORY_LIMIT: usize = 5; + +/// Fetches one snapshot of Hacker News top stories and submits it to a flow. +#[derive(Debug, Clone, Copy)] +pub struct HnPoller { + story_limit: usize, +} + +impl HnPoller { + pub fn new(story_limit: usize) -> Self { + Self { story_limit } + } + + pub fn poll_once( + &self, + engine: &Engine, + spec: RunSpec, + created_by: &str, + ) -> Result> { + let payload = fetch_top_stories()?; + self.poll_payload(engine, spec, created_by, &payload) + } + + /// Submits a recorded top-stories response without performing network I/O. + pub fn poll_payload( + &self, + engine: &Engine, + spec: RunSpec, + created_by: &str, + payload: &str, + ) -> Result> { + let story_ids: Vec = + serde_json::from_str(payload).context("parse HN top stories response")?; + story_ids + .into_iter() + .take(self.story_limit) + .map(|id| { + engine.submit_event( + spec.clone(), + Event { + event_type: "hn.story_posted".into(), + payload: json!({"id": id, "type": "story"}), + key: None, + }, + created_by, + ) + }) + .collect() + } +} + +impl Default for HnPoller { + fn default() -> Self { + Self::new(DEFAULT_STORY_LIMIT) + } +} + +fn fetch_top_stories() -> Result { + let output = Command::new("curl") + .args(["--fail", "--silent", "--show-error", TOP_STORIES_URL]) + .output() + .context("fetch HN top stories")?; + if !output.status.success() { + bail!( + "fetch HN top stories failed: {}", + String::from_utf8_lossy(&output.stderr).trim() + ); + } + String::from_utf8(output.stdout).context("HN top stories response was not UTF-8") +} diff --git a/kernel/relayflowd/src/lib.rs b/kernel/relayflowd/src/lib.rs index a90a98c37..5b2ce4013 100644 --- a/kernel/relayflowd/src/lib.rs +++ b/kernel/relayflowd/src/lib.rs @@ -5,6 +5,6 @@ pub mod server; pub mod worker; pub use engine::{ - DriveOptions, Engine, OutOfBandCompletion, RunOutcome, RunSnapshot, RunStatus, + DriveOptions, Engine, HnPoller, OutOfBandCompletion, RunOutcome, RunSnapshot, RunStatus, StepSnapshot, StepStatus, }; diff --git a/kernel/relayflowd/tests/hn_poller.rs b/kernel/relayflowd/tests/hn_poller.rs new file mode 100644 index 000000000..8a30f761c --- /dev/null +++ b/kernel/relayflowd/tests/hn_poller.rs @@ -0,0 +1,59 @@ +use relayflowd::{Engine, HnPoller}; +use relayflowd_core::{EntryType, RunSpec}; + +fn hn_monitor_spec() -> RunSpec { + let value = serde_json::from_str(include_str!(concat!( + env!("CARGO_MANIFEST_DIR"), + "/../../testdata/hn-monitor.spec.canonical.json" + ))) + .unwrap(); + RunSpec::parse(&value).unwrap() +} + +#[test] +fn recorded_top_stories_are_submitted_with_limit_and_deduped() { + let directory = tempfile::tempdir().unwrap(); + let engine = Engine::new(directory.path()); + let poller = HnPoller::new(2); + let payload = "[41380628, 41378954, 41370000]"; + + let first = poller + .poll_payload(&engine, hn_monitor_spec(), "hn-poller", payload) + .unwrap(); + assert_eq!(first.len(), 2); + assert!(first + .iter() + .all(|outcome| outcome.matched && !outcome.deduped)); + + let entries = engine + .journal_entries(&first[0].run.as_ref().unwrap().run_id, 1, 100) + .unwrap(); + let received = entries + .iter() + .find(|entry| entry.entry_type == EntryType::EventReceived) + .unwrap(); + assert_eq!(received.payload["event"]["type"], "hn.story_posted"); + assert_eq!(received.payload["event"]["payload"]["id"], 41380628); + assert_eq!(received.payload["event"]["payload"]["type"], "story"); + + let second = poller + .poll_payload(&engine, hn_monitor_spec(), "hn-poller", payload) + .unwrap(); + assert_eq!(second.len(), 2); + assert!(second + .iter() + .all(|outcome| outcome.matched && outcome.deduped)); + assert!(second.iter().all(|outcome| outcome.run.is_none())); +} + +#[test] +fn malformed_top_stories_payload_fails_closed() { + let directory = tempfile::tempdir().unwrap(); + let engine = Engine::new(directory.path()); + + let error = HnPoller::default() + .poll_payload(&engine, hn_monitor_spec(), "hn-poller", "not JSON") + .unwrap_err(); + + assert!(error.to_string().contains("parse HN top stories response")); +}