use std::collections::HashMap;
use anyhow::{Context, Result};
use reqwest::Client;
use serde::Deserialize;
use super::GcpCredentials;
#[derive(Debug, Clone)]
pub struct RemoteCallPattern {
pub remote_service_name: String,
pub remote_endpoint: String,
pub observed_count: u32,
pub p99_latency_ms: Option<f64>,
pub local_callers: Vec<String>,
}
pub struct GcpTraceProvider {
project_id: String,
credentials: GcpCredentials,
}
impl GcpTraceProvider {
pub fn is_available() -> bool {
GcpCredentials::is_available()
}
pub fn from_env() -> Option<Self> {
let credentials = GcpCredentials::from_env()?;
let project_id = credentials.project_id.clone();
Some(Self {
project_id,
credentials,
})
}
pub async fn fetch_remote_calls(
&self,
client: &Client,
lookback_minutes: u32,
page_size: u32,
) -> Result<Vec<RemoteCallPattern>> {
let token = self
.credentials
.access_token(client)
.await
.context("Failed to acquire GCP access token for Cloud Trace")?;
let traces = self
.list_traces(client, &token, lookback_minutes, page_size)
.await?;
Ok(extract_remote_call_patterns(traces))
}
async fn list_traces(
&self,
client: &Client,
token: &str,
lookback_minutes: u32,
page_size: u32,
) -> Result<Vec<TraceV1>> {
let end_time = chrono_like_now_rfc3339();
let start_time = rfc3339_minus_minutes(&end_time, lookback_minutes);
let url = format!(
"https://cloudtrace.googleapis.com/v1/projects/{}/traces",
self.project_id
);
let mut all_traces: Vec<TraceV1> = Vec::new();
let mut page_token: Option<String> = None;
loop {
let mut req = client.get(&url).bearer_auth(token).query(&[
("view", "COMPLETE"),
("pageSize", &page_size.to_string()),
("startTime", &start_time),
("endTime", &end_time),
]);
if let Some(ref tok) = page_token {
req = req.query(&[("pageToken", tok.as_str())]);
}
let resp = req
.send()
.await
.context("Failed to call Cloud Trace v1 ListTraces API")?;
if !resp.status().is_success() {
let status = resp.status();
if status.as_u16() == 403 || status.as_u16() == 404 {
break;
}
let body = resp.text().await.unwrap_or_default();
anyhow::bail!("Cloud Trace ListTraces error: {} — {}", status, body);
}
let response: ListTracesResponse = resp
.json()
.await
.context("Failed to parse Cloud Trace response")?;
all_traces.extend(response.traces.unwrap_or_default());
match response.next_page_token {
Some(t) if !t.is_empty() && all_traces.len() < page_size as usize * 3 => {
page_token = Some(t);
}
_ => break,
}
}
Ok(all_traces)
}
}
fn extract_remote_call_patterns(traces: Vec<TraceV1>) -> Vec<RemoteCallPattern> {
let mut aggregation: HashMap<String, (u32, Vec<f64>, Vec<String>, String)> = HashMap::new();
for trace in &traces {
let span_map: HashMap<&str, &TraceSpanV1> = trace
.spans
.iter()
.map(|s| (s.span_id.as_str(), s))
.collect();
let own_hosts: std::collections::HashSet<String> = trace
.spans
.iter()
.filter(|s| {
s.labels
.get("/component")
.map(|c| c == "AppServer")
.unwrap_or(false)
})
.filter_map(|s| {
s.labels
.get("/http/host")
.or_else(|| s.labels.get("http.host"))
.map(|h| h.split(':').next().unwrap_or(h).to_lowercase())
})
.collect();
for span in &trace.spans {
let is_explicit_client = span
.kind
.as_deref()
.map(|k| k == "RPC_CLIENT" || k == "CLIENT")
.unwrap_or(false);
let is_external_http = if !is_explicit_client {
let host = span
.labels
.get("/http/host")
.or_else(|| span.labels.get("http.host"))
.map(|h| h.split(':').next().unwrap_or(h).to_lowercase());
host.as_ref().map(|h| !own_hosts.contains(h)).unwrap_or(false)
&& host.is_some()
&& span.labels.get("/component").map(|c| c != "AppServer").unwrap_or(true)
} else {
false
};
if !is_explicit_client && !is_external_http {
continue;
}
let remote_name = infer_remote_service_name(span);
if remote_name.is_empty() {
continue;
}
let remote_endpoint = span
.labels
.get("/http/url")
.or_else(|| span.labels.get("http.url"))
.map(|u| normalize_endpoint(u))
.unwrap_or_default();
let local_caller = span
.parent_span_id
.as_deref()
.and_then(|pid| span_map.get(pid))
.map(|p| p.name.clone())
.unwrap_or_default();
let latency_ms = compute_span_latency_ms(span);
let entry = aggregation
.entry(remote_name.clone())
.or_insert_with(|| (0, Vec::new(), Vec::new(), remote_endpoint.clone()));
entry.0 += 1;
if let Some(lat) = latency_ms {
entry.1.push(lat);
}
if !local_caller.is_empty() && !entry.2.contains(&local_caller) {
entry.2.push(local_caller);
}
if entry.3.is_empty() && !remote_endpoint.is_empty() {
entry.3 = remote_endpoint;
}
}
}
aggregation
.into_iter()
.map(|(name, (count, latencies, callers, endpoint))| {
let p99 = compute_p99(&latencies);
RemoteCallPattern {
remote_service_name: name,
remote_endpoint: endpoint,
observed_count: count,
p99_latency_ms: p99,
local_callers: callers,
}
})
.collect()
}
fn infer_remote_service_name(span: &TraceSpanV1) -> String {
if let Some(svc) = span
.labels
.get("peer.service")
.or_else(|| span.labels.get("peer_service"))
{
return svc.clone();
}
if let Some(host) = span
.labels
.get("/http/host")
.or_else(|| span.labels.get("http.host"))
{
return host_to_service_name(host);
}
if let Some(peer) = span.labels.get("net.peer.name") {
return host_to_service_name(peer);
}
let name = &span.name;
if let Some(rest) = name.strip_prefix("Sent.") {
return rest
.split('/')
.next()
.unwrap_or(rest)
.split('.')
.next_back()
.unwrap_or(rest)
.to_string();
}
if name.starts_with("grpc.") || name.contains("grpc/") {
if let Some(parts) = name.split('/').nth(1) {
return parts.split('.').next_back().unwrap_or(parts).to_string();
}
}
if let Some(url) = span
.labels
.get("/http/url")
.or_else(|| span.labels.get("http.url"))
{
if let Some(host) = extract_host_from_url(url) {
return host_to_service_name(&host);
}
}
String::new()
}
fn host_to_service_name(host: &str) -> String {
let host = host.split(':').next().unwrap_or(host);
let host = host.split(".svc.").next().unwrap_or(host);
const PUBLIC_TLDS: &[&str] = &[
".googleapis.com",
".google.com",
".github.com",
".github.io",
".amazonaws.com",
".azure.com",
".cloudflare.com",
".run.app", ".com",
".io",
".dev",
".net",
".org",
".app",
];
if PUBLIC_TLDS.iter().any(|tld| host.ends_with(tld)) {
return host.to_string();
}
host.split('.').next().unwrap_or(host).to_string()
}
fn extract_host_from_url(url: &str) -> Option<String> {
let after_scheme = url.split("://").nth(1)?;
let host_and_more = after_scheme.split('/').next()?;
Some(host_and_more.to_string())
}
fn normalize_endpoint(url: &str) -> String {
if let Some(after_scheme) = url.split("://").nth(1) {
let scheme = url.split("://").next().unwrap_or("https");
let host = after_scheme.split('/').next().unwrap_or(after_scheme);
return format!("{}://{}", scheme, host);
}
url.to_string()
}
fn compute_span_latency_ms(span: &TraceSpanV1) -> Option<f64> {
let start = parse_rfc3339_ms(&span.start_time)?;
let end = parse_rfc3339_ms(&span.end_time)?;
if end >= start {
Some((end - start) as f64)
} else {
None
}
}
fn compute_p99(latencies: &[f64]) -> Option<f64> {
if latencies.is_empty() {
return None;
}
let mut sorted = latencies.to_vec();
sorted.sort_by(|a, b| a.partial_cmp(b).unwrap_or(std::cmp::Ordering::Equal));
let idx = ((sorted.len() as f64) * 0.99).ceil() as usize;
sorted.get(idx.saturating_sub(1)).copied()
}
fn chrono_like_now_rfc3339() -> String {
use std::time::{SystemTime, UNIX_EPOCH};
let secs = SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap_or_default()
.as_secs();
unix_secs_to_rfc3339(secs)
}
fn rfc3339_minus_minutes(rfc: &str, minutes: u32) -> String {
if let Some(secs) = parse_rfc3339_epoch_secs(rfc) {
return unix_secs_to_rfc3339(secs.saturating_sub((minutes as u64) * 60));
}
rfc.to_string()
}
fn unix_secs_to_rfc3339(secs: u64) -> String {
let s = secs % 60;
let m = (secs / 60) % 60;
let h = (secs / 3600) % 24;
let days = secs / 86400;
let (year, month, day) = days_to_ymd(days);
format!(
"{:04}-{:02}-{:02}T{:02}:{:02}:{:02}Z",
year, month, day, h, m, s
)
}
fn days_to_ymd(days: u64) -> (u64, u64, u64) {
let mut remaining = days;
let mut year = 1970u64;
loop {
let year_days = if is_leap(year) { 366 } else { 365 };
if remaining < year_days {
break;
}
remaining -= year_days;
year += 1;
}
let month_days: [u64; 12] = if is_leap(year) {
[31, 29, 31, 30, 31, 30, 31, 31, 30, 31, 30, 31]
} else {
[31, 28, 31, 30, 31, 30, 31, 31, 30, 31, 30, 31]
};
let mut month = 1u64;
for &md in &month_days {
if remaining < md {
break;
}
remaining -= md;
month += 1;
}
(year, month, remaining + 1)
}
fn is_leap(year: u64) -> bool {
(year % 4 == 0 && year % 100 != 0) || year % 400 == 0
}
fn parse_rfc3339_epoch_secs(s: &str) -> Option<u64> {
let s = s.trim_end_matches('Z');
let parts: Vec<&str> = s.split('T').collect();
if parts.len() != 2 {
return None;
}
let date_parts: Vec<u64> = parts[0].split('-').filter_map(|p| p.parse().ok()).collect();
let time_parts: Vec<u64> = parts[1].split(':').filter_map(|p| p.parse().ok()).collect();
if date_parts.len() < 3 || time_parts.len() < 3 {
return None;
}
let (y, mo, d) = (date_parts[0], date_parts[1], date_parts[2]);
let (h, mi, sec) = (time_parts[0], time_parts[1], time_parts[2]);
let mut days = 0u64;
for yr in 1970..y {
days += if is_leap(yr) { 366 } else { 365 };
}
let month_days: [u64; 12] = if is_leap(y) {
[31, 29, 31, 30, 31, 30, 31, 31, 30, 31, 30, 31]
} else {
[31, 28, 31, 30, 31, 30, 31, 31, 30, 31, 30, 31]
};
for mi_idx in 0..(mo as usize - 1) {
days += month_days[mi_idx];
}
days += d - 1;
Some(days * 86400 + h * 3600 + mi * 60 + sec)
}
fn parse_rfc3339_ms(s: &str) -> Option<u64> {
let s_trimmed = s.trim_end_matches('Z');
let (main_part, frac_part) = if let Some(dot_pos) = s_trimmed.rfind('.') {
(&s_trimmed[..dot_pos], &s_trimmed[dot_pos + 1..])
} else {
(s_trimmed, "")
};
let epoch_secs = parse_rfc3339_epoch_secs(&format!("{}Z", main_part))?;
let frac_ms = if frac_part.is_empty() {
0u64
} else {
let digits = &frac_part[..frac_part.len().min(3)];
let val: u64 = digits.parse().unwrap_or(0);
val * 10u64.pow(3u32.saturating_sub(digits.len() as u32))
};
Some(epoch_secs * 1000 + frac_ms)
}
#[derive(Debug, Deserialize, Default)]
struct ListTracesResponse {
#[serde(default)]
traces: Option<Vec<TraceV1>>,
#[serde(rename = "nextPageToken", default)]
next_page_token: Option<String>,
}
#[derive(Debug, Deserialize)]
struct TraceV1 {
#[serde(rename = "traceId", default)]
_trace_id: String,
#[serde(default)]
spans: Vec<TraceSpanV1>,
}
#[derive(Debug, Deserialize)]
#[serde(rename_all = "camelCase")]
struct TraceSpanV1 {
#[serde(default)]
span_id: String,
#[serde(default)]
kind: Option<String>,
#[serde(default)]
name: String,
#[serde(default)]
start_time: String,
#[serde(default)]
end_time: String,
#[serde(default)]
parent_span_id: Option<String>,
#[serde(default)]
labels: HashMap<String, String>,
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn host_to_service_name_strips_k8s_fqdn() {
assert_eq!(
host_to_service_name("inventory-svc.prod.svc.cluster.local:8080"),
"inventory-svc"
);
}
#[test]
fn host_to_service_name_keeps_gcp_apis() {
assert_eq!(
host_to_service_name("pubsub.googleapis.com"),
"pubsub.googleapis.com"
);
}
#[test]
fn host_to_service_name_strips_port() {
assert_eq!(host_to_service_name("payments-svc:9090"), "payments-svc");
}
#[test]
fn normalize_endpoint_strips_path() {
assert_eq!(
normalize_endpoint("https://inventory-svc:8080/api/v1/check"),
"https://inventory-svc:8080"
);
}
#[test]
fn extract_host_from_url_works() {
assert_eq!(
extract_host_from_url("https://payments-svc:9090/pay"),
Some("payments-svc:9090".to_string())
);
}
#[test]
fn unix_secs_to_rfc3339_round_trips() {
let secs = parse_rfc3339_epoch_secs("2026-04-15T00:00:00Z").unwrap();
let back = unix_secs_to_rfc3339(secs);
assert_eq!(back, "2026-04-15T00:00:00Z");
}
#[test]
fn rfc3339_minus_minutes_works() {
let base = "2026-04-15T01:00:00Z";
let result = rfc3339_minus_minutes(base, 60);
assert_eq!(result, "2026-04-15T00:00:00Z");
}
#[test]
fn parse_rfc3339_ms_with_fractional() {
let ms = parse_rfc3339_ms("2026-01-01T00:00:01.500Z").unwrap();
let base_ms = parse_rfc3339_ms("2026-01-01T00:00:01Z").unwrap();
assert_eq!(ms - base_ms, 500);
}
fn make_span(
id: &str,
kind: Option<&str>,
name: &str,
host: Option<&str>,
url: Option<&str>,
component: Option<&str>,
parent: Option<&str>,
) -> TraceSpanV1 {
let mut labels = HashMap::new();
if let Some(h) = host {
labels.insert("/http/host".to_string(), h.to_string());
}
if let Some(u) = url {
labels.insert("/http/url".to_string(), u.to_string());
}
if let Some(c) = component {
labels.insert("/component".to_string(), c.to_string());
}
TraceSpanV1 {
span_id: id.to_string(),
kind: kind.map(|s| s.to_string()),
name: name.to_string(),
start_time: "2026-04-15T00:00:00Z".to_string(),
end_time: "2026-04-15T00:00:00.100Z".to_string(),
parent_span_id: parent.map(|s| s.to_string()),
labels,
}
}
#[test]
fn infer_remote_service_name_from_peer_service_label() {
let mut labels = HashMap::new();
labels.insert("peer.service".to_string(), "inventory-service".to_string());
let span = TraceSpanV1 {
span_id: "1".to_string(),
kind: Some("RPC_CLIENT".to_string()),
name: "HTTP POST".to_string(),
start_time: "2026-04-15T00:00:00Z".to_string(),
end_time: "2026-04-15T00:00:00.100Z".to_string(),
parent_span_id: None,
labels,
};
assert_eq!(infer_remote_service_name(&span), "inventory-service");
}
#[test]
fn infer_remote_service_name_from_sent_prefix() {
let span = TraceSpanV1 {
span_id: "2".to_string(),
kind: Some("RPC_CLIENT".to_string()),
name: "Sent.payments.PaymentService/Charge".to_string(),
start_time: "2026-04-15T00:00:00Z".to_string(),
end_time: "2026-04-15T00:00:00.200Z".to_string(),
parent_span_id: None,
labels: HashMap::new(),
};
assert_eq!(infer_remote_service_name(&span), "PaymentService");
}
#[test]
fn extract_remote_call_patterns_aggregates_by_explicit_kind() {
let mut labels = HashMap::new();
labels.insert("peer.service".to_string(), "inventory".to_string());
labels.insert(
"/http/url".to_string(),
"https://inventory-svc:8080/check".to_string(),
);
let spans = vec![
TraceSpanV1 {
span_id: "10".to_string(),
kind: Some("RPC_CLIENT".to_string()),
name: "GET /check".to_string(),
start_time: "2026-04-15T00:00:00Z".to_string(),
end_time: "2026-04-15T00:00:00.050Z".to_string(),
parent_span_id: None,
labels: labels.clone(),
},
TraceSpanV1 {
span_id: "11".to_string(),
kind: Some("RPC_CLIENT".to_string()),
name: "GET /check".to_string(),
start_time: "2026-04-15T00:00:01Z".to_string(),
end_time: "2026-04-15T00:00:01.200Z".to_string(),
parent_span_id: None,
labels,
},
];
let traces = vec![TraceV1 {
_trace_id: "abc".to_string(),
spans,
}];
let patterns = extract_remote_call_patterns(traces);
assert_eq!(patterns.len(), 1);
assert_eq!(patterns[0].remote_service_name, "inventory");
assert_eq!(patterns[0].observed_count, 2);
assert!(patterns[0].p99_latency_ms.is_some());
}
#[test]
fn extract_patterns_via_host_inference_excludes_own_service() {
let inbound = make_span(
"1",
None,
"/",
Some("my-service.run.app"),
Some("https://my-service.run.app/"),
Some("AppServer"),
None,
);
let outbound = make_span(
"2",
None,
"GET",
Some("api.github.com"),
Some("https://api.github.com/"),
None,
Some("1"),
);
let traces = vec![TraceV1 {
_trace_id: "t1".to_string(),
spans: vec![inbound, outbound],
}];
let patterns = extract_remote_call_patterns(traces);
assert_eq!(patterns.len(), 1);
assert_eq!(patterns[0].remote_service_name, "api.github.com");
}
#[test]
fn extract_patterns_does_not_emit_own_service_as_remote() {
let inbound = make_span(
"1",
None,
"/",
Some("my-service.run.app"),
Some("https://my-service.run.app/"),
Some("AppServer"),
None,
);
let traces = vec![TraceV1 {
_trace_id: "t2".to_string(),
spans: vec![inbound],
}];
let patterns = extract_remote_call_patterns(traces);
assert!(patterns.is_empty());
}
#[test]
fn empty_response_deserializes() {
let json = "{}";
let resp: ListTracesResponse = serde_json::from_str(json).unwrap();
assert!(resp.traces.is_none());
assert!(resp.next_page_token.is_none());
}
#[test]
fn span_id_as_string_deserializes() {
let json = r#"{"spanId": "8686581962470036554", "name": "/"}"#;
let span: TraceSpanV1 = serde_json::from_str(json).unwrap();
assert_eq!(span.span_id, "8686581962470036554");
}
#[test]
fn is_available_does_not_panic() {
let _ = GcpTraceProvider::is_available();
}
}