Keyboard shortcuts

Press or to navigate between chapters

Press S or / to search in the book

Press ? to show this help

Press Esc to hide this help

Streaming events

Most finite reads should start with stream_events. The SDK opens a short-lived request, validates and deduplicates matching events across the selected relays, and yields each result as soon as it arrives. The application can begin work immediately and does not need to retain the complete result set in memory.

This is different from a live subscription. A finite stream normally ends at end-of-stored-events (EOSE), when all relay streams end, or when its timeout expires. A live subscription remains registered until it is closed.

For a pool request, EOSE happens per relay. One relay may finish immediately, another may yield events first, and a third may fail or reach the timeout. The SDK merges those independent flows into one stream while preserving the relay URL on each item.

The following example reads up to 20 text notes per relay and applies a ten-second upper bound:

Rust
use std::time::Duration;

use nostr_sdk::prelude::*;

async fn stream_events() -> Result<(), Box<dyn std::error::Error>> {
    let client = Client::default();
    client.add_relay("wss://relay.damus.io").await?;
    client.connect().await;

    let filter = Filter::new().kind(Kind::TextNote).limit(20);
    let mut stream = client
        .stream_events(filter)
        .timeout(Duration::from_secs(10))
        .await?;

    while let Some((relay_url, result)) = stream.next().await {
        match result {
            Ok(event) => println!("{relay_url}: {}", event.as_json()),
            Err(error) => eprintln!("{relay_url}: {error}"),
        }
    }

    Ok(())
}
Python
import asyncio
from datetime import timedelta

from nostr_sdk import Client, Filter, Kind, KindStandard, RelayUrl, ReqTarget

async def stream_events() -> None:
    client = Client()
    await client.add_relay(RelayUrl.parse("wss://relay.damus.io"))
    await client.connect()

    filter = Filter().kind(Kind.from_std(KindStandard.TEXT_NOTE)).limit(20)
    stream = await client.stream_events(
        ReqTarget.auto([filter]), timeout=timedelta(seconds=10)
    )

    while item := await stream.next():
        if item.event is not None:
            print(f"{item.relay_url}: {item.event.as_json()}")
        elif item.error is not None:
            print(f"{item.relay_url}: {item.error}")

JavaScript
Node.js
import {
    Client,
    Filter,
    Kind,
    KindStandard,
    RelayUrl,
    ReqTarget,
} from "@nostrdevkit/nostr-sdk-node";

async function streamEvents() {
    const client = new Client();
    await client.addRelay(RelayUrl.parse("wss://relay.damus.io"));
    await client.connect();

    const filter = new Filter()
        .kind(Kind.fromStd(KindStandard.TextNote))
        .limit(20n);
    const stream = await client.streamEvents(
        ReqTarget.auto([filter]),
        undefined,
        10_000,
    );

    while (true) {
        const item = await stream.next();
        if (!item) break;

        if (item.event) {
            console.log(`${item.relayUrl}: ${item.event.asJson()}`);
        } else if (item.error) {
            console.error(`${item.relayUrl}: ${item.error}`);
        }
    }

}
Web
import {
    Client,
    Filter,
    Kind,
    KindStandard,
    RelayUrl,
    ReqTarget,
    uniffiInitAsync,
} from "@nostrdevkit/nostr-sdk-web";

await uniffiInitAsync();

async function streamEvents() {
    const client = new Client();
    await client.addRelay(RelayUrl.parse("wss://relay.damus.io"));
    await client.connect();

    const filter = new Filter()
        .kind(Kind.fromStd(KindStandard.TextNote))
        .limit(20n);
    const stream = await client.streamEvents(
        ReqTarget.auto([filter]),
        undefined,
        10_000,
    );

    while (true) {
        const item = await stream.next();
        if (!item) break;

        if (item.event) {
            console.log(`${item.relayUrl}: ${item.event.asJson()}`);
        } else if (item.error) {
            console.error(`${item.relayUrl}: ${item.error}`);
        }
    }

}
React Native
import {
    Client,
    Filter,
    Kind,
    KindStandard,
    RelayUrl,
    ReqTarget,
} from "@nostrdevkit/nostr-sdk-react-native";

export async function streamEvents() {
    const client = new Client();
    await client.addRelay(RelayUrl.parse("wss://relay.damus.io"));
    await client.connect();

    const filter = new Filter()
        .kind(Kind.fromStd(KindStandard.TextNote))
        .limit(20n);
    const stream = await client.streamEvents(
        ReqTarget.auto([filter]),
        undefined,
        10_000,
    );

    while (true) {
        const item = await stream.next();
        if (!item) break;

        if (item.event) {
            console.log(`${item.relayUrl}: ${item.event.asJson()}`);
        } else if (item.error) {
            console.error(`${item.relayUrl}: ${item.error}`);
        }
    }

}
Kotlin
import java.time.Duration
import kotlinx.coroutines.runBlocking
import org.nostrdevkit.sdk.*

suspend fun streamEvents() {
    val client = Client()
    client.addRelay(RelayUrl.parse("wss://relay.damus.io"))
    client.connect()

    val filter = Filter()
        .kind(Kind.fromStd(KindStandard.TEXT_NOTE))
        .limit(20u)
    val stream = client.streamEvents(
        ReqTarget.auto(listOf(filter)),
        timeout = Duration.ofSeconds(10),
    )

    while (true) {
        val item = stream.next() ?: break
        val event = item.event
        val error = item.error
        when {
            event != null -> println("${item.relayUrl}: ${event.asJson()}")
            error != null -> System.err.println("${item.relayUrl}: $error")
        }
    }

}
Swift
import Foundation
import NostrSDK

func streamEvents() async throws {
    let client = Client()
    let relay = try RelayUrl.parse(url: "wss://relay.damus.io")
    _ = try await client.addRelay(url: relay)
    await client.connect()

    let filter = Filter()
        .kind(kind: Kind.fromStd(e: .textNote))
        .limit(limit: 20)
    let stream = try await client.streamEvents(
        target: ReqTarget.auto(filters: [filter]),
        timeout: 10.0
    )

    while let item = await stream.next() {
        if let event = item.event {
            print("\(item.relayUrl): \(try event.asJson())")
        } else if let error = item.error {
            print("\(item.relayUrl): \(error)")
        }
    }

}
C#
using Nostr.Sdk;

public static class ReadExample
{

    public static async Task StreamEvents()
    {
        var client = new Client();
        await client.AddRelay(RelayUrl.Parse("wss://relay.damus.io"));
        await client.Connect();

        var filter = new Filter()
            .Kind(Kind.FromStd(KindStandard.TextNote))
            .Limit(20);
        using var stream = await client.StreamEvents(
            ReqTarget.Auto([filter]),
            timeout: TimeSpan.FromSeconds(10)
        );

        while (await stream.Next() is { } item)
        {
            if (item.Event is { } @event)
            {
                Console.WriteLine($"{item.RelayUrl}: {@event.AsJson()}");
            }
            else if (item.Error is { } error)
            {
                Console.Error.WriteLine($"{item.RelayUrl}: {error}");
            }
        }

    }
}

ReqTarget.auto resolves eligible relays, the timeout bounds the operation, and next yields merged relay outcomes. The loop can process the first event without waiting for every relay or buffering the complete result set.

Stream items

Each item is one relay outcome. Rust returns a relay URL with Result<Event, _>; the bindings expose the URL and either an event or an error.

Process successful events even if another relay reports an error. Events are validated and deduplicated by ID before they reach the client stream, so seeing the same event on several relays does not require application-level deduplication for that request. The relay URL remains useful for diagnostics and for features whose policy depends on where an event was observed.

The stream completes when next returns no item. EOSE applies to the queried relays; it is not a claim that no other relay has matching events.

Request targets and termination

The filter describes what should match; the request target describes where it should be requested. ReqTarget.auto uses relays with read capability and can incorporate relay discovery when gossip is enabled. single and manual targets are available when the application needs exact relay selection or different filters per relay.

EOSE is the default successful completion policy; the timeout is the upper bound when completion does not arrive. Other exit policies are available in the API reference. Stop consuming and release the stream when the caller has enough data.

Use collecting events when the next operation genuinely requires the complete result set.