Observation
Pulls a bounded Java subscription with one active stream consumer. A full buffer yields a Delivery.Gap in this stream before the events that remain. FS2 demand does not provide tmux backpressure. A subscription does not reconnect: attach again and read a snapshot.
Source and package
-
Description#
The stream itself suspends the fiber, not a platform thread, while idle:
EventSubscriptionis already single-consumer by construction at the Java layer, with a non-blockingpollplus a one-shotonReadywakeup. Each wait polls first — if something is already buffered, nothing suspends at all — and only then armsonReady, disarming it withclearReadyif the fiber is cancelled first. NoExecutionContextsized for blocking stream reads is needed, because no read here ever blocks a thread.
Members6 members#
Common operations3 members#
- value method The value inside an event delivery, and none for a gap.
- kept method The event a strict reader kept. A gap fails the read.
- UnknownCause class Java ended the subscription without exposing the reason.
Other members3 members#
- droppedCount method The cumulative overflow count. A
Delivery.Gapinstreamsays where it sits. - isClosed method Reports subscription closure, including after resource finalization.
- stream method Reads on
F's own fiber scheduler. Cancelling the stream releases its consumer slot and disarms any pending wakeup; releasing the observation discards its buffered events. Deliberate closure ends the stream. A client that ended the subscription fails it with that cause.