Skip to main content

photon_runtime/admin/
snapshot.rs

1//! Compose [`AdminSnapshot`] from registries and storage checkpoint reads.
2//!
3//! Snapshots are metadata-only (topics, handlers, checkpoints). They never include
4//! event `actor_json` / `payload_json` or raw DLQ error bodies.
5
6use photon_backend::{
7    redact_credentials_in_text, sanitize_error_message, shard_storage_key, BackendCapabilities,
8    HandlerDescriptor, HandlerRegistry, PhotonError, ShardConfig, TopicDescriptor,
9};
10
11use crate::Photon;
12
13use super::types::{
14    AdminBackendSummary, AdminCheckpointSummary, AdminHandlerSummary, AdminSnapshot,
15    AdminTopicSummary,
16};
17
18/// Build a full admin snapshot for the given runtime (async checkpoint reads).
19///
20/// # Errors
21///
22/// Returns an error if a checkpoint load fails. Error text is sanitized before return.
23pub async fn collect_admin_snapshot(photon: &Photon) -> photon_backend::Result<AdminSnapshot> {
24    let topics = collect_topics(photon);
25    let handlers = collect_handlers();
26    let backend = collect_backend(photon.backend_capabilities());
27    let checkpoints = collect_checkpoints(photon)
28        .await
29        .map_err(|e| sanitize_admin_error(&e))?;
30
31    Ok(AdminSnapshot {
32        backend,
33        topics,
34        handlers,
35        checkpoints,
36    })
37}
38
39fn sanitize_admin_error(err: &PhotonError) -> PhotonError {
40    PhotonError::Internal(sanitize_error_message(&redact_credentials_in_text(
41        &err.to_string(),
42    )))
43}
44
45fn collect_topics(photon: &Photon) -> Vec<AdminTopicSummary> {
46    photon
47        .registry()
48        .sorted_topic_names()
49        .into_iter()
50        .filter_map(|name| photon.registry().get(name).map(topic_from_descriptor))
51        .collect()
52}
53
54fn topic_from_descriptor(descriptor: &TopicDescriptor) -> AdminTopicSummary {
55    AdminTopicSummary {
56        topic_name: descriptor.topic_name.to_string(),
57        keyed_by: descriptor.keyed_by.map(str::to_string),
58        schema_json: descriptor
59            .schema_value()
60            .unwrap_or_else(|| serde_json::json!({})),
61    }
62}
63
64fn collect_handlers() -> Vec<AdminHandlerSummary> {
65    let registry = HandlerRegistry::auto_discover();
66    registry.iter().map(handler_from_descriptor).collect()
67}
68
69fn handler_from_descriptor(handler: &HandlerDescriptor) -> AdminHandlerSummary {
70    let is_group = handler.is_consumer_group();
71    AdminHandlerSummary {
72        topic_name: handler.topic_name.to_string(),
73        subscription_name: if is_group || handler.subscription_name.is_empty() {
74            None
75        } else {
76            Some(handler.subscription_name.to_string())
77        },
78        consumer_group: handler.consumer_group.map(str::to_string),
79        registry_key: handler.registry_key.to_string(),
80        mode: if is_group {
81            "consumer_group".to_string()
82        } else {
83            "durable".to_string()
84        },
85    }
86}
87
88fn collect_backend(caps: BackendCapabilities) -> AdminBackendSummary {
89    AdminBackendSummary {
90        telemetry_label: caps.telemetry_label.to_string(),
91        supports_get_event: caps.supports_get_event,
92        max_replay_window_secs: caps.max_replay_window.map(|d| d.as_secs()),
93    }
94}
95
96async fn collect_checkpoints(
97    photon: &Photon,
98) -> photon_backend::Result<Vec<AdminCheckpointSummary>> {
99    let registry = HandlerRegistry::auto_discover();
100    let mut out = Vec::new();
101
102    for handler in registry.iter() {
103        if handler.is_consumer_group() {
104            let group_id = handler.consumer_group.unwrap_or("");
105            let shard_count = group_shard_count(photon, handler);
106            for shard_id in 0..shard_count {
107                let shard_key = shard_storage_key(shard_id);
108                let last_seq = photon
109                    .get_checkpoint_seq(group_id, handler.topic_name, Some(&shard_key))
110                    .await?;
111                out.push(AdminCheckpointSummary {
112                    subscription_name: group_id.to_string(),
113                    topic_name: handler.topic_name.to_string(),
114                    topic_key: Some(shard_key),
115                    last_seq,
116                });
117            }
118        } else {
119            let last_seq = photon
120                .get_checkpoint_seq(handler.subscription_name, handler.topic_name, None)
121                .await?;
122            out.push(AdminCheckpointSummary {
123                subscription_name: handler.subscription_name.to_string(),
124                topic_name: handler.topic_name.to_string(),
125                topic_key: None,
126                last_seq,
127            });
128        }
129    }
130
131    Ok(out)
132}
133
134fn group_shard_count(photon: &Photon, handler: &'static HandlerDescriptor) -> u32 {
135    handler
136        .group_shard_count
137        .or_else(|| {
138            photon
139                .registry()
140                .get(handler.topic_name)
141                .and_then(|d| d.shard_config)
142                .map(|c| c.shard_count)
143        })
144        .unwrap_or(ShardConfig::DEFAULT_SHARD_COUNT)
145}
146
147#[cfg(test)]
148mod tests {
149    use super::*;
150
151    #[test]
152    fn sanitize_admin_error_redacts_secrets_happy_path() {
153        let err = PhotonError::Internal(
154            "checkpoint load failed password=hunter2 nats://u:p@host:4222".into(),
155        );
156        let safe = sanitize_admin_error(&err).to_string();
157        assert!(safe.contains("[redacted]"), "safe: {safe}");
158        assert!(!safe.contains("hunter2"), "safe: {safe}");
159        assert!(!safe.contains("u:p@"), "safe: {safe}");
160    }
161
162    #[test]
163    fn sanitize_admin_error_preserves_benign_detail_sad_path() {
164        let err = PhotonError::Internal("checkpoint missing".into());
165        let safe = sanitize_admin_error(&err).to_string();
166        assert!(safe.contains("checkpoint missing"), "safe: {safe}");
167    }
168}