{"record":{"id":"721e520690ecac88","repo":"risingwavelabs/risingwave","slug":"illegal-message-id-string","errorCode":null,"errorMessage":"illegal message id string {}","messagePattern":"illegal message id string (.+?)","errorType":"validation","errorClass":"ConnectorError","httpStatus":null,"severity":"error","filePath":"src/connector/src/source/pulsar/source/reader.rs","lineNumber":156,"sourceCode":"pub struct PulsarBrokerReader {\n    #[expect(dead_code)]\n    pulsar: Pulsar<TokioExecutor>,\n    consumer: Consumer<Vec<u8>, TokioExecutor>,\n    split: PulsarSplit,\n    split_id: SplitId,\n    parser_config: ParserConfig,\n    source_ctx: SourceContextRef,\n\n    // for filter out already read messages\n    already_read_offset: Option<PulsarFilterOffset>,\n}\n\n// {ledger_id}:{entry_id}:{partition}:{batch_index}\nfn parse_message_id(id: &str) -> ConnectorResult<MessageIdData> {\n    let splits = id.split(':').collect_vec();\n\n    if splits.len() < 2 || splits.len() > 4 {\n        bail!(\"illegal message id string {}\", id);\n    }\n\n    let ledger_id = splits[0].parse::<u64>().context(\"illegal ledger id\")?;\n    let entry_id = splits[1].parse::<u64>().context(\"illegal entry id\")?;\n\n    let mut message_id = MessageIdData {\n        ledger_id,\n        entry_id,\n        partition: None,\n        batch_index: None,\n        ack_set: vec![],\n        batch_size: None,\n        first_chunk_message_id: None,\n    };\n\n    if splits.len() > 2 {\n        let partition = splits[2].parse::<i32>().context(\"illegal partition\")?;\n        message_id.partition = Some(partition);","sourceCodeStart":138,"sourceCodeEnd":174,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/connector/src/source/pulsar/source/reader.rs#L138-L174","documentation":"Pulsar message IDs are expected in the '{ledger_id}:{entry_id}:{partition}:{batch_index}' format (partition and batch_index optional). parse_message_id fails when the string does not contain 2–4 colon-separated segments, which would otherwise panic or corrupt the MessageIdData.","triggerScenarios":"Calling parse_message_id with a start-offset string that has fewer than 2 or more than 4 ':'-separated parts — e.g. '12345', 'a:b:c:d:e', or an empty string.","commonSituations":"Typing a Pulsar message id by hand into source start-offset config; copying an id from another system; including extra separators or pasting with whitespace/suffixes.","solutions":["Provide the message id as at least 'ledger:entry', e.g. '3:17' or '3:17:-1:-1'.","Verify the id was copied fully from Pulsar output (ledger_id:entry_id:partition:batch_index).","Strip whitespace, surrounding brackets, or trailing characters before passing the id.","Note the numeric segments must parse as u64; a malformed number fails with 'illegal ledger id'/'illegal entry id' context."],"exampleFix":"// before\nparse_message_id(\"1043-29\")?; // wrong separator\n// after\nparse_message_id(\"1043:29\")?; // ledger_id:entry_id","handlingStrategy":"validation","validationCode":"fn is_valid_message_id(id: &str) -> bool {\n    let parts: Vec<&str> = id.split(':').collect();\n    (2..=4).contains(&parts.len()) && parts[..2].iter().all(|p| p.parse::<u64>().is_ok())\n}","typeGuard":"fn parse_pulsar_message_id(id: &str) -> Option<(u64, u64)> {\n    let mut it = id.split(':');\n    Some((it.next()?.parse().ok()?, it.next()?.parse().ok()?))\n}","tryCatchPattern":"match parse_message_id(id) {\n    Err(e) if e.to_string().contains(\"illegal message id\") => eprintln!(\"expected ledger:entry[:partition][:batch_index]\"),\n    other => other,\n}","preventionTips":["Copy message ids verbatim from Pulsar output; they use ':' separators.","Validate the id with a regex like ^\\d+:\\d+(:-?\\d+){0,2}$ before use.","Trim whitespace and brackets when extracting ids from logs."],"tags":["pulsar","parsing","format","source"],"backgroundTag":"invalid-argument-format","analyzedSha":"6469eb736d691e8e9b8a419a57edd6429ca77417","analyzedAt":"2026-09-11T21:06:21.487Z","contentChangedAt":"2026-09-11T21:06:21.487Z","schemaVersion":2},"datasetVersion":"2026-09-23T08:17:48.524Z"}