1use std::io::Write;
2
3use serde::Serialize;
4use serde_json::to_writer;
5use shardline_index::{
6 AsyncIndexStore, LocalIndexStore, PostgresIndexStore, PostgresRecordStore, RecordTraversal,
7 xet_hash_hex_string,
8};
9use shardline_storage::{ObjectMetadata, ObjectPrefix};
10
11use crate::{
12 ServerConfig, ServerError,
13 object_store::{ServerObjectStore, object_store_from_config},
14 ops_record_store::OpsRecordStore,
15 overflow::{checked_add, checked_increment},
16 postgres_backend::connect_postgres_metadata_pool,
17 record_store::LocalRecordStore,
18};
19
20#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
22pub struct BackupManifestReport {
23 pub manifest_version: u64,
25 pub metadata_backend: String,
27 pub object_backend: String,
29 pub object_count: u64,
31 pub object_bytes: u64,
33 pub latest_records: u64,
35 pub version_records: u64,
37 pub reconstruction_rows: u64,
39 pub dedupe_shard_mappings: u64,
41 pub quarantine_candidates: u64,
43 pub retention_holds: u64,
45 pub webhook_deliveries: u64,
47 pub provider_repository_states: u64,
49}
50
51impl BackupManifestReport {
52 fn new(metadata_backend: &str, object_backend: &str) -> Self {
53 Self {
54 manifest_version: 1,
55 metadata_backend: metadata_backend.to_owned(),
56 object_backend: object_backend.to_owned(),
57 object_count: 0,
58 object_bytes: 0,
59 latest_records: 0,
60 version_records: 0,
61 reconstruction_rows: 0,
62 dedupe_shard_mappings: 0,
63 quarantine_candidates: 0,
64 retention_holds: 0,
65 webhook_deliveries: 0,
66 provider_repository_states: 0,
67 }
68 }
69}
70
71#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
72struct BackupManifestObjectEntry {
73 key: String,
74 length: u64,
75 checksum: Option<String>,
76}
77
78impl BackupManifestObjectEntry {
79 fn from_metadata(metadata: &ObjectMetadata) -> Self {
80 Self {
81 key: metadata.key().as_str().to_owned(),
82 length: metadata.length(),
83 checksum: metadata.checksum().map(xet_hash_hex_string),
84 }
85 }
86}
87
88pub async fn write_backup_manifest<Writer>(
99 config: ServerConfig,
100 writer: Writer,
101) -> Result<BackupManifestReport, ServerError>
102where
103 Writer: Write,
104{
105 let object_store = object_store_from_config(&config)?;
106 let object_backend = object_store.backend_name();
107
108 if let Some(index_postgres_url) = config.index_postgres_url() {
109 let pool = connect_postgres_metadata_pool(index_postgres_url, 4)?;
110 let index_store = PostgresIndexStore::new(pool.clone());
111 let record_store = PostgresRecordStore::new(pool);
112 let mut report = BackupManifestReport::new("postgres", object_backend);
113 collect_metadata_counts(&record_store, &index_store, &mut report).await?;
114 write_manifest_body(writer, &object_store, report)
115 } else {
116 let index_store = LocalIndexStore::open(config.root_dir().to_path_buf());
117 let record_store = LocalRecordStore::open(config.root_dir().to_path_buf());
118 let mut report = BackupManifestReport::new("local", object_backend);
119 collect_metadata_counts(&record_store, &index_store, &mut report).await?;
120 write_manifest_body(writer, &object_store, report)
121 }
122}
123
124async fn collect_metadata_counts<RecordAdapter, IndexAdapter>(
125 record_store: &RecordAdapter,
126 index_store: &IndexAdapter,
127 report: &mut BackupManifestReport,
128) -> Result<(), ServerError>
129where
130 RecordAdapter: OpsRecordStore + Sync,
131 RecordAdapter::Error: Into<ServerError>,
132 IndexAdapter: AsyncIndexStore + Sync,
133 IndexAdapter::Error: Into<ServerError>,
134{
135 RecordTraversal::visit_latest_records(record_store, |_entry| {
136 report.latest_records = checked_increment(report.latest_records)?;
137 Ok::<(), ServerError>(())
138 })
139 .await?;
140
141 RecordTraversal::visit_version_records(record_store, |_entry| {
142 report.version_records = checked_increment(report.version_records)?;
143 Ok::<(), ServerError>(())
144 })
145 .await?;
146
147 report.reconstruction_rows = u64::try_from(
148 index_store
149 .list_reconstruction_file_ids()
150 .await
151 .map_err(Into::into)?
152 .len(),
153 )?;
154
155 index_store
156 .visit_dedupe_shard_mappings(|_mapping| {
157 report.dedupe_shard_mappings = checked_increment(report.dedupe_shard_mappings)?;
158 Ok::<(), ServerError>(())
159 })
160 .await?;
161 index_store
162 .visit_quarantine_candidates(|_candidate| {
163 report.quarantine_candidates = checked_increment(report.quarantine_candidates)?;
164 Ok::<(), ServerError>(())
165 })
166 .await?;
167 index_store
168 .visit_retention_holds(|_hold| {
169 report.retention_holds = checked_increment(report.retention_holds)?;
170 Ok::<(), ServerError>(())
171 })
172 .await?;
173 index_store
174 .visit_webhook_deliveries(|_delivery| {
175 report.webhook_deliveries = checked_increment(report.webhook_deliveries)?;
176 Ok::<(), ServerError>(())
177 })
178 .await?;
179 index_store
180 .visit_provider_repository_states(|_state| {
181 report.provider_repository_states =
182 checked_increment(report.provider_repository_states)?;
183 Ok::<(), ServerError>(())
184 })
185 .await?;
186
187 Ok(())
188}
189
190fn write_manifest_body<Writer>(
191 mut writer: Writer,
192 object_store: &ServerObjectStore,
193 mut report: BackupManifestReport,
194) -> Result<BackupManifestReport, ServerError>
195where
196 Writer: Write,
197{
198 let prefix = ObjectPrefix::parse("")?;
199 let mut first_field = true;
200
201 writer.write_all(b"{")?;
202 write_named_value(
203 &mut writer,
204 "manifest_version",
205 &report.manifest_version,
206 &mut first_field,
207 )?;
208 write_named_value(
209 &mut writer,
210 "metadata_backend",
211 &report.metadata_backend,
212 &mut first_field,
213 )?;
214 write_named_value(
215 &mut writer,
216 "object_backend",
217 &report.object_backend,
218 &mut first_field,
219 )?;
220 write_named_value(
221 &mut writer,
222 "latest_records",
223 &report.latest_records,
224 &mut first_field,
225 )?;
226 write_named_value(
227 &mut writer,
228 "version_records",
229 &report.version_records,
230 &mut first_field,
231 )?;
232 write_named_value(
233 &mut writer,
234 "reconstruction_rows",
235 &report.reconstruction_rows,
236 &mut first_field,
237 )?;
238 write_named_value(
239 &mut writer,
240 "dedupe_shard_mappings",
241 &report.dedupe_shard_mappings,
242 &mut first_field,
243 )?;
244 write_named_value(
245 &mut writer,
246 "quarantine_candidates",
247 &report.quarantine_candidates,
248 &mut first_field,
249 )?;
250 write_named_value(
251 &mut writer,
252 "retention_holds",
253 &report.retention_holds,
254 &mut first_field,
255 )?;
256 write_named_value(
257 &mut writer,
258 "webhook_deliveries",
259 &report.webhook_deliveries,
260 &mut first_field,
261 )?;
262 write_named_value(
263 &mut writer,
264 "provider_repository_states",
265 &report.provider_repository_states,
266 &mut first_field,
267 )?;
268 write_field_name(&mut writer, "objects", &mut first_field)?;
269 writer.write_all(b"[")?;
270
271 let mut first_object = true;
272 crate::object_store::visit_object_prefix(object_store, &prefix, |metadata| {
273 if first_object {
274 first_object = false;
275 } else {
276 writer.write_all(b",")?;
277 }
278
279 let entry = BackupManifestObjectEntry::from_metadata(&metadata);
280 report.object_count = checked_increment(report.object_count)?;
281 report.object_bytes = checked_add(report.object_bytes, entry.length)?;
282 to_writer(&mut writer, &entry)?;
283 Ok(())
284 })?;
285
286 writer.write_all(b"]")?;
287 write_named_value(
288 &mut writer,
289 "object_count",
290 &report.object_count,
291 &mut first_field,
292 )?;
293 write_named_value(
294 &mut writer,
295 "object_bytes",
296 &report.object_bytes,
297 &mut first_field,
298 )?;
299 writer.write_all(b"}\n")?;
300
301 Ok(report)
302}
303
304fn write_named_value<Writer, Value>(
305 writer: &mut Writer,
306 name: &str,
307 value: &Value,
308 first_field: &mut bool,
309) -> Result<(), ServerError>
310where
311 Writer: Write,
312 Value: Serialize,
313{
314 write_field_name(writer, name, first_field)?;
315 to_writer(writer, value)?;
316 Ok(())
317}
318
319fn write_field_name<Writer>(
320 writer: &mut Writer,
321 name: &str,
322 first_field: &mut bool,
323) -> Result<(), ServerError>
324where
325 Writer: Write,
326{
327 if *first_field {
328 *first_field = false;
329 } else {
330 writer.write_all(b",")?;
331 }
332
333 to_writer(&mut *writer, name)?;
334 writer.write_all(b":")?;
335 Ok(())
336}