Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

56 changes: 55 additions & 1 deletion crates/adapters/src/transport/kafka/ft/input.rs
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,9 @@ use feldera_adapterlib::transport::{
use feldera_sqllib::{ByteArray, SqlString, Timestamp, Variant};
use feldera_types::config::FtModel;
use feldera_types::program_schema::Relation;
use feldera_types::transport::kafka::{KafkaInputConfig, KafkaStartFromConfig};
use feldera_types::transport::kafka::{
CompiledHeaderFilter, KafkaInputConfig, KafkaStartFromConfig,
};
use itertools::Itertools;
use rdkafka::client::OAuthToken;
use rdkafka::config::RDKafkaLogLevel;
Expand Down Expand Up @@ -66,6 +68,10 @@ const METADATA_TIMEOUT: Duration = Duration::from_secs(10);
// to the worker thread.
const ERROR_BUFFER_SIZE: usize = 1000;

/// A Kafka message's headers borrowed as `(key, value)` pairs for header
/// filtering. A `None` value denotes a header present with a null value.
type HeaderPairs<'a> = SmallVec<[(&'a str, Option<&'a [u8]>); 8]>;

pub struct KafkaFtInputEndpoint {
config: Arc<KafkaInputConfig>,
}
Expand Down Expand Up @@ -380,6 +386,16 @@ impl KafkaFtInputReaderInner {
})
.collect::<Vec<_>>();

// Compile the header filter once and share it across partition
// receivers. A malformed filter is already rejected by
// `KafkaInputConfig::validate`, so this normally cannot fail.
let header_filter = config
.header_filter
.as_ref()
.map(|filter| filter.compile())
.transpose()?
.map(Arc::new);

// Split every partition away as its own separate queue.
let mut receivers = BTreeMap::new();

Expand All @@ -405,6 +421,7 @@ impl KafkaFtInputReaderInner {
queue,
next_offset,
&config,
header_filter.clone(),
unparker,
));
receivers.insert(partition, receiver.clone());
Expand Down Expand Up @@ -999,6 +1016,10 @@ struct PartitionReceiver {
config: KafkaInputConfig,
metadata_requested: bool,

/// Compiled header filter, shared across partitions. Messages that do not
/// match are dropped before parsing. `None` admits all messages.
header_filter: Option<Arc<CompiledHeaderFilter>>,

/// The maximum message offset that we want to receive, used as follows:
///
/// - `i64::MIN`, the initial value, disables receiving messages entirely.
Expand Down Expand Up @@ -1041,6 +1062,7 @@ impl PartitionReceiver {
queue: PartitionQueue<KafkaFtInputContext>,
next_offset: i64,
config: &KafkaInputConfig,
header_filter: Option<Arc<CompiledHeaderFilter>>,
unparker: Unparker,
) -> Self {
let metadata_requested = config.metadata_requested();
Expand All @@ -1057,6 +1079,7 @@ impl PartitionReceiver {
fatal_error: AtomicBool::new(false),
config: config.clone(),
metadata_requested,
header_filter,
unparker,
}
}
Expand Down Expand Up @@ -1112,6 +1135,27 @@ impl PartitionReceiver {
self.next_offset.load(Ordering::Relaxed)
}

/// Returns `true` if `message` satisfies the configured header filter, or if
/// no filter is configured. Messages for which this returns `false` are
/// dropped before parsing.
fn header_filter_admits(&self, message: &BorrowedMessage<'_>) -> bool {
let Some(filter) = &self.header_filter else {
return true;
};

// Borrow the message's headers as (key, value) pairs; nothing is copied.
let headers: HeaderPairs = match message.headers() {
Some(headers) => (0..headers.count())
.map(|i| {
let header = headers.get(i);
(header.key, header.value)
})
.collect(),
None => SmallVec::new(),
};
filter.matches(&headers)
}

/// Create record metadata from Kafka message containing only properties specified in the connector config.
fn create_metadata(&self, message: &BorrowedMessage<'_>) -> Option<ConnectorMetadata> {
if !self.metadata_requested {
Expand Down Expand Up @@ -1189,6 +1233,16 @@ impl PartitionReceiver {
let next_offset = self.next_offset();
if offset >= next_offset {
self.next_offset.store(offset + 1, Ordering::Relaxed);

// Drop messages rejected by the header filter before parsing,
// so they never enter the offset ranges or the step hash.
// This runs after advancing `next_offset` (the consume
// position still moves forward) and is re-applied identically
// on replay, keeping ranges and hashes deterministic.
if !self.header_filter_admits(&message) {
return;
}

let timestamp = message.timestamp().to_millis().unwrap_or(i64::MIN);
let payload = message.payload().unwrap_or(&[]);
let metadata = self.create_metadata(&message);
Expand Down
185 changes: 179 additions & 6 deletions crates/adapters/src/transport/kafka/ft/test.rs
Original file line number Diff line number Diff line change
Expand Up @@ -136,14 +136,37 @@ fn create_reader(
DummyInputReceiver,
Box<dyn InputReader>,
) {
let config = serde_json::from_value(json!({
"name": "kafka_input",
"config": {
create_reader_config(topic, resume_info, synchronize_partitions, None)
}

/// Like [`create_reader`], but with an optional `header_filter` added to the
/// connector configuration.
fn create_reader_config(
topic: &str,
resume_info: Option<JsonValue>,
synchronize_partitions: bool,
header_filter: Option<JsonValue>,
) -> (
Box<dyn TransportInputEndpoint>,
DummyInputReceiver,
Box<dyn InputReader>,
) {
let mut inner_config = json!({
"topic": topic,
"log_level": "debug",
"start_from": "earliest",
"synchronize_partitions": synchronize_partitions,
}
"start_from": "earliest",
"synchronize_partitions": synchronize_partitions,
});
if let Some(header_filter) = header_filter {
inner_config
.as_object_mut()
.unwrap()
.insert("header_filter".to_string(), header_filter);
}

let config = serde_json::from_value(json!({
"name": "kafka_input",
"config": inner_config,
}))
.unwrap();

Expand All @@ -165,6 +188,23 @@ fn create_reader(
(endpoint, receiver, reader)
}

/// Drain and return the payloads flushed to `receiver` so far, in order.
fn take_flushed(receiver: &DummyInputReceiver) -> Vec<String> {
mem::take(&mut *receiver.inner.flushed.lock().unwrap())
}

/// Build the `headers` argument for [`TestProducer::send_message`] from
/// `(key, value)` pairs, where a `None` value denotes a header with a null
/// value.
fn headers(pairs: &[(&str, Option<&[u8]>)]) -> Option<BTreeMap<String, Option<Vec<u8>>>> {
Some(
pairs
.iter()
.map(|(key, value)| (key.to_string(), value.map(|v| v.to_vec())))
.collect(),
)
}

#[test]
fn single_input() {
test_input("single_input_ft", &[10]);
Expand Down Expand Up @@ -469,6 +509,139 @@ fn test_input(topic: &str, batch_sizes: &[u32]) {
}
}

/// A configured header filter drops non-matching messages before parsing, and
/// the drop decision is deterministic across replay: a fresh reader replaying
/// the checkpointed offset range reproduces exactly the admitted records.
///
/// To confirm this test catches a regression, make `header_filter_admits`
/// always return `true`; the connector then buffers all 6 messages instead of
/// 3 and both the live and the replay expectations fail.
#[test]
fn test_input_header_filter() {
init_test_logger();

let topic = "kafka_header_filter_ft";
let _kafka_resources = KafkaResources::create_topics(&[(topic, 1)]);

// Admit only messages whose `keep` header equals `yes`.
let filter = json!({"header": {"name": "keep", "pattern": "yes"}});

let (endpoint, receiver, reader) =
create_reader_config(topic, None, false, Some(filter.clone()));
reader.extend();

// Offsets 0..6. `keep=yes` at 0, 2, 4 are admitted; 1, 3 (`keep=no`) and 5
// (no header) are dropped. The trailing drop at offset 5 lies beyond the
// recorded range and must never be replayed.
let producer = TestProducer::new();
producer.send_message(b"m0", topic, headers(&[("keep", Some(b"yes"))]));
producer.send_message(b"m1", topic, headers(&[("keep", Some(b"no"))]));
producer.send_message(b"m2", topic, headers(&[("keep", Some(b"yes"))]));
producer.send_message(b"m3", topic, headers(&[("keep", Some(b"no"))]));
producer.send_message(b"m4", topic, headers(&[("keep", Some(b"yes"))]));
producer.send_message(b"m5", topic, None);

// Only the 3 admitted messages are buffered.
receiver.expect_buffering(3);
reader.queue(false);

// The offset range spans the first through the last admitted offset (0..5);
// dropped offsets inside and beyond the range do not change it.
let metadata = Metadata {
offsets: vec![0..5],
};
receiver.expect(vec![ConsumerCall::Extended {
num_records: 3,
metadata: serde_json::to_value(&metadata).unwrap(),
}]);
assert_eq!(take_flushed(&receiver), vec!["m0", "m2", "m4"]);

drop(endpoint);
drop(receiver);
drop(reader);

// Replay the recorded range with the same filter. The filter is re-applied
// during replay, so exactly the admitted records reappear, in order.
let (_endpoint, receiver, reader) = create_reader_config(topic, None, false, Some(filter));
receiver.inner.drop_buffered.store(true, Ordering::Release);
reader.replay(serde_json::to_value(&metadata).unwrap(), RmpValue::Nil);
receiver.expect(vec![ConsumerCall::Replayed { num_records: 3 }]);
assert_eq!(take_flushed(&receiver), vec!["m0", "m2", "m4"]);
}

/// A boolean header filter (`and`/`or`/`not`) drops the right messages in the
/// real connector, exercising leading, middle, and trailing drops, and its
/// decisions are re-applied identically on replay. The filter reads headers
/// directly, so it works even though `include_headers` is not set.
#[test]
fn test_input_header_filter_boolean() {
init_test_logger();

let topic = "kafka_header_filter_boolean_ft";
let _kafka_resources = KafkaResources::create_topics(&[(topic, 1)]);

// Admit iff (env is prod or staging) and not (skip == drop).
let filter = json!({
"and": [

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

We should standardize on the language for these expressions, in the push-down PR we are talking about JsonLogic https://crates.io/crates/jsonlogic-rs

But JsonLogic does not have regexes, these are very difficult across languages... sigh

{"or": [
{"header": {"name": "env", "pattern": "prod"}},
{"header": {"name": "env", "pattern": "staging"}}
]},
{"not": {"header": {"name": "skip", "pattern": "drop"}}}
]
});

let (endpoint, receiver, reader) =
create_reader_config(topic, None, false, Some(filter.clone()));
reader.extend();

let producer = TestProducer::new();
// 0: env=dev -> drop (fails `or`); leading drop
// 1: env=prod -> admit
// 2: env=staging, skip=drop -> drop (fails `not`); middle drop
// 3: env=staging -> admit
// 4: env=prod, skip=drop -> drop (fails `not`); trailing drop
producer.send_message(b"m0", topic, headers(&[("env", Some(b"dev"))]));
producer.send_message(b"m1", topic, headers(&[("env", Some(b"prod"))]));
producer.send_message(
b"m2",
topic,
headers(&[("env", Some(b"staging")), ("skip", Some(b"drop"))]),
);
producer.send_message(b"m3", topic, headers(&[("env", Some(b"staging"))]));
producer.send_message(
b"m4",
topic,
headers(&[("env", Some(b"prod")), ("skip", Some(b"drop"))]),
);

// Admitted: offsets 1 and 3.
receiver.expect_buffering(2);
reader.queue(false);

// The range starts at the first admitted offset (1) and ends after the last
// (3), skipping the leading drop at offset 0.
let metadata = Metadata {
offsets: vec![1..4],
};
receiver.expect(vec![ConsumerCall::Extended {
num_records: 2,
metadata: serde_json::to_value(&metadata).unwrap(),
}]);
assert_eq!(take_flushed(&receiver), vec!["m1", "m3"]);

drop(endpoint);
drop(receiver);
drop(reader);

// Replaying the recorded range reproduces the admitted records.
let (_endpoint, receiver, reader) = create_reader_config(topic, None, false, Some(filter));
receiver.inner.drop_buffered.store(true, Ordering::Release);
reader.replay(serde_json::to_value(&metadata).unwrap(), RmpValue::Nil);
receiver.expect(vec![ConsumerCall::Replayed { num_records: 2 }]);
assert_eq!(take_flushed(&receiver), vec!["m1", "m3"]);
}

#[derive(Debug, PartialEq)]
enum ConsumerCall {
ParseErrors,
Expand Down
1 change: 1 addition & 0 deletions crates/feldera-types/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -40,3 +40,4 @@ clap = { workspace = true }
[dev-dependencies]
csv = { workspace = true }
tempfile = { workspace = true }
serde_yaml = { workspace = true }

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

yech

Loading