{"record":{"id":"b591746d34b3ddce","repo":"risingwavelabs/risingwave","slug":"should-exist","errorCode":null,"errorMessage":"should exist","messagePattern":"should exist","errorType":"panic","errorClass":null,"httpStatus":null,"severity":"error","filePath":"src/frontend/src/handler/create_sink.rs","lineNumber":946,"sourceCode":"\n    Ok(SinkCreateMode::Replace { original_sink })\n}\n\npub fn fetch_incoming_sinks(\n    session: &Arc<SessionImpl>,\n    table: &TableCatalog,\n) -> Result<Vec<Arc<SinkCatalog>>> {\n    let reader = session.env().catalog_reader().read_guard();\n    let schema = reader.get_schema_by_id(table.database_id, table.schema_id)?;\n    let Some(incoming_sinks) = schema.table_incoming_sinks(table.id) else {\n        return Ok(vec![]);\n    };\n    let mut sinks = vec![];\n    for sink_id in incoming_sinks {\n        sinks.push(\n            schema\n                .get_sink_by_id(*sink_id)\n                .expect(\"should exist\")\n                .clone(),\n        );\n    }\n    Ok(sinks)\n}\n\nfn derive_sink_to_table_expr(\n    sink_schema: &Schema,\n    idx: usize,\n    target_type: &DataType,\n) -> Result<ExprImpl> {\n    let input_type = &sink_schema.fields()[idx].data_type;\n\n    if !target_type.equals_datatype(input_type) {\n        bail!(\n            \"column type mismatch: {:?} vs {:?}, column name: {:?}\",\n            target_type,\n            input_type,","sourceCodeStart":928,"sourceCodeEnd":964,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/frontend/src/handler/create_sink.rs#L928-L964","documentation":"`fetch_incoming_sinks` calls `expect(\"should exist\")` after looking up each incoming sink ID in the target schema's catalog. The ID came from the schema's own `incoming_sinks` list, so the sink must exist; if it does not, the catalog metadata is inconsistent (e.g. during concurrent DROP SINK or after a partial deletion), and the handler panics instead of returning an error.","triggerScenarios":"Listing/fetching downstream sinks of a table (`fetch_incoming_sinks`, used by e.g. table drop checks or SHOW) while a sink referencing the table was just dropped, leaving a stale entry in `incoming_sinks`; meta/catalog cache out of sync with the schema catalog.","commonSituations":"Race between `DROP SINK` and operations that enumerate incoming sinks (like `DROP TABLE` dependency checks); catalog cache staleness in the frontend after meta node failover.","solutions":["Retry the operation — transient races with DROP SINK usually resolve once catalogs resync.","Ensure sinks are dropped/created outside of concurrent table operations in scripts.","Check catalog consistency (`SHOW SINKS`) and recreate dangling sink references if any.","Internal fix: return a `catalog` error for missing sink IDs instead of `expect`."],"exampleFix":"// before\nsinks.push(schema.get_sink_by_id(*sink_id).expect(\"should exist\").clone());\n// after\nlet sink = schema.get_sink_by_id(*sink_id)\n    .ok_or_else(|| anyhow!(\"sink {:?} not found in catalog\", sink_id))?;\nsinks.push(sink.clone());","handlingStrategy":"retry","validationCode":null,"typeGuard":null,"tryCatchPattern":"// retry transient catalog races\nfor (let i = 0; i < 3; i++) {\n  try { return await fetchIncomingSinks(tableId); }\n  catch (e) {\n    if (String(e).includes('should exist') && i < 2) { await sleep(500); continue; }\n    throw e;\n  }\n}","preventionTips":["Avoid running DROP SINK concurrently with table dependency checks/show operations.","Serialize catalog mutations in deployment scripts.","Restart frontend/meta if catalog caches are suspected stale."],"tags":["catalog","sink","race-condition","panic"],"backgroundTag":"record-not-found","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"}