Polling an API, fanning out, and rescheduling¶
The how-to guides each cover one piece of dura in isolation. This page
walks through a single, more realistic script that combines several of
them: it polls a public API on a schedule, fans out per-item work with
its own retries, and waits for that work to finish before summarizing it.
The scenario: keep an eye on Hacker News's
front page, fetch and process any story we haven't seen yet, and print a
digest once a batch is done. The whole script is at
examples/poll_hacker_news.py
and runs as-is, no API key required.
The setup¶
import json
import urllib.request
from datetime import timedelta
from urllib.parse import urlparse
from dura import DurableEngine, RetryStrategy, run_workers
API = "https://hacker-news.firebaseio.com/v0"
STORIES_PER_POLL = 5
POLL_INTERVAL = timedelta(seconds=30) # a real poller would use minutes, not seconds
def fetch_json(url):
with urllib.request.urlopen(url, timeout=10) as response:
return json.load(response)
fetch_json is the only thing in this script that isn't dura: a small
wrapper around the standard library's urllib.request, since the point
of this example is what happens around a network call, not how to make
one.
Polling and fanning out¶
def poll_top_stories(engine, task):
run_number = task.params.get("run_number", 0) + 1
batch_id = task.task_id
all_ids = fetch_json(f"{API}/topstories.json")
unseen = [
story_id
for story_id in all_ids
if engine.get_state(namespace="hn:seen", key=str(story_id), default=None)
is None
]
batch = unseen[:STORIES_PER_POLL]
if batch:
# Set the countdown before spawning any child, so a fetch_story that
# finishes unusually fast never decrements a counter that isn't
# written yet.
engine.set_state(
namespace=f"batch:{batch_id}", key="remaining", value=len(batch)
)
for story_id in batch:
engine.spawn_task(
name="fetch_story",
params={"story_id": story_id, "batch_id": batch_id},
idempotency_key=f"fetch_story:{story_id}",
max_attempts=3,
retry=RetryStrategy(kind="fixed", base_seconds=5),
)
engine.spawn_task(
name="summarize_batch",
params={"batch_id": batch_id},
idempotency_key=f"summarize_batch:{batch_id}",
)
print(f"poll #{run_number}: fanned out {len(batch)} new stories")
else:
print(f"poll #{run_number}: nothing new")
# Reschedule the next poll before returning. A fresh idempotency_key,
# derived from a strictly-incrementing run_number, keeps this chain from
# forking or stalling.
engine.spawn_task(
name="poll_top_stories",
params={"run_number": run_number},
available_after=POLL_INTERVAL,
idempotency_key=f"poll_top_stories:{run_number}",
)
return {"fanned_out": len(batch)}
poll_top_stories does three things, and the order matters:
- It figures out which of the current top stories it hasn't already
processed, using a
hn:seendurable state entry as a dedup set. See how to use durable state. - It spawns one
fetch_storytask per new story, each with its ownidempotency_keyand retry strategy, since a flaky network call should only cost that one story a retry, not the whole poll. See how to chain and fan out tasks and how to configure retries. The countdown used for fan-in is written before any child is spawned, so afetch_storythat finishes unusually fast never decrements a counter that isn't there yet. - Before returning, it spawns its own next run,
POLL_INTERVALlater, with a freshidempotency_keyderived from a strictly-incrementingrun_number. That's what makes it recurring. See how to schedule recurring tasks.
Fetching and processing each story¶
def fetch_story(engine, task):
story_id = task.params["story_id"]
batch_id = task.params["batch_id"]
item = engine.checkpoint(
task_id=task.task_id,
step_name="fetch",
fn=lambda: fetch_json(f"{API}/item/{story_id}.json"),
)
digest = engine.checkpoint(
task_id=task.task_id,
step_name="process",
fn=lambda: {
"title": item.get("title", "(no title)"),
"score": item.get("score", 0),
"domain": urlparse(item.get("url", "")).netloc or "news.ycombinator.com",
},
)
engine.set_state(
namespace=f"batch:{batch_id}:digests", key=str(story_id), value=digest
)
engine.set_state(namespace="hn:seen", key=str(story_id), value=True)
remaining = engine.update_state(
namespace=f"batch:{batch_id}", key="remaining", fn=lambda n: n - 1
)
if remaining == 0:
engine.emit_event(event_name=f"batch_done:{batch_id}")
return digest
Fetching and processing are two separate checkpoints, not one, so a crash between them doesn't repeat the network call just to redo a bit of local computation. See how to checkpoint steps.
Once a story is processed, fetch_story writes its digest to durable
state (so summarize_batch can read it back later, from a different
task) and marks the story seen (so no later poll fans it out again).
Then it decrements the batch's countdown, and if this was the last story
in the batch, emits an event. Nothing is polling for that countdown to
hit zero: whichever fetch_story happens to be the last one just says so.
Waiting for the batch to finish¶
def summarize_batch(engine, task):
batch_id = task.params["batch_id"]
engine.wait_for_event(
run_id=task.run_id,
task_id=task.task_id,
step_name="wait",
event_name=f"batch_done:{batch_id}",
timeout_secs=600,
)
digests = engine.list_state(f"batch:{batch_id}:digests")
print(f"batch {batch_id[:8]} done, {len(digests)} stories:")
for digest in digests.values():
print(f" [{digest['score']:>4}] {digest['title']} ({digest['domain']})")
return {"summarized": len(digests)}
summarize_batch is spawned once per poll, right alongside the
fetch_story tasks, and immediately parks itself on that same event,
however many stories are in the batch or however long they take. See
how to wait for events. Once the event
fires, it reads back every digest written to durable state and prints
them.
Wiring it together¶
engine = DurableEngine("poll_hacker_news.db")
engine.spawn_task(
name="poll_top_stories",
params={"run_number": 0},
idempotency_key="poll_top_stories:0",
)
run_workers(
engine,
handlers={
"poll_top_stories": poll_top_stories,
"fetch_story": fetch_story,
"summarize_batch": summarize_batch,
},
worker_count=4,
)
One run_workers pool with a handful of threads runs all three task
types: the poll, the fetches, and the summary. dura doesn't need to
know they're related to each other, that relationship lives entirely in
the batch_id each one carries in its params. See how to run a worker
pool.
Run it¶
poll #1: fanned out 5 new stories
batch f6fd05ac done, 5 stories:
[ 127] A faster way to calculate the day of the week (www.benjoffe.com)
[ 819] OpenRouter is joining Stripe (openrouter.ai)
[ 609] Go 1.27 (go.dev)
[ 179] Turns are Better than Radians (2022) (www.computerenhance.com)
[ 100] Windows brings out the Rorschach test in everyone (devblogs.microsoft.com)
Leave it running and it polls again every POLL_INTERVAL, picking up a
fresh batch of stories each time, since the ones already seen are
excluded up front. Kill it, Ctrl+C or kill -9, at any point and run
it again: whatever had already completed, a fetch, a checkpoint, a
countdown, an emitted event, stays exactly as it was, and only the work
that hadn't finished yet gets picked back up.