Skip to content
Open
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: 0 additions & 1 deletion bus/core/src/channeled_agentbus.rs
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,6 @@ use crate::AppendRequest;
use crate::AppendResponse;
use crate::BlockingPollRequest;
use crate::BlockingPollResponse;
use crate::BusId;
use crate::BusResult;
use crate::CheckTailRequest;
use crate::CheckTailResponse;
Expand Down
110 changes: 107 additions & 3 deletions bus/impls/simple/src/in_memory_agentbus_state.rs
Original file line number Diff line number Diff line change
Expand Up @@ -158,6 +158,12 @@ impl InMemoryAgentBusState {
payload_types: &Option<Vec<i32>>,
end_position: i64,
) -> Result<(Vec<BusEntry>, i64)> {
if start_position < 0 {
return Err(anyhow::anyhow!(
"start_position {} must not be negative",
start_position
));
}
if start_position > end_position {
return Err(anyhow::anyhow!(
"start_position {} is beyond end_position {}",
Expand All @@ -175,6 +181,10 @@ impl InMemoryAgentBusState {
};

let start_idx = start_position as usize;
if start_idx >= bus.len() {
return Ok((vec![], end_position));
}

let end_idx = (end_position as usize).min(bus.len());
let slice = &bus[start_idx..end_idx];

Expand All @@ -193,13 +203,107 @@ impl InMemoryAgentBusState {
let next = entry
.header
.as_ref()
.expect("entry should have header")
.log_position
+ 1;
.map(|h| h.log_position + 1)
.unwrap_or(end_position);
return Ok((entries, next));
}
}

Ok((entries, end_position))
}
}

#[cfg(test)]
mod tests {
use super::*;
use agent_bus_proto_rust::agent_bus::intention::Intention as IntentionEnum;
use agent_bus_proto_rust::agent_bus::payload::Payload as PayloadEnum;

fn make_test_entry(position: i64, content: &str) -> BusEntry {
BusEntry {
header: Some(Header {
log_position: position,
rt_timestamp_ms: 1000 + position,
}),
payload: Some(Payload {
payload: Some(PayloadEnum::Intention(Intention {
intention: Some(IntentionEnum::StringIntention(content.to_string())),
..Default::default()
})),
}),
}
}

#[test]
fn test_read_filtered_entries_beyond_tail() {
let mut state = InMemoryAgentBusState::new();
state.buses.insert(
"test-bus".to_string(),
vec![make_test_entry(0, "e0"), make_test_entry(1, "e1")],
);

// Reading beyond tail should not panic and return empty vec with end_position
let res = state.read_filtered_entries("test-bus", 5, 10, &None, 10);
assert!(res.is_ok());
let (entries, next) = res.unwrap();
assert!(entries.is_empty());
assert_eq!(next, 10);

// Reading exactly at tail
let res = state.read_filtered_entries("test-bus", 2, 10, &None, 10);
assert!(res.is_ok());
let (entries, next) = res.unwrap();
assert!(entries.is_empty());
assert_eq!(next, 10);
}

#[test]
fn test_read_filtered_entries_empty_bus() {
let mut state = InMemoryAgentBusState::new();
state.buses.insert("empty-bus".to_string(), vec![]);

let res = state.read_filtered_entries("empty-bus", 0, 10, &None, 5);
assert!(res.is_ok());
let (entries, next) = res.unwrap();
assert!(entries.is_empty());
assert_eq!(next, 5);

let res2 = state.read_filtered_entries("empty-bus", 2, 10, &None, 5);
assert!(res2.is_ok());
let (entries2, next2) = res2.unwrap();
assert!(entries2.is_empty());
assert_eq!(next2, 5);
}

#[test]
fn test_read_filtered_entries_negative_start() {
let state = InMemoryAgentBusState::new();
let res = state.read_filtered_entries("test-bus", -1, 10, &None, 5);
assert!(res.is_err());
}

#[test]
fn test_read_filtered_entries_normal_pagination() {
let mut state = InMemoryAgentBusState::new();
state.buses.insert(
"test-bus".to_string(),
vec![
make_test_entry(0, "e0"),
make_test_entry(1, "e1"),
make_test_entry(2, "e2"),
],
);

let (entries, next) = state
.read_filtered_entries("test-bus", 0, 2, &None, 3)
.unwrap();
assert_eq!(entries.len(), 2);
assert_eq!(next, 2);

let (entries, next) = state
.read_filtered_entries("test-bus", next, 2, &None, 3)
.unwrap();
assert_eq!(entries.len(), 1);
assert_eq!(next, 3);
}
}
63 changes: 63 additions & 0 deletions bus/tests/src/scenarios/test_scenarios.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1384,6 +1384,69 @@ mod defs {
Ok(())
}

/// Reading beyond tail or reading an empty bus with non-zero bounds
/// must return empty entries without panicking and advance next_start_position appropriately.
#[scenario]
pub async fn run_test_read_next_beyond_tail<F: AgentBusTestFixture>(
fixture: &F,
) -> anyhow::Result<()> {
let env = fixture.get_env();
let bus = fixture.create_impl();
let bus_id: String = format!("bus-{}", env.with_rng(|rng| rng.random::<u64>()));

// Case 1: Empty bus, read range [0, 5) and [5, 10)
let resp1 = bus
.read_next(ReadNextRequest {
agent_bus_id: bus_id.clone(),
bus_id: Some(BusId {
agent_bus_id: bus_id.clone(),
}),
start_log_position: 0,
end_log_position: 5,
max_entries: 10,
filter: None,
})
.await?;
assert!(resp1.entries.is_empty());
assert_eq!(resp1.next_start_position, 5);

let resp2 = bus
.read_next(ReadNextRequest {
agent_bus_id: bus_id.clone(),
bus_id: Some(BusId {
agent_bus_id: bus_id.clone(),
}),
start_log_position: 5,
end_log_position: 10,
max_entries: 10,
filter: None,
})
.await?;
assert!(resp2.entries.is_empty());
assert_eq!(resp2.next_start_position, 10);

// Case 2: Append 2 entries, then read [5, 10) beyond tail
append_string_intention(&bus, bus_id.clone(), "entry-0".to_string()).await;
append_string_intention(&bus, bus_id.clone(), "entry-1".to_string()).await;

let resp3 = bus
.read_next(ReadNextRequest {
agent_bus_id: bus_id.clone(),
bus_id: Some(BusId {
agent_bus_id: bus_id.clone(),
}),
start_log_position: 5,
end_log_position: 10,
max_entries: 10,
filter: None,
})
.await?;
assert!(resp3.entries.is_empty());
assert_eq!(resp3.next_start_position, 10);

Ok(())
}

/// Multi-type filter: filter for intentions + votes, verify commits/aborts excluded.
#[scenario]
pub async fn run_test_read_next_multi_type_filter<F: AgentBusTestFixture>(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,6 @@

use std::marker::PhantomData;

use agent_bus_proto_rust::agent_bus::BusId;
use agentbus_core::client::AgentBusClient;
use logact_commit_service_api::CommitError;
use logact_commit_service_api::CommitIntentionCommand;
Expand Down