use coarsetime::Clock;
use super::range_index_manager__migration::PublishMigratedIndexResult;
#[inline]
fn now_ns() -> u64 {
Clock::now_since_epoch().as_u64()
}
#[derive(Debug, Clone)]
pub struct StreamActivity {
started_ns: u64,
chunk_size: usize,
file_size_bytes: i64,
chunk_count: usize,
total_bytes_enqueued: i64,
error: Option<String>,
}
impl StreamActivity {
pub fn start_activity(chunk_size: usize) -> Self {
Self {
started_ns: now_ns(),
chunk_size,
file_size_bytes: 0,
chunk_count: 0,
total_bytes_enqueued: 0,
error: None,
}
}
pub fn on_file_length(&mut self, file_bytes: i64) {
self.file_size_bytes = file_bytes;
}
pub fn on_chunk_enqueued(&mut self, bytes: usize) {
self.chunk_count += 1;
self.total_bytes_enqueued += bytes as i64;
}
pub fn on_error(&mut self, error: &str) {
self.error.get_or_insert_with(|| error.to_string());
}
pub fn end_and_log(&self, key: &[u8]) {
let total_ticks = now_ns().saturating_sub(self.started_ns);
log::info!(
"RangeIndexReplicationStreamActivity: key={key} isError={is_error} errorStr={error} chunkSize={chunk_size} fileSizeBytes={file_size_bytes} chunkCount={chunk_count} totalBytesEnqueued={total_bytes_enqueued} totalTicks={total_ticks}",
key = String::from_utf8_lossy(key),
is_error = self.error.is_some(),
error = self.error.as_deref().unwrap_or(""),
chunk_size = self.chunk_size,
file_size_bytes = self.file_size_bytes,
chunk_count = self.chunk_count,
total_bytes_enqueued = self.total_bytes_enqueued,
total_ticks = total_ticks,
);
}
}
#[derive(Debug, Clone)]
pub struct ReassemblyActivity {
started_ns: u64,
chunk_count: usize,
total_bytes_received: i64,
publish_result: Option<PublishMigratedIndexResult>,
}
impl ReassemblyActivity {
pub fn start_activity() -> Self {
Self {
started_ns: now_ns(),
chunk_count: 0,
total_bytes_received: 0,
publish_result: None,
}
}
pub fn on_chunk_received(&mut self, chunk_length: usize) {
self.chunk_count += 1;
self.total_bytes_received += chunk_length as i64;
}
pub fn on_publish_result(&mut self, result: PublishMigratedIndexResult) {
self.publish_result = Some(result);
}
pub fn end_and_log(&self, key: &[u8], reason: &str) {
let total_ticks = now_ns().saturating_sub(self.started_ns);
let publish_result_text = self
.publish_result
.as_ref()
.map_or("n/a".to_string(), |r| r.to_string());
log::info!(
"RangeIndexReplicationReassemblyActivity: key={key} reason={reason} publishResult={publish_result} chunkCount={chunk_count} totalBytesReceived={total_bytes_received} totalTicks={total_ticks}",
key = String::from_utf8_lossy(key),
reason = reason,
publish_result = publish_result_text,
chunk_count = self.chunk_count,
total_bytes_received = self.total_bytes_received,
total_ticks = total_ticks,
);
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn stream_activity_counters_and_first_error_wins() {
let mut a = StreamActivity::start_activity(4096);
a.on_file_length(100_000);
a.on_chunk_enqueued(4096);
a.on_chunk_enqueued(2048);
assert_eq!(a.chunk_count, 2);
assert_eq!(a.total_bytes_enqueued, 6144);
a.on_error("ZeroLengthChunkFromReader");
a.on_error("SecondErrorIgnored");
assert_eq!(a.chunk_size, 4096);
assert_eq!(a.error.as_deref(), Some("ZeroLengthChunkFromReader"));
a.end_and_log(b"idx-key");
}
#[test]
fn stream_activity_success_path_logs() {
let mut a = StreamActivity::start_activity(1024);
a.on_file_length(2048);
a.on_chunk_enqueued(1024);
a.on_chunk_enqueued(1024);
assert!(a.error.is_none());
a.end_and_log(b"single-chunk");
}
#[test]
fn reassembly_activity_counters_and_publish_result() {
let mut a = ReassemblyActivity::start_activity();
a.on_chunk_received(512);
a.on_chunk_received(512);
a.on_chunk_received(47);
assert_eq!(a.chunk_count, 3);
assert_eq!(a.total_bytes_received, 1071);
a.end_and_log(b"k", "ChunkProcessingError");
a.on_publish_result(PublishMigratedIndexResult::Success);
a.end_and_log(b"k", "Complete");
}
}