Skip to main content

architect_sdk/
store.rs

1//! _sys_* table DDL and config persistence. All _sys_* tables live in a schema named from `ARCHITECT_SCHEMA` env (default `architect`).
2
3use crate::db::{pool::Pool, Dialect};
4use crate::error::AppError;
5use chrono::{DateTime, Utc};
6use std::collections::HashMap;
7
8/// Schema name for _sys_* tables. From env `ARCHITECT_SCHEMA`, default `architect`. Must be a valid PostgreSQL identifier.
9pub fn architect_schema() -> String {
10    std::env::var("ARCHITECT_SCHEMA").unwrap_or_else(|_| "architect".into())
11}
12
13/// Returns schema-qualified table name for _sys_* tables (e.g. "architect._sys_schemas").
14pub fn qualified_sys_table(table: &str) -> String {
15    format!("{}.{}", architect_schema(), table)
16}
17
18/// Config tables (each row is keyed by id + package_id). Excludes _sys_packages.
19const CONFIG_TABLES: &[&str] = &[
20    "_sys_schemas",
21    "_sys_enums",
22    "_sys_tables",
23    "_sys_columns",
24    "_sys_indexes",
25    "_sys_relationships",
26    "_sys_api_entities",
27    "_sys_kv_stores",
28    "_sys_reports",
29];
30
31/// Package id used when config is posted directly (no package install). Ensures (id, package_id) is unique per package.
32pub const DEFAULT_PACKAGE_ID: &str = "_default";
33
34/// Create schema from `ARCHITECT_SCHEMA` env if not exists, then _sys_* tables.
35/// Config tables have (id, package_id) as composite primary key; _sys_packages has id only.
36pub async fn ensure_sys_tables(pool: &Pool, dialect: &dyn Dialect) -> Result<(), AppError> {
37    let schema = architect_schema();
38    if dialect.supports_schemas() {
39        sqlx::query(&format!("CREATE SCHEMA IF NOT EXISTS {}", schema))
40            .execute(pool)
41            .await?;
42    }
43
44    for table in CONFIG_TABLES {
45        let q_table = qualified_sys_table(table);
46        let ddl = format!(
47            "CREATE TABLE IF NOT EXISTS {} (\
48                id TEXT NOT NULL, \
49                package_id TEXT NOT NULL, \
50                payload {} NOT NULL, \
51                updated_at {} NOT NULL DEFAULT {}, \
52                version BIGINT NOT NULL DEFAULT 1, \
53                PRIMARY KEY (id, package_id)\
54            )",
55            q_table,
56            dialect.sys_json_type(),
57            dialect.sys_timestamp_type(),
58            dialect.now_fn(),
59        );
60        sqlx::query(&ddl).execute(pool).await?;
61        let alter_version = format!(
62            "ALTER TABLE {} ADD COLUMN IF NOT EXISTS version BIGINT NOT NULL DEFAULT 1",
63            q_table
64        );
65        let _ = sqlx::query(&alter_version).execute(pool).await;
66        let alter_package = format!(
67            "ALTER TABLE {} ADD COLUMN IF NOT EXISTS package_id TEXT NOT NULL DEFAULT '{}'",
68            q_table, DEFAULT_PACKAGE_ID
69        );
70        let _ = sqlx::query(&alter_package).execute(pool).await;
71
72        let history_table = qualified_sys_table(&format!("{}_history", table));
73        let history_ddl = format!(
74            "CREATE TABLE IF NOT EXISTS {} (\
75                id TEXT NOT NULL, \
76                package_id TEXT NOT NULL, \
77                payload {} NOT NULL, \
78                version BIGINT NOT NULL, \
79                created_at {} NOT NULL DEFAULT {}, \
80                PRIMARY KEY (id, package_id, version)\
81            )",
82            history_table,
83            dialect.sys_json_type(),
84            dialect.sys_timestamp_type(),
85            dialect.now_fn(),
86        );
87        sqlx::query(&history_ddl).execute(pool).await?;
88        let alter_history_package = format!(
89            "ALTER TABLE {} ADD COLUMN IF NOT EXISTS package_id TEXT NOT NULL DEFAULT '{}'",
90            history_table, DEFAULT_PACKAGE_ID
91        );
92        let _ = sqlx::query(&alter_history_package).execute(pool).await;
93    }
94
95    let q_packages = qualified_sys_table("_sys_packages");
96    let packages_ddl = format!(
97        "CREATE TABLE IF NOT EXISTS {} (\
98            id TEXT PRIMARY KEY, \
99            payload {} NOT NULL, \
100            updated_at {} NOT NULL DEFAULT {}, \
101            version BIGINT NOT NULL DEFAULT 1, \
102            semantic_version TEXT\
103        )",
104        q_packages,
105        dialect.sys_json_type(),
106        dialect.sys_timestamp_type(),
107        dialect.now_fn(),
108    );
109    sqlx::query(&packages_ddl).execute(pool).await?;
110    let alter_pkg_semver = format!(
111        "ALTER TABLE {} ADD COLUMN IF NOT EXISTS semantic_version TEXT",
112        q_packages
113    );
114    let _ = sqlx::query(&alter_pkg_semver).execute(pool).await;
115    let q_packages_history = qualified_sys_table("_sys_packages_history");
116    let packages_history_ddl = format!(
117        "CREATE TABLE IF NOT EXISTS {} (\
118            id TEXT NOT NULL, \
119            payload {} NOT NULL, \
120            version BIGINT NOT NULL, \
121            created_at {} NOT NULL DEFAULT {}, \
122            semantic_version TEXT, \
123            PRIMARY KEY (id, version)\
124        )",
125        q_packages_history,
126        dialect.sys_json_type(),
127        dialect.sys_timestamp_type(),
128        dialect.now_fn(),
129    );
130    sqlx::query(&packages_history_ddl).execute(pool).await?;
131    let alter_pkg_hist_semver = format!(
132        "ALTER TABLE {} ADD COLUMN IF NOT EXISTS semantic_version TEXT",
133        q_packages_history
134    );
135    let _ = sqlx::query(&alter_pkg_hist_semver).execute(pool).await;
136    // Migrate to surrogate PK so multiple uninstalls of the same package (same id/version) don't violate uniqueness
137    let add_history_id = format!(
138        "ALTER TABLE {} ADD COLUMN IF NOT EXISTS history_id {}",
139        q_packages_history,
140        dialect.sys_bigserial_type()
141    );
142    let _ = sqlx::query(&add_history_id).execute(pool).await;
143    let drop_old_pk = format!(
144        "ALTER TABLE {} DROP CONSTRAINT IF EXISTS _sys_packages_history_pkey",
145        q_packages_history
146    );
147    let _ = sqlx::query(&drop_old_pk).execute(pool).await;
148    if dialect.name() == "postgres" {
149        let add_new_pk_cond = format!(
150            "DO $$ BEGIN IF NOT EXISTS (SELECT 1 FROM pg_constraint WHERE conname = '_sys_packages_history_history_id_pkey') THEN \
151             ALTER TABLE {} ADD CONSTRAINT _sys_packages_history_history_id_pkey PRIMARY KEY (history_id); END IF; END $$",
152            q_packages_history
153        );
154        let _ = sqlx::query(&add_new_pk_cond).execute(pool).await;
155    }
156
157    let q_tenants = qualified_sys_table("_sys_tenants");
158    let tenants_ddl = format!(
159        "CREATE TABLE IF NOT EXISTS {} (\
160            id TEXT PRIMARY KEY, \
161            strategy TEXT NOT NULL, \
162            database_url TEXT, \
163            updated_at {} NOT NULL DEFAULT {}, \
164            comment TEXT\
165        )",
166        q_tenants,
167        dialect.sys_timestamp_type(),
168        dialect.now_fn(),
169    );
170    sqlx::query(&tenants_ddl).execute(pool).await?;
171    let drop_schema_name = format!(
172        "ALTER TABLE {} DROP COLUMN IF EXISTS schema_name",
173        q_tenants
174    );
175    let _ = sqlx::query(&drop_schema_name).execute(pool).await;
176
177    // Auto-provision the Platform Admin tenant (sole writer of `global` shared tables) so global
178    // tables are writable out of the box. Idempotent; only meaningful where RLS is supported.
179    if dialect.supports_rls() {
180        let platform_id = crate::tenant::platform_tenant_id();
181        let insert_platform = format!(
182            "INSERT INTO {} (id, strategy, database_url, updated_at, comment) \
183             VALUES ('{}', 'rls', NULL, {}, 'Platform Admin (auto-provisioned): sole writer of global tables') \
184             ON CONFLICT (id) DO NOTHING",
185            q_tenants,
186            platform_id.replace('\'', "''"),
187            dialect.now_fn(),
188        );
189        sqlx::query(&insert_platform).execute(pool).await?;
190    }
191
192    let q_kv_data = qualified_sys_table("_sys_kv_data");
193    let kv_data_ddl = format!(
194        "CREATE TABLE IF NOT EXISTS {} (\
195            tenant_id TEXT NOT NULL, \
196            package_id TEXT NOT NULL, \
197            namespace TEXT NOT NULL, \
198            key TEXT NOT NULL, \
199            value {} NOT NULL, \
200            updated_at {} NOT NULL DEFAULT {}, \
201            PRIMARY KEY (tenant_id, package_id, namespace, key)\
202        )",
203        q_kv_data,
204        dialect.sys_json_type(),
205        dialect.sys_timestamp_type(),
206        dialect.now_fn(),
207    );
208    sqlx::query(&kv_data_ddl).execute(pool).await?;
209    // Migrate existing tables that had no tenant_id: add column and new PK.
210    let alter_kv_tenant = format!(
211        "ALTER TABLE {} ADD COLUMN IF NOT EXISTS tenant_id TEXT NOT NULL DEFAULT '_shared'",
212        q_kv_data
213    );
214    let _ = sqlx::query(&alter_kv_tenant).execute(pool).await;
215    let drop_pk = format!(
216        "ALTER TABLE {} DROP CONSTRAINT IF EXISTS _sys_kv_data_pkey",
217        q_kv_data
218    );
219    let _ = sqlx::query(&drop_pk).execute(pool).await;
220    let add_pk = format!(
221        "ALTER TABLE {} ADD PRIMARY KEY (tenant_id, package_id, namespace, key)",
222        q_kv_data
223    );
224    let _ = sqlx::query(&add_pk).execute(pool).await;
225    // Ensure value column is JSON type (for existing tables that had value as text).
226    if dialect.name() == "postgres" {
227        let alter_value_json = format!(
228            "ALTER TABLE {} ALTER COLUMN value TYPE JSONB USING value::jsonb",
229            q_kv_data
230        );
231        let _ = sqlx::query(&alter_value_json).execute(pool).await;
232    }
233
234    ensure_migration_tables(pool, dialect).await?;
235
236    Ok(())
237}
238
239/// Reserved `_sys_kv_data` namespace holding cached report results. Not a configured KV store, so
240/// the normal namespace-existence check is bypassed for these rows.
241pub const REPORT_CACHE_NAMESPACE: &str = "__report_cache__";
242
243/// Fetch a cached report envelope (the JSON value stored under the reserved cache namespace), or
244/// `None` when absent. Expiry and request-match verification are the caller's responsibility — the
245/// TTL lives as a field inside the returned envelope, not as a column. The lookup key is a bounded
246/// hash, so a collision (astronomically unlikely) surfaces as an envelope that fails the caller's
247/// verification and is treated as a miss, never a wrong hit.
248pub async fn report_cache_get(
249    pool: &Pool,
250    tenant_id: &str,
251    package_id: &str,
252    cache_key: &str,
253) -> Result<Option<serde_json::Value>, AppError> {
254    let q = qualified_sys_table("_sys_kv_data");
255    let sql = format!(
256        "SELECT value FROM {} WHERE tenant_id = $1 AND package_id = $2 AND namespace = $3 AND key = $4",
257        q
258    );
259    let row: Option<(serde_json::Value,)> = sqlx::query_as(&sql)
260        .bind(tenant_id)
261        .bind(package_id)
262        .bind(REPORT_CACHE_NAMESPACE)
263        .bind(cache_key)
264        .fetch_optional(pool)
265        .await
266        .map_err(AppError::Db)?;
267    Ok(row.map(|(v,)| v))
268}
269
270/// Upsert a cached report envelope into `_sys_kv_data` under the reserved cache namespace. Uses the
271/// UPDATE-then-INSERT pattern (like `kv_put`) to stay dialect-portable.
272pub async fn report_cache_put(
273    pool: &Pool,
274    dialect: &dyn Dialect,
275    tenant_id: &str,
276    package_id: &str,
277    cache_key: &str,
278    envelope: &serde_json::Value,
279) -> Result<(), AppError> {
280    let q = qualified_sys_table("_sys_kv_data");
281    let update_sql = format!(
282        "UPDATE {} SET value = $5, updated_at = {} \
283         WHERE tenant_id = $1 AND package_id = $2 AND namespace = $3 AND key = $4",
284        q,
285        dialect.now_fn()
286    );
287    let res = sqlx::query(&update_sql)
288        .bind(tenant_id)
289        .bind(package_id)
290        .bind(REPORT_CACHE_NAMESPACE)
291        .bind(cache_key)
292        .bind(envelope)
293        .execute(pool)
294        .await
295        .map_err(AppError::Db)?;
296    if res.rows_affected() == 0 {
297        let insert_sql = format!(
298            "INSERT INTO {} (tenant_id, package_id, namespace, key, value) VALUES ($1, $2, $3, $4, $5)",
299            q
300        );
301        sqlx::query(&insert_sql)
302            .bind(tenant_id)
303            .bind(package_id)
304            .bind(REPORT_CACHE_NAMESPACE)
305            .bind(cache_key)
306            .bind(envelope)
307            .execute(pool)
308            .await
309            .map_err(AppError::Db)?;
310    }
311    Ok(())
312}
313
314/// Create _sys_migration_plans and _sys_migration_audit tables if they don't exist.
315async fn ensure_migration_tables(pool: &Pool, dialect: &dyn Dialect) -> Result<(), AppError> {
316    let q_plans = qualified_sys_table("_sys_migration_plans");
317    let expires_at_col = match dialect.default_now_plus_hours(24) {
318        Some(expr) => format!(
319            "expires_at {} NOT NULL DEFAULT {}",
320            dialect.sys_timestamp_type(),
321            expr
322        ),
323        None => format!("expires_at {}", dialect.sys_timestamp_type()),
324    };
325    sqlx::query(&format!(
326        "CREATE TABLE IF NOT EXISTS {} (\
327            id TEXT PRIMARY KEY, \
328            package_id TEXT NOT NULL, \
329            tenant_id TEXT NOT NULL, \
330            from_version TEXT, \
331            to_version TEXT NOT NULL, \
332            plan_json {} NOT NULL, \
333            zip_bytes {} NOT NULL, \
334            status TEXT NOT NULL DEFAULT 'pending', \
335            created_at {} NOT NULL DEFAULT {}, \
336            {}, \
337            applied_at {}\
338        )",
339        q_plans,
340        dialect.sys_json_type(),
341        dialect.sys_bytes_type(),
342        dialect.sys_timestamp_type(),
343        dialect.now_fn(),
344        expires_at_col,
345        dialect.sys_timestamp_type(),
346    ))
347    .execute(pool)
348    .await?;
349
350    let q_audit = qualified_sys_table("_sys_migration_audit");
351    sqlx::query(&format!(
352        "CREATE TABLE IF NOT EXISTS {} (\
353            id {} PRIMARY KEY, \
354            migration_plan_id TEXT NOT NULL, \
355            package_id TEXT NOT NULL, \
356            tenant_id TEXT NOT NULL, \
357            from_version TEXT, \
358            to_version TEXT NOT NULL, \
359            step_number INT NOT NULL, \
360            operation TEXT NOT NULL, \
361            schema_name TEXT NOT NULL, \
362            table_name TEXT, \
363            object_name TEXT NOT NULL, \
364            object_type TEXT NOT NULL, \
365            description TEXT NOT NULL, \
366            ddl TEXT, \
367            safety TEXT NOT NULL, \
368            risk TEXT NOT NULL, \
369            status TEXT NOT NULL, \
370            error_message TEXT, \
371            executed_at {} NOT NULL DEFAULT {}\
372        )",
373        q_audit,
374        dialect.sys_bigserial_type(),
375        dialect.sys_timestamp_type(),
376        dialect.now_fn(),
377    ))
378    .execute(pool)
379    .await?;
380
381    Ok(())
382}
383
384/// Row returned from _sys_migration_plans.
385pub struct MigrationPlanRow {
386    pub id: String,
387    pub package_id: String,
388    pub tenant_id: String,
389    pub from_version: Option<String>,
390    pub to_version: String,
391    pub plan_json: serde_json::Value,
392    pub zip_bytes: Vec<u8>,
393    pub status: String,
394    pub created_at: DateTime<Utc>,
395    pub expires_at: DateTime<Utc>,
396    pub applied_at: Option<DateTime<Utc>>,
397}
398
399/// Persist a migration plan (zip bytes + serialized steps) for later confirmation.
400#[allow(clippy::too_many_arguments)]
401pub async fn save_migration_plan(
402    pool: &Pool,
403    id: &str,
404    package_id: &str,
405    tenant_id: &str,
406    from_version: Option<&str>,
407    to_version: &str,
408    plan_json: &serde_json::Value,
409    zip_bytes: &[u8],
410) -> Result<(), AppError> {
411    let q = qualified_sys_table("_sys_migration_plans");
412    sqlx::query(&format!(
413        "INSERT INTO {} (id, package_id, tenant_id, from_version, to_version, plan_json, zip_bytes, status, created_at, expires_at) \
414         VALUES ($1, $2, $3, $4, $5, $6, $7, 'pending', NOW(), NOW() + INTERVAL '24 hours')",
415        q
416    ))
417    .bind(id)
418    .bind(package_id)
419    .bind(tenant_id)
420    .bind(from_version)
421    .bind(to_version)
422    .bind(plan_json)
423    .bind(zip_bytes)
424    .execute(pool)
425    .await?;
426    Ok(())
427}
428
429/// Fetch a migration plan by id, or None if not found.
430pub async fn get_migration_plan(
431    pool: &Pool,
432    id: &str,
433) -> Result<Option<MigrationPlanRow>, AppError> {
434    let q = qualified_sys_table("_sys_migration_plans");
435    #[allow(clippy::type_complexity)]
436    let row: Option<(String, String, String, Option<String>, String, serde_json::Value, Vec<u8>, String, DateTime<Utc>, DateTime<Utc>, Option<DateTime<Utc>>)> =
437        sqlx::query_as(&format!(
438            "SELECT id, package_id, tenant_id, from_version, to_version, plan_json, zip_bytes, status, created_at, expires_at, applied_at FROM {} WHERE id = $1",
439            q
440        ))
441        .bind(id)
442        .fetch_optional(pool)
443        .await
444        .map_err(AppError::Db)?;
445    Ok(row.map(
446        |(
447            id,
448            package_id,
449            tenant_id,
450            from_version,
451            to_version,
452            plan_json,
453            zip_bytes,
454            status,
455            created_at,
456            expires_at,
457            applied_at,
458        )| {
459            MigrationPlanRow {
460                id,
461                package_id,
462                tenant_id,
463                from_version,
464                to_version,
465                plan_json,
466                zip_bytes,
467                status,
468                created_at,
469                expires_at,
470                applied_at,
471            }
472        },
473    ))
474}
475
476/// Atomically mark a migration plan as applied. Returns false if already applied or not found.
477pub async fn mark_migration_plan_applied(pool: &Pool, id: &str) -> Result<bool, AppError> {
478    let q = qualified_sys_table("_sys_migration_plans");
479    let result = sqlx::query(&format!(
480        "UPDATE {} SET status = 'applied', applied_at = NOW() WHERE id = $1 AND status = 'pending'",
481        q
482    ))
483    .bind(id)
484    .execute(pool)
485    .await?;
486    Ok(result.rows_affected() > 0)
487}
488
489/// Append one audit record for a migration step execution.
490#[allow(clippy::too_many_arguments)]
491pub async fn insert_migration_audit(
492    pool: &Pool,
493    migration_plan_id: &str,
494    package_id: &str,
495    tenant_id: &str,
496    from_version: Option<&str>,
497    to_version: &str,
498    step_number: i32,
499    operation: &str,
500    schema_name: &str,
501    table_name: Option<&str>,
502    object_name: &str,
503    object_type: &str,
504    description: &str,
505    ddl: Option<&str>,
506    safety: &str,
507    risk: &str,
508    status: &str,
509    error_message: Option<&str>,
510) -> Result<(), AppError> {
511    let q = qualified_sys_table("_sys_migration_audit");
512    sqlx::query(&format!(
513        "INSERT INTO {} (migration_plan_id, package_id, tenant_id, from_version, to_version, step_number, operation, schema_name, table_name, object_name, object_type, description, ddl, safety, risk, status, error_message, executed_at) \
514         VALUES ($1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11,$12,$13,$14,$15,$16,$17,NOW())",
515        q
516    ))
517    .bind(migration_plan_id)
518    .bind(package_id)
519    .bind(tenant_id)
520    .bind(from_version)
521    .bind(to_version)
522    .bind(step_number)
523    .bind(operation)
524    .bind(schema_name)
525    .bind(table_name)
526    .bind(object_name)
527    .bind(object_type)
528    .bind(description)
529    .bind(ddl)
530    .bind(safety)
531    .bind(risk)
532    .bind(status)
533    .bind(error_message)
534    .execute(pool)
535    .await?;
536    Ok(())
537}
538
539/// Resolve the storage id for a config record. For api_entities, entity_id is used when id is absent.
540fn config_record_id(table: &str, rec: &serde_json::Value) -> Result<String, AppError> {
541    let id = rec.get("id").and_then(|v| v.as_str());
542    let entity_id = rec.get("entity_id").and_then(|v| v.as_str());
543    match (table, id, entity_id) {
544        ("_sys_api_entities", None, Some(eid)) => Ok(eid.to_string()),
545        (_, Some(id), _) => Ok(id.to_string()),
546        _ => Err(AppError::BadRequest(
547            "each config record must have an 'id' field (or 'entity_id' for api_entities)".into(),
548        )),
549    }
550}
551
552/// Deep-compare incoming records with current stored payloads (by id).
553/// Returns true if they are identical (same ids and same payload per id).
554fn config_payloads_unchanged(
555    table: &str,
556    current: &HashMap<String, serde_json::Value>,
557    records: &[serde_json::Value],
558) -> Result<bool, AppError> {
559    if current.len() != records.len() {
560        return Ok(false);
561    }
562    for rec in records {
563        let id = config_record_id(table, rec)?;
564        match current.get(&id) {
565            None => return Ok(false),
566            Some(existing) if existing != rec => return Ok(false),
567            Some(_) => {}
568        }
569    }
570    Ok(true)
571}
572
573/// Replace all rows for a config type for one package: copy current (for this package_id) to history, delete, insert with new version.
574/// If incoming payloads are deep-equal to current, no write is performed and no new version is created.
575/// Returns (count inserted, version). Call within transaction for atomicity.
576pub async fn replace_config_rows(
577    tx: &mut crate::db::pool::Connection,
578    table: &str,
579    package_id: &str,
580    records: &[serde_json::Value],
581) -> Result<(u64, i64), AppError> {
582    let q_table = qualified_sys_table(table);
583    let current_version: (Option<i64>,) = sqlx::query_as(&format!(
584        "SELECT COALESCE(MAX(version), 0) FROM {} WHERE package_id = $1",
585        q_table
586    ))
587    .bind(package_id)
588    .fetch_one(&mut *tx)
589    .await
590    .map_err(AppError::Db)?;
591    let current_version = current_version.0.unwrap_or(0);
592
593    let rows: Vec<(String, serde_json::Value)> = sqlx::query_as(&format!(
594        "SELECT id, payload FROM {} WHERE package_id = $1",
595        q_table
596    ))
597    .bind(package_id)
598    .fetch_all(&mut *tx)
599    .await
600    .map_err(AppError::Db)?;
601    let current: HashMap<String, serde_json::Value> = rows.into_iter().collect();
602
603    if config_payloads_unchanged(table, &current, records)? {
604        return Ok((0, current_version));
605    }
606
607    let history_table = qualified_sys_table(&format!("{}_history", table));
608    let new_version = current_version + 1;
609
610    sqlx::query(&format!(
611        "INSERT INTO {} (id, package_id, payload, version, created_at) SELECT id, package_id, payload, version, updated_at FROM {} WHERE package_id = $1",
612        history_table, q_table
613    ))
614    .bind(package_id)
615    .execute(&mut *tx)
616    .await?;
617
618    sqlx::query(&format!("DELETE FROM {} WHERE package_id = $1", q_table))
619        .bind(package_id)
620        .execute(&mut *tx)
621        .await?;
622
623    let mut count = 0u64;
624    for rec in records {
625        let id = config_record_id(table, rec)?;
626        sqlx::query(&format!(
627            "INSERT INTO {} (id, package_id, payload, updated_at, version) VALUES ($1, $2, $3, NOW(), $4)",
628            q_table
629        ))
630        .bind(&id)
631        .bind(package_id)
632        .bind(rec)
633        .bind(new_version)
634        .execute(&mut *tx)
635        .await?;
636        count += 1;
637    }
638    Ok((count, new_version))
639}
640
641const PACKAGES_TABLE: &str = "_sys_packages";
642const PACKAGES_HISTORY_TABLE: &str = "_sys_packages_history";
643
644pub struct PackageRow {
645    pub id: String,
646    pub payload: serde_json::Value,
647    pub version: i64,
648    pub updated_at: DateTime<Utc>,
649    pub semantic_version: Option<String>,
650}
651
652/// List all rows from _sys_packages ordered by id.
653pub async fn list_packages(pool: &Pool) -> Result<Vec<PackageRow>, AppError> {
654    let q = qualified_sys_table(PACKAGES_TABLE);
655    #[allow(clippy::type_complexity)]
656    let rows: Vec<(
657        String,
658        serde_json::Value,
659        i64,
660        DateTime<Utc>,
661        Option<String>,
662    )> = sqlx::query_as(&format!(
663        "SELECT id, payload, version, updated_at, semantic_version FROM {} ORDER BY id",
664        q
665    ))
666    .fetch_all(pool)
667    .await
668    .map_err(AppError::Db)?;
669    Ok(rows
670        .into_iter()
671        .map(
672            |(id, payload, version, updated_at, semantic_version)| PackageRow {
673                id,
674                payload,
675                version,
676                updated_at,
677                semantic_version,
678            },
679        )
680        .collect())
681}
682
683/// Fetch a single package row by id, or None if not installed.
684pub async fn get_package(pool: &Pool, id: &str) -> Result<Option<PackageRow>, AppError> {
685    let q = qualified_sys_table(PACKAGES_TABLE);
686    #[allow(clippy::type_complexity)]
687    let row: Option<(
688        String,
689        serde_json::Value,
690        i64,
691        DateTime<Utc>,
692        Option<String>,
693    )> = sqlx::query_as(&format!(
694        "SELECT id, payload, version, updated_at, semantic_version FROM {} WHERE id = $1",
695        q
696    ))
697    .bind(id)
698    .fetch_optional(pool)
699    .await
700    .map_err(AppError::Db)?;
701    Ok(row.map(
702        |(id, payload, version, updated_at, semantic_version)| PackageRow {
703            id,
704            payload,
705            version,
706            updated_at,
707            semantic_version,
708        },
709    ))
710}
711
712/// Count rows in a config table for a given package.
713pub async fn count_package_kind(
714    pool: &Pool,
715    kind: &str,
716    package_id: &str,
717) -> Result<i64, AppError> {
718    let table = sys_table_for_kind(kind)
719        .ok_or_else(|| AppError::BadRequest(format!("unknown config kind: {}", kind)))?;
720    let q = qualified_sys_table(table);
721    let (count,): (i64,) =
722        sqlx::query_as(&format!("SELECT COUNT(*) FROM {} WHERE package_id = $1", q))
723            .bind(package_id)
724            .fetch_one(pool)
725            .await
726            .map_err(AppError::Db)?;
727    Ok(count)
728}
729
730/// List all package ids from _sys_packages (what is installed in the DB). Used to generate OpenAPI spec from _sys_* config.
731pub async fn list_package_ids(pool: &Pool) -> Result<Vec<String>, AppError> {
732    let q = qualified_sys_table(PACKAGES_TABLE);
733    let rows: Vec<(String,)> = sqlx::query_as(&format!("SELECT id FROM {} ORDER BY id", q))
734        .fetch_all(pool)
735        .await
736        .map_err(AppError::Db)?;
737    Ok(rows.into_iter().map(|(id,)| id).collect())
738}
739
740/// Upsert one package row by id: copy current to history if exists, then insert or replace with new payload.
741/// Semantic version is read from payload.version (e.g. manifest "version": "1.0.0"). Version is only incremented when semantic_version changes.
742pub async fn upsert_package(
743    pool: &Pool,
744    id: &str,
745    payload: &serde_json::Value,
746) -> Result<i64, AppError> {
747    let semantic_version = payload
748        .get("version")
749        .and_then(serde_json::Value::as_str)
750        .map(String::from)
751        .unwrap_or_default();
752
753    let q_packages = qualified_sys_table(PACKAGES_TABLE);
754    let q_packages_history = qualified_sys_table(PACKAGES_HISTORY_TABLE);
755    let mut tx = pool.begin().await?;
756    let current: Option<(serde_json::Value, i64, Option<String>)> = sqlx::query_as(&format!(
757        "SELECT payload, version, semantic_version FROM {} WHERE id = $1",
758        q_packages
759    ))
760    .bind(id)
761    .fetch_optional(&mut *tx)
762    .await
763    .map_err(AppError::Db)?;
764
765    let new_version = match &current {
766        Some((_, v, Some(ref old_semver))) if *old_semver == semantic_version => *v,
767        Some((_, v, _)) => v + 1,
768        None => 1,
769    };
770
771    if let Some((old_payload, old_version, old_semver)) = current {
772        sqlx::query(&format!(
773            "INSERT INTO {} (id, payload, version, created_at, semantic_version) VALUES ($1, $2, $3, NOW(), $4)",
774            q_packages_history
775        ))
776        .bind(id)
777        .bind(old_payload)
778        .bind(old_version)
779        .bind(old_semver)
780        .execute(&mut *tx)
781        .await?;
782    }
783
784    sqlx::query(&format!("DELETE FROM {} WHERE id = $1", q_packages))
785        .bind(id)
786        .execute(&mut *tx)
787        .await?;
788
789    let semver_param: Option<&str> = if semantic_version.is_empty() {
790        None
791    } else {
792        Some(semantic_version.as_str())
793    };
794    sqlx::query(&format!(
795        "INSERT INTO {} (id, payload, updated_at, version, semantic_version) VALUES ($1, $2, NOW(), $3, $4)",
796        q_packages
797    ))
798    .bind(id)
799    .bind(payload)
800    .bind(new_version)
801    .bind(semver_param)
802    .execute(&mut *tx)
803    .await?;
804
805    tx.commit().await?;
806    Ok(new_version)
807}
808
809/// Delete all config rows and KV data for a package, then remove the package record.
810/// Copies the current package row to _sys_packages_history before delete. Call after reverting migrations on the tenant DB.
811pub async fn delete_package_and_config(pool: &Pool, package_id: &str) -> Result<(), AppError> {
812    let q_packages = qualified_sys_table(PACKAGES_TABLE);
813    let q_packages_history = qualified_sys_table(PACKAGES_HISTORY_TABLE);
814    let q_kv_data = qualified_sys_table("_sys_kv_data");
815
816    let mut tx = pool.begin().await?;
817
818    // Copy current package row to history (if exists)
819    let current: Option<(serde_json::Value, i64, Option<String>)> = sqlx::query_as(&format!(
820        "SELECT payload, version, semantic_version FROM {} WHERE id = $1",
821        q_packages
822    ))
823    .bind(package_id)
824    .fetch_optional(&mut *tx)
825    .await
826    .map_err(AppError::Db)?;
827
828    if let Some((payload, version, semantic_version)) = current {
829        sqlx::query(&format!(
830            "INSERT INTO {} (id, payload, version, created_at, semantic_version) VALUES ($1, $2, $3, NOW(), $4)",
831            q_packages_history
832        ))
833        .bind(package_id)
834        .bind(payload)
835        .bind(version)
836        .bind(semantic_version)
837        .execute(&mut *tx)
838        .await?;
839    }
840
841    // Delete from each config table and its history (by package_id)
842    for table in CONFIG_TABLES {
843        let q_table = qualified_sys_table(table);
844        sqlx::query(&format!("DELETE FROM {} WHERE package_id = $1", q_table))
845            .bind(package_id)
846            .execute(&mut *tx)
847            .await?;
848        let history_table = qualified_sys_table(&format!("{}_history", table));
849        sqlx::query(&format!(
850            "DELETE FROM {} WHERE package_id = $1",
851            history_table
852        ))
853        .bind(package_id)
854        .execute(&mut *tx)
855        .await?;
856    }
857
858    // Delete KV data for this package
859    sqlx::query(&format!("DELETE FROM {} WHERE package_id = $1", q_kv_data))
860        .bind(package_id)
861        .execute(&mut *tx)
862        .await?;
863
864    // Delete package row
865    sqlx::query(&format!("DELETE FROM {} WHERE id = $1", q_packages))
866        .bind(package_id)
867        .execute(&mut *tx)
868        .await?;
869
870    tx.commit().await?;
871    Ok(())
872}
873
874/// Create a connection pool for the compiled-in dialect.
875///
876/// This is the recommended way to build a pool in consumer binaries — it uses the correct
877/// pool type for whichever dialect feature is active without requiring `#[cfg(feature = ...)]`
878/// in caller code.
879pub async fn create_pool(database_url: &str, max_connections: u32) -> Result<Pool, AppError> {
880    #[cfg(feature = "postgres")]
881    return sqlx::postgres::PgPoolOptions::new()
882        .max_connections(max_connections)
883        .connect(database_url)
884        .await
885        .map_err(AppError::Db);
886
887    #[cfg(feature = "mysql")]
888    return sqlx::mysql::MySqlPoolOptions::new()
889        .max_connections(max_connections)
890        .connect(database_url)
891        .await
892        .map_err(AppError::Db);
893
894    #[cfg(feature = "sqlite")]
895    return sqlx::sqlite::SqlitePoolOptions::new()
896        .max_connections(max_connections)
897        .connect(database_url)
898        .await
899        .map_err(AppError::Db);
900
901    #[cfg(not(any(feature = "postgres", feature = "mysql", feature = "sqlite")))]
902    Err(AppError::BadRequest(
903        "No database dialect feature enabled. Enable one of: postgres, mysql, sqlite.".into(),
904    ))
905}
906
907/// Ensure the database in `database_url` exists; create it if not.
908///
909/// - **Postgres**: connects to the admin `postgres` database and runs `CREATE DATABASE`.
910/// - **SQLite**: the database file is created automatically on first connect — this is a no-op.
911/// - **MySQL**: database auto-creation is not supported here; create the database manually or
912///   rely on connection-string options like `createDatabaseIfNotExist=true`.
913pub async fn ensure_database_exists(database_url: &str) -> Result<(), AppError> {
914    #[cfg(feature = "postgres")]
915    {
916        use sqlx::ConnectOptions as _;
917        use std::str::FromStr;
918        let (admin_url, db_name) = parse_db_name_from_url(database_url)?;
919        if db_name.is_empty() || db_name == "postgres" {
920            return Ok(());
921        }
922        let opts = sqlx::postgres::PgConnectOptions::from_str(&admin_url)
923            .map_err(|e| AppError::BadRequest(format!("invalid DATABASE_URL: {}", e)))?;
924        let mut conn: sqlx::PgConnection = opts.connect().await.map_err(AppError::Db)?;
925        let exists: (bool,) =
926            sqlx::query_as("SELECT EXISTS(SELECT 1 FROM pg_database WHERE datname = $1)")
927                .bind(&db_name)
928                .fetch_one(&mut conn)
929                .await
930                .map_err(AppError::Db)?;
931        if !exists.0 {
932            let quoted = quote_ident(&db_name);
933            sqlx::query(&format!("CREATE DATABASE {}", quoted))
934                .execute(&mut conn)
935                .await
936                .map_err(AppError::Db)?;
937        }
938    }
939    #[cfg(not(feature = "postgres"))]
940    let _ = database_url;
941    Ok(())
942}
943
944#[cfg(feature = "postgres")]
945fn parse_db_name_from_url(url: &str) -> Result<(String, String), AppError> {
946    let path_start = url
947        .rfind('/')
948        .ok_or_else(|| AppError::BadRequest("DATABASE_URL: no path".into()))?
949        + 1;
950    let path_and_query = url.get(path_start..).unwrap_or("");
951    let db_name = path_and_query.split('?').next().unwrap_or("").trim();
952    let base = url.get(..path_start).unwrap_or(url);
953    let admin_url = format!("{}postgres", base);
954    Ok((admin_url, db_name.to_string()))
955}
956
957#[cfg(feature = "postgres")]
958fn quote_ident(name: &str) -> String {
959    format!("\"{}\"", name.replace('\\', "\\\\").replace('"', "\\\""))
960}
961
962pub fn sys_table_for_kind(kind: &str) -> Option<&'static str> {
963    match kind {
964        "schemas" => Some("_sys_schemas"),
965        "enums" => Some("_sys_enums"),
966        "tables" => Some("_sys_tables"),
967        "columns" => Some("_sys_columns"),
968        "indexes" => Some("_sys_indexes"),
969        "relationships" => Some("_sys_relationships"),
970        "api_entities" => Some("_sys_api_entities"),
971        "kv_stores" => Some("_sys_kv_stores"),
972        "reports" => Some("_sys_reports"),
973        _ => None,
974    }
975}