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 CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -35,3 +35,4 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
instead of panicking the decoder.
- SlateDB compaction skips a stored batch it cannot inflate, with a warning,
instead of abandoning the whole maintenance pass.
- ListOffsets by timestamp now answers offset -1 and timestamp -1 with error NONE when no record has a timestamp at or after the target, as Apache Kafka does. The broker answered offset 0 before. That made the Java consumer's `offsetsForTimes()` throw `IllegalArgumentException: Invalid negative timestamp`, and it sent a client that seeks to the returned offset back to the start of the partition. A partition answered with an error code now also carries offset -1 instead of 0.
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.

142 changes: 136 additions & 6 deletions nisshi-broker/tests/it/list_offsets.rs
Original file line number Diff line number Diff line change
Expand Up @@ -442,9 +442,12 @@ pub async fn multiple_record(broker: Broker) -> Result<()> {
.is_some_and(|timestamp| timestamp > 0)
);

// Other partitions are empty: no record matches this timestamp, so this is
// a not-found result, not an empty-partition result. Offset is -1, matching
// Kafka's buildErrorResponse default.
for partition in partitions[1..].iter() {
assert_eq!(i16::from(ErrorCode::None), partition.error_code);
assert_eq!(Some(0), partition.offset);
assert_eq!(Some(-1), partition.offset);
assert_eq!(Some(-1), partition.timestamp);
}

Expand Down Expand Up @@ -492,9 +495,12 @@ pub async fn multiple_record(broker: Broker) -> Result<()> {
.is_some_and(|timestamp| timestamp > 0)
);

// Other partitions are empty: no record matches this timestamp, so this is
// a not-found result, not an empty-partition result. Offset is -1, matching
// Kafka's buildErrorResponse default.
for partition in partitions[1..].iter() {
assert_eq!(i16::from(ErrorCode::None), partition.error_code);
assert_eq!(Some(0), partition.offset);
assert_eq!(Some(-1), partition.offset);
assert_eq!(Some(-1), partition.timestamp);
}

Expand Down Expand Up @@ -542,9 +548,56 @@ pub async fn multiple_record(broker: Broker) -> Result<()> {
.is_some_and(|timestamp| timestamp > 0)
);

// Other partitions are empty: no record matches this timestamp, so this is
// a not-found result, not an empty-partition result. Offset is -1, matching
// Kafka's buildErrorResponse default.
for partition in partitions[1..].iter() {
assert_eq!(i16::from(ErrorCode::None), partition.error_code);
assert_eq!(Some(0), partition.offset);
assert_eq!(Some(-1), partition.offset);
assert_eq!(Some(-1), partition.timestamp);
}

debug!(phase = "after last record");
let timestamp = ListOffset::Timestamp(third + Duration::from_secs(1))
.try_into()
.inspect(|after| debug!(?after))?;

let response = broker
.serve(RequestInput {
request: ListOffsetsRequest::default()
.isolation_level(isolation)
.replica_id(replica_id)
.topics(Some(
[ListOffsetsTopic::default()
.name(topic_name.into())
.partitions(Some(
(0..num_partitions)
.map(|partition_index| {
ListOffsetsPartition::default()
.partition_index(partition_index)
.max_num_offsets(max_num_offsets)
.timestamp(timestamp)
.current_leader_epoch(Some(current_leader_epoch))
})
.collect::<Vec<_>>(),
))]
.into(),
)),
extensions: extensions.clone(),
})
.await?;

let topics = response.topics.as_deref().unwrap_or_default();
assert_eq!(1, topics.len());
assert_eq!(topic_name, topics[0].name);
let partitions = topics[0].partitions.as_deref().unwrap_or_default();
assert_eq!(num_partitions as usize, partitions.len());

// Partition 0 holds records, but none at or after the target: Kafka
// answers offset -1 and timestamp -1 with no error.
for partition in partitions {
assert_eq!(i16::from(ErrorCode::None), partition.error_code);
assert_eq!(Some(-1), partition.offset);
assert_eq!(Some(-1), partition.timestamp);
}

Expand Down Expand Up @@ -638,8 +691,11 @@ where

assert!(!items.is_empty());

// The topic has no records at all, so no record matches this timestamp:
// it's a not-found result at the raw storage layer (the -1 default is
// applied by the broker service, not the storage backend itself).
for (_toptition, response) in items {
assert_eq!(Some(0), response.offset);
assert_eq!(None, response.offset);
assert_eq!(None, response.timestamp);
}

Expand Down Expand Up @@ -695,7 +751,10 @@ where
.inspect(|offset| debug!(?offset))?
);

let after = SystemTime::now();
// A Timestamp lookup matches a record at or after the target, at
// millisecond precision. `after` is one second past now, so it is
// later than the record's timestamp.
let after = SystemTime::now() + Duration::from_secs(1);
debug!(after = to_timestamp(&after)?);

let offsets = [(topition.clone(), ListOffset::Latest)];
Expand Down Expand Up @@ -733,8 +792,12 @@ where
.list_offsets(IsolationLevel::ReadUncommitted, &offsets[..])
.await?;

// No record's timestamp is >= `after`: not-found at the raw storage layer
// (the -1 default is applied by the broker service, not the storage
// backend itself).
assert_eq!(1, responses.len());
assert_eq!(Some(0), responses[0].1.offset);
assert_eq!(None, responses[0].1.offset);
assert_eq!(None, responses[0].1.timestamp);

Ok(())
}
Expand Down Expand Up @@ -931,6 +994,73 @@ mod lite {
}
}

#[cfg(feature = "turso")]
mod turso {
use super::*;
use nisshi_storage::ArcDynStorage;

async fn storage_container(
cluster: impl Into<String> + Clone,
node: i32,
) -> Result<ArcDynStorage> {
common::storage_container(
StorageType::Turso,
cluster,
node,
Url::parse("tcp://127.0.0.1/")?,
None,
)
.await
}

#[ignore]
#[tokio::test]
async fn multiple_record() -> Result<()> {
let _guard = init_tracing()?;

let cluster_id = Uuid::now_v7();
let broker_id = rng().random_range(0..i32::MAX);

let sc = storage_container(cluster_id, broker_id).await?;
register_broker(cluster_id, broker_id, &sc).await?;

let broker = broker(sc)?;
super::multiple_record(broker).await
}

#[ignore]
#[tokio::test]
async fn new_topic() -> Result<()> {
let _guard = init_tracing()?;

let cluster_id = Uuid::now_v7();
let broker_id = rng().random_range(0..i32::MAX);

super::new_topic(
cluster_id,
broker_id,
storage_container(cluster_id, broker_id).await?,
)
.await
}

#[ignore]
#[tokio::test]
async fn single_record() -> Result<()> {
let _guard = init_tracing()?;

let cluster_id = Uuid::now_v7();
let broker_id = rng().random_range(0..i32::MAX);

super::single_record(
cluster_id,
broker_id,
storage_container(cluster_id, broker_id).await?,
)
.await
}
}

#[cfg(feature = "slatedb")]
mod slatedb {
use super::*;
Expand Down
145 changes: 145 additions & 0 deletions nisshi-broker/tests/it/storage_list_offsets.rs
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,8 @@
// See the License for the specific language governing permissions and
// limitations under the License.

use std::time::SystemTime;

use crate::common::{
alphanumeric_string, init_tracing, lite_storage, memory_storage, postgres_storage,
slate_storage,
Expand Down Expand Up @@ -107,6 +109,93 @@ async fn simple(storage: impl Storage + Clone, broker_id: i32) -> Result<()> {
Ok(())
}

/// A Timestamp lookup that matches no record (here: an empty topic) is a
/// not-found result, not an empty-partition result. At the wire level, Kafka's
/// `buildErrorResponse` (KafkaApis.scala) sends offset=-1/timestamp=-1 with
/// error_code=NONE for this case, and `ListOffsetsService` must apply that
/// default rather than the Earliest/Latest empty-partition default of 0.
/// Kafka sends leader epoch -1 here; nisshi sends 0, a known difference that
/// clients ignore when the offset is -1.
async fn timestamp_no_match(storage: impl Storage + Clone, broker_id: i32) -> Result<()> {
let extensions = Extensions::default();

let create_topic = CreateTopicsService {
storage: storage.clone(),
};

let topic = &alphanumeric_string(15)[..];

let num_partitions = 1;
let replication_factor = 0;

{
let response = create_topic
.serve(RequestInput {
request: CreateTopicsRequest::default()
.validate_only(Some(false))
.topics(Some(
[CreatableTopic::default()
.name(topic.into())
.num_partitions(num_partitions)
.replication_factor(replication_factor)
.assignments(Some([].into()))
.configs(Some([].into()))]
.into(),
)),
extensions: extensions.clone(),
})
.await?;

let topics = response.topics.as_deref().unwrap_or_default();
assert_eq!(1, topics.len());
assert_eq!(ErrorCode::None, ErrorCode::try_from(topics[0].error_code)?);
}

let service = ListOffsetsService {
storage: storage.clone(),
};

let response = service
.serve(RequestInput {
request: ListOffsetsRequest::default()
.isolation_level(Some(IsolationLevel::ReadUncommitted.into()))
.replica_id(broker_id)
.topics(Some(
[ListOffsetsTopic::default()
.name(topic.into())
.partitions(Some(
[ListOffsetsPartition::default()
.current_leader_epoch(Some(-1))
.max_num_offsets(Some(3))
.partition_index(0)
.timestamp(ListOffset::Timestamp(SystemTime::now()).try_into()?)]
.into(),
))]
.into(),
)),
extensions: extensions.clone(),
})
.await?;

let topics = response.topics.as_deref().unwrap_or_default();
assert_eq!(1, topics.len());
assert_eq!(topic, topics[0].name);

let partitions = topics[0].partitions.as_deref().unwrap_or_default();
assert_eq!(1, partitions.len());
assert_eq!(0, partitions[0].partition_index);
assert!(partitions[0].old_style_offsets.is_none());
assert_eq!(
ErrorCode::None,
ErrorCode::try_from(partitions[0].error_code)?
);
assert_eq!(Some(-1), partitions[0].timestamp);
assert_eq!(Some(-1), partitions[0].offset);
assert_eq!(Some(0), partitions[0].leader_epoch);
Comment thread
solace-aross marked this conversation as resolved.

Ok(())
}

#[cfg(feature = "dynostore")]
mod in_memory {
use super::*;
Expand All @@ -131,6 +220,20 @@ mod in_memory {

Ok(())
}

#[tokio::test]
async fn timestamp_no_match() -> Result<()> {
let _guard = init_tracing()?;

let cluster_id = Uuid::now_v7();
let broker_id = rng().random_range(0..i32::MAX);

let storage = storage_container(cluster_id, broker_id).await?;

super::timestamp_no_match(storage, broker_id).await?;

Ok(())
}
}

#[cfg(feature = "libsql")]
Expand All @@ -157,6 +260,20 @@ mod lite {

Ok(())
}

#[tokio::test]
async fn timestamp_no_match() -> Result<()> {
let _guard = init_tracing()?;

let cluster_id = Uuid::now_v7();
let broker_id = rng().random_range(0..i32::MAX);

let storage = storage_container(cluster_id, broker_id).await?;

super::timestamp_no_match(storage, broker_id).await?;

Ok(())
}
}

#[cfg(feature = "slatedb")]
Expand All @@ -183,6 +300,20 @@ mod slatedb {

Ok(())
}

#[tokio::test]
async fn timestamp_no_match() -> Result<()> {
let _guard = init_tracing()?;

let cluster_id = Uuid::now_v7();
let broker_id = rng().random_range(0..i32::MAX);

let storage = storage_container(cluster_id, broker_id).await?;

super::timestamp_no_match(storage, broker_id).await?;

Ok(())
}
}

#[cfg(feature = "postgres")]
Expand All @@ -209,4 +340,18 @@ mod pg {

Ok(())
}

#[tokio::test]
async fn timestamp_no_match() -> Result<()> {
let _guard = init_tracing()?;

let cluster_id = Uuid::now_v7();
let broker_id = rng().random_range(0..i32::MAX);

let storage = storage_container(cluster_id, broker_id).await?;

super::timestamp_no_match(storage, broker_id).await?;

Ok(())
}
}
Loading
Loading