{"record":{"id":"af56cb2679636965","repo":"nautechsystems/nautilus_trader","slug":"stream-receiver-already-taken-af56cb","errorCode":null,"errorMessage":"Stream receiver already taken","messagePattern":"Stream receiver already taken","errorType":"exception","errorClass":null,"httpStatus":null,"severity":"error","filePath":"crates/infrastructure/src/redis/msgbus.rs","lineNumber":407,"sourceCode":"            });\n        });\n\n        log::debug!(\"Closed\");\n    }\n}\n\nimpl RedisMessageBusBacking {\n    /// Retrieves the Redis stream receiver for this message bus instance.\n    ///\n    /// # Errors\n    ///\n    /// Returns an error if the stream receiver has already been taken.\n    pub fn get_stream_receiver(\n        &mut self,\n    ) -> anyhow::Result<tokio::sync::mpsc::Receiver<BusMessage>> {\n        self.stream_rx\n            .take()\n            .ok_or_else(|| anyhow::anyhow!(\"Stream receiver already taken\"))\n    }\n\n    /// Streams messages arriving on the stream receiver channel.\n    pub fn stream(\n        mut stream_rx: tokio::sync::mpsc::Receiver<BusMessage>,\n    ) -> impl Stream<Item = BusMessage> + 'static {\n        async_stream::stream! {\n            while let Some(msg) = stream_rx.recv().await {\n                yield msg;\n            }\n        }\n    }\n\n    pub async fn close_async(&mut self) {\n        await_handle(self.pub_handle.take(), MSGBUS_PUBLISH).await;\n        await_handle(self.stream_handle.take(), MSGBUS_STREAM).await;\n        await_handle(self.heartbeat_handle.take(), MSGBUS_HEARTBEAT).await;\n    }","sourceCodeStart":389,"sourceCodeEnd":425,"githubUrl":"https://github.com/nautechsystems/nautilus_trader/blob/18893faf8b356be3320add8de2f861b0b647cf06/crates/infrastructure/src/redis/msgbus.rs#L389-L425","documentation":"`get_stream_receiver` hands out the `mpsc::Receiver<BusMessage>` stored in the Redis message bus subscriber via `Option::take` — it can only be given out once. This error means the receiver was already taken (typically by `take_receiver`/`stream` startup), so a second consumer cannot be attached.","triggerScenarios":"Calling `get_stream_receiver()` (or `take_receiver`) twice on the same RedisBusSubscriber/MessageBus instance, e.g. calling `stream()` again after already starting the stream, or two components both trying to consume the same Redis stream subscription.","commonSituations":"Accidentally initializing the message bus stream twice (e.g. on reconnect logic without checking prior state); sharing one subscriber between two consumers; re-running an async setup task after a partial failure where the receiver was already consumed.","solutions":["Call `get_stream_receiver` exactly once per subscriber, storing the returned Receiver for the stream loop","Before re-initializing, check whether the subscriber already started streaming and skip the second call","If a second consumer is needed, create a separate subscriber/subscription instead of reusing the same receiver","On reconnect, recreate the subscriber object rather than re-taking its receiver"],"exampleFix":"// before\nlet rx1 = subscriber.get_stream_receiver()?;\nlet rx2 = subscriber.get_stream_receiver()?; // Err: Stream receiver already taken\n// after\nlet mut rx_opt = Some(subscriber.get_stream_receiver()?);\nif let Some(rx) = rx_opt.take() {\n    // consume rx exactly once\n}","handlingStrategy":"type-guard","validationCode":"// only attempt to take the receiver if streaming has not started\nif !subscriber.is_streaming() {\n    let rx = subscriber.get_stream_receiver()?;\n    // start stream(rx) once\n}","typeGuard":"// Rust\nfn try_get_receiver(sub: &mut RedisBusSubscriber) -> Option<tokio::sync::mpsc::Receiver<BusMessage>> {\n    sub.get_stream_receiver().ok()\n}","tryCatchPattern":"// Rust\nlet rx = match subscriber.get_stream_receiver() {\n    Ok(rx) => rx,\n    Err(e) if e.to_string().contains(\"already taken\") => {\n        log::debug!(\"stream already started; reusing existing consumer\");\n        return Ok(()); // not an error for idempotent setup\n    }\n    Err(e) => return Err(e),\n};","preventionTips":["Take the receiver exactly once, at a single well-known initialization point","Make stream initialization idempotent (skip if already streaming)","Never share one subscriber's receiver between multiple consumers","On reconnect, recreate the subscriber instead of re-taking its receiver"],"tags":["redis","msgbus","stream","channel","state"],"backgroundTag":"stream-receiver-already-taken","analyzedSha":"18893faf8b356be3320add8de2f861b0b647cf06","analyzedAt":"2026-09-08T20:49:34.690Z","contentChangedAt":"2026-09-08T20:49:34.690Z","schemaVersion":2},"datasetVersion":"2026-09-14T05:17:10.506Z"}