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:
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(())
}
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}")
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}`);
}
}
}
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}`);
}
}
}
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}`);
}
}
}
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")
}
}
}
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)")
}
}
}
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.