Firehose class final

A durable firehose consumer: it connects, hands each message to a handler, reconnects with exponential backoff when the connection drops, and persists how far it got so a restart resumes instead of skipping ahead to the live edge.

Why a handler and not a Stream

Exposing a Stream would be the more idiomatic Dart shape, but it cannot carry the cursor correctly: a stream tells this class when a message was delivered, never when the consumer finished with it. Advancing on delivery loses everything that was in flight when the process died — the exact loss a persisted cursor exists to prevent. The completion of the future returned by FirehoseMessageHandler is the only evidence that a message was handled, so start takes a handler and the delivery guarantee is at-least-once: after a crash, messages inside the unflushed window are replayed. Handlers should be idempotent.

Callers who do not care about the cursor can keep using atproto.sync.subscribeReposAsMessages() directly, which is unchanged.

Example

final firehose = Firehose(
  connect: (cursor) async {
    final subscription = await atproto.sync.subscribeReposAsMessages(
      cursor: cursor,
    );

    return FirehoseConnection(
      subscription.data.stream,
      close: subscription.data.close,
    );
  },
  cursorStore: MyCursorStore(),
);

await firehose.start((message) async {
  if (message.isCommit) await index(message.commit!);
});

Constructors

Firehose({required FirehoseConnector connect, CursorStore? cursorStore, Duration initialBackoff = const Duration(seconds: 1), Duration maxBackoff = const Duration(minutes: 1), Duration healthyConnectionThreshold = const Duration(seconds: 30), double jitter = 0.2, int flushEveryEvents = defaultFlushEveryEvents, Duration flushEveryInterval = defaultFlushEveryInterval, void onError(Object error, StackTrace stackTrace)?, Random? random})
Returns the new instance of Firehose.

Properties

flushEveryEvents → int
How many handled messages may pass between cursor writes.
final
flushEveryInterval → Duration
How long may pass between cursor writes.
final
hashCode → int
The hash code for this object.
no setterinherited
healthyConnectionThreshold → Duration
How long a connection must stay up before it counts as healthy and the consecutive-failure counter is reset.
final
initialBackoff → Duration
The reconnect delay after the first failure; doubles per consecutive failure up to maxBackoff.
final
jitter → double
The fraction of the base delay added at random, in 0.0..1.0. Keeps a fleet of instances from reconnecting in lockstep after a shared outage.
final
maxBackoff → Duration
final
runtimeType → Type
A representation of the runtime type of the object.
no setterinherited

Methods

noSuchMethod(Invocation invocation) → dynamic
Invoked when a nonexistent method or property is accessed.
inherited
start(FirehoseMessageHandler onMessage) → Future<void>
Runs until stop is called: connects, hands every message to onMessage, and reconnects with exponential backoff when the connection fails or closes.
stop() → Future<void>
Stops the consumer, tears the live connection down and writes out the pending cursor, after which start returns.
toString() → String
A string representation of this object.
inherited

Operators

operator ==(Object other) → bool
The equality operator.
inherited