1use crate::{MetadataExecutor, PgMetaError};
2use ahri_tre_libpq_oauth::{LibpqOAuthConnection, LibpqOAuthRow};
3use ahri_tre_types::{
4 DatastoreCreationFailureCode, DatastoreIdentityBinding, DatastoreLakeLocation,
5 DatastoreLifecycleState, LakeCatalogCredentialMode, NewDatastoreIdentityBinding,
6};
7use chrono::{DateTime, Utc};
8use postgres::Row;
9use uuid::Uuid;
10
11pub(crate) const CREATE_DATASTORE_IDENTITY_TABLE_SQL: &str = "
12CREATE TABLE IF NOT EXISTS datastore_identity (
13 datastore_id UUID PRIMARY KEY,
14 datastore_name TEXT NOT NULL,
15 ducklake_catalog_database TEXT NOT NULL,
16 ducklake_catalog_schema TEXT NULL,
17 lake_path TEXT NOT NULL,
18 storage_policy_id TEXT NULL,
19 lake_location_json TEXT NULL,
20 ducklake_encryption BOOLEAN NOT NULL,
21 lake_catalog_credential_mode TEXT NOT NULL,
22 lake_catalog_role_name TEXT NULL,
23 managed_secret_ref TEXT NULL,
24 credential_version INTEGER NOT NULL DEFAULT 1,
25 credential_last_rotated_at TIMESTAMPTZ NULL,
26 created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
27 creating_tool_version TEXT NULL,
28 binding_fingerprint CHAR(64) NOT NULL,
29 lifecycle_state TEXT NOT NULL,
30 failure_code TEXT NULL,
31 failure_summary TEXT NULL,
32 CONSTRAINT ck_datastore_identity_lifecycle_state
33 CHECK (lifecycle_state IN ('creating', 'ready', 'failed', 'retired')),
34 CONSTRAINT ck_datastore_identity_credential_mode
35 CHECK (lake_catalog_credential_mode = 'managed_local'),
36 CONSTRAINT ck_datastore_identity_fingerprint
37 CHECK (binding_fingerprint ~ '^[0-9a-f]{64}$'),
38 CONSTRAINT ck_datastore_identity_failure_code
39 CHECK (failure_code IS NULL OR failure_code IN (
40 'postgresql_setup_failed', 'credential_setup_failed',
41 'lake_setup_failed', 'verification_failed', 'internal_failed'
42 )),
43 CONSTRAINT ck_datastore_identity_failure_summary
44 CHECK (failure_summary IS NULL OR octet_length(failure_summary) <= 512),
45 CONSTRAINT ck_datastore_identity_failure_lifecycle
46 CHECK (
47 (lifecycle_state = 'failed' AND failure_code IS NOT NULL)
48 OR (lifecycle_state <> 'failed' AND failure_code IS NULL AND failure_summary IS NULL)
49 ),
50 CONSTRAINT ck_datastore_identity_name
51 CHECK (length(btrim(datastore_name)) > 0),
52 CONSTRAINT ck_datastore_identity_catalog_database
53 CHECK (length(btrim(ducklake_catalog_database)) > 0),
54 CONSTRAINT ck_datastore_identity_lake_path
55 CHECK (length(btrim(lake_path)) > 0)
56)
57";
58
59pub(crate) const CREATE_DATASTORE_IDENTITY_SINGLETON_INDEX_SQL: &str = "
60CREATE UNIQUE INDEX IF NOT EXISTS ux_datastore_identity_singleton
61 ON datastore_identity ((true))
62";
63
64pub(crate) const CREATE_DATASTORE_IDENTITY_IMMUTABILITY_FUNCTION_SQL: &str = "
65CREATE OR REPLACE FUNCTION tre_prevent_ready_datastore_identity_mutation()
66RETURNS trigger
67LANGUAGE plpgsql
68AS $$
69BEGIN
70 IF TG_OP = 'DELETE' THEN
71 RAISE EXCEPTION 'datastore identity binding is immutable';
72 END IF;
73 IF NEW.datastore_id <> OLD.datastore_id
74 OR NEW.datastore_name <> OLD.datastore_name
75 OR NEW.ducklake_catalog_database <> OLD.ducklake_catalog_database
76 OR NEW.ducklake_catalog_schema IS DISTINCT FROM OLD.ducklake_catalog_schema
77 OR NEW.storage_policy_id IS DISTINCT FROM OLD.storage_policy_id
78 OR NEW.lake_location_json IS DISTINCT FROM OLD.lake_location_json
79 OR NEW.ducklake_encryption <> OLD.ducklake_encryption
80 OR NEW.lake_catalog_credential_mode <> OLD.lake_catalog_credential_mode
81 OR NEW.lake_catalog_role_name IS DISTINCT FROM OLD.lake_catalog_role_name
82 OR NEW.created_at <> OLD.created_at
83 OR NEW.creating_tool_version IS DISTINCT FROM OLD.creating_tool_version THEN
84 RAISE EXCEPTION 'datastore identity binding is immutable';
85 END IF;
86 IF OLD.lifecycle_state = 'ready' AND NEW.lifecycle_state <> OLD.lifecycle_state THEN
87 RAISE EXCEPTION 'ready datastore identity lifecycle is immutable';
88 END IF;
89 RETURN NEW;
90END;
91$$
92";
93
94pub(crate) const DROP_DATASTORE_IDENTITY_IMMUTABILITY_TRIGGER_SQL: &str = "
95DROP TRIGGER IF EXISTS trg_datastore_identity_ready_immutable ON datastore_identity;
96";
97
98pub(crate) const CREATE_DATASTORE_IDENTITY_IMMUTABILITY_TRIGGER_SQL: &str = "
99CREATE TRIGGER trg_datastore_identity_ready_immutable
100 BEFORE UPDATE OR DELETE ON datastore_identity
101 FOR EACH ROW
102 EXECUTE FUNCTION tre_prevent_ready_datastore_identity_mutation()
103";
104
105pub fn create_datastore_identity_table<E>(executor: &mut E) -> Result<(), PgMetaError>
106where
107 E: MetadataExecutor,
108{
109 executor.execute_command(
110 "create datastore identity table",
111 CREATE_DATASTORE_IDENTITY_TABLE_SQL,
112 &[],
113 )?;
114 executor.execute_command(
115 "create datastore identity singleton index",
116 CREATE_DATASTORE_IDENTITY_SINGLETON_INDEX_SQL,
117 &[],
118 )?;
119 executor.execute_command(
120 "create datastore identity immutability function",
121 CREATE_DATASTORE_IDENTITY_IMMUTABILITY_FUNCTION_SQL,
122 &[],
123 )?;
124 executor.execute_command(
125 "drop datastore identity immutability trigger",
126 DROP_DATASTORE_IDENTITY_IMMUTABILITY_TRIGGER_SQL,
127 &[],
128 )?;
129 executor.execute_command(
130 "create datastore identity immutability trigger",
131 CREATE_DATASTORE_IDENTITY_IMMUTABILITY_TRIGGER_SQL,
132 &[],
133 )?;
134 Ok(())
135}
136
137pub fn insert_datastore_identity_binding<E>(
138 executor: &mut E,
139 binding: &NewDatastoreIdentityBinding,
140) -> Result<DatastoreIdentityBinding, PgMetaError>
141where
142 E: MetadataExecutor<Row = Row>,
143{
144 let fingerprint = binding.binding_fingerprint();
145 let credential_mode = binding.lake_catalog_credential_mode.as_str().to_string();
146 let lifecycle_state = binding.lifecycle_state.as_str().to_string();
147 let lake_location_json = binding
148 .lake_location
149 .as_ref()
150 .map(serde_json::to_string)
151 .transpose()
152 .map_err(|error| PgMetaError::Decode {
153 field: "lake_location",
154 value: "<structured>".to_string(),
155 message: error.to_string(),
156 })?;
157 let row = executor.query_one(
158 "insert datastore identity binding",
159 "
160 INSERT INTO datastore_identity (
161 datastore_id,
162 datastore_name,
163 ducklake_catalog_database,
164 ducklake_catalog_schema,
165 lake_path,
166 storage_policy_id,
167 lake_location_json,
168 ducklake_encryption,
169 lake_catalog_credential_mode,
170 lake_catalog_role_name,
171 managed_secret_ref,
172 credential_version,
173 credential_last_rotated_at,
174 creating_tool_version,
175 binding_fingerprint,
176 lifecycle_state,
177 failure_code,
178 failure_summary
179 )
180 VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16, $17, $18)
181 RETURNING
182 datastore_id,
183 datastore_name,
184 ducklake_catalog_database,
185 ducklake_catalog_schema,
186 lake_path,
187 storage_policy_id,
188 lake_location_json,
189 ducklake_encryption,
190 lake_catalog_credential_mode,
191 lake_catalog_role_name,
192 managed_secret_ref,
193 credential_version,
194 credential_last_rotated_at,
195 created_at,
196 creating_tool_version,
197 binding_fingerprint,
198 lifecycle_state,
199 failure_code,
200 failure_summary
201 ",
202 &[
203 &binding.datastore_id,
204 &binding.datastore_name,
205 &binding.ducklake_catalog_database,
206 &binding.ducklake_catalog_schema,
207 &binding.lake_path,
208 &binding.storage_policy_id,
209 &lake_location_json,
210 &binding.ducklake_encryption,
211 &credential_mode,
212 &binding.lake_catalog_role_name,
213 &binding.managed_secret_ref,
214 &binding.credential_version,
215 &binding.credential_last_rotated_at,
216 &binding.creating_tool_version,
217 &fingerprint,
218 &lifecycle_state,
219 &binding
220 .failure_code
221 .map(DatastoreCreationFailureCode::as_str),
222 &binding.failure_summary,
223 ],
224 )?;
225 datastore_identity_binding_from_row(&row)
226}
227
228pub(crate) const MARK_DATASTORE_IDENTITY_READY_SQL: &str = "
229 UPDATE datastore_identity
230 SET lifecycle_state = 'ready',
231 credential_last_rotated_at = $2,
232 failure_code = NULL,
233 failure_summary = NULL
234 WHERE datastore_id = $1
235 AND lifecycle_state = 'creating'
236 RETURNING
237 datastore_id,
238 datastore_name,
239 ducklake_catalog_database,
240 ducklake_catalog_schema,
241 lake_path,
242 storage_policy_id,
243 lake_location_json,
244 ducklake_encryption,
245 lake_catalog_credential_mode,
246 lake_catalog_role_name,
247 managed_secret_ref,
248 credential_version,
249 credential_last_rotated_at,
250 created_at,
251 creating_tool_version,
252 binding_fingerprint,
253 lifecycle_state,
254 failure_code,
255 failure_summary
256 ";
257
258pub(crate) const MARK_DATASTORE_IDENTITY_FAILED_SQL: &str = "
259 UPDATE datastore_identity
260 SET lifecycle_state = 'failed',
261 failure_code = $2,
262 failure_summary = $3
263 WHERE datastore_id = $1
264 AND lifecycle_state = 'creating'
265 RETURNING
266 datastore_id,
267 datastore_name,
268 ducklake_catalog_database,
269 ducklake_catalog_schema,
270 lake_path,
271 storage_policy_id,
272 lake_location_json,
273 ducklake_encryption,
274 lake_catalog_credential_mode,
275 lake_catalog_role_name,
276 managed_secret_ref,
277 credential_version,
278 credential_last_rotated_at,
279 created_at,
280 creating_tool_version,
281 binding_fingerprint,
282 lifecycle_state,
283 failure_code,
284 failure_summary
285 ";
286
287pub fn mark_datastore_identity_ready<E>(
288 executor: &mut E,
289 datastore_id: Uuid,
290 credential_created_at: DateTime<Utc>,
291) -> Result<DatastoreIdentityBinding, PgMetaError>
292where
293 E: MetadataExecutor<Row = Row>,
294{
295 let row = executor.query_optional(
296 "mark datastore identity ready",
297 MARK_DATASTORE_IDENTITY_READY_SQL,
298 &[&datastore_id, &credential_created_at],
299 )?;
300 row.as_ref()
301 .ok_or_else(|| PgMetaError::Decode {
302 field: "lifecycle_state",
303 value: "<not creating>".to_string(),
304 message: "Datastore binding is not eligible to become ready".to_string(),
305 })
306 .and_then(datastore_identity_binding_from_row)
307}
308
309pub fn mark_datastore_identity_failed<E>(
310 executor: &mut E,
311 datastore_id: Uuid,
312 failure_code: DatastoreCreationFailureCode,
313 safe_summary: Option<&str>,
314) -> Result<DatastoreIdentityBinding, PgMetaError>
315where
316 E: MetadataExecutor<Row = Row>,
317{
318 if safe_summary.is_some_and(|summary| summary.len() > 512) {
319 return Err(PgMetaError::Decode {
320 field: "failure_summary",
321 value: "<bounded summary>".to_string(),
322 message: "Datastore creation failure summary exceeds 512 bytes".to_string(),
323 });
324 }
325 let code = failure_code.as_str();
326 let row = executor.query_optional(
327 "mark datastore identity failed",
328 MARK_DATASTORE_IDENTITY_FAILED_SQL,
329 &[&datastore_id, &code, &safe_summary],
330 )?;
331 row.as_ref()
332 .ok_or_else(|| PgMetaError::Decode {
333 field: "lifecycle_state",
334 value: "<not creating>".to_string(),
335 message: "Datastore binding is not eligible to become failed".to_string(),
336 })
337 .and_then(datastore_identity_binding_from_row)
338}
339
340pub fn update_datastore_identity_credential<E>(
341 executor: &mut E,
342 datastore_id: Uuid,
343 managed_secret_ref: &str,
344 credential_version: i32,
345 credential_last_rotated_at: DateTime<Utc>,
346 binding_fingerprint: &str,
347) -> Result<DatastoreIdentityBinding, PgMetaError>
348where
349 E: MetadataExecutor<Row = Row>,
350{
351 let row = executor.query_one(
352 "update datastore identity credential metadata",
353 "
354 UPDATE datastore_identity
355 SET managed_secret_ref = $2,
356 credential_version = $3,
357 credential_last_rotated_at = $4,
358 binding_fingerprint = $5
359 WHERE datastore_id = $1
360 AND lifecycle_state = 'ready'
361 RETURNING
362 datastore_id,
363 datastore_name,
364 ducklake_catalog_database,
365 ducklake_catalog_schema,
366 lake_path,
367 storage_policy_id,
368 lake_location_json,
369 ducklake_encryption,
370 lake_catalog_credential_mode,
371 lake_catalog_role_name,
372 managed_secret_ref,
373 credential_version,
374 credential_last_rotated_at,
375 created_at,
376 creating_tool_version,
377 binding_fingerprint,
378 lifecycle_state,
379 failure_code,
380 failure_summary
381 ",
382 &[
383 &datastore_id,
384 &managed_secret_ref,
385 &credential_version,
386 &credential_last_rotated_at,
387 &binding_fingerprint,
388 ],
389 )?;
390 datastore_identity_binding_from_row(&row)
391}
392
393pub fn advance_ready_datastore_identity_lake_path<E>(
394 executor: &mut E,
395 datastore_id: Uuid,
396 old_lake_path: &str,
397 new_lake_path: &str,
398) -> Result<DatastoreIdentityBinding, PgMetaError>
399where
400 E: MetadataExecutor<Row = Row>,
401{
402 executor.with_transaction(
403 "advance ready datastore identity Lake location",
404 |transaction| {
405 transaction.execute_command(
406 "refresh datastore identity immutability function",
407 CREATE_DATASTORE_IDENTITY_IMMUTABILITY_FUNCTION_SQL,
408 &[],
409 )?;
410
411 let row = transaction.query_optional(
412 "read ready datastore identity binding for Lake location cutover",
413 READY_DATASTORE_IDENTITY_BINDING_FOR_UPDATE_SQL,
414 &[&datastore_id],
415 )?;
416 let Some(row) = row else {
417 return Err(PgMetaError::Decode {
418 field: "datastore_id",
419 value: datastore_id.to_string(),
420 message: "ready datastore identity binding was not found for Lake location cutover"
421 .to_string(),
422 });
423 };
424 let binding = datastore_identity_binding_from_row(&row)?;
425 if binding.lake_path != old_lake_path {
426 return Err(PgMetaError::Decode {
427 field: "lake_path",
428 value: binding.lake_path,
429 message: format!("expected current ready binding Lake location {old_lake_path}"),
430 });
431 }
432
433 let mut advanced = binding.clone();
434 advanced.lake_path = new_lake_path.to_string();
435 let advanced_fingerprint = advanced.computed_binding_fingerprint();
436
437 let row = transaction.query_optional(
438 "advance ready datastore identity Lake location",
439 ADVANCE_READY_DATASTORE_IDENTITY_LAKE_PATH_SQL,
440 &[
441 &datastore_id,
442 &old_lake_path,
443 &binding.binding_fingerprint,
444 &new_lake_path,
445 &advanced_fingerprint,
446 ],
447 )?;
448 let Some(row) = row else {
449 return Err(PgMetaError::Decode {
450 field: "binding_fingerprint",
451 value: binding.binding_fingerprint,
452 message:
453 "ready datastore identity binding changed before Lake location cutover could apply"
454 .to_string(),
455 });
456 };
457 datastore_identity_binding_from_row(&row)
458 },
459 )
460}
461
462pub fn read_ready_datastore_identity_bindings<E>(
463 executor: &mut E,
464) -> Result<Vec<DatastoreIdentityBinding>, PgMetaError>
465where
466 E: MetadataExecutor<Row = Row>,
467{
468 let rows = executor.query_many(
469 "read ready datastore identity bindings",
470 "
471 SELECT
472 datastore_id,
473 datastore_name,
474 ducklake_catalog_database,
475 ducklake_catalog_schema,
476 lake_path,
477 storage_policy_id,
478 lake_location_json,
479 ducklake_encryption,
480 lake_catalog_credential_mode,
481 lake_catalog_role_name,
482 managed_secret_ref,
483 credential_version,
484 credential_last_rotated_at,
485 created_at,
486 creating_tool_version,
487 binding_fingerprint,
488 lifecycle_state,
489 failure_code,
490 failure_summary
491 FROM datastore_identity
492 WHERE lifecycle_state = 'ready'
493 ORDER BY created_at, datastore_id
494 ",
495 &[],
496 )?;
497 rows.iter()
498 .map(datastore_identity_binding_from_row)
499 .collect()
500}
501
502pub fn read_datastore_identity_bindings<E>(
503 executor: &mut E,
504) -> Result<Vec<DatastoreIdentityBinding>, PgMetaError>
505where
506 E: MetadataExecutor<Row = Row>,
507{
508 let rows = executor.query_many(
509 "read datastore identity bindings",
510 DATASTORE_IDENTITY_BINDINGS_SQL,
511 &[],
512 )?;
513 rows.iter()
514 .map(datastore_identity_binding_from_row)
515 .collect()
516}
517
518pub fn datastore_identity_table_exists<E>(executor: &mut E) -> Result<bool, PgMetaError>
519where
520 E: MetadataExecutor<Row = Row>,
521{
522 let row = executor.query_one(
523 "check datastore identity table",
524 "SELECT to_regclass('public.datastore_identity')::text",
525 &[],
526 )?;
527 Ok(row.get::<_, Option<String>>(0).is_some())
528}
529
530pub fn read_ready_datastore_identity_bindings_oauth(
531 executor: &mut LibpqOAuthConnection,
532) -> Result<Vec<DatastoreIdentityBinding>, PgMetaError> {
533 let rows = executor.query_many(
534 "read ready datastore identity bindings",
535 READY_DATASTORE_IDENTITY_BINDINGS_TEXT_SQL,
536 &[],
537 )?;
538 rows.iter()
539 .map(datastore_identity_binding_from_oauth_row)
540 .collect()
541}
542
543pub fn read_datastore_identity_bindings_oauth(
544 executor: &mut LibpqOAuthConnection,
545) -> Result<Vec<DatastoreIdentityBinding>, PgMetaError> {
546 let rows = executor.query_many(
547 "read datastore identity bindings",
548 DATASTORE_IDENTITY_BINDINGS_TEXT_SQL,
549 &[],
550 )?;
551 rows.iter()
552 .map(datastore_identity_binding_from_oauth_row)
553 .collect()
554}
555
556pub fn datastore_identity_table_exists_oauth(
557 executor: &mut LibpqOAuthConnection,
558) -> Result<bool, PgMetaError> {
559 let row = executor.query_one(
560 "check datastore identity table",
561 "SELECT to_regclass('public.datastore_identity')::text",
562 &[],
563 )?;
564 Ok(row.get(0).is_some())
565}
566
567fn datastore_identity_binding_from_row(row: &Row) -> Result<DatastoreIdentityBinding, PgMetaError> {
568 let credential_mode: String = row.get("lake_catalog_credential_mode");
569 let lifecycle_state: String = row.get("lifecycle_state");
570 Ok(DatastoreIdentityBinding {
571 datastore_id: row.get::<_, Uuid>("datastore_id"),
572 datastore_name: row.get("datastore_name"),
573 ducklake_catalog_database: row.get("ducklake_catalog_database"),
574 ducklake_catalog_schema: row.get("ducklake_catalog_schema"),
575 lake_path: row.get("lake_path"),
576 storage_policy_id: row.get("storage_policy_id"),
577 lake_location: parse_lake_location(row.get("lake_location_json"))?,
578 ducklake_encryption: row.get("ducklake_encryption"),
579 lake_catalog_credential_mode: LakeCatalogCredentialMode::try_from(credential_mode.as_str())
580 .map_err(|message| PgMetaError::Decode {
581 field: "lake_catalog_credential_mode",
582 value: credential_mode,
583 message,
584 })?,
585 lake_catalog_role_name: row.get("lake_catalog_role_name"),
586 managed_secret_ref: row.get("managed_secret_ref"),
587 credential_version: row.get("credential_version"),
588 credential_last_rotated_at: row
589 .get::<_, Option<DateTime<Utc>>>("credential_last_rotated_at"),
590 created_at: row.get("created_at"),
591 creating_tool_version: row.get("creating_tool_version"),
592 binding_fingerprint: row.get("binding_fingerprint"),
593 lifecycle_state: DatastoreLifecycleState::try_from(lifecycle_state.as_str()).map_err(
594 |message| PgMetaError::Decode {
595 field: "lifecycle_state",
596 value: lifecycle_state,
597 message,
598 },
599 )?,
600 failure_code: row
601 .get::<_, Option<String>>("failure_code")
602 .map(|code| DatastoreCreationFailureCode::try_from(code.as_str()))
603 .transpose()
604 .map_err(|message| PgMetaError::Decode {
605 field: "failure_code",
606 value: "<closed code>".to_string(),
607 message,
608 })?,
609 failure_summary: row.get("failure_summary"),
610 })
611}
612
613const DATASTORE_IDENTITY_BINDINGS_SQL: &str = "
614 SELECT
615 datastore_id,
616 datastore_name,
617 ducklake_catalog_database,
618 ducklake_catalog_schema,
619 lake_path,
620 storage_policy_id,
621 lake_location_json,
622 ducklake_encryption,
623 lake_catalog_credential_mode,
624 lake_catalog_role_name,
625 managed_secret_ref,
626 credential_version,
627 credential_last_rotated_at,
628 created_at,
629 creating_tool_version,
630 binding_fingerprint,
631 lifecycle_state,
632 failure_code,
633 failure_summary
634 FROM datastore_identity
635 ORDER BY created_at, datastore_id
636 ";
637
638const READY_DATASTORE_IDENTITY_BINDING_FOR_UPDATE_SQL: &str = "
639 SELECT
640 datastore_id,
641 datastore_name,
642 ducklake_catalog_database,
643 ducklake_catalog_schema,
644 lake_path,
645 storage_policy_id,
646 lake_location_json,
647 ducklake_encryption,
648 lake_catalog_credential_mode,
649 lake_catalog_role_name,
650 managed_secret_ref,
651 credential_version,
652 credential_last_rotated_at,
653 created_at,
654 creating_tool_version,
655 binding_fingerprint,
656 lifecycle_state,
657 failure_code,
658 failure_summary
659 FROM datastore_identity
660 WHERE datastore_id = $1
661 AND lifecycle_state = 'ready'
662 FOR UPDATE
663 ";
664
665const ADVANCE_READY_DATASTORE_IDENTITY_LAKE_PATH_SQL: &str = "
666 UPDATE datastore_identity
667 SET lake_path = $4,
668 binding_fingerprint = $5
669 WHERE datastore_id = $1
670 AND lifecycle_state = 'ready'
671 AND lake_path = $2
672 AND binding_fingerprint = $3
673 RETURNING
674 datastore_id,
675 datastore_name,
676 ducklake_catalog_database,
677 ducklake_catalog_schema,
678 lake_path,
679 storage_policy_id,
680 lake_location_json,
681 ducklake_encryption,
682 lake_catalog_credential_mode,
683 lake_catalog_role_name,
684 managed_secret_ref,
685 credential_version,
686 credential_last_rotated_at,
687 created_at,
688 creating_tool_version,
689 binding_fingerprint,
690 lifecycle_state,
691 failure_code,
692 failure_summary
693 ";
694
695const READY_DATASTORE_IDENTITY_BINDINGS_TEXT_SQL: &str = "
696 SELECT
697 datastore_id::text,
698 datastore_name,
699 ducklake_catalog_database,
700 ducklake_catalog_schema,
701 lake_path,
702 storage_policy_id,
703 lake_location_json,
704 ducklake_encryption::text,
705 lake_catalog_credential_mode,
706 lake_catalog_role_name,
707 managed_secret_ref,
708 credential_version::text,
709 CASE
710 WHEN credential_last_rotated_at IS NULL THEN NULL
711 ELSE to_char(credential_last_rotated_at AT TIME ZONE 'UTC', 'YYYY-MM-DD\"T\"HH24:MI:SS.US\"Z\"')
712 END,
713 to_char(created_at AT TIME ZONE 'UTC', 'YYYY-MM-DD\"T\"HH24:MI:SS.US\"Z\"'),
714 creating_tool_version,
715 binding_fingerprint,
716 lifecycle_state,
717 failure_code,
718 failure_summary
719 FROM datastore_identity
720 WHERE lifecycle_state = 'ready'
721 ORDER BY created_at, datastore_id
722 ";
723
724const DATASTORE_IDENTITY_BINDINGS_TEXT_SQL: &str = "
725 SELECT
726 datastore_id::text,
727 datastore_name,
728 ducklake_catalog_database,
729 ducklake_catalog_schema,
730 lake_path,
731 storage_policy_id,
732 lake_location_json,
733 ducklake_encryption::text,
734 lake_catalog_credential_mode,
735 lake_catalog_role_name,
736 managed_secret_ref,
737 credential_version::text,
738 CASE
739 WHEN credential_last_rotated_at IS NULL THEN NULL
740 ELSE to_char(credential_last_rotated_at AT TIME ZONE 'UTC', 'YYYY-MM-DD\"T\"HH24:MI:SS.US\"Z\"')
741 END,
742 to_char(created_at AT TIME ZONE 'UTC', 'YYYY-MM-DD\"T\"HH24:MI:SS.US\"Z\"'),
743 creating_tool_version,
744 binding_fingerprint,
745 lifecycle_state,
746 failure_code,
747 failure_summary
748 FROM datastore_identity
749 ORDER BY created_at, datastore_id
750 ";
751
752fn datastore_identity_binding_from_oauth_row(
753 row: &LibpqOAuthRow,
754) -> Result<DatastoreIdentityBinding, PgMetaError> {
755 let datastore_id = parse_uuid(required_text(row, 0, "datastore_id")?, "datastore_id")?;
756 let datastore_name = required_text(row, 1, "datastore_name")?.to_string();
757 let ducklake_catalog_database = required_text(row, 2, "ducklake_catalog_database")?.to_string();
758 let ducklake_catalog_schema = optional_text(row, 3).map(str::to_string);
759 let lake_path = required_text(row, 4, "lake_path")?.to_string();
760 let storage_policy_id = optional_text(row, 5).map(str::to_string);
761 let lake_location = optional_text(row, 6)
762 .map(|value| parse_lake_location(Some(value.to_string())))
763 .transpose()?
764 .flatten();
765 let ducklake_encryption = parse_bool(
766 required_text(row, 7, "ducklake_encryption")?,
767 "ducklake_encryption",
768 )?;
769 let credential_mode = required_text(row, 8, "lake_catalog_credential_mode")?;
770 let lake_catalog_credential_mode = LakeCatalogCredentialMode::try_from(credential_mode)
771 .map_err(|message| PgMetaError::Decode {
772 field: "lake_catalog_credential_mode",
773 value: credential_mode.to_string(),
774 message,
775 })?;
776 let lake_catalog_role_name = optional_text(row, 9).map(str::to_string);
777 let managed_secret_ref = optional_text(row, 10).map(str::to_string);
778 let credential_version = parse_i32(
779 required_text(row, 11, "credential_version")?,
780 "credential_version",
781 )?;
782 let credential_last_rotated_at = optional_text(row, 12)
783 .map(|value| parse_datetime(value, "credential_last_rotated_at"))
784 .transpose()?;
785 let created_at = parse_datetime(required_text(row, 13, "created_at")?, "created_at")?;
786 let creating_tool_version = optional_text(row, 14).map(str::to_string);
787 let binding_fingerprint = required_text(row, 15, "binding_fingerprint")?.to_string();
788 let lifecycle_state = required_text(row, 16, "lifecycle_state")?;
789 let lifecycle_state =
790 DatastoreLifecycleState::try_from(lifecycle_state).map_err(|message| {
791 PgMetaError::Decode {
792 field: "lifecycle_state",
793 value: lifecycle_state.to_string(),
794 message,
795 }
796 })?;
797 let failure_code = optional_text(row, 17)
798 .map(DatastoreCreationFailureCode::try_from)
799 .transpose()
800 .map_err(|message| PgMetaError::Decode {
801 field: "failure_code",
802 value: "<closed code>".to_string(),
803 message,
804 })?;
805 let failure_summary = optional_text(row, 18).map(str::to_string);
806
807 Ok(DatastoreIdentityBinding {
808 datastore_id,
809 datastore_name,
810 ducklake_catalog_database,
811 ducklake_catalog_schema,
812 lake_path,
813 storage_policy_id,
814 lake_location,
815 ducklake_encryption,
816 lake_catalog_credential_mode,
817 lake_catalog_role_name,
818 managed_secret_ref,
819 credential_version,
820 credential_last_rotated_at,
821 created_at,
822 creating_tool_version,
823 binding_fingerprint,
824 lifecycle_state,
825 failure_code,
826 failure_summary,
827 })
828}
829
830fn parse_lake_location(
831 encoded: Option<String>,
832) -> Result<Option<DatastoreLakeLocation>, PgMetaError> {
833 encoded
834 .map(|value| {
835 serde_json::from_str(&value).map_err(|error| PgMetaError::Decode {
836 field: "lake_location_json",
837 value: "<structured>".to_string(),
838 message: error.to_string(),
839 })
840 })
841 .transpose()
842}
843
844fn optional_text(row: &LibpqOAuthRow, index: usize) -> Option<&str> {
845 row.get(index).filter(|value| !value.trim().is_empty())
846}
847
848fn required_text<'a>(
849 row: &'a LibpqOAuthRow,
850 index: usize,
851 field: &'static str,
852) -> Result<&'a str, PgMetaError> {
853 optional_text(row, index).ok_or_else(|| PgMetaError::Decode {
854 field,
855 value: "<null>".to_string(),
856 message: "expected non-empty text value".to_string(),
857 })
858}
859
860fn parse_uuid(value: &str, field: &'static str) -> Result<Uuid, PgMetaError> {
861 Uuid::parse_str(value).map_err(|error| PgMetaError::Decode {
862 field,
863 value: value.to_string(),
864 message: error.to_string(),
865 })
866}
867
868fn parse_bool(value: &str, field: &'static str) -> Result<bool, PgMetaError> {
869 value.parse::<bool>().map_err(|error| PgMetaError::Decode {
870 field,
871 value: value.to_string(),
872 message: error.to_string(),
873 })
874}
875
876fn parse_i32(value: &str, field: &'static str) -> Result<i32, PgMetaError> {
877 value.parse::<i32>().map_err(|error| PgMetaError::Decode {
878 field,
879 value: value.to_string(),
880 message: error.to_string(),
881 })
882}
883
884fn parse_datetime(value: &str, field: &'static str) -> Result<DateTime<Utc>, PgMetaError> {
885 DateTime::parse_from_rfc3339(value)
886 .map(|value| value.with_timezone(&Utc))
887 .map_err(|error| PgMetaError::Decode {
888 field,
889 value: value.to_string(),
890 message: error.to_string(),
891 })
892}
893
894#[cfg(test)]
895mod tests {
896 use super::{
897 ADVANCE_READY_DATASTORE_IDENTITY_LAKE_PATH_SQL,
898 CREATE_DATASTORE_IDENTITY_IMMUTABILITY_FUNCTION_SQL,
899 CREATE_DATASTORE_IDENTITY_IMMUTABILITY_TRIGGER_SQL,
900 CREATE_DATASTORE_IDENTITY_SINGLETON_INDEX_SQL, CREATE_DATASTORE_IDENTITY_TABLE_SQL,
901 MARK_DATASTORE_IDENTITY_FAILED_SQL, MARK_DATASTORE_IDENTITY_READY_SQL,
902 READY_DATASTORE_IDENTITY_BINDING_FOR_UPDATE_SQL,
903 };
904
905 #[test]
906 fn identity_table_ddl_records_safe_binding_facts() {
907 assert!(CREATE_DATASTORE_IDENTITY_TABLE_SQL.contains("datastore_id UUID PRIMARY KEY"));
908 assert!(CREATE_DATASTORE_IDENTITY_TABLE_SQL.contains("ducklake_catalog_database"));
909 assert!(CREATE_DATASTORE_IDENTITY_TABLE_SQL.contains("managed_secret_ref"));
910 assert!(CREATE_DATASTORE_IDENTITY_TABLE_SQL.contains("binding_fingerprint CHAR(64)"));
911 assert!(
912 CREATE_DATASTORE_IDENTITY_SINGLETON_INDEX_SQL
913 .contains("ux_datastore_identity_singleton")
914 );
915 assert!(!CREATE_DATASTORE_IDENTITY_SINGLETON_INDEX_SQL.contains("WHERE"));
916 assert!(CREATE_DATASTORE_IDENTITY_TABLE_SQL.contains("failure_code TEXT NULL"));
917 assert!(
918 CREATE_DATASTORE_IDENTITY_TABLE_SQL.contains("octet_length(failure_summary) <= 512")
919 );
920 assert!(
921 CREATE_DATASTORE_IDENTITY_IMMUTABILITY_FUNCTION_SQL
922 .contains("datastore identity binding is immutable")
923 );
924 assert!(
925 !CREATE_DATASTORE_IDENTITY_IMMUTABILITY_FUNCTION_SQL
926 .contains("NEW.lake_path = OLD.lake_path")
927 );
928 assert!(
929 CREATE_DATASTORE_IDENTITY_IMMUTABILITY_TRIGGER_SQL
930 .contains("trg_datastore_identity_ready_immutable")
931 );
932 assert!(!CREATE_DATASTORE_IDENTITY_TABLE_SQL.contains("password"));
933 }
934
935 #[test]
936 fn creation_lifecycle_advances_only_the_existing_creating_binding() {
937 for sql in [
938 MARK_DATASTORE_IDENTITY_READY_SQL,
939 MARK_DATASTORE_IDENTITY_FAILED_SQL,
940 ] {
941 assert!(sql.contains("WHERE datastore_id = $1"));
942 assert!(sql.contains("AND lifecycle_state = 'creating'"));
943 assert!(sql.contains("RETURNING"));
944 }
945 assert!(MARK_DATASTORE_IDENTITY_READY_SQL.contains("SET lifecycle_state = 'ready'"));
946 assert!(MARK_DATASTORE_IDENTITY_FAILED_SQL.contains("SET lifecycle_state = 'failed'"));
947 assert!(MARK_DATASTORE_IDENTITY_FAILED_SQL.contains("failure_code = $2"));
948 assert!(MARK_DATASTORE_IDENTITY_FAILED_SQL.contains("failure_summary = $3"));
949 }
950
951 #[test]
952 fn lake_path_cutover_sql_rechecks_ready_binding_identity() {
953 assert!(READY_DATASTORE_IDENTITY_BINDING_FOR_UPDATE_SQL.contains("FOR UPDATE"));
954 assert!(ADVANCE_READY_DATASTORE_IDENTITY_LAKE_PATH_SQL.contains("SET lake_path = $4"));
955 assert!(
956 ADVANCE_READY_DATASTORE_IDENTITY_LAKE_PATH_SQL.contains("AND binding_fingerprint = $3")
957 );
958 assert!(
959 ADVANCE_READY_DATASTORE_IDENTITY_LAKE_PATH_SQL.contains("binding_fingerprint = $5")
960 );
961 }
962}