{"record":{"id":"9438ac13e99fc377","repo":"crowdsecurity/crowdsec","slug":"cannot-register-stream-consumer-w","errorCode":null,"errorMessage":"cannot register stream consumer: %w","messagePattern":"cannot register stream consumer: %w","errorType":"exception","errorClass":null,"httpStatus":null,"severity":"error","filePath":"pkg/acquisition/modules/kinesis/run.go","lineNumber":147,"sourceCode":"\t\t\treturn nil\n\t\t}\n\n\t\ttime.Sleep(time.Millisecond * 200 * time.Duration(i+1))\n\t\ts.logger.Debugf(\"Waiting for consumer registration %d\", i)\n\t}\n\n\treturn fmt.Errorf(\"consumer %s is not active after %d tries\", consumerARN, maxTries)\n}\n\nfunc (s *Source) RegisterConsumer(ctx context.Context) (*kinesis.RegisterStreamConsumerOutput, error) {\n\ts.logger.Debugf(\"Registering consumer %s\", s.Config.ConsumerName)\n\n\tstreamConsumer, err := s.kClient.RegisterStreamConsumer(ctx, &kinesis.RegisterStreamConsumerInput{\n\t\t\tConsumerName: aws.String(s.Config.ConsumerName),\n\t\t\tStreamARN:    aws.String(s.Config.StreamARN),\n\t\t})\n\tif err != nil {\n\t\treturn nil, fmt.Errorf(\"cannot register stream consumer: %w\", err)\n\t}\n\n\terr = s.WaitForConsumerRegistration(ctx, *streamConsumer.Consumer.ConsumerARN)\n\tif err != nil {\n\t\treturn nil, fmt.Errorf(\"timeout while waiting for consumer to be active: %w\", err)\n\t}\n\n\treturn streamConsumer, nil\n}\n\nfunc (s *Source) ParseAndPushRecords(records []kinTypes.Record, out chan pipeline.Event, logger *log.Entry, shardID string) {\n\tfor _, record := range records {\n\t\tif s.Config.StreamARN != \"\" {\n\t\t\tif s.metricsLevel != metrics.AcquisitionMetricsLevelNone {\n\t\t\t\tmetrics.KinesisDataSourceLinesReadShards.With(prometheus.Labels{\"stream\": s.Config.StreamARN, \"shard\": shardID}).Inc()\n\t\t\t\tmetrics.KinesisDataSourceLinesRead.With(prometheus.Labels{\"stream\": s.Config.StreamARN, \"datasource_type\": ModuleName, \"acquis_type\": s.Config.Labels[\"type\"]}).Inc()\n\t\t\t}\n\t\t} else {","sourceCodeStart":129,"sourceCodeEnd":165,"githubUrl":"https://github.com/crowdsecurity/crowdsec/blob/909b5157986a2b2c2163300fdaef5ed01289f7d2/pkg/acquisition/modules/kinesis/run.go#L129-L165","documentation":"RegisterConsumer calls kinesis RegisterStreamConsumer for an EFO (enhanced fan-out) consumer; the AWS API call failed (permissions, invalid ARN, stream state, throttling), so the consumer could not be registered and EnhancedRead cannot proceed.","triggerScenarios":"RegisterStreamConsumer returns ResourceNotFoundException (stream doesn't exist), LimitExceededException (20 consumers per stream max), InvalidArgumentException (bad consumer name/ARN), AccessDeniedException, or a throttling/network error.","commonSituations":"Stream deleted or ARN typo; already 20 EFO consumers registered on the stream; IAM policy missing kinesis:RegisterStreamConsumer; running against localstack without EFO support.","solutions":["Verify stream ARN, region and consumer name in the DSN","Check IAM permissions for kinesis:RegisterStreamConsumer","Ensure the consumer name is valid (≤128 chars, alphanumerics plus _.=+@-)"],"exampleFix":"// IAM policy addition\n{\"Effect\": \"Allow\", \"Action\": [\"kinesis:RegisterStreamConsumer\", \"kinesis:DescribeStreamConsumer\", \"kinesis:DeregisterStreamConsumer\"], \"Resource\": \"*\"}","handlingStrategy":"try-catch","validationCode":"// Pre-flight checks before RegisterStreamConsumer\n_, err := client.DescribeStreamSummary(ctx, &kinesis.DescribeStreamSummaryInput{StreamName: aws.String(name)})\nif err != nil { return fmt.Errorf(\"stream missing or undescribable: %w\", err) }\nvar re = regexp.MustCompile(`^[a-zA-Z0-9_.-]{1,128}$`)\nif !re.MatchString(consumerName) { return errors.New(\"invalid consumer_name\") }","typeGuard":"var nf *kinTypes.ResourceNotFoundException\nif errors.As(err, &nf) { /* stream or consumer ARN wrong — fix config */ }\nvar lim *kinTypes.LimitExceededException\nif errors.As(err, &lim) { /* 20-consumer cap: deregister stale ones */ }","tryCatchPattern":"out, err := client.RegisterStreamConsumer(ctx, in)\nif err != nil {\n    var exists *kinTypes.ResourceInUseException\n    if errors.As(err, &exists) { /* consumer already registered — proceed to describe */ }\n    return nil, fmt.Errorf(\"cannot register stream consumer: %w\", err)\n}","preventionTips":["Verify stream_arn with a describe call at config validation time.","Keep EFO consumer count below the 20-per-stream quota; schedule deregistration of stale consumers.","Add kinesis:RegisterStreamConsumer to the IAM role.","When testing against localstack, confirm EFO endpoints are supported by that version."],"tags":["aws","kinesis","efo","permissions"],"backgroundTag":"resource-not-found","analyzedSha":"909b5157986a2b2c2163300fdaef5ed01289f7d2","analyzedAt":"2026-09-06T12:27:26.012Z","contentChangedAt":"2026-09-06T12:27:26.012Z","schemaVersion":2},"datasetVersion":"2026-09-14T00:17:10.932Z"}