Skip to main content

ahri_tre_pgmeta/
schema.rs

1use crate::{MetadataExecutor, PgMetaError};
2use ahri_tre_libpq_oauth::LibpqOAuthRow;
3use postgres::Row;
4
5use crate::adapter::M0001_INITIAL_SCHEMA_VERSION;
6use crate::read_models::METADATA_TABLES;
7
8/// Latest released-shape datastore schema supported by this unreleased baseline.
9pub const CURRENT_DATASTORE_SCHEMA_VERSION: &str = "m0005_governance_custodians";
10pub(crate) const SEMANTIC_INSTANCES_SCHEMA_VERSION: &str = "m0004_semantic_instances";
11pub(crate) const SEMANTIC_PROVENANCE_SCHEMA_VERSION: &str = "m0003_semantic_provenance";
12pub(crate) const GOVERNANCE_SCHEMA_VERSION: &str = "m0002_governance_disclosure";
13
14/// Deployment authority retained for the stable migration API, with an explicit
15/// Governance appointment required for the preserving disclosure upgrade.
16pub struct AccessRequestAuthority<'a> {
17    _role_name: &'a str,
18    _password: &'a str,
19    _oauth_group_role: &'a str,
20    governance: Option<crate::governance::GovernanceAppointment>,
21}
22
23impl<'a> AccessRequestAuthority<'a> {
24    pub fn new(role_name: &'a str, password: &'a str, oauth_group_role: &'a str) -> Self {
25        Self {
26            _role_name: role_name,
27            _password: password,
28            _oauth_group_role: oauth_group_role,
29            governance: None,
30        }
31    }
32    pub fn with_governance(
33        mut self,
34        appointment: crate::governance::GovernanceAppointment,
35    ) -> Self {
36        self.governance = Some(appointment);
37        self
38    }
39}
40
41#[derive(Debug, Clone, PartialEq, Eq)]
42pub struct DatastoreSchemaStatus {
43    pub current_version: Option<String>,
44    pub target_version: String,
45    pub status: DatastoreSchemaCompatibility,
46    pub migration_history_present: bool,
47    pub applied_versions: Vec<String>,
48    pub missing_tables: Vec<String>,
49    pub database: String,
50    pub user: String,
51    pub server_version: String,
52}
53
54#[derive(Debug, Clone, PartialEq, Eq)]
55pub struct DatastoreMigrationPlan {
56    pub status: DatastoreSchemaCompatibility,
57    pub current_version: Option<String>,
58    pub target_version: String,
59    pub steps: Vec<DatastoreMigrationStep>,
60    pub blocked_reason: Option<String>,
61}
62
63#[derive(Debug, Clone, PartialEq, Eq)]
64pub struct DatastoreMigrationApplyReport {
65    pub before: DatastoreMigrationPlan,
66    pub after: DatastoreSchemaStatus,
67    pub applied_steps: Vec<DatastoreMigrationStep>,
68    pub blocked_reason: Option<String>,
69}
70
71#[derive(Debug, Clone, PartialEq, Eq)]
72pub struct DatastoreMigrationStep {
73    pub version: String,
74    pub description: String,
75    pub scope: DatastoreMigrationScope,
76    pub transactional: bool,
77    pub action: DatastoreMigrationAction,
78}
79
80#[derive(Debug, Clone, Copy, PartialEq, Eq)]
81pub enum DatastoreMigrationScope {
82    Metadata,
83}
84
85impl DatastoreMigrationScope {
86    pub const fn as_str(self) -> &'static str {
87        match self {
88            Self::Metadata => "metadata",
89        }
90    }
91}
92
93#[derive(Debug, Clone, Copy, PartialEq, Eq)]
94pub enum DatastoreMigrationAction {
95    Apply,
96    Skipped,
97    Failed,
98    Blocked,
99}
100
101impl DatastoreMigrationAction {
102    pub const fn as_str(self) -> &'static str {
103        match self {
104            Self::Apply => "apply",
105            Self::Skipped => "skipped",
106            Self::Failed => "failed",
107            Self::Blocked => "blocked",
108        }
109    }
110}
111
112#[derive(Debug, Clone, Copy, PartialEq, Eq)]
113pub enum DatastoreSchemaCompatibility {
114    Current,
115    Pending,
116    MissingHistory,
117    Dirty,
118    TooNew,
119    Unsupported,
120}
121
122impl DatastoreSchemaCompatibility {
123    pub const fn as_str(self) -> &'static str {
124        match self {
125            Self::Current => "current",
126            Self::Pending => "pending",
127            Self::MissingHistory => "missing_history",
128            Self::Dirty => "dirty",
129            Self::TooNew => "too_new",
130            Self::Unsupported => "unsupported",
131        }
132    }
133}
134
135pub fn datastore_schema_status<E>(executor: &mut E) -> Result<DatastoreSchemaStatus, PgMetaError>
136where
137    E: MetadataExecutor,
138    E::Row: SchemaStatusRow,
139{
140    let identity = executor.query_one(
141        "datastore schema status identity",
142        "SELECT current_database()::text, current_user::text, current_setting('server_version')::text",
143        &[],
144    )?;
145    let database = identity.required_text(0, "current_database")?;
146    let user = identity.required_text(1, "current_user")?;
147    let server_version = identity.required_text(2, "server_version")?;
148
149    let migration_history_present = table_exists(executor, "ahri_tre_schema_migrations")?;
150    let missing_tables = missing_canonical_tables(executor)?;
151    let mut applied_versions = if migration_history_present {
152        executor
153            .query_many(
154                "datastore schema migration history",
155                "SELECT version::text FROM ahri_tre_schema_migrations ORDER BY applied_at, version",
156                &[],
157            )?
158            .iter()
159            .map(|row| row.required_text(0, "version"))
160            .collect::<Result<Vec<_>, _>>()?
161    } else {
162        Vec::new()
163    };
164    applied_versions.sort();
165
166    let status = classify_schema_status(
167        migration_history_present,
168        &applied_versions,
169        &missing_tables,
170    );
171    let current_version = if applied_versions
172        .iter()
173        .any(|version| version == CURRENT_DATASTORE_SCHEMA_VERSION)
174    {
175        Some(CURRENT_DATASTORE_SCHEMA_VERSION.to_string())
176    } else {
177        applied_versions.last().cloned()
178    };
179
180    Ok(DatastoreSchemaStatus {
181        current_version,
182        target_version: CURRENT_DATASTORE_SCHEMA_VERSION.to_string(),
183        status,
184        migration_history_present,
185        applied_versions,
186        missing_tables,
187        database,
188        user,
189        server_version,
190    })
191}
192
193pub fn plan_datastore_schema_migration<E>(
194    executor: &mut E,
195) -> Result<DatastoreMigrationPlan, PgMetaError>
196where
197    E: MetadataExecutor,
198    E::Row: SchemaStatusRow,
199{
200    let status = datastore_schema_status(executor)?;
201    Ok(plan_from_status(&status))
202}
203
204pub fn apply_datastore_schema_migration<E>(
205    executor: &mut E,
206    authority: AccessRequestAuthority<'_>,
207) -> Result<DatastoreMigrationApplyReport, PgMetaError>
208where
209    E: MetadataExecutor,
210    E::Row: SchemaStatusRow,
211{
212    let before_status = datastore_schema_status(executor)?;
213    let before = plan_from_status(&before_status);
214    let mut blocked_reason = before.blocked_reason.clone();
215    let mut applied_steps = Vec::new();
216    if before.status == DatastoreSchemaCompatibility::Pending {
217        if let Some(appointment) = authority.governance {
218            crate::governance::GovernanceUpgrade::prepare(executor, &appointment)?;
219            applied_steps = before.steps.clone();
220        } else {
221            blocked_reason = Some("explicit Governance OIDC appointment is required".into());
222        }
223    }
224    let after = datastore_schema_status(executor)?;
225    Ok(DatastoreMigrationApplyReport {
226        before,
227        after,
228        applied_steps,
229        blocked_reason,
230    })
231}
232
233fn plan_from_status(status: &DatastoreSchemaStatus) -> DatastoreMigrationPlan {
234    let blocked_reason = match status.status {
235        DatastoreSchemaCompatibility::Current => None,
236        DatastoreSchemaCompatibility::Pending => None,
237        DatastoreSchemaCompatibility::MissingHistory => Some(
238            "migration history is missing; unreleased legacy schemas are not adoptable".to_string(),
239        ),
240        DatastoreSchemaCompatibility::Dirty => {
241            Some("migration history is dirty or malformed".to_string())
242        }
243        DatastoreSchemaCompatibility::TooNew => {
244            Some("datastore schema was migrated by a newer binary".to_string())
245        }
246        DatastoreSchemaCompatibility::Unsupported => {
247            Some("datastore database does not match the version-1 baseline".to_string())
248        }
249    };
250    DatastoreMigrationPlan {
251        status: status.status,
252        current_version: status.current_version.clone(),
253        target_version: status.target_version.clone(),
254        steps: if status.status == DatastoreSchemaCompatibility::Pending {
255            vec![DatastoreMigrationStep {
256            version:CURRENT_DATASTORE_SCHEMA_VERSION.into(),description:"Preserving Governance and semantic provenance; existing unproven semantic fields remain unknown".into(),scope:DatastoreMigrationScope::Metadata,transactional:true,action:DatastoreMigrationAction::Apply
257        }]
258        } else {
259            Vec::new()
260        },
261        blocked_reason,
262    }
263}
264
265fn classify_schema_status(
266    migration_history_present: bool,
267    applied_versions: &[String],
268    missing_tables: &[String],
269) -> DatastoreSchemaCompatibility {
270    if !migration_history_present {
271        return DatastoreSchemaCompatibility::MissingHistory;
272    }
273    if applied_versions.is_empty() {
274        return DatastoreSchemaCompatibility::Dirty;
275    }
276    if !missing_tables.is_empty() {
277        return DatastoreSchemaCompatibility::Unsupported;
278    }
279    let supported = [
280        M0001_INITIAL_SCHEMA_VERSION,
281        GOVERNANCE_SCHEMA_VERSION,
282        SEMANTIC_PROVENANCE_SCHEMA_VERSION,
283        SEMANTIC_INSTANCES_SCHEMA_VERSION,
284        CURRENT_DATASTORE_SCHEMA_VERSION,
285    ];
286    if applied_versions == supported {
287        return DatastoreSchemaCompatibility::Current;
288    }
289    if applied_versions.len() < supported.len()
290        && applied_versions == &supported[..applied_versions.len()]
291    {
292        return DatastoreSchemaCompatibility::Pending;
293    }
294    if applied_versions.iter().any(|version| {
295        version.starts_with('m') && version.as_str() > CURRENT_DATASTORE_SCHEMA_VERSION
296    }) {
297        return DatastoreSchemaCompatibility::TooNew;
298    }
299    DatastoreSchemaCompatibility::Unsupported
300}
301
302fn table_exists<E>(executor: &mut E, table: &str) -> Result<bool, PgMetaError>
303where
304    E: MetadataExecutor,
305    E::Row: SchemaStatusRow,
306{
307    // The unreleased operation baseline includes indexed, typed creation time.
308    // Report an older shape as an incomplete canonical table; never repair it
309    // from a workflow or claim it is the current baseline.
310    if table == "operations" {
311        let row = executor.query_one(
312            "datastore operation history schema",
313            "SELECT CASE WHEN EXISTS (SELECT 1 FROM pg_attribute WHERE attrelid=to_regclass('public.operations') AND attname='created_at' AND atttypid='timestamptz'::regtype AND attnotnull AND NOT attisdropped) AND EXISTS (SELECT 1 FROM pg_index WHERE indexrelid=to_regclass('public.operation_history_actor_created_id') AND indrelid=to_regclass('public.operations') AND indisvalid) AND to_regprocedure('public.operation_now()') IS NOT NULL THEN 'operations'::text END",
314            &[],
315        )?;
316        return Ok(row.optional_text(0).is_some());
317    }
318    let row = executor.query_one(
319        "datastore schema table existence",
320        "SELECT to_regclass($1)::text",
321        &[&format!("public.{table}")],
322    )?;
323    Ok(row.optional_text(0).is_some())
324}
325
326fn missing_canonical_tables<E>(executor: &mut E) -> Result<Vec<String>, PgMetaError>
327where
328    E: MetadataExecutor,
329    E::Row: SchemaStatusRow,
330{
331    let mut missing = Vec::new();
332    for table in METADATA_TABLES {
333        if !table_exists(executor, table)? {
334            missing.push((*table).to_string());
335        }
336    }
337    for table in ["ahri_tre_schema_migrations", "datastore_identity"] {
338        if !table_exists(executor, table)? && !missing.iter().any(|missing| missing == table) {
339            missing.push(table.to_string());
340        }
341    }
342    missing.sort();
343    Ok(missing)
344}
345
346pub trait SchemaStatusRow {
347    fn optional_text(&self, index: usize) -> Option<String>;
348
349    fn required_text(&self, index: usize, field: &'static str) -> Result<String, PgMetaError> {
350        self.optional_text(index)
351            .ok_or_else(|| PgMetaError::Decode {
352                field,
353                value: "<null>".to_string(),
354                message: "expected text value".to_string(),
355            })
356    }
357}
358
359impl SchemaStatusRow for Row {
360    fn optional_text(&self, index: usize) -> Option<String> {
361        self.get::<_, Option<String>>(index)
362    }
363}
364
365impl SchemaStatusRow for LibpqOAuthRow {
366    fn optional_text(&self, index: usize) -> Option<String> {
367        self.get(index).map(str::to_string)
368    }
369}
370
371#[cfg(test)]
372mod tests {
373    use super::{
374        CURRENT_DATASTORE_SCHEMA_VERSION, DatastoreSchemaCompatibility, classify_schema_status,
375    };
376
377    #[test]
378    fn preserving_upgrade_has_distinct_pending_and_current_histories() {
379        assert_eq!(
380            classify_schema_status(
381                true,
382                &[super::M0001_INITIAL_SCHEMA_VERSION.to_string()],
383                &[]
384            ),
385            DatastoreSchemaCompatibility::Pending
386        );
387        assert_eq!(
388            classify_schema_status(
389                true,
390                &[
391                    super::M0001_INITIAL_SCHEMA_VERSION.to_string(),
392                    super::GOVERNANCE_SCHEMA_VERSION.to_string(),
393                    super::SEMANTIC_PROVENANCE_SCHEMA_VERSION.to_string(),
394                    super::SEMANTIC_INSTANCES_SCHEMA_VERSION.to_string(),
395                    CURRENT_DATASTORE_SCHEMA_VERSION.to_string()
396                ],
397                &[],
398            ),
399            DatastoreSchemaCompatibility::Current
400        );
401    }
402
403    #[test]
404    fn unreleased_legacy_schema_without_history_is_not_adoptable() {
405        assert_eq!(
406            classify_schema_status(false, &[], &[]),
407            DatastoreSchemaCompatibility::MissingHistory
408        );
409        assert_eq!(
410            classify_schema_status(true, &["p1_3_full_metadata_schema".to_string()], &[]),
411            DatastoreSchemaCompatibility::Unsupported
412        );
413    }
414}