Skip to main content

ahri_tre_pgmeta/
datastore_identity.rs

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}