Skip to content
Closed
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 kernel/relayflowd/src/engine.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
79 changes: 79 additions & 0 deletions kernel/relayflowd/src/engine/hn_poller.rs
Original file line number Diff line number Diff line change
@@ -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<C: Clock>(
&self,
engine: &Engine<C>,
spec: RunSpec,
created_by: &str,
) -> Result<Vec<EventSubmitOutcome>> {
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<C: Clock>(
&self,
engine: &Engine<C>,
spec: RunSpec,
created_by: &str,
payload: &str,
) -> Result<Vec<EventSubmitOutcome>> {
let story_ids: Vec<u64> =
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<String> {
let output = Command::new("curl")
.args(["--fail", "--silent", "--show-error", TOP_STORIES_URL])
.output()
Comment on lines +68 to +70

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Badge Keep provider polling out of the kernel

When HnPoller::poll_once is called, relayflowd now performs Hacker News–specific network I/O by spawning curl, coupling the durable kernel to both an external provider and a host-installed executable. The parent already exports sdk/src/hn-poller.ts, which submits these events through event.submit; remove this kernel module and keep polling on that SDK/protocol surface.

AGENTS.md reference: AGENTS.md:L11-L15

Useful? React with 👍 / 👎.

.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")
}
2 changes: 1 addition & 1 deletion kernel/relayflowd/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
};
59 changes: 59 additions & 0 deletions kernel/relayflowd/tests/hn_poller.rs
Original file line number Diff line number Diff line change
@@ -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"));
}