use tailtriage_core::{RequestEvent, Run};
use super::{
analyze_run_internal, AnalyzeOptions, DiagnosisKind, Report, SignalCoverageStatus,
TemporalSegment,
};
use crate::route;
const TEMPORAL_RUNTIME_ATTRIBUTION_WARNING: &str = "Runtime and in-flight evidence is sparse in this segment after timestamp filtering; executor/blocking attribution is limited.";
pub(super) const TEMPORAL_SUSPECT_SHIFT_WARNING: &str = "Temporal segments show different primary suspects; inspect temporal_segments before acting on the global suspect.";
pub(super) const TEMPORAL_P95_SHIFT_WARNING: &str =
"Temporal segments show a large p95 latency shift between early and late requests.";
pub(super) const TEMPORAL_OVERLAP_ATTRIBUTION_WARNING: &str = "Segment windows overlap under concurrent requests; timestamp-filtered runtime/in-flight attribution is approximate.";
pub(super) const TEMPORAL_WALL_CLOCK_FALLBACK_WARNING: &str = "Temporal segment used wall-clock timestamp fallback; attribution is approximate for artifacts without complete run-relative timing.";
#[derive(Clone, Copy)]
struct SegmentWindow {
unix_start: u64,
unix_finish: u64,
run_relative: Option<(u64, u64)>,
}
fn all_requests_have_run_relative_start(requests: &[RequestEvent]) -> bool {
requests
.iter()
.all(|request| request.started_at_run_us.is_some())
}
fn sort_requests_for_temporal_segments(requests: &mut [RequestEvent]) {
if all_requests_have_run_relative_start(requests) {
requests.sort_by(|a, b| {
a.started_at_run_us
.unwrap()
.cmp(&b.started_at_run_us.unwrap())
.then_with(|| a.started_at_unix_ms.cmp(&b.started_at_unix_ms))
.then_with(|| a.request_id.cmp(&b.request_id))
});
} else {
requests.sort_by(|a, b| {
a.started_at_unix_ms
.cmp(&b.started_at_unix_ms)
.then_with(|| a.request_id.cmp(&b.request_id))
});
}
}
fn segment_run_relative_window(requests: &[RequestEvent]) -> Option<(u64, u64)> {
if !requests
.iter()
.all(|request| request.started_at_run_us.is_some() && request.finished_at_run_us.is_some())
{
return None;
}
let start = requests
.iter()
.filter_map(|request| request.started_at_run_us)
.min()?;
let finish = requests
.iter()
.filter_map(|request| request.finished_at_run_us)
.max()?;
Some((start, finish))
}
fn segment_unix_window(requests: &[RequestEvent]) -> Option<(u64, u64)> {
let start = requests
.iter()
.map(|request| request.started_at_unix_ms)
.min()?;
let finish = requests
.iter()
.map(|request| request.finished_at_unix_ms)
.max()?;
Some((start, finish))
}
fn retain_segment_sample(
at_unix_ms: u64,
at_run_us: Option<u64>,
window: SegmentWindow,
) -> (bool, bool) {
if let (Some((start, finish)), Some(at)) = (window.run_relative, at_run_us) {
return (at >= start && at <= finish, false);
}
let retained = at_unix_ms >= window.unix_start && at_unix_ms <= window.unix_finish;
let used_unix_fallback = retained && window.run_relative.is_some() && at_run_us.is_none();
(retained, used_unix_fallback)
}
fn filter_segment_runtime_and_inflight(run: &mut Run, source: &Run, window: SegmentWindow) -> bool {
let mut used_unix_fallback = false;
run.runtime_snapshots = source
.runtime_snapshots
.iter()
.filter(|snapshot| {
let (retained, used_fallback) =
retain_segment_sample(snapshot.at_unix_ms, snapshot.at_run_us, window);
used_unix_fallback |= used_fallback;
retained
})
.cloned()
.collect();
run.inflight = source
.inflight
.iter()
.filter(|snapshot| {
let (retained, used_fallback) =
retain_segment_sample(snapshot.at_unix_ms, snapshot.at_run_us, window);
used_unix_fallback |= used_fallback;
retained
})
.cloned()
.collect();
used_unix_fallback
}
fn filtered_run_for_temporal_segment(
run: &Run,
request_ids: &[String],
window: SegmentWindow,
) -> (Run, bool) {
let mut filtered = route::filtered_run_for_route(run, request_ids);
let used_unix_fallback = filter_segment_runtime_and_inflight(&mut filtered, run, window);
(filtered, used_unix_fallback)
}
fn temporal_segment_from_report(
name: &str,
analyzed: Report,
start: Option<u64>,
finish: Option<u64>,
) -> TemporalSegment {
TemporalSegment {
name: name.to_string(),
request_count: analyzed.request_count,
started_at_unix_ms: start,
finished_at_unix_ms: finish,
p50_latency_us: analyzed.p50_latency_us,
p95_latency_us: analyzed.p95_latency_us,
p99_latency_us: analyzed.p99_latency_us,
p95_queue_share_permille: analyzed.p95_queue_share_permille,
p95_service_share_permille: analyzed.p95_service_share_permille,
evidence_quality: analyzed.evidence_quality,
primary_suspect: analyzed.primary_suspect,
secondary_suspects: analyzed.secondary_suspects,
warnings: analyzed.warnings,
}
}
pub(super) fn temporal_segments(
run: &Run,
global_warnings: &mut Vec<String>,
options: &AnalyzeOptions,
) -> Vec<TemporalSegment> {
if run.requests.len() < options.temporal.min_request_count {
return vec![];
}
let mut requests = run.requests.clone();
sort_requests_for_temporal_segments(&mut requests);
let split = requests.len() / 2;
let (early, late) = requests.split_at(split);
if early.len() < options.temporal.min_segment_request_count
|| late.len() < options.temporal.min_segment_request_count
{
return vec![];
}
let build = |name: &str, seg: &[RequestEvent]| {
let ids: Vec<String> = seg.iter().map(|r| r.request_id.clone()).collect();
let Some((unix_start, unix_finish)) = segment_unix_window(seg) else {
let analyzed = analyze_run_internal(&route::filtered_run_for_route(run, &ids), options);
return temporal_segment_from_report(name, analyzed, None, None);
};
let start = Some(unix_start);
let finish = Some(unix_finish);
let run_relative_window = segment_run_relative_window(seg);
let window = SegmentWindow {
unix_start,
unix_finish,
run_relative: run_relative_window,
};
let (filtered, used_snapshot_unix_fallback) =
filtered_run_for_temporal_segment(run, &ids, window);
let mut analyzed = analyze_run_internal(&filtered, options);
if run_relative_window.is_none() || used_snapshot_unix_fallback {
analyzed
.warnings
.push(TEMPORAL_WALL_CLOCK_FALLBACK_WARNING.to_string());
}
let sparse_runtime =
analyzed.evidence_quality.runtime_snapshots != SignalCoverageStatus::Present;
let sparse_inflight =
analyzed.evidence_quality.inflight_snapshots != SignalCoverageStatus::Present;
if matches!(
analyzed.primary_suspect.kind,
DiagnosisKind::ExecutorPressureSuspected | DiagnosisKind::BlockingPoolPressure
) && (sparse_runtime || sparse_inflight)
{
analyzed
.warnings
.push(TEMPORAL_RUNTIME_ATTRIBUTION_WARNING.to_string());
}
temporal_segment_from_report(name, analyzed, start, finish)
};
let mut early_seg = build("early", early);
let mut late_seg = build("late", late);
let p95_shift =
has_material_p95_shift(early_seg.p95_latency_us, late_seg.p95_latency_us, options);
let queue_move = has_material_share_shift(
early_seg.p95_queue_share_permille,
late_seg.p95_queue_share_permille,
options,
);
let service_move = has_material_share_shift(
early_seg.p95_service_share_permille,
late_seg.p95_service_share_permille,
options,
);
let movement = TemporalMovement {
p95_shift,
queue_move,
service_move,
};
let suspect_shift = has_material_suspect_shift(&early_seg, &late_seg, movement, options);
let material = has_material_temporal_signal(suspect_shift, movement, options);
if !material {
return vec![];
}
if options.temporal.emit_on_suspect_shift && suspect_shift {
global_warnings.push(TEMPORAL_SUSPECT_SHIFT_WARNING.to_string());
}
if p95_shift {
global_warnings.push(TEMPORAL_P95_SHIFT_WARNING.to_string());
}
apply_temporal_overlap_attribution_warning(&mut early_seg, &mut late_seg);
vec![early_seg, late_seg]
}
fn has_material_share_shift(
left: Option<u64>,
right: Option<u64>,
options: &AnalyzeOptions,
) -> bool {
matches!((left, right), (Some(a), Some(b)) if a.abs_diff(b) >= options.temporal.share_shift_permille)
}
fn has_runtime_sparse_temporal_evidence(early: &TemporalSegment, late: &TemporalSegment) -> bool {
early.evidence_quality.runtime_snapshots != SignalCoverageStatus::Present
|| early.evidence_quality.inflight_snapshots != SignalCoverageStatus::Present
|| late.evidence_quality.runtime_snapshots != SignalCoverageStatus::Present
|| late.evidence_quality.inflight_snapshots != SignalCoverageStatus::Present
}
#[derive(Clone, Copy)]
struct TemporalMovement {
p95_shift: bool,
queue_move: bool,
service_move: bool,
}
fn is_runtime_dependent_suspect_shift(early: &TemporalSegment, late: &TemporalSegment) -> bool {
matches!(
(&early.primary_suspect.kind, &late.primary_suspect.kind),
(
DiagnosisKind::ExecutorPressureSuspected | DiagnosisKind::BlockingPoolPressure,
_
) | (
_,
DiagnosisKind::ExecutorPressureSuspected | DiagnosisKind::BlockingPoolPressure
)
)
}
fn has_material_suspect_shift(
early: &TemporalSegment,
late: &TemporalSegment,
movement: TemporalMovement,
options: &AnalyzeOptions,
) -> bool {
let suspect_shift_raw = early.primary_suspect.kind != late.primary_suspect.kind;
let runtime_sparse = has_runtime_sparse_temporal_evidence(early, late);
let runtime_dependent_shift = is_runtime_dependent_suspect_shift(early, late);
suspect_shift_raw
&& (!options
.temporal
.suppress_runtime_sparse_suspect_shift_without_supporting_movement
|| !runtime_sparse
|| !runtime_dependent_shift
|| movement.p95_shift
|| movement.queue_move
|| movement.service_move)
}
fn has_material_temporal_signal(
suspect_shift: bool,
movement: TemporalMovement,
options: &AnalyzeOptions,
) -> bool {
(options.temporal.emit_on_suspect_shift && suspect_shift)
|| movement.p95_shift
|| movement.queue_move
|| movement.service_move
}
pub(super) fn apply_temporal_overlap_attribution_warning(
early_seg: &mut TemporalSegment,
late_seg: &mut TemporalSegment,
) {
let windows_overlap = matches!(
(
early_seg.started_at_unix_ms,
early_seg.finished_at_unix_ms,
late_seg.started_at_unix_ms,
late_seg.finished_at_unix_ms,
),
(Some(_), Some(early_finish), Some(late_start), Some(_)) if early_finish >= late_start
);
let has_segment_runtime_or_inflight_samples = early_seg.evidence_quality.runtime_snapshot_count
> 0
|| early_seg.evidence_quality.inflight_snapshot_count > 0
|| late_seg.evidence_quality.runtime_snapshot_count > 0
|| late_seg.evidence_quality.inflight_snapshot_count > 0;
if windows_overlap && has_segment_runtime_or_inflight_samples {
early_seg
.warnings
.push(TEMPORAL_OVERLAP_ATTRIBUTION_WARNING.to_string());
late_seg
.warnings
.push(TEMPORAL_OVERLAP_ATTRIBUTION_WARNING.to_string());
}
}
pub(super) fn has_material_p95_shift(
left: Option<u64>,
right: Option<u64>,
options: &AnalyzeOptions,
) -> bool {
let (Some(a), Some(b)) = (left, right) else {
return false;
};
let lower = a.min(b);
let higher = a.max(b);
if lower == 0 {
return false;
}
higher.saturating_mul(options.temporal.p95_shift_ratio_denominator)
>= lower.saturating_mul(options.temporal.p95_shift_ratio_numerator)
}