{"record":{"id":"b641df2e30eb38ba","repo":"pola-rs/polars","slug":"not-implemented-b641df","errorCode":null,"errorMessage":"not implemented","messagePattern":"not implemented","errorType":"panic","errorClass":null,"httpStatus":null,"severity":"error","filePath":"crates/polars-arrow/src/io/ipc/read/flight.rs","lineNumber":372,"sourceCode":"                    )\n                    .map(Some)\n                } else {\n                    // Needed to memory map.\n                    let arrow_data = Arc::new(msg.arrow_data);\n                    unsafe {\n                        mmap_record(\n                            &self.md.schema,\n                            &self.md.ipc_schema.fields,\n                            arrow_data,\n                            batch,\n                            0,\n                            &self.dictionaries,\n                        )\n                        .map(Some)\n                    }\n                }\n            },\n            _ => unimplemented!(),\n        }\n    }\n}\n\npub struct FlightstreamConsumer<S: Stream<Item = PolarsResult<EncodedData>> + Unpin> {\n    inner: FlightConsumer,\n    stream: S,\n}\n\nimpl<S: Stream<Item = PolarsResult<EncodedData>> + Unpin> FlightstreamConsumer<S> {\n    pub async fn new(mut stream: S) -> PolarsResult<Self> {\n        let Some(first) = stream.next().await else {\n            polars_bail!(ComputeError: \"expected the schema\")\n        };\n        let first = first?;\n\n        Ok(FlightstreamConsumer {\n            inner: FlightConsumer::new(first)?,","sourceCodeStart":354,"sourceCodeEnd":390,"githubUrl":"https://github.com/pola-rs/polars/blob/df599052daf96e7a9cc30a3b0c6bd25d6947e3c0/crates/polars-arrow/src/io/ipc/read/flight.rs#L354-L390","documentation":"FlightConsumer::consume (crates/polars-arrow/src/io/ipc/read/flight.rs) decodes one Arrow IPC stream message per call and only handles Schema (rejected as unexpected), DictionaryBatch and RecordBatch headers. Any other MessageHeaderRef variant — flatbuffer union tag NONE (0) or a message kind this reader predates — falls into _ => unimplemented!() at line 372 and panics mid-stream. FlightstreamConsumer is built on top of it, so async flight streams hit the same arm.","triggerScenarios":"consumer.consume(msg) where MessageRef::read_as_root(...).header() yields something other than the three known kinds: truncated/corrupt flatbuffers, or a producer emitting message types this reader was never written to process.","commonSituations":"Arrow Flight / Flight SQL streams from newer servers; custom EncodedData producers; network truncation or proxy corruption producing garbage message headers.","solutions":["Validate the header kind before calling consume and skip or error on unknown messages","Pin producer and consumer to compatible Arrow IPC versions","Upgrade polars-arrow once the message kind is supported","Treat unknown headers as a corrupt stream: surface an error and abort the session instead of unwinding"],"exampleFix":"// before\nlet batch = consumer.consume(msg)?; // panics on unknown MessageHeader\n\n// after\nlet header = arrow_format::ipc::MessageRef::read_as_root(&msg.ipc_message)\n    .map_err(|e| polars_err!(oos = OutOfSpecKind::InvalidFlatbufferMessage(e)))?\n    .header()\n    .map_err(|e| polars_err!(oos = OutOfSpecKind::InvalidFlatbufferHeader(e)))?;\nmatch header {\n    Some(MessageHeaderRef::DictionaryBatch(_)) | Some(MessageHeaderRef::RecordBatch(_)) =>\n        consumer.consume(msg),\n    _ => polars_bail!(ComputeError: \"unsupported IPC message header in flight stream\"),\n}","handlingStrategy":"validation","validationCode":"let header = arrow_format::ipc::MessageRef::read_as_root(&msg.ipc_message)\n    .map_err(|e| polars_err!(oos = OutOfSpecKind::InvalidFlatbufferMessage(e)))?\n    .header()\n    .map_err(|e| polars_err!(oos = OutOfSpecKind::InvalidFlatbufferHeader(e)))?;\nmatch header {\n    Some(MessageHeaderRef::DictionaryBatch(_)) | Some(MessageHeaderRef::RecordBatch(_)) => consumer.consume(msg),\n    _ => polars_bail!(ComputeError: \"unsupported IPC message header in flight stream\"),\n}","typeGuard":null,"tryCatchPattern":"let res = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| consumer.consume(msg)));\nlet batch = match res {\n    Ok(v) => v?,\n    Err(_) => polars_bail!(ComputeError: \"flight stream contained an unsupported message kind; treat stream as corrupt\"),\n};","preventionTips":["Pin flight client and server to compatible Arrow IPC versions","Validate message headers before dispatching to FlightConsumer","Terminate the session with a clear error on unknown message kinds instead of retrying"],"tags":["rust","polars-arrow","flight","ipc","flatbuffers","panic"],"backgroundTag":null,"analyzedSha":"df599052daf96e7a9cc30a3b0c6bd25d6947e3c0","analyzedAt":"2026-08-16T12:10:03.978Z","schemaVersion":2},"datasetVersion":"2026-08-16T13:17:31.715Z"}