photon_runtime/admin/
snapshot.rs1use 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
18pub 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}