use arrow_array::builder::{
BinaryBuilder, FixedSizeBinaryBuilder, Int32Builder, StringBuilder, TimestampNanosecondBuilder,
UInt16Builder, UInt32Builder,
};
use arrow_array::{ArrayRef, RecordBatch};
use prost::Message;
use std::sync::Arc;
use mira_proto::collector::logs::v1::ExportLogsServiceRequest;
use mira_proto::common::v1::any_value::Value;
use crate::attrs::{AttrsBuilder, DictColumn, ResourceScope, resource_kv, scope_kv};
use crate::error::Result;
use crate::schema::LOGS;
use crate::signal::{Sealed, Sidecars, SignalBuilder};
pub struct LogsBuilder {
id: UInt32Builder,
time: TimestampNanosecondBuilder,
observed: TimestampNanosecondBuilder,
sev_num: Int32Builder,
sev_text: DictColumn,
event_name: DictColumn,
body: StringBuilder,
body_ser: BinaryBuilder,
trace_id: FixedSizeBinaryBuilder,
span_id: FixedSizeBinaryBuilder,
flags: UInt32Builder,
dropped: UInt32Builder,
resource_id: UInt16Builder,
scope_id: UInt16Builder,
log_attrs: AttrsBuilder,
rs: ResourceScope,
next_id: u32,
min_ts: i64,
max_ts: i64,
}
impl Default for LogsBuilder {
fn default() -> Self {
Self::new()
}
}
impl LogsBuilder {
pub fn new() -> Self {
Self {
id: UInt32Builder::new(),
time: TimestampNanosecondBuilder::new(),
observed: TimestampNanosecondBuilder::new(),
sev_num: Int32Builder::new(),
sev_text: DictColumn::new("logs.severity_text"),
event_name: DictColumn::new("logs.event_name"),
body: StringBuilder::new(),
body_ser: BinaryBuilder::new(),
trace_id: FixedSizeBinaryBuilder::new(16),
span_id: FixedSizeBinaryBuilder::new(8),
flags: UInt32Builder::new(),
dropped: UInt32Builder::new(),
resource_id: UInt16Builder::new(),
scope_id: UInt16Builder::new(),
log_attrs: AttrsBuilder::new("log_attrs.key"),
rs: ResourceScope::new(),
next_id: 0,
min_ts: i64::MAX,
max_ts: i64::MIN,
}
}
pub fn num_rows(&self) -> usize {
self.next_id as usize
}
pub fn is_empty(&self) -> bool {
self.next_id == 0
}
pub fn approx_bytes(&self) -> usize {
self.next_id as usize * 64
+ (self.log_attrs.len() + self.rs.len()) * 48
+ self.body.values_slice().len()
+ self.body_ser.values_slice().len()
+ self.log_attrs.heap_bytes()
+ self.rs.heap_bytes()
}
pub fn has_headroom_for(&self, req: &ExportLogsServiceRequest) -> bool {
let (mut resources, mut scopes, mut records) = (0usize, 0usize, 0usize);
let (mut res_kv, mut sc_kv, mut log_kv) = (0usize, 0usize, 0usize);
for rl in &req.resource_logs {
resources += 1;
res_kv += resource_kv(rl.resource.as_ref());
for sl in &rl.scope_logs {
scopes += 1;
sc_kv += scope_kv(sl.scope.as_ref());
records += sl.log_records.len();
log_kv += sl
.log_records
.iter()
.map(|r| r.attributes.len())
.sum::<usize>();
}
}
self.rs.has_headroom(resources, scopes, res_kv, sc_kv)
&& self.log_attrs.has_headroom(log_kv)
&& self.sev_text.has_headroom(records)
&& self.event_name.has_headroom(records)
}
pub fn append_request(&mut self, req: &ExportLogsServiceRequest) -> Result<usize> {
let mut added = 0;
for rl in &req.resource_logs {
let rid = self.rs.resource(rl.resource.as_ref())?;
for sl in &rl.scope_logs {
let sid = self.rs.scope(sl.scope.as_ref())?;
for rec in &sl.log_records {
self.sev_text.append(&rec.severity_text)?;
self.event_name.append(&rec.event_name)?;
let id = self.next_id;
self.next_id += 1;
let observed = nanos(rec.observed_time_unix_nano);
let t = match nanos(rec.time_unix_nano) {
0 if observed != 0 => observed,
0 => now_nanos(),
t => t,
};
self.min_ts = self.min_ts.min(t);
self.max_ts = self.max_ts.max(t);
self.id.append_value(id);
self.time.append_value(t);
if observed != 0 {
self.observed.append_value(observed);
} else {
self.observed.append_null();
}
self.sev_num.append_value(rec.severity_number);
match rec.body.as_ref().and_then(|b| b.value.as_ref()) {
Some(Value::StringValue(s)) => {
self.body.append_value(s);
self.body_ser.append_null();
}
Some(_) => {
self.body.append_null();
self.body_ser
.append_value(rec.body.as_ref().unwrap().encode_to_vec());
}
None => {
self.body.append_null();
self.body_ser.append_null();
}
}
append_fixed(&mut self.trace_id, &rec.trace_id, 16)?;
append_fixed(&mut self.span_id, &rec.span_id, 8)?;
self.flags.append_value(rec.flags);
self.dropped.append_value(rec.dropped_attributes_count);
self.resource_id.append_value(rid);
self.scope_id.append_value(sid);
self.log_attrs.append_all(id, &rec.attributes)?;
added += 1;
}
}
}
Ok(added)
}
pub fn finish(&mut self) -> Result<Sealed> {
let out = self.seal(Sidecars::Build);
*self = Self::new();
out
}
fn seal(&self, sidecars: Sidecars) -> Result<Sealed> {
let trace_ids = self.trace_id.finish_cloned();
let trace_idx = match sidecars {
Sidecars::Build => crate::bloom::build(&trace_ids),
Sidecars::Skip => None,
};
let cols: Vec<ArrayRef> = vec![
Arc::new(self.id.finish_cloned()),
Arc::new(self.time.finish_cloned()),
Arc::new(self.observed.finish_cloned()),
Arc::new(self.sev_num.finish_cloned()),
self.sev_text.finish(),
self.event_name.finish(),
Arc::new(self.body.finish_cloned()),
Arc::new(self.body_ser.finish_cloned()),
Arc::new(trace_ids),
Arc::new(self.span_id.finish_cloned()),
Arc::new(self.flags.finish_cloned()),
Arc::new(self.dropped.finish_cloned()),
Arc::new(self.resource_id.finish_cloned()),
Arc::new(self.scope_id.finish_cloned()),
];
let mut tables = vec![
("logs", RecordBatch::try_new(LOGS.clone(), cols)?),
("log_attrs", self.log_attrs.finish()?),
];
tables.extend(self.rs.finish()?);
Ok(Sealed::with(
sidecars,
self.next_id as usize,
tables,
self.min_ts,
self.max_ts,
)
.with_sidecar(crate::bloom::TRACE_IDX, trace_idx))
}
}
impl SignalBuilder for LogsBuilder {
type Request = ExportLogsServiceRequest;
const SIGNAL: &'static str = "logs";
fn has_headroom_for(&self, req: &Self::Request) -> bool {
LogsBuilder::has_headroom_for(self, req)
}
fn append_request(&mut self, req: &Self::Request) -> Result<usize> {
LogsBuilder::append_request(self, req)
}
fn approx_bytes(&self) -> usize {
LogsBuilder::approx_bytes(self)
}
fn is_empty(&self) -> bool {
LogsBuilder::is_empty(self)
}
fn finish(&mut self) -> Result<Sealed> {
LogsBuilder::finish(self)
}
fn snapshot(&self) -> Result<Sealed> {
self.seal(Sidecars::Skip)
}
}
pub(crate) fn nanos(t: u64) -> i64 {
i64::try_from(t).unwrap_or(0)
}
fn now_nanos() -> i64 {
let since_epoch = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default();
i64::try_from(since_epoch.as_nanos()).unwrap_or(i64::MAX)
}
pub(crate) fn append_fixed(b: &mut FixedSizeBinaryBuilder, v: &[u8], width: usize) -> Result<()> {
if v.len() == width {
b.append_value(v)?;
} else {
b.append_null();
}
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
use arrow_array::{Array, TimestampNanosecondArray};
use mira_proto::logs::v1::{LogRecord, ResourceLogs, ScopeLogs};
fn request(records: Vec<LogRecord>) -> ExportLogsServiceRequest {
ExportLogsServiceRequest {
resource_logs: vec![ResourceLogs {
scope_logs: vec![ScopeLogs {
log_records: records,
..Default::default()
}],
..Default::default()
}],
}
}
fn col(s: &Sealed, name: &str) -> TimestampNanosecondArray {
s.table("logs")
.unwrap()
.column_by_name(name)
.unwrap()
.as_any()
.downcast_ref::<TimestampNanosecondArray>()
.unwrap()
.clone()
}
#[test]
fn a_timestamp_past_i64_still_publishes_a_block_the_catalog_can_see() {
let mut b = LogsBuilder::new();
b.append_request(&request(vec![
LogRecord {
time_unix_nano: u64::MAX,
observed_time_unix_nano: u64::MAX,
..Default::default()
},
LogRecord {
time_unix_nano: 5_000,
..Default::default()
},
]))
.unwrap();
let sealed = b.finish().unwrap();
assert_eq!(sealed.min_ts, 5_000);
assert!(sealed.min_ts >= 0 && sealed.max_ts >= sealed.min_ts);
assert!(col(&sealed, "time_unix_nano").value(0) > 5_000);
assert!(col(&sealed, "observed_time_unix_nano").is_null(0));
let root = std::env::temp_dir().join(format!("mira-ts-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&root);
let published = crate::block::publish(&root, "logs", 7, 1, 0, &sealed).unwrap();
assert_eq!(crate::block::scan(&root, "logs").unwrap(), vec![published]);
assert_eq!(crate::block::expire(&root, "logs", i64::MAX).unwrap(), 1);
let _ = std::fs::remove_dir_all(&root);
}
#[test]
fn a_record_with_no_clock_falls_back_to_observed_then_to_receipt_time() {
let before = now_nanos();
let mut b = LogsBuilder::new();
assert!(b.is_empty() && b.num_rows() == 0, "nothing appended yet");
b.append_request(&request(vec![
LogRecord {
time_unix_nano: 3_000,
observed_time_unix_nano: 4_000,
..Default::default()
},
LogRecord {
observed_time_unix_nano: 4_000,
..Default::default()
},
LogRecord::default(),
]))
.unwrap();
assert_eq!(b.num_rows(), 3);
assert!(!b.is_empty());
let sealed = b.finish().unwrap();
assert_eq!(sealed.table("logs").unwrap().num_rows(), 3);
let time = col(&sealed, "time_unix_nano");
assert_eq!(time.value(0), 3_000, "a real clock wins");
assert_eq!(time.value(1), 4_000, "then the receiver's own reading");
assert!(time.value(2) >= before, "then the moment we took delivery");
assert_eq!(sealed.min_ts, 3_000);
assert_eq!(sealed.max_ts, time.value(2));
let observed = col(&sealed, "observed_time_unix_nano");
assert_eq!(observed.value(0), 4_000);
assert!(observed.is_null(2), "the fallback is not written back");
}
}