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        supports_list_events: caps.supports_list_events,
93        max_replay_window_secs: caps.max_replay_window.map(|d| d.as_secs()),
94    }
95}
96
97async fn collect_checkpoints(
98    photon: &Photon,
99) -> photon_backend::Result<Vec<AdminCheckpointSummary>> {
100    let registry = HandlerRegistry::auto_discover();
101    let mut out = Vec::new();
102
103    for handler in registry.iter() {
104        if handler.is_consumer_group() {
105            let group_id = handler.consumer_group.unwrap_or("");
106            let shard_count = group_shard_count(photon, handler);
107            for shard_id in 0..shard_count {
108                let shard_key = shard_storage_key(shard_id);
109                let last_seq = photon
110                    .get_checkpoint_seq(group_id, handler.topic_name, Some(&shard_key))
111                    .await?;
112                out.push(AdminCheckpointSummary {
113                    subscription_name: group_id.to_string(),
114                    topic_name: handler.topic_name.to_string(),
115                    topic_key: Some(shard_key),
116                    last_seq,
117                });
118            }
119        } else {
120            let last_seq = photon
121                .get_checkpoint_seq(handler.subscription_name, handler.topic_name, None)
122                .await?;
123            out.push(AdminCheckpointSummary {
124                subscription_name: handler.subscription_name.to_string(),
125                topic_name: handler.topic_name.to_string(),
126                topic_key: None,
127                last_seq,
128            });
129        }
130    }
131
132    Ok(out)
133}
134
135fn group_shard_count(photon: &Photon, handler: &'static HandlerDescriptor) -> u32 {
136    handler
137        .group_shard_count
138        .or_else(|| {
139            photon
140                .registry()
141                .get(handler.topic_name)
142                .and_then(|d| d.shard_config)
143                .map(|c| c.shard_count)
144        })
145        .unwrap_or(ShardConfig::DEFAULT_SHARD_COUNT)
146}
147
148#[cfg(test)]
149mod tests {
150    use super::*;
151
152    #[test]
153    fn sanitize_admin_error_redacts_secrets_happy_path() {
154        let err = PhotonError::Internal(
155            "checkpoint load failed password=hunter2 nats://u:p@host:4222".into(),
156        );
157        let safe = sanitize_admin_error(&err).to_string();
158        assert!(safe.contains("[redacted]"), "safe: {safe}");
159        assert!(!safe.contains("hunter2"), "safe: {safe}");
160        assert!(!safe.contains("u:p@"), "safe: {safe}");
161    }
162
163    #[test]
164    fn sanitize_admin_error_preserves_benign_detail_sad_path() {
165        let err = PhotonError::Internal("checkpoint missing".into());
166        let safe = sanitize_admin_error(&err).to_string();
167        assert!(safe.contains("checkpoint missing"), "safe: {safe}");
168    }
169}