use std::io::{Cursor, Read};
use crate::{
context::{read_build_event_context, BuildEventContext},
error::MuninError,
field_flags::BuildEventArgsFieldFlags,
header::{open_binlog, BinlogHeader},
nvl_table::{NameValueListTable, NameValuePair},
primitives::{read_7bit_count, read_7bit_int},
reader::{dispatch_event, ArchiveEntry, BinlogEvent},
record_kind::BinaryLogRecordKind,
string_table::StringTable,
};
#[derive(Debug, Clone)]
pub struct EventMeta {
pub record_kind: BinaryLogRecordKind,
pub byte_offset: u64,
pub payload_len: usize,
pub context: Option<BuildEventContext>,
}
#[derive(Debug, Clone)]
struct IndexEntry {
meta: EventMeta,
payload: Vec<u8>,
}
#[derive(Debug)]
pub struct BinlogIndex {
header: BinlogHeader,
strings: StringTable,
nvl_table: NameValueListTable,
archives: Vec<Vec<u8>>,
entries: Vec<IndexEntry>,
}
impl BinlogIndex {
pub fn open(reader: impl Read) -> Result<Self, MuninError> {
let (header, mut gz_reader) = open_binlog(reader)?;
let version = header.file_format_version;
let mut strings = StringTable::new();
let mut nvl_table = NameValueListTable::new();
let mut archives = Vec::new();
let mut entries = Vec::new();
let mut offset: u64 = 8;
loop {
let kind_start = offset;
let (kind_raw, kind_bytes) = read_7bit_int_counted(&mut gz_reader)?;
offset += kind_bytes;
if kind_raw == BinaryLogRecordKind::EndOfFile as i32 {
break;
}
let (record_length, len_bytes) = read_7bit_int_counted(&mut gz_reader)?;
if record_length < 0 {
return Err(MuninError::InvalidFormat(format!(
"negative record length: {record_length}"
)));
}
let record_length = record_length as usize;
if record_length > crate::primitives::MAX_BINLOG_FIELD_LEN {
return Err(MuninError::InvalidFormat(format!(
"record length too large: {record_length} (max {})",
crate::primitives::MAX_BINLOG_FIELD_LEN
)));
}
offset += len_bytes;
let mut payload = vec![0u8; record_length];
gz_reader.read_exact(&mut payload)?;
offset += record_length as u64;
let kind = BinaryLogRecordKind::from_raw(kind_raw);
match kind {
Some(BinaryLogRecordKind::String) => {
let s = String::from_utf8(payload).map_err(|_| MuninError::InvalidUtf8)?;
strings.add(s);
}
Some(BinaryLogRecordKind::NameValueList) => {
let mut cursor = Cursor::new(&payload);
let count = read_7bit_count(&mut cursor, "name-value list count")?;
let mut pairs = Vec::with_capacity(count);
for _ in 0..count {
let key_index = read_7bit_int(&mut cursor)?;
let value_index = read_7bit_int(&mut cursor)?;
pairs.push(NameValuePair {
key_index,
value_index,
});
}
nvl_table.add(pairs);
}
Some(BinaryLogRecordKind::ProjectImportArchive) => {
archives.push(payload);
}
Some(record_kind) if !record_kind.is_auxiliary() => {
let context = extract_context(&payload, version);
let meta = EventMeta {
record_kind,
byte_offset: kind_start,
payload_len: record_length,
context,
};
entries.push(IndexEntry { meta, payload });
}
_ => {}
}
}
Ok(Self {
header,
strings,
nvl_table,
archives,
entries,
})
}
pub fn header(&self) -> &BinlogHeader {
&self.header
}
pub fn strings(&self) -> &StringTable {
&self.strings
}
pub fn nvl_table(&self) -> &NameValueListTable {
&self.nvl_table
}
pub fn archives(&self) -> &[Vec<u8>] {
&self.archives
}
pub fn len(&self) -> usize {
self.entries.len()
}
pub fn is_empty(&self) -> bool {
self.entries.is_empty()
}
pub fn meta(&self, index: usize) -> Option<&EventMeta> {
self.entries.get(index).map(|e| &e.meta)
}
pub fn iter_meta(&self) -> impl Iterator<Item = (usize, &EventMeta)> {
self.entries.iter().enumerate().map(|(i, e)| (i, &e.meta))
}
pub fn get(&self, index: usize) -> Result<Option<BinlogEvent>, MuninError> {
let entry = match self.entries.get(index) {
Some(e) => e,
None => return Ok(None),
};
let mut cursor = Cursor::new(&entry.payload);
let event = dispatch_event(
&mut cursor,
entry.meta.record_kind,
&self.strings,
&self.nvl_table,
self.header.file_format_version,
)?;
Ok(Some(event))
}
pub fn get_all(&self) -> Result<Vec<BinlogEvent>, MuninError> {
let mut events = Vec::with_capacity(self.entries.len());
for i in 0..self.entries.len() {
if let Some(event) = self.get(i)? {
events.push(event);
}
}
Ok(events)
}
pub fn indices_by_kind(&self, kind: BinaryLogRecordKind) -> Vec<usize> {
self.entries
.iter()
.enumerate()
.filter(|(_, e)| e.meta.record_kind == kind)
.map(|(i, _)| i)
.collect()
}
pub fn indices_by_project_context(&self, project_context_id: i32) -> Vec<usize> {
self.entries
.iter()
.enumerate()
.filter(|(_, e)| {
e.meta
.context
.is_some_and(|c| c.project_context_id == project_context_id)
})
.map(|(i, _)| i)
.collect()
}
pub fn indices_by_target_id(&self, target_id: i32) -> Vec<usize> {
self.entries
.iter()
.enumerate()
.filter(|(_, e)| e.meta.context.is_some_and(|c| c.target_id == target_id))
.map(|(i, _)| i)
.collect()
}
pub fn indices_by_task_id(&self, task_id: i32) -> Vec<usize> {
self.entries
.iter()
.enumerate()
.filter(|(_, e)| e.meta.context.is_some_and(|c| c.task_id == task_id))
.map(|(i, _)| i)
.collect()
}
pub fn query(
&self,
kind: Option<BinaryLogRecordKind>,
project_context_id: Option<i32>,
target_id: Option<i32>,
task_id: Option<i32>,
) -> Vec<usize> {
self.entries
.iter()
.enumerate()
.filter(|(_, e)| {
if let Some(k) = kind {
if e.meta.record_kind != k {
return false;
}
}
if let Some(pid) = project_context_id {
if e.meta.context.is_none_or(|c| c.project_context_id != pid) {
return false;
}
}
if let Some(tid) = target_id {
if e.meta.context.is_none_or(|c| c.target_id != tid) {
return false;
}
}
if let Some(tsk) = task_id {
if e.meta.context.is_none_or(|c| c.task_id != tsk) {
return false;
}
}
true
})
.map(|(i, _)| i)
.collect()
}
pub fn extract_archives(&self) -> Result<Vec<ArchiveEntry>, MuninError> {
let mut entries = Vec::new();
for archive_bytes in &self.archives {
let cursor = Cursor::new(archive_bytes);
let mut archive = zip::ZipArchive::new(cursor).map_err(|e| {
MuninError::InvalidFormat(format!("invalid ProjectImportArchive zip: {e}"))
})?;
for i in 0..archive.len() {
let mut file = archive.by_index(i).map_err(|e| {
MuninError::InvalidFormat(format!("cannot read zip entry {i}: {e}"))
})?;
if file.is_dir() {
continue;
}
let path = file.name().to_string();
let mut contents = String::new();
if file.read_to_string(&mut contents).is_ok() {
entries.push(ArchiveEntry { path, contents });
}
}
}
Ok(entries)
}
}
fn read_7bit_int_counted(reader: &mut impl Read) -> Result<(i32, u64), MuninError> {
let mut result: u32 = 0;
let mut shift: u32 = 0;
let mut buf = [0u8; 1];
let mut count: u64 = 0;
for _ in 0..5 {
reader.read_exact(&mut buf)?;
count += 1;
let byte = buf[0];
result |= ((byte & 0x7F) as u32) << shift;
if byte & 0x80 == 0 {
return Ok((result as i32, count));
}
shift += 7;
}
Err(MuninError::OverlongVarInt)
}
fn extract_context(payload: &[u8], file_format_version: i32) -> Option<BuildEventContext> {
let mut cursor = Cursor::new(payload);
let flags_raw = read_7bit_int(&mut cursor).ok()? as u32;
let flags = BuildEventArgsFieldFlags::from_raw(flags_raw);
if flags.contains(BuildEventArgsFieldFlags::MESSAGE) {
let _ = read_7bit_int(&mut cursor).ok()?;
}
if flags.contains(BuildEventArgsFieldFlags::BUILD_EVENT_CONTEXT) {
read_build_event_context(&mut cursor, file_format_version).ok()
} else {
None
}
}
#[cfg(test)]
mod tests;