skip to content

Observable Streams

Sources that push values to subscribers: the contract they obey, whether each subscriber gets its own run, and when work starts and stops. Cold-versus-hot explains most surprising reactive bugs.

on this pageshow

explore

questions

18

A hand-written stream source emits a failure signal and then keeps pushing values — which rule does that break?

level: juniorimportance: must knowfreq 70%

answer

  1. the sequence has a grammar
  2. endings are signals too
  3. how many endings are allowed
  4. what may follow the ending
  5. consumer releases state at the ending

basics

~20 s

The signal grammar: a run carries any number of value signals and then at most one terminal signal, completion or failure, never both. A failure is that ending, not another value, so the source must go silent after it.

solid answer

~50 s

A push-based source speaks a fixed grammar: zero or more **value** signals, then **at most one terminal** signal — either completion or failure — and nothing at all afterwards. The adapter treats the failure as if it were a warning and keeps draining its buffer, but the failure already ended the sequence. The consumer is entitled to release its per-run state the moment the terminal arrives, so anything delivered later hits a stage that has already torn itself down, or is silently dropped by a defensive stage and the data loss never surfaces. Note what the rule does not say: it caps terminals at one, it does not require one, so a source that runs forever is fine. The fix is a latched flag at the emission boundary that turns every later signal into a no-op.

code

pseudocode · 9 lines
pseudocode
function runSource(consumer):
    try:
        for each row in rows:
            consumer.value(row)
        consumer.completed()
    catch err:
        consumer.failure(err)          // the sequence ends here
        for each row in leftovers:     // defect: it already ended
            consumer.value(row)

go deeper

for a junior

Recall the shape of a run: any number of values, then at most one ending, and silence after it. Know that a failure is an ending rather than a strange value.

for a middle

Explain why the prohibition sits on the source: the consumer releases per-run state at the terminal, so a later signal reaches a stage that has torn itself down or is dropped invisibly.

for a senior

Show how you would enforce it in code you accept from other teams — one emission boundary, a latched flag, no error path that can reach the consumer after another path ended the run.

for a principal

Frame it as an interoperability guarantee: stages written by different people compose only because the ending is agreed, and every locally convenient exception to it pushes cost into every consumer.

## The two kinds of signal A push-based source talks to a consumer with exactly two kinds of signal. A **value signal** carries one element of the sequence. A **terminal signal** says the sequence is over, and comes in two flavours: **completion**, meaning the source produced everything it had, and **failure**, meaning the source cannot continue and is handing over a reason instead of a value. The whole grammar is one line: *any number of value signals, then at most one terminal signal, and nothing after it.* As a pattern, `value* (completion | failure)?`. Three rules fall straight out of it: - A run carries **at most one** terminal signal — not one of each kind, not two completions. - A failure **is** the ending, not an element the consumer reads past. - After the terminal signal the source is **silent on that subscription**, permanently. | signal kind | how many per run | may anything follow it | |---|---|---| | value | zero or more | yes — more values, or the terminal | | completion | at most one | no | | failure | at most one | no | ## Why the reviewed adapter is broken The adapter catches an error, reports it as a failure, and then drains whatever it had already buffered. Read as a protocol rather than as a pile of callbacks, that is two messages after a goodbye. Four things go wrong, none of them at the line that caused them: - The consumer may have **released its per-run state** at the failure — an accumulator, a handle it was holding, an open resource. A later value arrives at a stage that no longer has anywhere to put it. - A stage that follows the contract defensively **drops** the stray signals. That is the worse outcome, because the violation is now invisible: values vanish and nothing reports it. - If the leftovers are pushed from the worker that was unwinding the error, an exception thrown downstream surfaces **inside the source's own error path**, where nobody is looking for it. - The symptom appears far from the cause, usually in whichever stage happened to hold state, which is why this defect is expensive to find and cheap to prevent. ## What the consumer is entitled to do at the ending The grammar is what makes a consumer's job finite. On the terminal signal a consumer may: 1. Emit whatever it had accumulated and forget it — a stage that counts, batches or folds knows it will never be asked to fold again. 2. Release everything tied to this run: buffers, timers, the handle it was given when it attached. 3. Treat any later signal as a defect in the source rather than as data, because the contract says one cannot legitimately arrive. That third point is the reason the rule is stated as a prohibition on the source rather than as advice to the consumer. If a consumer had to stay ready for stragglers, it could never release anything, and the ending would carry no information. ## Exactly-one is required; ending at all is not The two halves of the rule are often collapsed and they are different. A source is **forbidden** to deliver a second terminal, and **not obliged** to deliver a first one: a source over a live feed may run for the lifetime of the process. That has a practical consequence — a consumer that needs an ending has to impose one itself, and from the outside a source that will never end is indistinguishable from one that has stalled, except by waiting a chosen amount of time and deciding. One related case is worth naming because it comes up in exactly this review. If the source finishes cleanly, signals completion, and *then* its own cleanup step fails, it may not signal that failure: the sequence has ended and accepts nothing more. The failure has to be reported through some channel outside the stream. Trying to squeeze it into the ended sequence is the same violation wearing a sympathetic motive. ## Reading it as a review checklist For a hand-written source, four questions settle conformance on this rule alone: - Is there **one** place where signals leave the source, or several scattered through the code? - Does that place carry a **latched flag** that is set by the first terminal and consulted by every later signal? - Can any error path reach the consumer **after** another path has already ended the run? - Does the failure signal carry a **reason**, rather than being a completion with a log line beside it? A source that answers those four cleanly cannot break this rule, whatever else it gets wrong.

  • Is a source that never emits a terminal signal violating this contract?
    No. The grammar caps terminal signals at one; it does not require one. A source over a continuing feed may run as long as the process does. A consumer that needs an ending must impose it itself, and from outside, endless and stalled look identical until you pick a waiting time.
  • May one run deliver both a completion and a failure?
    No — at most one terminal signal of either kind, whichever the source reaches first. If a source completes and then its cleanup fails, the sequence has already ended and cannot carry the failure; it has to be reported outside the stream, not pushed into a run that is over.

A sequence is a letter: any number of lines, then one sign-off. A line added after the sign-off does not extend the letter — it is a second letter nobody agreed to read.

saying these in an interview costs you the question

  • Treats a failure as a value the consumer can keep reading past
  • Says a source may complete after it has already failed
  • Thinks stray signals are harmless because consumers drop them
  • Assumes every stream must eventually reach an ending
  • Leaves cleanup waiting for more values after the terminal signal
  • Signals a post-completion cleanup error into the ended sequence
open as a page

A chat screen builds a message stream but no request is sent — what act starts the work, and what does it return?

level: juniorimportance: must knowfreq 72%

basics

~20 s

Subscribing starts the work; building a stream only describes it. A subscription call attaches a subscriber to the source, triggers whatever the source does to produce values, and hands back a handle the caller uses to cancel that run.

open as a page

What has happened when a source declared to carry at most one value completes without emitting a value?

level: juniorimportance: must knowfreq 65%

basics

~20 s

Nothing was found and nothing failed: empty completion is a third outcome beside a value and a failure. The caller must decide what absence means here - a default, a fallback source, or an error it raises itself.

open as a page

Three dashboard panels subscribe to one source built around a query, and that query runs three times - why?

level: middleimportance: must knowfreq 70%

basics

~10 s

The source is cold: it describes work instead of sharing a sequence that is already running, so every subscriber starts a fresh, independent run. Three panels subscribed, so the query executed three times.

open as a page

Two workers in a hand-written stream source push values to one subscriber concurrently — why does the contract forbid that?

level: middleimportance: must knowfreq 58%

basics

~20 s

Signals must be delivered one at a time, with each one seeing what the previous one wrote. Every stage downstream is written on that promise and keeps its accumulated state unsynchronised, so overlapping deliveries corrupt it.

open as a page

A user closes a conversation screen and its message subscription is cancelled — what travels upstream, and when does the source stop?

level: middleimportance: must knowfreq 62%

basics

~20 s

A cancel request travels up the same chain the values came down: each stage stops asking its upstream, forwards the request and tears itself down. The source stops at the next point where it observes the request, not instantly.

open as a page

When you narrow a many-valued source to one that carries at most one value, what must the conversion decide?

level: middleimportance: must knowfreq 52%

basics

~20 s

Two things: what happens to values after the first - ignored, or treated as a contract violation - and what zero values means, since requiring exactly one turns absence into a failure while taking the first leaves it an empty completion.

open as a page

A dashboard panel that subscribes to a live shared feed after it starts shows nothing until the next update - why?

level: middleimportance: should knowfreq 54%

basics

~20 s

A live source runs independently of its subscribers and hands each one only what arrives while it is attached. Everything emitted before the panel subscribed was delivered to whoever was attached then and kept nowhere, so the panel waits for the next value.

open as a page

Why is a push-based stream source called a protocol between three roles rather than a producer invoking a consumer callback?

level: middleimportance: should knowfreq 46%

basics

~20 s

Because three parties are named and bound by rules: a source, a consumer, and a per-subscription handle created when they attach. A bare callback has no ending, no separate failure channel, no identity for one run, and no way to answer back.

open as a page

A message source opens a connection when subscribed — on which endings must its teardown run, and which ending is forgotten?

level: middleimportance: should knowfreq 50%

basics

~20 s

A subscription ends in one of three ways: normal completion, failure, or cancellation. Teardown must run on all three, and cancellation is the one teams forget, because it is the only ending the source did not choose.

open as a page

A directory exposes lookup by identifier beside search by name - what does declaring each result's cardinality commit you to?

level: middleimportance: should knowfreq 44%

basics

~20 s

Declaring at most one promises no caller ever sees a second value but obliges every caller to handle absence; declaring many promises nothing about count. The signature makes that decision once, for every caller that ever reads it.

open as a page

In a multicast source with reference counting, what happens to the upstream query when the last dashboard panel closes?

level: seniorimportance: should knowfreq 47%

basics

~20 s

The shared stage cancels its single upstream subscription when the subscriber count falls to zero, stopping the query and usually discarding its accumulated state. The next panel that opens starts a brand new upstream run rather than rejoining the old one.

open as a page

Closed chat screens keep receiving pushed messages and memory grows with every open-and-close cycle — how do you diagnose and fix it?

level: seniorimportance: should knowfreq 58%

basics

~20 s

Subscriptions are never cancelled, so each closed screen leaves a live one behind. The source holds the subscriber, the subscriber holds the screen, and nothing can be released. Fix it by tying every subscription's disposal to the lifetime of the owner that created it.

open as a page

A stage maps each search result through an at-most-one lookup and emits fewer values than it received, with no failure - why?

level: seniorimportance: should knowfreq 38%

basics

~20 s

Every inner lookup that completes empty contributes zero values, so n inputs produce between zero and n outputs. Absence is a normal ending rather than a failure, so the misses vanish silently unless the empty case is given a value of its own.

open as a page

Your platform publishes one dashboard data source to several product teams - would you expose it per-subscriber or multicast, and why?

level: principalimportance: should knowfreq 38%

basics

~20 s

Decide it as a published contract, not an implementation detail. Per-subscriber re-execution keeps consumers independent and parameterisable; one multicast sequence collapses cost but couples every consumer to one set of parameters, one failure and one lifecycle.

open as a page

When a shared dashboard feed replays its last values to a late panel, what can the replayed value mislead that panel about?

level: seniorimportance: nice to knowfreq 32%

basics

~20 s

Its age and its pacing. A replayed value arrives at subscription time carrying no record of when it was produced, so a panel that treats arrival as observation renders an old reading as current and computes rates from a burst that never happened at that speed.

open as a page

Your platform accepts hand-written stream sources from other teams — what must an automated conformance check assert about each?

level: seniorimportance: nice to knowfreq 30%

basics

~20 s

Subscribe with a recording consumer and assert the grammar: the handle arrives before any value, at most one terminal signal, nothing after it, a failure carries a reason, and no two deliveries overlap. The last one can only be sampled under stress.

open as a page

A message subscription stops delivering values — how can the code tell a cancelled run from a completed one, and why does the difference matter?

level: seniorimportance: nice to knowfreq 30%

basics

~20 s

Completion arrives as a terminal signal from the source; cancellation is a request the subscriber sent upstream, after which it is simply not told anything. So a completion handler never fires on a cancel, while an any-ending hook fires on both and cannot distinguish them by itself.

open as a page