Skip to main content

shardline_server/
backup.rs

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/// Adapter-neutral backup manifest summary.
21#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
22pub struct BackupManifestReport {
23    /// Stable manifest format version.
24    pub manifest_version: u64,
25    /// Metadata backend used by the deployment.
26    pub metadata_backend: String,
27    /// Object backend used by the deployment.
28    pub object_backend: String,
29    /// Number of object-store entries written to the manifest.
30    pub object_count: u64,
31    /// Total bytes reported by object-store metadata.
32    pub object_bytes: u64,
33    /// Number of visible latest records in metadata storage.
34    pub latest_records: u64,
35    /// Number of immutable version records in metadata storage.
36    pub version_records: u64,
37    /// Number of durable reconstruction rows in metadata storage.
38    pub reconstruction_rows: u64,
39    /// Number of retained dedupe-shard mappings in metadata storage.
40    pub dedupe_shard_mappings: u64,
41    /// Number of durable quarantine candidates in metadata storage.
42    pub quarantine_candidates: u64,
43    /// Number of durable retention holds in metadata storage.
44    pub retention_holds: u64,
45    /// Number of processed provider webhook delivery claims in metadata storage.
46    pub webhook_deliveries: u64,
47    /// Number of provider repository lifecycle states in metadata storage.
48    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
88/// Writes an adapter-neutral backup manifest for the configured deployment.
89///
90/// The manifest inventories object metadata and durable index state without reading object
91/// bodies. Payload bytes remain in the configured object store; operators should combine
92/// this manifest with the storage backend's native backup or replication mechanism.
93///
94/// # Errors
95///
96/// Returns [`ServerError`] when metadata enumeration, object inventory, or manifest
97/// writing fails.
98pub 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}