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