@@ -136,14 +136,37 @@ fn create_reader(
136136 DummyInputReceiver ,
137137 Box < dyn InputReader > ,
138138) {
139- let config = serde_json:: from_value ( json ! ( {
140- "name" : "kafka_input" ,
141- "config" : {
139+ create_reader_config ( topic, resume_info, synchronize_partitions, None )
140+ }
141+
142+ /// Like [`create_reader`], but with an optional `header_filter` added to the
143+ /// connector configuration.
144+ fn create_reader_config (
145+ topic : & str ,
146+ resume_info : Option < JsonValue > ,
147+ synchronize_partitions : bool ,
148+ header_filter : Option < JsonValue > ,
149+ ) -> (
150+ Box < dyn TransportInputEndpoint > ,
151+ DummyInputReceiver ,
152+ Box < dyn InputReader > ,
153+ ) {
154+ let mut inner_config = json ! ( {
142155 "topic" : topic,
143156 "log_level" : "debug" ,
144- "start_from" : "earliest" ,
145- "synchronize_partitions" : synchronize_partitions,
146- }
157+ "start_from" : "earliest" ,
158+ "synchronize_partitions" : synchronize_partitions,
159+ } ) ;
160+ if let Some ( header_filter) = header_filter {
161+ inner_config
162+ . as_object_mut ( )
163+ . unwrap ( )
164+ . insert ( "header_filter" . to_string ( ) , header_filter) ;
165+ }
166+
167+ let config = serde_json:: from_value ( json ! ( {
168+ "name" : "kafka_input" ,
169+ "config" : inner_config,
147170 } ) )
148171 . unwrap ( ) ;
149172
@@ -165,6 +188,23 @@ fn create_reader(
165188 ( endpoint, receiver, reader)
166189}
167190
191+ /// Drain and return the payloads flushed to `receiver` so far, in order.
192+ fn take_flushed ( receiver : & DummyInputReceiver ) -> Vec < String > {
193+ mem:: take ( & mut * receiver. inner . flushed . lock ( ) . unwrap ( ) )
194+ }
195+
196+ /// Build the `headers` argument for [`TestProducer::send_message`] from
197+ /// `(key, value)` pairs, where a `None` value denotes a header with a null
198+ /// value.
199+ fn headers ( pairs : & [ ( & str , Option < & [ u8 ] > ) ] ) -> Option < BTreeMap < String , Option < Vec < u8 > > > > {
200+ Some (
201+ pairs
202+ . iter ( )
203+ . map ( |( key, value) | ( key. to_string ( ) , value. map ( |v| v. to_vec ( ) ) ) )
204+ . collect ( ) ,
205+ )
206+ }
207+
168208#[ test]
169209fn single_input ( ) {
170210 test_input ( "single_input_ft" , & [ 10 ] ) ;
@@ -469,6 +509,139 @@ fn test_input(topic: &str, batch_sizes: &[u32]) {
469509 }
470510}
471511
512+ /// A configured header filter drops non-matching messages before parsing, and
513+ /// the drop decision is deterministic across replay: a fresh reader replaying
514+ /// the checkpointed offset range reproduces exactly the admitted records.
515+ ///
516+ /// To confirm this test catches a regression, make `header_filter_admits`
517+ /// always return `true`; the connector then buffers all 6 messages instead of
518+ /// 3 and both the live and the replay expectations fail.
519+ #[ test]
520+ fn test_input_header_filter ( ) {
521+ init_test_logger ( ) ;
522+
523+ let topic = "kafka_header_filter_ft" ;
524+ let _kafka_resources = KafkaResources :: create_topics ( & [ ( topic, 1 ) ] ) ;
525+
526+ // Admit only messages whose `keep` header equals `yes`.
527+ let filter = json ! ( { "header" : { "name" : "keep" , "pattern" : "yes" } } ) ;
528+
529+ let ( endpoint, receiver, reader) =
530+ create_reader_config ( topic, None , false , Some ( filter. clone ( ) ) ) ;
531+ reader. extend ( ) ;
532+
533+ // Offsets 0..6. `keep=yes` at 0, 2, 4 are admitted; 1, 3 (`keep=no`) and 5
534+ // (no header) are dropped. The trailing drop at offset 5 lies beyond the
535+ // recorded range and must never be replayed.
536+ let producer = TestProducer :: new ( ) ;
537+ producer. send_message ( b"m0" , topic, headers ( & [ ( "keep" , Some ( b"yes" ) ) ] ) ) ;
538+ producer. send_message ( b"m1" , topic, headers ( & [ ( "keep" , Some ( b"no" ) ) ] ) ) ;
539+ producer. send_message ( b"m2" , topic, headers ( & [ ( "keep" , Some ( b"yes" ) ) ] ) ) ;
540+ producer. send_message ( b"m3" , topic, headers ( & [ ( "keep" , Some ( b"no" ) ) ] ) ) ;
541+ producer. send_message ( b"m4" , topic, headers ( & [ ( "keep" , Some ( b"yes" ) ) ] ) ) ;
542+ producer. send_message ( b"m5" , topic, None ) ;
543+
544+ // Only the 3 admitted messages are buffered.
545+ receiver. expect_buffering ( 3 ) ;
546+ reader. queue ( false ) ;
547+
548+ // The offset range spans the first through the last admitted offset (0..5);
549+ // dropped offsets inside and beyond the range do not change it.
550+ let metadata = Metadata {
551+ offsets : vec ! [ 0 ..5 ] ,
552+ } ;
553+ receiver. expect ( vec ! [ ConsumerCall :: Extended {
554+ num_records: 3 ,
555+ metadata: serde_json:: to_value( & metadata) . unwrap( ) ,
556+ } ] ) ;
557+ assert_eq ! ( take_flushed( & receiver) , vec![ "m0" , "m2" , "m4" ] ) ;
558+
559+ drop ( endpoint) ;
560+ drop ( receiver) ;
561+ drop ( reader) ;
562+
563+ // Replay the recorded range with the same filter. The filter is re-applied
564+ // during replay, so exactly the admitted records reappear, in order.
565+ let ( _endpoint, receiver, reader) = create_reader_config ( topic, None , false , Some ( filter) ) ;
566+ receiver. inner . drop_buffered . store ( true , Ordering :: Release ) ;
567+ reader. replay ( serde_json:: to_value ( & metadata) . unwrap ( ) , RmpValue :: Nil ) ;
568+ receiver. expect ( vec ! [ ConsumerCall :: Replayed { num_records: 3 } ] ) ;
569+ assert_eq ! ( take_flushed( & receiver) , vec![ "m0" , "m2" , "m4" ] ) ;
570+ }
571+
572+ /// A boolean header filter (`and`/`or`/`not`) drops the right messages in the
573+ /// real connector, exercising leading, middle, and trailing drops, and its
574+ /// decisions are re-applied identically on replay. The filter reads headers
575+ /// directly, so it works even though `include_headers` is not set.
576+ #[ test]
577+ fn test_input_header_filter_boolean ( ) {
578+ init_test_logger ( ) ;
579+
580+ let topic = "kafka_header_filter_boolean_ft" ;
581+ let _kafka_resources = KafkaResources :: create_topics ( & [ ( topic, 1 ) ] ) ;
582+
583+ // Admit iff (env is prod or staging) and not (skip == drop).
584+ let filter = json ! ( {
585+ "and" : [
586+ { "or" : [
587+ { "header" : { "name" : "env" , "pattern" : "prod" } } ,
588+ { "header" : { "name" : "env" , "pattern" : "staging" } }
589+ ] } ,
590+ { "not" : { "header" : { "name" : "skip" , "pattern" : "drop" } } }
591+ ]
592+ } ) ;
593+
594+ let ( endpoint, receiver, reader) =
595+ create_reader_config ( topic, None , false , Some ( filter. clone ( ) ) ) ;
596+ reader. extend ( ) ;
597+
598+ let producer = TestProducer :: new ( ) ;
599+ // 0: env=dev -> drop (fails `or`); leading drop
600+ // 1: env=prod -> admit
601+ // 2: env=staging, skip=drop -> drop (fails `not`); middle drop
602+ // 3: env=staging -> admit
603+ // 4: env=prod, skip=drop -> drop (fails `not`); trailing drop
604+ producer. send_message ( b"m0" , topic, headers ( & [ ( "env" , Some ( b"dev" ) ) ] ) ) ;
605+ producer. send_message ( b"m1" , topic, headers ( & [ ( "env" , Some ( b"prod" ) ) ] ) ) ;
606+ producer. send_message (
607+ b"m2" ,
608+ topic,
609+ headers ( & [ ( "env" , Some ( b"staging" ) ) , ( "skip" , Some ( b"drop" ) ) ] ) ,
610+ ) ;
611+ producer. send_message ( b"m3" , topic, headers ( & [ ( "env" , Some ( b"staging" ) ) ] ) ) ;
612+ producer. send_message (
613+ b"m4" ,
614+ topic,
615+ headers ( & [ ( "env" , Some ( b"prod" ) ) , ( "skip" , Some ( b"drop" ) ) ] ) ,
616+ ) ;
617+
618+ // Admitted: offsets 1 and 3.
619+ receiver. expect_buffering ( 2 ) ;
620+ reader. queue ( false ) ;
621+
622+ // The range starts at the first admitted offset (1) and ends after the last
623+ // (3), skipping the leading drop at offset 0.
624+ let metadata = Metadata {
625+ offsets : vec ! [ 1 ..4 ] ,
626+ } ;
627+ receiver. expect ( vec ! [ ConsumerCall :: Extended {
628+ num_records: 2 ,
629+ metadata: serde_json:: to_value( & metadata) . unwrap( ) ,
630+ } ] ) ;
631+ assert_eq ! ( take_flushed( & receiver) , vec![ "m1" , "m3" ] ) ;
632+
633+ drop ( endpoint) ;
634+ drop ( receiver) ;
635+ drop ( reader) ;
636+
637+ // Replaying the recorded range reproduces the admitted records.
638+ let ( _endpoint, receiver, reader) = create_reader_config ( topic, None , false , Some ( filter) ) ;
639+ receiver. inner . drop_buffered . store ( true , Ordering :: Release ) ;
640+ reader. replay ( serde_json:: to_value ( & metadata) . unwrap ( ) , RmpValue :: Nil ) ;
641+ receiver. expect ( vec ! [ ConsumerCall :: Replayed { num_records: 2 } ] ) ;
642+ assert_eq ! ( take_flushed( & receiver) , vec![ "m1" , "m3" ] ) ;
643+ }
644+
472645#[ derive( Debug , PartialEq ) ]
473646enum ConsumerCall {
474647 ParseErrors ,
0 commit comments