Skip to content

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:

  1. It figures out which of the current top stories it hasn't already processed, using a hn:seen durable state entry as a dedup set. See how to use durable state.
  2. It spawns one fetch_story task per new story, each with its own idempotency_key and 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 a fetch_story that finishes unusually fast never decrements a counter that isn't there yet.
  3. Before returning, it spawns its own next run, POLL_INTERVAL later, with a fresh idempotency_key derived from a strictly-incrementing run_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

python poll_hacker_news.py
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.