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
8pub 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
14pub 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 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}