use crate::{
crash_info::{CrashInfo, CrashInfoBuilder, ErrorKind, SigInfo, Span, StackFrame},
runtime_callback::RuntimeStack,
shared::constants::*,
CrashtrackerConfiguration,
};
use anyhow::Context;
use serde::{Deserialize, Serialize};
use std::time::{Duration, Instant};
use tokio::io::AsyncBufReadExt;
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
struct RuntimeStackFrame {
#[serde(default, skip_serializing_if = "Option::is_none")]
line: Option<u32>,
#[serde(default, skip_serializing_if = "Option::is_none")]
column: Option<u32>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
function: Vec<u8>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
type_name: Vec<u8>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
file: Vec<u8>,
}
impl From<RuntimeStackFrame> for StackFrame {
fn from(value: RuntimeStackFrame) -> Self {
let mut stack_frame = StackFrame::new();
stack_frame.function = if value.function.is_empty() {
None
} else {
Some(String::from_utf8_lossy(&value.function).to_string())
};
stack_frame.type_name = if value.type_name.is_empty() {
None
} else {
Some(String::from_utf8_lossy(&value.type_name).to_string())
};
stack_frame.file = if value.file.is_empty() {
None
} else {
Some(String::from_utf8_lossy(&value.file).to_string())
};
stack_frame.line = value.line;
stack_frame.column = value.column;
stack_frame
}
}
#[derive(Debug)]
pub(crate) enum StdinState {
AdditionalTags,
Config,
Counters,
Done,
File(String, Vec<String>),
Metadata,
ProcInfo,
SigInfo,
SpanIds,
StackTrace,
TraceIds,
Ucontext,
Waiting,
RuntimeStackFrame(Vec<StackFrame>),
RuntimeStackString(Vec<String>),
}
fn process_line(
builder: &mut CrashInfoBuilder,
config: &mut Option<CrashtrackerConfiguration>,
line: &str,
state: StdinState,
) -> anyhow::Result<StdinState> {
let next = match state {
StdinState::AdditionalTags if line.starts_with(DD_CRASHTRACK_END_ADDITIONAL_TAGS) => {
StdinState::Waiting
}
StdinState::AdditionalTags => {
let additional_tags: Vec<String> = serde_json::from_str(line)?;
builder.with_experimental_additional_tags(additional_tags)?;
StdinState::AdditionalTags
}
StdinState::Config if line.starts_with(DD_CRASHTRACK_END_CONFIG) => StdinState::Waiting,
StdinState::Config => {
if config.is_some() {
eprintln!("Unexpected double config");
}
*config = Some(serde_json::from_str(line)?);
StdinState::Config
}
StdinState::Counters if line.starts_with(DD_CRASHTRACK_END_COUNTERS) => StdinState::Waiting,
StdinState::Counters => {
let v: serde_json::Value = serde_json::from_str(line)?;
let map = v.as_object().context("Expected map type value")?;
anyhow::ensure!(map.len() == 1);
let (key, val) = map
.iter()
.next()
.context("we know there is one value here")?;
let val = val.as_i64().context("Vals are ints")?;
builder.with_counter(key.clone(), val)?;
StdinState::Counters
}
StdinState::Done => {
builder.with_log_message(
format!("Unexpected line after crashreport is done: {line}"),
true,
)?;
StdinState::Done
}
StdinState::File(filename, lines) if line.starts_with(DD_CRASHTRACK_END_FILE) => {
builder.with_file_and_contents(filename, lines)?;
StdinState::Waiting
}
StdinState::File(name, mut contents) => {
contents.push(line.to_string());
StdinState::File(name, contents)
}
StdinState::Metadata if line.starts_with(DD_CRASHTRACK_END_METADATA) => StdinState::Waiting,
StdinState::Metadata => {
let metadata = serde_json::from_str(line)?;
builder.with_metadata(metadata)?;
StdinState::Metadata
}
StdinState::ProcInfo if line.starts_with(DD_CRASHTRACK_END_PROCINFO) => StdinState::Waiting,
StdinState::ProcInfo => {
let proc_info = serde_json::from_str(line)?;
builder.with_proc_info(proc_info)?;
StdinState::ProcInfo
}
StdinState::RuntimeStackFrame(frames)
if line.starts_with(DD_CRASHTRACK_END_RUNTIME_STACK_FRAME) =>
{
let runtime_stack = RuntimeStack {
format: "Datadog Runtime Callback 1.0".to_string(),
frames,
stacktrace_string: None,
};
builder.with_experimental_runtime_stack(runtime_stack)?;
StdinState::Waiting
}
StdinState::RuntimeStackFrame(mut frames) => {
let frame_json: RuntimeStackFrame = serde_json::from_str(line)?;
frames.push(frame_json.into());
StdinState::RuntimeStackFrame(frames)
}
StdinState::RuntimeStackString(lines)
if line.starts_with(DD_CRASHTRACK_END_RUNTIME_STACK_STRING) =>
{
let runtime_stack = RuntimeStack {
format: "Datadog Runtime Callback 1.0".to_string(),
frames: vec![],
stacktrace_string: Some(lines.join("\n")),
};
builder.with_experimental_runtime_stack(runtime_stack)?;
StdinState::Waiting
}
StdinState::RuntimeStackString(mut lines) => {
lines.push(line.to_string());
StdinState::RuntimeStackString(lines)
}
StdinState::SigInfo if line.starts_with(DD_CRASHTRACK_END_SIGINFO) => StdinState::Waiting,
StdinState::SigInfo => {
let sig_info: SigInfo = serde_json::from_str(line)?;
let message = format!(
"Process terminated with {:?} ({:?})",
sig_info.si_code_human_readable, sig_info.si_signo_human_readable
);
builder
.with_timestamp_now()?
.with_sig_info(sig_info)?
.with_incomplete(true)?
.with_message(message)?;
StdinState::SigInfo
}
StdinState::SpanIds if line.starts_with(DD_CRASHTRACK_END_SPAN_IDS) => StdinState::Waiting,
StdinState::SpanIds => {
let span_ids: Vec<Span> = serde_json::from_str(line)?;
builder.with_span_ids(span_ids)?;
StdinState::SpanIds
}
StdinState::StackTrace if line.starts_with(DD_CRASHTRACK_END_STACKTRACE) => {
builder.with_stack_set_complete()?;
StdinState::Waiting
}
StdinState::StackTrace => {
let frame = serde_json::from_str(line)?;
builder.with_stack_frame(frame, true)?;
StdinState::StackTrace
}
StdinState::TraceIds if line.starts_with(DD_CRASHTRACK_END_TRACE_IDS) => {
StdinState::Waiting
}
StdinState::TraceIds => {
let trace_ids: Vec<Span> = serde_json::from_str(line)?;
builder.with_trace_ids(trace_ids)?;
StdinState::TraceIds
}
StdinState::Ucontext if line.starts_with(DD_CRASHTRACK_END_UCONTEXT) => StdinState::Waiting,
StdinState::Ucontext => {
builder.with_experimental_ucontext(line.to_string())?;
StdinState::Ucontext
}
StdinState::Waiting if line.starts_with(DD_CRASHTRACK_BEGIN_ADDITIONAL_TAGS) => {
StdinState::AdditionalTags
}
StdinState::Waiting if line.starts_with(DD_CRASHTRACK_BEGIN_CONFIG) => StdinState::Config,
StdinState::Waiting if line.starts_with(DD_CRASHTRACK_BEGIN_COUNTERS) => {
StdinState::Counters
}
StdinState::Waiting if line.starts_with(DD_CRASHTRACK_BEGIN_FILE) => {
let (_, filename) = line.split_once(' ').unwrap_or(("", "MISSING_FILENAME"));
StdinState::File(filename.to_string(), vec![])
}
StdinState::Waiting if line.starts_with(DD_CRASHTRACK_BEGIN_METADATA) => {
StdinState::Metadata
}
StdinState::Waiting if line.starts_with(DD_CRASHTRACK_BEGIN_PROCINFO) => {
StdinState::ProcInfo
}
StdinState::Waiting if line.starts_with(DD_CRASHTRACK_BEGIN_SIGINFO) => StdinState::SigInfo,
StdinState::Waiting if line.starts_with(DD_CRASHTRACK_BEGIN_SPAN_IDS) => {
StdinState::SpanIds
}
StdinState::Waiting if line.starts_with(DD_CRASHTRACK_BEGIN_STACKTRACE) => {
StdinState::StackTrace
}
StdinState::Waiting if line.starts_with(DD_CRASHTRACK_BEGIN_RUNTIME_STACK_STRING) => {
StdinState::RuntimeStackString(vec![])
}
StdinState::Waiting if line.starts_with(DD_CRASHTRACK_BEGIN_RUNTIME_STACK_FRAME) => {
StdinState::RuntimeStackFrame(vec![])
}
StdinState::Waiting if line.starts_with(DD_CRASHTRACK_BEGIN_TRACE_IDS) => {
StdinState::TraceIds
}
StdinState::Waiting if line.starts_with(DD_CRASHTRACK_BEGIN_UCONTEXT) => {
StdinState::Ucontext
}
StdinState::Waiting if line.starts_with(DD_CRASHTRACK_DONE) => {
builder.with_incomplete(false)?;
StdinState::Done
}
StdinState::Waiting => {
builder.with_log_message(
format!("Unexpected line while receiving crashreport: {line}"),
true,
)?;
StdinState::Waiting
}
};
Ok(next)
}
pub(crate) async fn receive_report_from_stream(
timeout: Duration,
stream: impl AsyncBufReadExt + std::marker::Unpin,
) -> anyhow::Result<Option<(CrashtrackerConfiguration, CrashInfo)>> {
let mut builder = CrashInfoBuilder::new();
let mut stdin_state = StdinState::Waiting;
let mut config: Option<CrashtrackerConfiguration> = None;
let mut crash_ping_sent = false;
let mut lines = stream.lines();
let mut deadline = None;
let mut remaining_timeout = Duration::MAX;
loop {
if !crash_ping_sent && builder.is_ping_ready() {
if let Some(ref config_ref) = config {
let config_clone = config_ref.clone();
crash_ping_sent = true;
let crash_ping = builder.build_crash_ping()?;
tokio::task::spawn(async move {
if let Err(e) = crash_ping
.upload_to_endpoint_async(config_clone.endpoint())
.await
{
eprintln!("Failed to send crash ping: {e}");
}
});
} else {
eprintln!("No config found, skipping crash ping");
}
}
let next_line = tokio::time::timeout(remaining_timeout, lines.next_line()).await;
let Ok(next_line) = next_line else {
builder.with_log_message(format!("Timeout: {next_line:?}"), true)?;
break;
};
let Ok(next_line) = next_line else {
builder.with_log_message(format!("IO Error: {next_line:?}"), true)?;
break;
};
let Some(next_line) = next_line else { break };
match process_line(&mut builder, &mut config, &next_line, stdin_state) {
Ok(next_state) => {
stdin_state = next_state;
if matches!(stdin_state, StdinState::Done) {
break;
}
}
Err(e) => {
builder.with_log_message(
format!("Unable to process line: {next_line}. Error: {e}"),
true,
)?;
break;
}
}
if let Some(deadline) = deadline {
remaining_timeout = deadline - Instant::now()
} else {
deadline = Some(Instant::now() + timeout);
remaining_timeout = timeout;
}
}
if !builder.has_data() {
return Ok(None);
}
builder.with_kind(ErrorKind::UnixSignal)?;
let config = config.context("Missing crashtracker configuration")?;
for filename in config.additional_files() {
if let Err(e) = builder.with_file(filename.clone()) {
builder.with_log_message(e.to_string(), true)?;
}
}
let crash_info = builder.build()?;
Ok(Some((config, crash_info)))
}