diff --git a/bus/core/src/channeled_agentbus.rs b/bus/core/src/channeled_agentbus.rs index 158d72d..f3342c0 100644 --- a/bus/core/src/channeled_agentbus.rs +++ b/bus/core/src/channeled_agentbus.rs @@ -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; diff --git a/bus/impls/simple/src/in_memory_agentbus_state.rs b/bus/impls/simple/src/in_memory_agentbus_state.rs index 31d6dd2..d2060bc 100644 --- a/bus/impls/simple/src/in_memory_agentbus_state.rs +++ b/bus/impls/simple/src/in_memory_agentbus_state.rs @@ -158,6 +158,12 @@ impl InMemoryAgentBusState { payload_types: &Option>, end_position: i64, ) -> Result<(Vec, 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 {}", @@ -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]; @@ -193,9 +203,8 @@ 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)); } } @@ -203,3 +212,98 @@ impl InMemoryAgentBusState { 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); + } +} diff --git a/bus/tests/src/scenarios/test_scenarios.rs b/bus/tests/src/scenarios/test_scenarios.rs index b7daaf9..d4d4e92 100644 --- a/bus/tests/src/scenarios/test_scenarios.rs +++ b/bus/tests/src/scenarios/test_scenarios.rs @@ -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( + 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::())); + + // 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( diff --git a/logact/commit_service/tests/src/fixtures/grpc_bus_id_encoding.rs b/logact/commit_service/tests/src/fixtures/grpc_bus_id_encoding.rs index cef4f10..b829a6e 100644 --- a/logact/commit_service/tests/src/fixtures/grpc_bus_id_encoding.rs +++ b/logact/commit_service/tests/src/fixtures/grpc_bus_id_encoding.rs @@ -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;