Skip to main content

icydb_core/db/session/
catalog.rs

1//! Module: db::session::catalog
2//! Responsibility: session-owned catalog, schema-description, and storage
3//! observability surfaces.
4//! Does not own: schema reconciliation policy, query planning, or storage
5//! mutation.
6//! Boundary: converts accepted/generated schema authority into stable
7//! introspection DTOs at the session facade.
8
9#[cfg(feature = "sql")]
10use crate::db::schema::show_indexes_for_schema_info_with_runtime_state;
11use crate::{
12    db::{
13        DbSession, EntityCatalogCounts, EntityCatalogDescription, EntityIdentityDescription,
14        EntitySchemaDescription, IndexState, QueryError, SchemaApplicationTarget,
15        SchemaChangeJobId, SchemaChangeProgress, SchemaChangeReceipt, StorageReport,
16        StoreCatalogDescription,
17        commit::database_incarnation_id,
18        query::plan::VisibleIndexes,
19        schema::{
20            AcceptedEntityDescriptionMetadata, ConstraintValidationJob, SchemaInfo,
21            describe_accepted_entity_with_persisted_schema, describe_accepted_identity,
22        },
23    },
24    error::InternalError,
25    traits::CanisterKind,
26};
27use icydb_schema::{SchemaProposal, SchemaSubmissionKey, TargetDatabaseIdentity};
28
29#[cfg(feature = "migration")]
30use crate::db::{SchemaMigrationCommand, SchemaMigrationStatusPage, SchemaMigrationStatusRequest};
31
32impl<C: CanisterKind> DbSession<C> {
33    /// Require generated schema application to wait for an exact active migration
34    /// to finish. This does not change ordinary row-operation admission.
35    #[cfg(feature = "migration")]
36    pub fn ensure_generated_schema_application_admitted(
37        &self,
38        proposal: &SchemaProposal,
39    ) -> Result<(), InternalError> {
40        crate::db::schema::ensure_generated_schema_application_admitted(&self.db, proposal)
41    }
42
43    /// Execute one explicit metadata source-migration operation.
44    #[cfg(feature = "migration")]
45    pub fn migrate_schema(
46        &self,
47        proposal: &SchemaProposal,
48        command: SchemaMigrationCommand,
49    ) -> Result<SchemaMigrationStatusPage, InternalError> {
50        crate::db::schema::migrate_schema(&self.db, proposal, command)
51    }
52
53    /// Return one bounded source-migration status page.
54    #[cfg(feature = "migration")]
55    pub fn schema_migration_status(
56        &self,
57        proposal: &SchemaProposal,
58        request: &SchemaMigrationStatusRequest,
59    ) -> Result<SchemaMigrationStatusPage, InternalError> {
60        crate::db::schema::schema_migration_status(&self.db, proposal, request)
61    }
62
63    /// Apply one canonical generated proposal after proving it contains no
64    /// explicit removals.
65    #[doc(hidden)]
66    pub fn apply_generated_schema(
67        &self,
68        proposal: &SchemaProposal,
69    ) -> Result<SchemaChangeReceipt, InternalError> {
70        crate::db::schema::apply_generated_schema(&self.db, proposal)
71    }
72
73    /// Apply one exact source-keyed schema proposal through accepted catalog
74    /// authority and return its durable idempotent receipt.
75    pub fn apply_schema(
76        &self,
77        proposal: &SchemaProposal,
78    ) -> Result<SchemaChangeReceipt, InternalError> {
79        crate::db::schema::apply_schema(&self.db, proposal)
80    }
81
82    /// Issue the opaque database/store identities and exact accepted head used
83    /// to compose one optimistic schema proposal.
84    pub fn schema_application_target(&self) -> Result<SchemaApplicationTarget, InternalError> {
85        crate::db::schema::schema_application_target(&self.db)
86    }
87
88    /// Load one durable schema-application receipt by exact target and
89    /// submission identity.
90    pub fn schema_application_receipt(
91        &self,
92        database_identity: TargetDatabaseIdentity,
93        submission_key: &SchemaSubmissionKey,
94    ) -> Result<Option<SchemaChangeReceipt>, InternalError> {
95        crate::db::schema::schema_application_receipt(&self.db, database_identity, submission_key)
96    }
97
98    /// Advance one pending schema application by at most one bounded
99    /// activation step.
100    pub fn continue_schema_application(
101        &self,
102        job_id: SchemaChangeJobId,
103        acknowledged_receipt: Option<u64>,
104    ) -> Result<SchemaChangeProgress, InternalError> {
105        crate::db::schema::continue_schema_application(&self.db, job_id, acknowledged_receipt)
106    }
107
108    /// Abort one pending schema application after acknowledging any retained
109    /// finding page by exact sequence.
110    pub fn abort_schema_application(
111        &self,
112        job_id: SchemaChangeJobId,
113        acknowledged_receipt: Option<u64>,
114    ) -> Result<SchemaChangeProgress, InternalError> {
115        crate::db::schema::abort_schema_application(&self.db, job_id, acknowledged_receipt)
116    }
117
118    // Return one stable, human-readable index listing for one resolved
119    // store/accepted-schema pair, attaching the current runtime lifecycle state
120    // when the registry can resolve the backing store handle.
121    #[cfg(feature = "sql")]
122    pub(in crate::db) fn show_indexes_for_store_schema_info(
123        &self,
124        store_path: &str,
125        schema: &SchemaInfo,
126        snapshot: &crate::db::schema::PersistedSchemaSnapshot,
127    ) -> Vec<String> {
128        let runtime_state = self
129            .db
130            .with_store_registry(|registry| registry.try_get_store(store_path).ok())
131            .map(|store| store.index_state());
132
133        show_indexes_for_schema_info_with_runtime_state(schema, snapshot, runtime_state)
134    }
135
136    /// Return one stable list of accepted runtime entity catalog entries.
137    pub fn show_entities(&self) -> Result<Vec<EntityCatalogDescription>, InternalError> {
138        let runtime_entities = self.db.accepted_runtime_entities()?;
139        let mut entities = Vec::with_capacity(runtime_entities.len());
140
141        for runtime_entity in runtime_entities {
142            let store = self.db.recovered_store(runtime_entity.store_path())?;
143            let storage = store
144                .storage_capabilities()
145                .storage_mode()
146                .as_str()
147                .to_string();
148            let accepted = self.accepted_schema_catalog_context_for_runtime_entity(
149                runtime_entity.clone(),
150                store,
151            )?;
152            let snapshot = accepted.snapshot().persisted_snapshot();
153
154            entities.push(EntityCatalogDescription::new(
155                snapshot.entity_name().to_string(),
156                snapshot.entity_path().to_string(),
157                runtime_entity.store_path().to_string(),
158                storage,
159                EntityCatalogCounts::new(
160                    u32::try_from(snapshot.fields().len()).unwrap_or(u32::MAX),
161                    u32::try_from(snapshot.indexes().len()).unwrap_or(u32::MAX),
162                    u32::try_from(snapshot.relations().len()).unwrap_or(u32::MAX),
163                    snapshot.version().get(),
164                ),
165            ));
166        }
167
168        Ok(entities)
169    }
170
171    /// Return one stable list of runtime-registered stores.
172    #[must_use]
173    pub fn show_stores(&self) -> Vec<StoreCatalogDescription> {
174        self.db.runtime_store_catalog()
175    }
176
177    /// Return one stable list of runtime-registered stable-memory allocations.
178    #[must_use]
179    pub fn show_memory(&self) -> Vec<crate::db::MemoryCatalogDescription> {
180        self.db.runtime_memory_catalog()
181    }
182
183    // Resolve the exact secondary-index set that is visible to planner-owned
184    // query planning for one recovered store and accepted schema pair.
185    pub(in crate::db::session) fn visible_indexes_for_store_accepted_schema(
186        &self,
187        store_path: &str,
188        schema_info: &SchemaInfo,
189    ) -> Result<VisibleIndexes, QueryError> {
190        // Phase 1: resolve the recovered store state once at the session
191        // boundary so query/executor planning does not reopen lifecycle checks.
192        let store = self
193            .db
194            .recovered_store(store_path)
195            .map_err(QueryError::execute)?;
196        let state = store.index_state();
197        if state != IndexState::Ready {
198            return Ok(VisibleIndexes::none());
199        }
200        debug_assert_eq!(state, IndexState::Ready);
201
202        // Phase 2: planner-visible indexes are accepted schema contracts once
203        // the recovered store is query-visible.
204        let visible_indexes =
205            VisibleIndexes::accepted_schema_visible(schema_info).map_err(QueryError::execute)?;
206        debug_assert!(visible_indexes.accepted_field_path_contracts_are_consistent());
207        debug_assert!(visible_indexes.accepted_expression_contracts_are_consistent());
208        debug_assert_eq!(
209            visible_indexes.accepted_expression_index_count(),
210            Some(visible_indexes.accepted_expression_indexes().len()),
211        );
212
213        Ok(visible_indexes)
214    }
215
216    /// Return one schema description selected by an immutable authored source key.
217    pub fn try_describe_entity_by_source_key(
218        &self,
219        entity_source: &str,
220    ) -> Result<EntitySchemaDescription, InternalError> {
221        let catalog = self.accepted_schema_catalog_context_for_entity_source_key(entity_source)?;
222        self.describe_accepted_catalog(&catalog)
223    }
224
225    /// Return one schema description selected by its accepted display name.
226    pub fn try_describe_entity_by_name(
227        &self,
228        entity: &str,
229    ) -> Result<EntitySchemaDescription, InternalError> {
230        let catalog = self.accepted_schema_catalog_context_for_entity_name(Some(entity))?;
231        self.describe_accepted_catalog(&catalog)
232    }
233
234    fn describe_accepted_catalog(
235        &self,
236        catalog: &crate::db::session::AcceptedSchemaCatalogContext,
237    ) -> Result<EntitySchemaDescription, InternalError> {
238        let validation_jobs = self.constraint_validation_jobs_for_accepted_catalog(catalog)?;
239        let identity = self.identity_description_for_accepted_catalog(catalog)?;
240
241        describe_accepted_entity_with_persisted_schema(
242            catalog.snapshot(),
243            catalog.value_catalog_handle(),
244            validation_jobs.as_slice(),
245            AcceptedEntityDescriptionMetadata::new(
246                identity,
247                catalog.identity().entity_tag().value(),
248                catalog.fingerprint_method_version(),
249                catalog.fingerprint(),
250            ),
251            |target_path| catalog.relation_target_description(target_path),
252        )
253    }
254
255    pub(in crate::db::session) fn identity_description_for_accepted_catalog(
256        &self,
257        catalog: &crate::db::session::AcceptedSchemaCatalogContext,
258    ) -> Result<Option<EntityIdentityDescription>, InternalError> {
259        let Some(identity) = catalog.inspection_plan().identity_inspection() else {
260            return Ok(None);
261        };
262        let catalog_identity = catalog.identity();
263        let store = self.db.recovered_store(catalog_identity.store_path())?;
264        let incarnation = database_incarnation_id()?;
265        let high_water = store.with_schema(|schema_store| {
266            schema_store.identity_high_water_for_integrity(
267                incarnation,
268                catalog_identity.entity_tag(),
269                identity.field_id(),
270                identity.accepted_kind(),
271            )
272        })?;
273        describe_accepted_identity(identity, high_water).map(Some)
274    }
275
276    pub(in crate::db::session) fn constraint_validation_jobs_for_accepted_catalog(
277        &self,
278        catalog: &crate::db::session::AcceptedSchemaCatalogContext,
279    ) -> Result<Vec<ConstraintValidationJob>, InternalError> {
280        let identity = catalog.inspection_plan().identity();
281        let store = self.db.recovered_store(identity.store_path())?;
282        store.with_schema(|schema_store| {
283            let jobs = catalog
284                .snapshot()
285                .persisted_snapshot()
286                .constraint_activations()
287                .iter()
288                .map(|activation| {
289                    schema_store.constraint_validation_job(identity.entity_tag(), activation.id())
290                })
291                .collect::<Result<Vec<_>, InternalError>>()?;
292            jobs.into_iter()
293                .flatten()
294                .map(|job| {
295                    if job.entity_tag() != identity.entity_tag()
296                        || job.entity_path() != catalog.snapshot().entity_path()
297                    {
298                        return Err(InternalError::store_invariant());
299                    }
300                    Ok(job)
301                })
302                .collect()
303        })
304    }
305
306    /// Build one point-in-time storage report for observability endpoints.
307    pub fn storage_report(
308        &self,
309        name_to_path: &[(&'static str, &'static str)],
310    ) -> Result<StorageReport, InternalError> {
311        self.db.storage_report(name_to_path)
312    }
313
314    /// Build one point-in-time storage report using default entity-path labels.
315    pub fn storage_report_default(&self) -> Result<StorageReport, InternalError> {
316        self.db.storage_report_default()
317    }
318}