risingwavelabs/risingwave · error · ConnectorError

Unsupported source

Error message

Unsupported source: {:?}

What it means

The connector enumerator factory only recognizes a fixed set of source configs; get_source_list matches on the ConnectorProperties variant and bails with 'Unsupported source' for any variant that has no list-stream implementation. This indicates the caller constructed a source config the enumerator cannot enumerate.

Solutions

  1. Check which source type was configured and use one supported by the list/enumerator path (S3, GCS, POSIX FS, etc.)
  2. Update the match statement in get_source_list to handle the new variant if you are adding a connector
  3. Rebuild with the relevant `source-*` cargo feature enabled so the variant has an implementation

Example fix

// before
other => bail!("Unsupported source: {:?}", other),
// after
ConnectorProperties::S3(prop) => { /* build S3 list stream */ }
other => bail!("Unsupported source: {:?}", other),
Defensive patterns

Strategy: type-guard

Validate before calling

const supportedListSources = ['s3','gcs','posix_fs','opendal'];
function validateEnumerableSource(sourceType) {
  if (!supportedListSources.includes(sourceType)) throw new Error(`source ${sourceType} has no list/enumerator support`);
}

Type guard

const isEnumerableProps = (p) => ['S3','Gcs','PosixFs','Opendal'].some(v => p.variant === v);

Try / catch

match get_source_list(props) { Err(e) if String(e).contains("Unsupported source") => pick_supported_enumerator_or_static_assignment(), r => r }

Prevention

When it happens

Trigger: Calling get_source_list on a ConnectorConnector/properties enum holding a source type without enumerator support (e.g. a CDC or push-based source passed into the FS enumeration path), typically due to a wiring/config mistake rather than user input.

Common situations: Adding a new connector without updating the match in reader.rs, or misconfiguring a source so the wrong ConnectorProperties variant is built; also binary/feature builds where some enumerators are compiled out.

Related errors


AI-assisted analysis of risingwavelabs/risingwave@6469eb736d (2026-09-11). Data as JSON: /api/errors/e9ef5ce4add22019. Report an issue: GitHub.

Appendix: source

Thrown at src/connector/src/source/reader/reader.rs:119

                    &prop.s3_properties,
                    prop.assume_role,
                    prop.fs_common.compression_format,
                )?;
                Ok(build_opendal_fs_list_stream(lister, list_interval_sec))
            }
            ConnectorProperties::Azblob(prop) => {
                list_interval_sec = get_list_interval_sec(prop.fs_common.refresh_interval_sec);
                let lister: OpendalEnumerator<OpendalAzblob> =
                    OpendalEnumerator::new_azblob_source(*prop)?;
                Ok(build_opendal_fs_list_stream(lister, list_interval_sec))
            }
            ConnectorProperties::PosixFs(prop) => {
                list_interval_sec = get_list_interval_sec(prop.fs_common.refresh_interval_sec);
                let lister: OpendalEnumerator<OpendalPosixFs> =
                    OpendalEnumerator::new_posix_fs_source(*prop)?;
                Ok(build_opendal_fs_list_stream(lister, list_interval_sec))
            }
            other => bail!("Unsupported source: {:?}", other),
        }
    }

    /// Refer to `WaitCheckpointWorker` for more details.
    pub async fn create_wait_checkpoint_task(&self) -> ConnectorResult<Option<WaitCheckpointTask>> {
        Ok(match &self.config {
            ConnectorProperties::PostgresCdc(_) => Some(WaitCheckpointTask::CommitCdcOffset(None)),
            ConnectorProperties::GooglePubsub(prop) => Some(WaitCheckpointTask::AckPubsubMessage(
                prop.subscription_client().await?,
                vec![],
            )),
            ConnectorProperties::Nats(prop) => {
                match prop.nats_properties_consumer.get_ack_policy()? {
                    a @ AckPolicy::Explicit | a @ AckPolicy::All => {
                        Some(WaitCheckpointTask::AckNatsJetStream(
                            prop.common.build_context().await?,
                            vec![],
                            a,

View on GitHub (pinned to 6469eb736d)