1use crate::db::{pool::Pool, Dialect};
4use crate::error::AppError;
5use chrono::{DateTime, Utc};
6use std::collections::HashMap;
7
8pub fn architect_schema() -> String {
10 std::env::var("ARCHITECT_SCHEMA").unwrap_or_else(|_| "architect".into())
11}
12
13pub fn qualified_sys_table(table: &str) -> String {
15 format!("{}.{}", architect_schema(), table)
16}
17
18const 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
31pub const DEFAULT_PACKAGE_ID: &str = "_default";
33
34pub 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 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 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 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 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
239pub const REPORT_CACHE_NAMESPACE: &str = "__report_cache__";
242
243pub 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
270pub 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
314async 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
384pub 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#[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
429pub 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
476pub 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#[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
539fn 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
552fn 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
573pub 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, ¤t, 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
652pub 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
683pub 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
712pub 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
730pub 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
740pub 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 ¤t {
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
809pub 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 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 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 sqlx::query(&format!("DELETE FROM {} WHERE package_id = $1", q_kv_data))
860 .bind(package_id)
861 .execute(&mut *tx)
862 .await?;
863
864 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
874pub 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
907pub 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}