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 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}