1use std::collections::HashMap;
30
31use sea_query::{
32 Alias, Asterisk, Condition, Expr, Func, Order, PostgresQueryBuilder, Query, SqliteQueryBuilder,
33 Value as SeaValue,
34};
35use sea_query_binder::SqlxBinder;
36use sqlx::Row;
37
38use crate::db::{DbPool, pool_for_dispatched};
39use crate::migrate::{Column, ModelMeta};
40use crate::orm::SqlType;
41use crate::orm::write::{WriteError, json_to_sea_value, null_for};
42
43fn resolve_pool_dyn(meta: &crate::migrate::ModelMeta, op: crate::db::RouteOp) -> crate::db::DbPool {
46 let ctx = crate::db::route_context();
47 let r = crate::db::router::router();
48 let alias = match op {
49 crate::db::RouteOp::Read => r.db_for_read(meta, &ctx),
50 crate::db::RouteOp::Write => r.db_for_write(meta, &ctx),
51 };
52 pool_for_dispatched(alias.as_str()).clone()
53}
54
55#[derive(Debug)]
78pub enum DynError {
79 Write(WriteError),
84 Sqlx(sqlx::Error),
88}
89
90impl std::fmt::Display for DynError {
91 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
92 match self {
93 Self::Write(e) => write!(f, "{e}"),
94 Self::Sqlx(e) => write!(f, "{e}"),
95 }
96 }
97}
98
99impl std::error::Error for DynError {
100 fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
101 match self {
102 Self::Write(e) => Some(e),
103 Self::Sqlx(e) => Some(e),
104 }
105 }
106}
107
108impl From<sqlx::Error> for DynError {
109 fn from(e: sqlx::Error) -> Self {
110 Self::Sqlx(e)
111 }
112}
113
114impl From<WriteError> for DynError {
115 fn from(e: WriteError) -> Self {
116 Self::Write(e)
117 }
118}
119
120pub struct DynQuerySet<'a> {
127 meta: &'a ModelMeta,
128 where_clauses: Vec<Condition>,
133 order: Vec<(String, bool)>,
134 limit: Option<u64>,
135 offset: Option<u64>,
136 select_cols: Vec<String>,
137 with_deleted: bool,
138 only_deleted: bool,
139 hard_delete: bool,
140 select_related: Vec<String>,
148 allow_privileged: Vec<String>,
154}
155
156impl<'a> DynQuerySet<'a> {
157 pub fn for_meta(meta: &'a ModelMeta) -> Self {
161 let select_cols = meta.fields.iter().map(|c| c.name.clone()).collect();
162 Self {
163 meta,
164 where_clauses: Vec::new(),
165 order: Vec::new(),
166 limit: None,
167 offset: None,
168 select_cols,
169 with_deleted: false,
170 only_deleted: false,
171 hard_delete: false,
172 select_related: Vec::new(),
173 allow_privileged: Vec::new(),
174 }
175 }
176
177 pub fn allow_privileged(mut self, cols: &[&str]) -> Self {
195 self.allow_privileged
196 .extend(cols.iter().map(|c| c.to_string()));
197 self
198 }
199
200 pub fn with_deleted(mut self) -> Self {
203 self.with_deleted = true;
204 self
205 }
206
207 pub fn only_deleted(mut self) -> Self {
210 self.only_deleted = true;
211 self
212 }
213
214 pub fn hard_delete(mut self) -> Self {
216 self.hard_delete = true;
217 self
218 }
219
220 fn effective_where_clauses(&self) -> Vec<Condition> {
221 let mut clauses = self.where_clauses.clone();
222 if self.meta.soft_delete {
223 if self.only_deleted {
224 clauses
225 .push(Condition::all().add(Expr::col(Alias::new("deleted_at")).is_not_null()));
226 } else if !self.with_deleted {
227 clauses.push(Condition::all().add(Expr::col(Alias::new("deleted_at")).is_null()));
228 }
229 }
230 clauses
231 }
232
233 fn live_where_clauses(&self) -> Vec<Condition> {
234 let mut clauses = self.where_clauses.clone();
235 if self.meta.soft_delete {
236 clauses.push(Condition::all().add(Expr::col(Alias::new("deleted_at")).is_null()));
237 }
238 clauses
239 }
240
241 pub fn select_cols(mut self, cols: &[String]) -> Self {
245 let valid: Vec<String> = cols
246 .iter()
247 .filter(|n| self.meta.fields.iter().any(|c| &c.name == *n))
248 .cloned()
249 .collect();
250 if !valid.is_empty() {
251 self.select_cols = valid;
252 }
253 self
254 }
255
256 pub fn select_related_dyn(mut self, fields: &[String]) -> Self {
283 for name in fields {
284 let canonical = normalize_sr_token(name);
285 if validate_sr_chain(self.meta, &canonical).is_none() {
286 continue;
287 }
288 if !self.select_related.iter().any(|n| n == &canonical) {
289 self.select_related.push(canonical);
290 }
291 }
292 self
293 }
294
295 #[doc(hidden)]
298 pub fn select_related_fields(&self) -> &[String] {
299 &self.select_related
300 }
301
302 pub fn search(mut self, fields: &[String], term: &str) -> Self {
326 let term = term.trim();
327 if term.is_empty() {
328 return self;
329 }
330
331 let restricted = !fields.is_empty();
332 let as_int = term.parse::<i64>().ok();
333 let as_float = term.parse::<f64>().ok();
334 let as_bool = match term.to_ascii_lowercase().as_str() {
335 "true" => Some(true),
336 "false" => Some(false),
337 _ => None,
338 };
339 let like_pat = format!("%{}%", crate::orm::escape_like_literal(term)).to_uppercase();
342
343 let mut cond = Condition::any();
344 let mut added = 0;
345 for col in &self.meta.fields {
346 if restricted && !fields.iter().any(|f| f == &col.name) {
347 continue;
348 }
349 let predicate: Option<sea_query::SimpleExpr> = match col.ty {
350 SqlType::Text => Some(
351 Expr::expr(Func::upper(Expr::col(Alias::new(&col.name))))
352 .like(sea_query::LikeExpr::new(like_pat.clone()).escape('\\')),
353 ),
354 SqlType::SmallInt | SqlType::Integer | SqlType::BigInt | SqlType::ForeignKey => {
355 as_int.map(|n| Expr::col(Alias::new(&col.name)).eq(n))
356 }
357 SqlType::Real | SqlType::Double => {
358 as_float.map(|n| Expr::col(Alias::new(&col.name)).eq(n))
359 }
360 SqlType::Boolean => as_bool.map(|b| Expr::col(Alias::new(&col.name)).eq(b)),
361 _ => None,
362 };
363 if let Some(p) = predicate {
364 cond = cond.add(p);
365 added += 1;
366 }
367 }
368 if added > 0 {
369 self.where_clauses.push(cond);
370 }
371 self
372 }
373
374 pub fn filter_condition(mut self, cond: sea_query::Condition) -> Self {
380 self.where_clauses.push(cond);
381 self
382 }
383
384 pub fn filter_in_i64(mut self, col: &str, vals: &[i64]) -> Self {
387 if vals.is_empty() || !self.meta.fields.iter().any(|c| c.name == col) {
388 return self;
389 }
390 let cond = Condition::all().add(Expr::col(Alias::new(col)).is_in(vals.iter().copied()));
391 self.where_clauses.push(cond);
392 self
393 }
394
395 pub fn filter_m2m_contains_any(mut self, field_name: &str, child_ids: &[String]) -> Self {
419 if child_ids.is_empty() {
420 return self;
421 }
422 let Some(rel) = self
423 .meta
424 .m2m_relations
425 .iter()
426 .find(|r| r.field_name == field_name)
427 else {
428 return self;
429 };
430 let Some(pk_col) = self.meta.pk_column() else {
431 return self;
432 };
433 let target_pk_ty = crate::migrate::pk_meta_for_table(&rel.target_table)
444 .map(|(_, ty)| ty)
445 .unwrap_or(SqlType::BigInt);
446 let junction_table = format!("{}_{}", self.meta.table, rel.field_name);
447 let child_id_expr = Expr::col(Alias::new("child_id"));
448 let in_clause: sea_query::SimpleExpr = match target_pk_ty {
449 SqlType::Text | SqlType::Uuid => {
450 let bound: Vec<String> = child_ids
454 .iter()
455 .filter_map(|s| {
456 let s = s.trim();
457 if s.is_empty() {
458 None
459 } else {
460 Some(s.to_string())
461 }
462 })
463 .collect();
464 if bound.is_empty() {
465 return self;
466 }
467 child_id_expr.is_in(bound)
468 }
469 _ => {
470 let parsed: Vec<i64> = child_ids.iter().filter_map(|s| s.parse().ok()).collect();
474 if parsed.is_empty() {
475 return self;
476 }
477 child_id_expr.is_in(parsed)
478 }
479 };
480 let subq = Query::select()
481 .column(Alias::new("parent_id"))
482 .from(crate::db::router::schema_qualified_table(&junction_table))
483 .and_where(in_clause)
484 .to_owned();
485 let cond =
486 Condition::all().add(Expr::col(Alias::new(pk_col.name.clone())).in_subquery(subq));
487 self.where_clauses.push(cond);
488 self
489 }
490
491 pub fn filter_in_strings(mut self, col: &str, vals: &[String]) -> Self {
503 let Some(meta_col) = self.meta.fields.iter().find(|c| c.name == col) else {
504 return self;
505 };
506 if vals.is_empty() {
507 return self;
508 }
509 let expr = Expr::col(Alias::new(col));
510 let cond = match crate::migrate::fk_effective_type(meta_col) {
517 SqlType::SmallInt | SqlType::Integer => {
518 let parsed: Vec<i32> = vals.iter().filter_map(|s| s.parse().ok()).collect();
519 if parsed.is_empty() {
520 return self;
521 }
522 Condition::all().add(expr.is_in(parsed))
523 }
524 SqlType::BigInt | SqlType::ForeignKey => {
525 let parsed: Vec<i64> = vals.iter().filter_map(|s| s.parse().ok()).collect();
526 if parsed.is_empty() {
527 return self;
528 }
529 Condition::all().add(expr.is_in(parsed))
530 }
531 SqlType::Real | SqlType::Double => {
532 let parsed: Vec<f64> = vals.iter().filter_map(|s| s.parse().ok()).collect();
533 if parsed.is_empty() {
534 return self;
535 }
536 Condition::all().add(expr.is_in(parsed))
537 }
538 SqlType::Boolean => {
539 let parsed: Vec<bool> = vals
540 .iter()
541 .map(|s| matches!(s.as_str(), "true" | "on" | "1"))
542 .collect();
543 Condition::all().add(expr.is_in(parsed))
544 }
545 SqlType::Uuid => {
550 let parsed: Vec<uuid::Uuid> = vals
551 .iter()
552 .filter_map(|s| uuid::Uuid::parse_str(s).ok())
553 .collect();
554 if parsed.is_empty() {
555 return self;
556 }
557 Condition::all().add(expr.is_in(parsed))
558 }
559 _ => Condition::all().add(expr.is_in(vals.iter().map(|s| s.to_string()))),
560 };
561 self.where_clauses.push(cond);
562 self
563 }
564
565 pub fn filter_eq_string(mut self, col: &str, value: &str) -> Self {
569 let Some(meta_col) = self.meta.fields.iter().find(|c| c.name == col) else {
570 return self;
571 };
572 let expr = Expr::col(Alias::new(col));
573 let predicate = match crate::migrate::fk_effective_type(meta_col) {
576 SqlType::SmallInt | SqlType::Integer => value.parse::<i32>().ok().map(|v| expr.eq(v)),
577 SqlType::BigInt | SqlType::ForeignKey => value.parse::<i64>().ok().map(|v| expr.eq(v)),
578 SqlType::Real | SqlType::Double => value.parse::<f64>().ok().map(|v| expr.eq(v)),
579 SqlType::Boolean => {
580 let v = matches!(value, "true" | "on" | "1");
581 Some(expr.eq(v))
582 }
583 SqlType::Uuid => uuid::Uuid::parse_str(value).ok().map(|u| expr.eq(u)),
586 _ => Some(expr.eq(value.to_string())),
587 };
588 if let Some(p) = predicate {
589 self.where_clauses.push(Condition::all().add(p));
590 }
591 self
592 }
593
594 pub fn order_by_col(mut self, col: &str, descending: bool) -> Self {
597 if self.meta.fields.iter().any(|c| c.name == col) {
598 self.order.push((col.to_string(), descending));
599 }
600 self
601 }
602
603 pub fn limit(mut self, n: u64) -> Self {
605 self.limit = Some(n);
606 self
607 }
608
609 pub fn offset(mut self, n: u64) -> Self {
611 self.offset = Some(n);
612 self
613 }
614
615 pub async fn count(self) -> Result<i64, DynError> {
619 let mut q = Query::select();
620 q.from(crate::db::router::schema_qualified_table(&self.meta.table));
621 q.expr(Func::count(Expr::col(Asterisk)));
622 let where_clauses = self.effective_where_clauses();
623 for cond in &where_clauses {
624 q.cond_where(cond.clone());
625 }
626
627 match resolve_pool_dyn(self.meta, crate::db::RouteOp::Read) {
628 DbPool::Sqlite(pool) => {
629 let (sql, values) = q.build_sqlx(SqliteQueryBuilder);
630 let row = sqlx::query_with(&sql, values).fetch_one(&pool).await?;
631 Ok(row.try_get::<i64, _>(0)?)
632 }
633 DbPool::Postgres(pool) => {
634 let (sql, values) = q.build_sqlx(PostgresQueryBuilder);
635 let row = sqlx::query_with(&sql, values).fetch_one(&pool).await?;
636 Ok(row.try_get::<i64, _>(0)?)
637 }
638 }
639 }
640
641 pub async fn fetch_distinct_strings(self, col: &str) -> Result<Vec<String>, DynError> {
646 let Some(col_meta) = self.meta.fields.iter().find(|c| c.name == col) else {
647 return Ok(Vec::new());
648 };
649 let mut q = Query::select();
650 q.distinct();
651 q.from(crate::db::router::schema_qualified_table(&self.meta.table));
652 q.column(Alias::new(col));
653 let where_clauses = self.effective_where_clauses();
654 for cond in &where_clauses {
655 q.cond_where(cond.clone());
656 }
657 if let Some(n) = self.limit {
658 q.limit(n);
659 }
660
661 match resolve_pool_dyn(self.meta, crate::db::RouteOp::Read) {
662 DbPool::Sqlite(pool) => {
663 let (sql, values) = q.build_sqlx(SqliteQueryBuilder);
664 let rows = sqlx::query_with(&sql, values).fetch_all(&pool).await?;
665 let mut out = Vec::with_capacity(rows.len());
666 for row in rows {
667 out.push(decode_to_string(&row, col_meta)?);
668 }
669 Ok(out)
670 }
671 DbPool::Postgres(pool) => {
672 let (sql, values) = q.build_sqlx(PostgresQueryBuilder);
673 let rows = sqlx::query_with(&sql, values).fetch_all(&pool).await?;
674 let mut out = Vec::with_capacity(rows.len());
675 for row in rows {
676 out.push(decode_pg_to_string(&row, col_meta)?);
677 }
678 Ok(out)
679 }
680 }
681 }
682
683 pub async fn delete(self) -> Result<u64, DynError> {
692 if self.meta.soft_delete && !self.hard_delete {
693 return self.soft_delete_update().await;
694 }
695 let where_clauses = self.effective_where_clauses();
696 let parent_pks: Vec<serde_json::Value> = match self.meta.pk_column() {
700 Some(pk_col) => collect_parent_pks(&self.meta, pk_col, &where_clauses)
701 .await
702 .unwrap_or_default(),
703 None => Vec::new(),
704 };
705
706 let mut q = Query::delete();
707 q.from_table(crate::db::router::schema_qualified_table(&self.meta.table));
708 for cond in &where_clauses {
709 q.cond_where(cond.clone());
710 }
711
712 let rows_affected = match resolve_pool_dyn(self.meta, crate::db::RouteOp::Write) {
713 DbPool::Sqlite(pool) => {
714 let (sql, values) = q.build_sqlx(SqliteQueryBuilder);
715 let res = sqlx::query_with(&sql, values).execute(&pool).await?;
716 res.rows_affected()
717 }
718 DbPool::Postgres(pool) => {
719 let (sql, values) = q.build_sqlx(PostgresQueryBuilder);
720 let res = sqlx::query_with(&sql, values).execute(&pool).await?;
721 res.rows_affected()
722 }
723 };
724
725 crate::signals::emit_bulk_post_delete_by_table(&self.meta.table, parent_pks).await;
730 Ok(rows_affected)
731 }
732
733 async fn soft_delete_update(self) -> Result<u64, DynError> {
734 let where_clauses = self.live_where_clauses();
735 let parent_pks: Vec<serde_json::Value> = match self.meta.pk_column() {
736 Some(pk_col) => collect_parent_pks(self.meta, pk_col, &where_clauses)
737 .await
738 .unwrap_or_default(),
739 None => Vec::new(),
740 };
741
742 let mut q = Query::update();
743 q.table(crate::db::router::schema_qualified_table(&self.meta.table));
744 q.value(
745 Alias::new("deleted_at"),
746 sea_query::Value::ChronoDateTimeUtc(Some(Box::new(chrono::Utc::now()))),
747 );
748 for cond in &where_clauses {
749 q.cond_where(cond.clone());
750 }
751
752 let rows_affected = match resolve_pool_dyn(self.meta, crate::db::RouteOp::Write) {
753 DbPool::Sqlite(pool) => {
754 let (sql, values) = q.build_sqlx(SqliteQueryBuilder);
755 let res = sqlx::query_with(&sql, values).execute(&pool).await?;
756 res.rows_affected()
757 }
758 DbPool::Postgres(pool) => {
759 let (sql, values) = q.build_sqlx(PostgresQueryBuilder);
760 let res = sqlx::query_with(&sql, values).execute(&pool).await?;
761 res.rows_affected()
762 }
763 };
764
765 crate::signals::emit_bulk_post_delete_by_table(&self.meta.table, parent_pks).await;
766 Ok(rows_affected)
767 }
768
769 pub async fn restore(self) -> Result<u64, DynError> {
780 if !self.meta.soft_delete {
781 return Ok(0);
782 }
783 let mut where_clauses = self.where_clauses.clone();
787 where_clauses.push(Condition::all().add(Expr::col(Alias::new("deleted_at")).is_not_null()));
788
789 let parent_pks: Vec<serde_json::Value> = match self.meta.pk_column() {
790 Some(pk_col) => collect_parent_pks(self.meta, pk_col, &where_clauses)
791 .await
792 .unwrap_or_default(),
793 None => Vec::new(),
794 };
795
796 let mut q = Query::update();
797 q.table(crate::db::router::schema_qualified_table(&self.meta.table));
798 q.value(
799 Alias::new("deleted_at"),
800 sea_query::Value::ChronoDateTimeUtc(None),
801 );
802 for cond in &where_clauses {
803 q.cond_where(cond.clone());
804 }
805
806 let rows_affected = match resolve_pool_dyn(self.meta, crate::db::RouteOp::Write) {
807 DbPool::Sqlite(pool) => {
808 let (sql, values) = q.build_sqlx(SqliteQueryBuilder);
809 let res = sqlx::query_with(&sql, values).execute(&pool).await?;
810 res.rows_affected()
811 }
812 DbPool::Postgres(pool) => {
813 let (sql, values) = q.build_sqlx(PostgresQueryBuilder);
814 let res = sqlx::query_with(&sql, values).execute(&pool).await?;
815 res.rows_affected()
816 }
817 };
818
819 crate::signals::emit_bulk_post_save_by_table(&self.meta.table, parent_pks, false).await;
823 Ok(rows_affected)
824 }
825
826 pub async fn update_one(self, col: &str, value: &str) -> Result<u64, DynError> {
831 let Some(col_meta) = self.meta.fields.iter().find(|c| c.name == col) else {
832 return Ok(0);
833 };
834 let sea_value = match form_str_to_sea_value(col_meta, value) {
835 Ok(v) => v,
836 Err(e) => {
838 return Err(DynError::Write(WriteError::Validator {
839 field: col_meta.name.clone(),
840 message: e.to_string(),
841 }));
842 }
843 };
844
845 let mut q = Query::update();
846 q.table(crate::db::router::schema_qualified_table(&self.meta.table));
847 q.value(Alias::new(col), sea_value);
848 let where_clauses = self.effective_where_clauses();
849 for cond in &where_clauses {
850 q.cond_where(cond.clone());
851 }
852
853 match resolve_pool_dyn(self.meta, crate::db::RouteOp::Write) {
854 DbPool::Sqlite(pool) => {
855 let (sql, values) = q.build_sqlx(SqliteQueryBuilder);
856 let res = sqlx::query_with(&sql, values).execute(&pool).await?;
857 Ok(res.rows_affected())
858 }
859 DbPool::Postgres(pool) => {
860 let (sql, values) = q.build_sqlx(PostgresQueryBuilder);
861 let res = sqlx::query_with(&sql, values).execute(&pool).await?;
862 Ok(res.rows_affected())
863 }
864 }
865 }
866
867 pub async fn update_form(
874 self,
875 form: &HashMap<String, String>,
876 skip: &[String],
877 ) -> Result<u64, DynError> {
878 let Some(q) = self.build_update_form_query(form, skip)? else {
879 return Ok(0);
880 };
881
882 match resolve_pool_dyn(self.meta, crate::db::RouteOp::Write) {
883 DbPool::Sqlite(pool) => {
884 let (sql, values) = q.build_sqlx(SqliteQueryBuilder);
885 let res = sqlx::query_with(&sql, values).execute(&pool).await?;
886 Ok(res.rows_affected())
887 }
888 DbPool::Postgres(pool) => {
889 let (sql, values) = q.build_sqlx(PostgresQueryBuilder);
890 let res = sqlx::query_with(&sql, values).execute(&pool).await?;
891 Ok(res.rows_affected())
892 }
893 }
894 }
895
896 fn build_update_form_query(
904 &self,
905 form: &HashMap<String, String>,
906 skip: &[String],
907 ) -> Result<Option<sea_query::UpdateStatement>, DynError> {
908 let mut q = Query::update();
909 q.table(crate::db::router::schema_qualified_table(&self.meta.table));
910 let mut any = false;
911 for col in &self.meta.fields {
912 if col.primary_key || skip.iter().any(|s| s == &col.name) {
913 continue;
914 }
915 if is_unauthorized_privileged(col, &self.allow_privileged) {
919 continue;
920 }
921 if col.auto_now {
929 q.value(
930 Alias::new(&col.name),
931 crate::orm::write::now_for_column(col.ty),
932 );
933 any = true;
934 continue;
935 }
936 let Some(raw) = form.get(&col.name) else {
937 continue;
938 };
939 let sea_value = match form_str_to_sea_value(col, raw) {
940 Ok(v) => v,
941 Err(e) => {
947 return Err(DynError::Write(WriteError::Validator {
948 field: col.name.clone(),
949 message: e.to_string(),
950 }));
951 }
952 };
953 q.value(Alias::new(&col.name), sea_value);
954 any = true;
955 }
956 if !any {
957 return Ok(None);
958 }
959 let where_clauses = self.effective_where_clauses();
960 for cond in &where_clauses {
961 q.cond_where(cond.clone());
962 }
963 Ok(Some(q))
964 }
965
966 pub async fn update_form_in_tx(
974 self,
975 tx: &mut crate::db::Transaction,
976 form: &HashMap<String, String>,
977 skip: &[String],
978 ) -> Result<u64, DynError> {
979 let Some(q) = self.build_update_form_query(form, skip)? else {
980 return Ok(0);
981 };
982
983 match tx.backend_name() {
984 "sqlite" => {
985 let (sql, values) = q.build_sqlx(SqliteQueryBuilder);
986 let inner = tx.as_sqlite_mut().expect("sqlite backend_name");
987 let res = sqlx::query_with(&sql, values).execute(&mut **inner).await?;
988 Ok(res.rows_affected())
989 }
990 _ => {
991 let (sql, values) = q.build_sqlx(PostgresQueryBuilder);
992 let inner = tx.as_pg_mut().expect("postgres backend_name");
993 let res = sqlx::query_with(&sql, values).execute(&mut **inner).await?;
994 Ok(res.rows_affected())
995 }
996 }
997 }
998
999 pub async fn insert_form(
1005 self,
1006 form: &HashMap<String, String>,
1007 skip: &[String],
1008 ) -> Result<i64, DynError> {
1009 let Some(mut q) = self.build_insert_form_query(form, skip)? else {
1010 return Ok(0);
1011 };
1012
1013 match resolve_pool_dyn(self.meta, crate::db::RouteOp::Write) {
1014 DbPool::Sqlite(pool) => {
1015 let (sql, vals) = q.build_sqlx(SqliteQueryBuilder);
1016 let res = sqlx::query_with(&sql, vals).execute(&pool).await?;
1017 Ok(res.last_insert_rowid())
1018 }
1019 DbPool::Postgres(pool) => {
1020 let pk_name = self
1026 .meta
1027 .fields
1028 .iter()
1029 .find(|c| c.primary_key)
1030 .map(|c| c.name.clone());
1031 if let Some(pk) = pk_name {
1032 q.returning_col(Alias::new(&pk));
1033 let (sql, vals) = q.build_sqlx(PostgresQueryBuilder);
1034 let row = sqlx::query_with(&sql, vals).fetch_one(&pool).await?;
1035 Ok(row.try_get::<i64, _>(pk.as_str()).unwrap_or(0))
1036 } else {
1037 let (sql, vals) = q.build_sqlx(PostgresQueryBuilder);
1038 let _ = sqlx::query_with(&sql, vals).execute(&pool).await?;
1039 Ok(0)
1040 }
1041 }
1042 }
1043 }
1044
1045 fn build_insert_form_query(
1054 &self,
1055 form: &HashMap<String, String>,
1056 skip: &[String],
1057 ) -> Result<Option<sea_query::InsertStatement>, DynError> {
1058 let mut cols: Vec<&str> = Vec::new();
1059 let mut values: Vec<SeaValue> = Vec::new();
1060 for col in &self.meta.fields {
1061 if skip.iter().any(|s| s == &col.name) {
1062 continue;
1063 }
1064 if is_unauthorized_privileged(col, &self.allow_privileged) {
1069 continue;
1070 }
1071 if col.primary_key
1074 && matches!(
1075 col.ty,
1076 SqlType::Integer | SqlType::BigInt | SqlType::SmallInt
1077 )
1078 && form.get(&col.name).is_none_or(|v| v.is_empty())
1079 {
1080 continue;
1081 }
1082 if (col.auto_now_add || col.auto_now)
1090 && form.get(&col.name).is_none_or(|v| v.is_empty())
1091 {
1092 cols.push(&col.name);
1093 values.push(crate::orm::write::now_for_column(col.ty));
1094 continue;
1095 }
1096 let raw = form.get(&col.name).map(|s| s.as_str()).unwrap_or("");
1097 let sea_value = match form_str_to_sea_value(col, raw) {
1098 Ok(v) => v,
1099 Err(e) => {
1102 return Err(DynError::Write(WriteError::Validator {
1103 field: col.name.clone(),
1104 message: e.to_string(),
1105 }));
1106 }
1107 };
1108 cols.push(&col.name);
1109 values.push(sea_value);
1110 }
1111 if cols.is_empty() {
1112 return Ok(None);
1113 }
1114
1115 let mut q = Query::insert();
1116 q.into_table(crate::db::router::schema_qualified_table(&self.meta.table));
1117 q.columns(cols.iter().map(|c| Alias::new(*c)).collect::<Vec<_>>());
1118 let exprs: Vec<sea_query::SimpleExpr> = values.into_iter().map(Into::into).collect();
1119 q.values_panic(exprs);
1120 Ok(Some(q))
1121 }
1122
1123 pub async fn insert_form_in_tx(
1132 self,
1133 tx: &mut crate::db::Transaction,
1134 form: &HashMap<String, String>,
1135 skip: &[String],
1136 ) -> Result<i64, DynError> {
1137 let Some(mut q) = self.build_insert_form_query(form, skip)? else {
1138 return Ok(0);
1139 };
1140
1141 match tx.backend_name() {
1142 "sqlite" => {
1143 let (sql, vals) = q.build_sqlx(SqliteQueryBuilder);
1144 let inner = tx.as_sqlite_mut().expect("sqlite backend_name");
1145 let res = sqlx::query_with(&sql, vals).execute(&mut **inner).await?;
1146 Ok(res.last_insert_rowid())
1147 }
1148 _ => {
1149 let pk_name = self
1153 .meta
1154 .fields
1155 .iter()
1156 .find(|c| c.primary_key)
1157 .map(|c| c.name.clone());
1158 let inner = tx.as_pg_mut().expect("postgres backend_name");
1159 if let Some(pk) = pk_name {
1160 q.returning_col(Alias::new(&pk));
1161 let (sql, vals) = q.build_sqlx(PostgresQueryBuilder);
1162 let row = sqlx::query_with(&sql, vals).fetch_one(&mut **inner).await?;
1163 Ok(row.try_get::<i64, _>(pk.as_str()).unwrap_or(0))
1164 } else {
1165 let (sql, vals) = q.build_sqlx(PostgresQueryBuilder);
1166 let _ = sqlx::query_with(&sql, vals).execute(&mut **inner).await?;
1167 Ok(0)
1168 }
1169 }
1170 }
1171 }
1172
1173 pub async fn fetch_as_strings(self) -> Result<Vec<HashMap<String, String>>, DynError> {
1178 let mut q = Query::select();
1179 q.from(crate::db::router::schema_qualified_table(&self.meta.table));
1180 for c in &self.select_cols {
1181 q.column(Alias::new(c));
1182 }
1183 let where_clauses = self.effective_where_clauses();
1184 for cond in &where_clauses {
1185 q.cond_where(cond.clone());
1186 }
1187 for (col, descending) in &self.order {
1188 q.order_by(
1189 Alias::new(col),
1190 if *descending { Order::Desc } else { Order::Asc },
1191 );
1192 }
1193 if let Some(n) = self.limit {
1194 q.limit(n);
1195 }
1196 if let Some(n) = self.offset {
1197 q.offset(n);
1198 }
1199
1200 match resolve_pool_dyn(self.meta, crate::db::RouteOp::Read) {
1201 DbPool::Sqlite(pool) => {
1202 let (sql, values) = q.build_sqlx(SqliteQueryBuilder);
1203 let rows = sqlx::query_with(&sql, values).fetch_all(&pool).await?;
1204 let mut out: Vec<HashMap<String, String>> = Vec::with_capacity(rows.len());
1205 for row in rows {
1206 let mut entry = HashMap::new();
1207 for col_name in &self.select_cols {
1208 if let Some(col_meta) =
1209 self.meta.fields.iter().find(|c| &c.name == col_name)
1210 {
1211 let v = decode_to_string(&row, col_meta)?;
1212 entry.insert(col_name.clone(), v);
1213 }
1214 }
1215 out.push(entry);
1216 }
1217 Ok(out)
1218 }
1219 DbPool::Postgres(pool) => {
1220 let (sql, values) = q.build_sqlx(PostgresQueryBuilder);
1221 let rows = sqlx::query_with(&sql, values).fetch_all(&pool).await?;
1222 let mut out: Vec<HashMap<String, String>> = Vec::with_capacity(rows.len());
1223 for row in rows {
1224 let mut entry = HashMap::new();
1225 for col_name in &self.select_cols {
1226 if let Some(col_meta) =
1227 self.meta.fields.iter().find(|c| &c.name == col_name)
1228 {
1229 let v = decode_pg_to_string(&row, col_meta)?;
1230 entry.insert(col_name.clone(), v);
1231 }
1232 }
1233 out.push(entry);
1234 }
1235 Ok(out)
1236 }
1237 }
1238 }
1239
1240 pub async fn fetch_as_json(
1246 self,
1247 ) -> Result<Vec<serde_json::Map<String, serde_json::Value>>, DynError> {
1248 let mut q = Query::select();
1249 q.from(crate::db::router::schema_qualified_table(&self.meta.table));
1250 for c in &self.select_cols {
1251 q.column(Alias::new(c));
1252 }
1253 let where_clauses = self.effective_where_clauses();
1254 for cond in &where_clauses {
1255 q.cond_where(cond.clone());
1256 }
1257 for (col, descending) in &self.order {
1258 q.order_by(
1259 Alias::new(col),
1260 if *descending { Order::Desc } else { Order::Asc },
1261 );
1262 }
1263 if let Some(n) = self.limit {
1264 q.limit(n);
1265 }
1266 if let Some(n) = self.offset {
1267 q.offset(n);
1268 }
1269
1270 let pk_name = self
1271 .meta
1272 .pk_column()
1273 .map(|c| c.name.clone())
1274 .unwrap_or_default();
1275 let selected_cols: Vec<(&String, &Column)> = self
1276 .select_cols
1277 .iter()
1278 .filter_map(|col_name| {
1279 self.meta
1280 .fields
1281 .iter()
1282 .find(|c| &c.name == col_name)
1283 .map(|col| (col_name, col))
1284 })
1285 .collect();
1286 let mut out: Vec<serde_json::Map<String, serde_json::Value>> =
1287 match resolve_pool_dyn(self.meta, crate::db::RouteOp::Read) {
1288 DbPool::Sqlite(pool) => {
1289 let (sql, values) = q.build_sqlx(SqliteQueryBuilder);
1290 let rows = sqlx::query_with(&sql, values).fetch_all(&pool).await?;
1291 let mut out: Vec<serde_json::Map<String, serde_json::Value>> =
1292 Vec::with_capacity(rows.len());
1293 for row in rows {
1294 let mut entry = serde_json::Map::new();
1295 for (col_name, col_meta) in &selected_cols {
1296 entry.insert((*col_name).clone(), decode_to_json(&row, col_meta)?);
1297 }
1298 out.push(entry);
1299 }
1300 out
1301 }
1302 DbPool::Postgres(pool) => {
1303 let (sql, values) = q.build_sqlx(PostgresQueryBuilder);
1304 let rows = sqlx::query_with(&sql, values).fetch_all(&pool).await?;
1305 let mut out: Vec<serde_json::Map<String, serde_json::Value>> =
1306 Vec::with_capacity(rows.len());
1307 for row in rows {
1308 let mut entry = serde_json::Map::new();
1309 for (col_name, col_meta) in &selected_cols {
1310 entry.insert((*col_name).clone(), decode_pg_to_json(&row, col_meta)?);
1311 }
1312 out.push(entry);
1313 }
1314 out
1315 }
1316 };
1317
1318 if !self.meta.m2m_relations.is_empty() && !out.is_empty() {
1328 hydrate_m2m_batched(&self.meta, &pk_name, &mut out).await?;
1329 }
1330
1331 if !self.select_related.is_empty() && !out.is_empty() {
1340 hydrate_select_related_into(&self.meta, &self.select_related, &mut out).await?;
1341 }
1342 Ok(out)
1343 }
1344
1345 pub async fn first_as_json(
1348 mut self,
1349 ) -> Result<Option<serde_json::Map<String, serde_json::Value>>, DynError> {
1350 self.limit = Some(1);
1351 let mut rows = self.fetch_as_json().await?;
1352 Ok(rows.pop())
1353 }
1354
1355 pub async fn fetch_one_json_in_tx(
1365 self,
1366 tx: &mut crate::db::Transaction,
1367 ) -> Result<Option<serde_json::Map<String, serde_json::Value>>, DynError> {
1368 let mut q = Query::select();
1369 q.from(crate::db::router::schema_qualified_table(&self.meta.table));
1370 for c in &self.meta.fields {
1371 q.column(Alias::new(&c.name));
1372 }
1373 let where_clauses = self.effective_where_clauses();
1374 for cond in &where_clauses {
1375 q.cond_where(cond.clone());
1376 }
1377 q.limit(1);
1378
1379 let out = match tx.backend_name() {
1380 "sqlite" => {
1381 let (sql, values) = q.build_sqlx(SqliteQueryBuilder);
1382 let inner = tx.as_sqlite_mut().expect("sqlite backend_name");
1383 let row = sqlx::query_with(&sql, values)
1384 .fetch_optional(&mut **inner)
1385 .await?;
1386 match row {
1387 Some(row) => {
1388 let mut entry = serde_json::Map::new();
1389 for col in &self.meta.fields {
1390 entry.insert(col.name.clone(), decode_to_json(&row, col)?);
1391 }
1392 Some(entry)
1393 }
1394 None => None,
1395 }
1396 }
1397 _ => {
1398 let (sql, values) = q.build_sqlx(PostgresQueryBuilder);
1399 let inner = tx.as_pg_mut().expect("postgres backend_name");
1400 let row = sqlx::query_with(&sql, values)
1401 .fetch_optional(&mut **inner)
1402 .await?;
1403 match row {
1404 Some(row) => {
1405 let mut entry = serde_json::Map::new();
1406 for col in &self.meta.fields {
1407 entry.insert(col.name.clone(), decode_pg_to_json(&row, col)?);
1408 }
1409 Some(entry)
1410 }
1411 None => None,
1412 }
1413 }
1414 };
1415 Ok(out)
1416 }
1417
1418 pub async fn insert_json(
1426 self,
1427 body: &serde_json::Map<String, serde_json::Value>,
1428 ) -> Result<serde_json::Map<String, serde_json::Value>, crate::orm::write::WriteError> {
1429 use crate::orm::write::WriteError;
1430
1431 let body_owned: serde_json::Map<String, serde_json::Value>;
1434 let body: &serde_json::Map<String, serde_json::Value> =
1435 match normalise_insert_body(self.meta, body, &self.allow_privileged) {
1436 Some(owned) => {
1437 body_owned = owned;
1438 &body_owned
1439 }
1440 None => body,
1441 };
1442
1443 let validation_errors = crate::orm::validation::validate_on_create(self.meta, body).await;
1445 if !validation_errors.is_empty() {
1446 return Err(WriteError::Multiple {
1447 errors: validation_errors,
1448 });
1449 }
1450
1451 let InsertPlan {
1454 mut q,
1455 pk_name,
1456 pk_ty,
1457 } = build_insert_plan(self.meta, body)?;
1458
1459 crate::signals::emit_pre_save_by_table(
1465 &self.meta.table,
1466 serde_json::Value::Object(body.clone()),
1467 true,
1468 )
1469 .await;
1470
1471 let mut tx = match resolve_pool_dyn(self.meta, crate::db::RouteOp::Write) {
1482 DbPool::Sqlite(pool) => crate::db::begin_sqlite(&pool).await,
1483 DbPool::Postgres(pool) => crate::db::begin_pg(&pool).await,
1484 }?;
1485
1486 let mut out = match tx.backend_name() {
1487 "sqlite" => {
1488 let (sql, vals) = q.build_sqlx(SqliteQueryBuilder);
1489 let res = {
1490 let inner = tx.as_sqlite_mut().expect("sqlite backend_name");
1491 sqlx::query_with(&sql, vals)
1492 .execute(&mut **inner)
1493 .await
1494 .map_err(|e| classify_or_sqlx(e, body))?
1495 };
1496 let pk_pred = match pk_ty {
1500 SqlType::Integer | SqlType::BigInt | SqlType::SmallInt => {
1501 Expr::col(Alias::new(&pk_name)).eq(res.last_insert_rowid())
1502 }
1503 _ => {
1504 let supplied = body
1507 .get(&pk_name)
1508 .cloned()
1509 .unwrap_or(serde_json::Value::Null);
1510 let sea_value = crate::orm::write::json_to_sea_value(
1511 pk_ty, &supplied, false, &pk_name, None,
1512 )?;
1513 Expr::col(Alias::new(&pk_name)).eq(sea_value)
1514 }
1515 };
1516 let mut sel = Query::select();
1517 sel.from(crate::db::router::schema_qualified_table(&self.meta.table));
1518 for c in &self.meta.fields {
1519 sel.column(Alias::new(&c.name));
1520 }
1521 sel.cond_where(Condition::all().add(pk_pred));
1522 let (sel_sql, sel_vals) = sel.build_sqlx(SqliteQueryBuilder);
1523 let mut out = serde_json::Map::new();
1524 {
1525 let inner = tx.as_sqlite_mut().expect("sqlite backend_name");
1526 let row = sqlx::query_with(&sel_sql, sel_vals)
1527 .fetch_one(&mut **inner)
1528 .await?;
1529 for col in &self.meta.fields {
1530 out.insert(col.name.clone(), decode_to_json(&row, col)?);
1531 }
1532 }
1533 out
1534 }
1535 _ => {
1536 q.returning_all();
1541 let (sql, vals) = q.build_sqlx(PostgresQueryBuilder);
1542 let mut out = serde_json::Map::new();
1543 {
1544 let inner = tx.as_pg_mut().expect("postgres backend_name");
1545 let row = sqlx::query_with(&sql, vals)
1546 .fetch_one(&mut **inner)
1547 .await
1548 .map_err(|e| classify_or_sqlx(e, body))?;
1549 for col in &self.meta.fields {
1550 out.insert(col.name.clone(), decode_pg_to_json(&row, col)?);
1551 }
1552 }
1553 out
1554 }
1555 };
1556
1557 let pk_value = out.get(&pk_name).cloned();
1564 write_m2m_junctions_in_tx(self.meta, pk_value.as_ref(), body, &mut tx).await?;
1565 hydrate_m2m_into_tx(self.meta, pk_value.as_ref(), &mut out, &mut tx).await?;
1568
1569 tx.commit().await?;
1570
1571 crate::signals::emit_post_save_by_table(
1574 &self.meta.table,
1575 serde_json::Value::Object(out.clone()),
1576 true,
1577 )
1578 .await;
1579 Ok(out)
1580 }
1581
1582 pub async fn insert_json_in_tx(
1606 self,
1607 body: &serde_json::Map<String, serde_json::Value>,
1608 tx: &mut crate::db::Transaction,
1609 ) -> Result<serde_json::Map<String, serde_json::Value>, crate::orm::write::WriteError> {
1610 use crate::orm::write::WriteError;
1611
1612 let body_owned: serde_json::Map<String, serde_json::Value>;
1614 let body: &serde_json::Map<String, serde_json::Value> =
1615 match normalise_insert_body(self.meta, body, &self.allow_privileged) {
1616 Some(owned) => {
1617 body_owned = owned;
1618 &body_owned
1619 }
1620 None => body,
1621 };
1622
1623 let validation_errors =
1626 crate::orm::validation::validate_on_create_in_tx(self.meta, body, tx).await;
1627 if !validation_errors.is_empty() {
1628 return Err(WriteError::Multiple {
1629 errors: validation_errors,
1630 });
1631 }
1632
1633 let InsertPlan {
1635 mut q,
1636 pk_name,
1637 pk_ty,
1638 } = build_insert_plan(self.meta, body)?;
1639
1640 match tx.backend_name() {
1641 "sqlite" => {
1642 let (sql, vals) = q.build_sqlx(SqliteQueryBuilder);
1643 let res = {
1644 let inner = tx.as_sqlite_mut().expect("sqlite backend_name");
1645 sqlx::query_with(&sql, vals)
1646 .execute(&mut **inner)
1647 .await
1648 .map_err(|e| classify_or_sqlx(e, body))?
1649 };
1650 let pk_pred = match pk_ty {
1653 SqlType::Integer | SqlType::BigInt | SqlType::SmallInt => {
1654 Expr::col(Alias::new(&pk_name)).eq(res.last_insert_rowid())
1655 }
1656 _ => {
1657 let supplied = body
1658 .get(&pk_name)
1659 .cloned()
1660 .unwrap_or(serde_json::Value::Null);
1661 let sea_value = crate::orm::write::json_to_sea_value(
1662 pk_ty, &supplied, false, &pk_name, None,
1663 )?;
1664 Expr::col(Alias::new(&pk_name)).eq(sea_value)
1665 }
1666 };
1667 let mut sel = Query::select();
1668 sel.from(crate::db::router::schema_qualified_table(&self.meta.table));
1669 for c in &self.meta.fields {
1670 sel.column(Alias::new(&c.name));
1671 }
1672 sel.cond_where(Condition::all().add(pk_pred));
1673 let (sel_sql, sel_vals) = sel.build_sqlx(SqliteQueryBuilder);
1674 let mut out = serde_json::Map::new();
1675 {
1676 let inner = tx.as_sqlite_mut().expect("sqlite backend_name");
1677 let row = sqlx::query_with(&sel_sql, sel_vals)
1678 .fetch_one(&mut **inner)
1679 .await?;
1680 for col in &self.meta.fields {
1681 out.insert(col.name.clone(), decode_to_json(&row, col)?);
1682 }
1683 }
1684 let pk_value = out.get(&pk_name).cloned();
1686 write_m2m_junctions_in_tx(self.meta, pk_value.as_ref(), body, tx).await?;
1687 hydrate_m2m_into_tx(self.meta, pk_value.as_ref(), &mut out, tx).await?;
1688 Ok(out)
1689 }
1690 _ => {
1691 q.returning_all();
1692 let (sql, vals) = q.build_sqlx(PostgresQueryBuilder);
1693 let mut out = serde_json::Map::new();
1694 {
1695 let inner = tx.as_pg_mut().expect("postgres backend_name");
1696 let row = sqlx::query_with(&sql, vals)
1697 .fetch_one(&mut **inner)
1698 .await
1699 .map_err(|e| classify_or_sqlx(e, body))?;
1700 for col in &self.meta.fields {
1701 out.insert(col.name.clone(), decode_pg_to_json(&row, col)?);
1702 }
1703 }
1704 let pk_value = out.get(&pk_name).cloned();
1705 write_m2m_junctions_in_tx(self.meta, pk_value.as_ref(), body, tx).await?;
1706 hydrate_m2m_into_tx(self.meta, pk_value.as_ref(), &mut out, tx).await?;
1707 Ok(out)
1708 }
1709 }
1710 }
1711
1712 pub async fn update_json_in_tx(
1723 self,
1724 body: &serde_json::Map<String, serde_json::Value>,
1725 tx: &mut crate::db::Transaction,
1726 ) -> Result<u64, crate::orm::write::WriteError> {
1727 use crate::orm::write::WriteError;
1728
1729 let body_owned: serde_json::Map<String, serde_json::Value>;
1732 let body: &serde_json::Map<String, serde_json::Value> =
1733 match normalise_update_body(self.meta, body, &self.allow_privileged) {
1734 Some(owned) => {
1735 body_owned = owned;
1736 &body_owned
1737 }
1738 None => body,
1739 };
1740
1741 let validation_errors =
1745 crate::orm::validation::validate_on_update_in_tx(self.meta, body, tx).await;
1746 if !validation_errors.is_empty() {
1747 return Err(WriteError::Multiple {
1748 errors: validation_errors,
1749 });
1750 }
1751
1752 let mut q = Query::update();
1753 q.table(crate::db::router::schema_qualified_table(&self.meta.table));
1754 let mut any = false;
1755 for col in &self.meta.fields {
1756 if col.primary_key {
1757 continue;
1758 }
1759 let Some(json) = body.get(&col.name) else {
1760 if col.auto_now {
1761 let now_value = crate::orm::write::now_for_column(col.ty);
1762 q.value(Alias::new(&col.name), now_value);
1763 any = true;
1764 }
1765 continue;
1766 };
1767 validate_numeric_bounds(col, json)?;
1768 if let (Some(fmt), Some(s)) = (col.text_format.as_deref(), json.as_str()) {
1769 if let Err(e) = crate::orm::validators::validate_text_format(fmt, s) {
1770 return Err(WriteError::Validator {
1771 field: col.name.clone(),
1772 message: e.to_string(),
1773 });
1774 }
1775 }
1776 let normalized_json = normalize_json_for_col(col, json);
1780 let json = normalized_json.as_ref().unwrap_or(json);
1781 let sealed = crate::orm::write::seal_masked_json(col, json)?;
1784 let sea_value = crate::orm::write::json_to_sea_value(
1785 col.ty,
1786 sealed.as_ref().unwrap_or(json),
1787 col.nullable,
1788 &col.name,
1789 fk_target_pk_sql_type(col),
1790 )?;
1791 q.value(Alias::new(&col.name), sea_value);
1792 any = true;
1793 }
1794 let touches_m2m = self
1795 .meta
1796 .m2m_relations
1797 .iter()
1798 .any(|r| body.contains_key(&r.field_name));
1799 if !any && !touches_m2m {
1800 return Ok(0);
1801 }
1802 let where_clauses = self.effective_where_clauses();
1803 for cond in &where_clauses {
1804 q.cond_where(cond.clone());
1805 }
1806
1807 let parent_pks: Vec<serde_json::Value> = match self.meta.pk_column() {
1811 Some(pk_col) => collect_parent_pks_in_tx(self.meta, pk_col, &where_clauses, tx).await?,
1812 None => Vec::new(),
1813 };
1814
1815 if any {
1816 match tx.backend_name() {
1817 "sqlite" => {
1818 let (sql, values) = q.build_sqlx(SqliteQueryBuilder);
1819 let inner = tx.as_sqlite_mut().expect("sqlite backend_name");
1820 sqlx::query_with(&sql, values)
1821 .execute(&mut **inner)
1822 .await
1823 .map_err(|e| classify_or_sqlx(e, body))?;
1824 }
1825 _ => {
1826 let (sql, values) = q.build_sqlx(PostgresQueryBuilder);
1827 let inner = tx.as_pg_mut().expect("postgres backend_name");
1828 sqlx::query_with(&sql, values)
1829 .execute(&mut **inner)
1830 .await
1831 .map_err(|e| classify_or_sqlx(e, body))?;
1832 }
1833 }
1834 }
1835 for pk in &parent_pks {
1836 write_m2m_junctions_in_tx(self.meta, Some(pk), body, tx).await?;
1837 }
1838 Ok(parent_pks.len().max(if any { 1 } else { 0 }) as u64)
1839 }
1840
1841 pub async fn delete_in_tx(self, tx: &mut crate::db::Transaction) -> Result<u64, DynError> {
1849 let soft = self.meta.soft_delete && !self.hard_delete;
1850 let where_clauses = if soft {
1851 self.live_where_clauses()
1852 } else {
1853 self.effective_where_clauses()
1854 };
1855
1856 let table = crate::db::router::schema_qualified_table(&self.meta.table);
1860 let build = |is_sqlite: bool| {
1861 if soft {
1862 let mut u = Query::update();
1863 u.table(table.clone());
1864 u.value(
1865 Alias::new("deleted_at"),
1866 sea_query::Value::ChronoDateTimeUtc(Some(Box::new(chrono::Utc::now()))),
1867 );
1868 for cond in &where_clauses {
1869 u.cond_where(cond.clone());
1870 }
1871 if is_sqlite {
1872 u.build_sqlx(SqliteQueryBuilder)
1873 } else {
1874 u.build_sqlx(PostgresQueryBuilder)
1875 }
1876 } else {
1877 let mut d = Query::delete();
1878 d.from_table(table.clone());
1879 for cond in &where_clauses {
1880 d.cond_where(cond.clone());
1881 }
1882 if is_sqlite {
1883 d.build_sqlx(SqliteQueryBuilder)
1884 } else {
1885 d.build_sqlx(PostgresQueryBuilder)
1886 }
1887 }
1888 };
1889
1890 let rows_affected = match tx.backend_name() {
1891 "sqlite" => {
1892 let (sql, values) = build(true);
1893 let inner = tx.as_sqlite_mut().expect("sqlite backend_name");
1894 sqlx::query_with(&sql, values)
1895 .execute(&mut **inner)
1896 .await?
1897 .rows_affected()
1898 }
1899 _ => {
1900 let (sql, values) = build(false);
1901 let inner = tx.as_pg_mut().expect("postgres backend_name");
1902 sqlx::query_with(&sql, values)
1903 .execute(&mut **inner)
1904 .await?
1905 .rows_affected()
1906 }
1907 };
1908 Ok(rows_affected)
1909 }
1910
1911 pub async fn update_json(
1915 self,
1916 body: &serde_json::Map<String, serde_json::Value>,
1917 ) -> Result<u64, crate::orm::write::WriteError> {
1918 use crate::orm::write::WriteError;
1919
1920 let body_owned: serde_json::Map<String, serde_json::Value>;
1927 let body: &serde_json::Map<String, serde_json::Value> =
1928 match normalise_update_body(self.meta, body, &self.allow_privileged) {
1929 Some(owned) => {
1930 body_owned = owned;
1931 &body_owned
1932 }
1933 None => body,
1934 };
1935
1936 let validation_errors = crate::orm::validation::validate_on_update(&self.meta, body).await;
1942 if !validation_errors.is_empty() {
1943 return Err(WriteError::Multiple {
1944 errors: validation_errors,
1945 });
1946 }
1947
1948 let mut q = Query::update();
1949 q.table(crate::db::router::schema_qualified_table(&self.meta.table));
1950 let mut any = false;
1951 for col in &self.meta.fields {
1952 if col.primary_key {
1953 continue;
1954 }
1955 let Some(json) = body.get(&col.name) else {
1956 if col.auto_now {
1961 let now_value = crate::orm::write::now_for_column(col.ty);
1962 q.value(Alias::new(&col.name), now_value);
1963 any = true;
1964 }
1965 continue;
1966 };
1967 validate_numeric_bounds(col, json)?;
1968 if let (Some(fmt), Some(s)) = (col.text_format.as_deref(), json.as_str()) {
1971 if let Err(e) = crate::orm::validators::validate_text_format(fmt, s) {
1972 return Err(WriteError::Validator {
1973 field: col.name.clone(),
1974 message: e.to_string(),
1975 });
1976 }
1977 }
1978 let normalized_json = normalize_json_for_col(col, json);
1982 let json = normalized_json.as_ref().unwrap_or(json);
1983 let sealed = crate::orm::write::seal_masked_json(col, json)?;
1986 let sea_value = crate::orm::write::json_to_sea_value(
1987 col.ty,
1988 sealed.as_ref().unwrap_or(json),
1989 col.nullable,
1990 &col.name,
1991 fk_target_pk_sql_type(col),
1992 )?;
1993 q.value(Alias::new(&col.name), sea_value);
1994 any = true;
1995 }
1996 let touches_m2m = self
2001 .meta
2002 .m2m_relations
2003 .iter()
2004 .any(|r| body.contains_key(&r.field_name));
2005 if !any && !touches_m2m {
2006 return Ok(0);
2007 }
2008 let where_clauses = self.effective_where_clauses();
2009 for cond in &where_clauses {
2010 q.cond_where(cond.clone());
2011 }
2012
2013 let mut tx = match resolve_pool_dyn(self.meta, crate::db::RouteOp::Write) {
2018 DbPool::Sqlite(pool) => crate::db::begin_sqlite(&pool).await,
2019 DbPool::Postgres(pool) => crate::db::begin_pg(&pool).await,
2020 }?;
2021
2022 let parent_pks: Vec<serde_json::Value> = match self.meta.pk_column() {
2036 Some(pk_col) => {
2037 collect_parent_pks_in_tx(self.meta, pk_col, &where_clauses, &mut tx).await?
2038 }
2039 None => Vec::new(),
2040 };
2041
2042 let rows_affected = if any {
2047 match tx.backend_name() {
2048 "sqlite" => {
2049 let (sql, values) = q.build_sqlx(SqliteQueryBuilder);
2050 let inner = tx.as_sqlite_mut().expect("sqlite backend_name");
2051 sqlx::query_with(&sql, values)
2052 .execute(&mut **inner)
2053 .await
2054 .map_err(|e| classify_or_sqlx(e, body))?
2055 .rows_affected()
2056 }
2057 _ => {
2058 let (sql, values) = q.build_sqlx(PostgresQueryBuilder);
2059 let inner = tx.as_pg_mut().expect("postgres backend_name");
2060 sqlx::query_with(&sql, values)
2061 .execute(&mut **inner)
2062 .await
2063 .map_err(|e| classify_or_sqlx(e, body))?
2064 .rows_affected()
2065 }
2066 }
2067 } else {
2068 0
2069 };
2070
2071 for pk in &parent_pks {
2072 write_m2m_junctions_in_tx(self.meta, Some(pk), body, &mut tx).await?;
2073 }
2074
2075 tx.commit().await?;
2076
2077 crate::signals::emit_bulk_post_save_by_table(&self.meta.table, parent_pks.clone(), false)
2082 .await;
2083
2084 if any {
2089 Ok(rows_affected)
2090 } else {
2091 Ok(parent_pks.len() as u64)
2092 }
2093 }
2094}
2095
2096pub fn decode_to_string(
2102 row: &sqlx::sqlite::SqliteRow,
2103 col: &Column,
2104) -> Result<String, sqlx::Error> {
2105 use chrono::{DateTime, NaiveDate, NaiveTime, Utc};
2106 use serde_json::Value;
2107 use uuid::Uuid;
2108
2109 let name = col.name.as_str();
2110 if col.nullable {
2111 return Ok(match col.ty {
2112 SqlType::SmallInt | SqlType::Integer => row
2113 .try_get::<Option<i32>, _>(name)?
2114 .map_or(String::new(), |v| v.to_string()),
2115 SqlType::BigInt => row
2116 .try_get::<Option<i64>, _>(name)?
2117 .map_or(String::new(), |v| v.to_string()),
2118 SqlType::Real => row
2119 .try_get::<Option<f32>, _>(name)?
2120 .map_or(String::new(), |v| v.to_string()),
2121 SqlType::Double => row
2122 .try_get::<Option<f64>, _>(name)?
2123 .map_or(String::new(), |v| v.to_string()),
2124 SqlType::Boolean => row
2125 .try_get::<Option<bool>, _>(name)?
2126 .map_or(String::new(), |v| {
2127 if v { "true" } else { "false" }.to_string()
2128 }),
2129 SqlType::Text => row.try_get::<Option<String>, _>(name)?.unwrap_or_default(),
2130 SqlType::Date => row
2131 .try_get::<Option<NaiveDate>, _>(name)?
2132 .map_or(String::new(), |v| v.to_string()),
2133 SqlType::Time => row
2134 .try_get::<Option<NaiveTime>, _>(name)?
2135 .map_or(String::new(), |v| v.to_string()),
2136 SqlType::Timestamptz => row
2137 .try_get::<Option<DateTime<Utc>>, _>(name)?
2138 .map_or(String::new(), |v| v.to_rfc3339()),
2139 SqlType::Uuid => row
2140 .try_get::<Option<Uuid>, _>(name)?
2141 .map_or(String::new(), |v| v.to_string()),
2142 SqlType::Json => row
2143 .try_get::<Option<Value>, _>(name)?
2144 .map_or(String::new(), |v| v.to_string()),
2145 SqlType::Array(_) => panic_array_unsupported(&col.name),
2146 SqlType::Inet
2147 | SqlType::Cidr
2148 | SqlType::MacAddr
2149 | SqlType::Xml
2150 | SqlType::Ltree
2151 | SqlType::Bit
2152 | SqlType::FullText => panic_pg_only_unsupported(&col.name),
2153 SqlType::ForeignKey => match fk_target_pk_sql_type(col) {
2157 Some(SqlType::Text) => row.try_get::<Option<String>, _>(name)?.unwrap_or_default(),
2158 Some(SqlType::Uuid) => row
2159 .try_get::<Option<Uuid>, _>(name)?
2160 .map_or(String::new(), |v| v.to_string()),
2161 _ => row
2162 .try_get::<Option<i64>, _>(name)?
2163 .map_or(String::new(), |v| v.to_string()),
2164 },
2165 SqlType::Bytes => row
2166 .try_get::<Option<Vec<u8>>, _>(name)?
2167 .map_or(String::new(), |b| hex_encode(&b)),
2168 SqlType::Decimal => panic_pg_only_unsupported(&col.name),
2169 });
2170 }
2171 Ok(match col.ty {
2172 SqlType::SmallInt | SqlType::Integer => row.try_get::<i32, _>(name)?.to_string(),
2173 SqlType::BigInt => row.try_get::<i64, _>(name)?.to_string(),
2174 SqlType::Real => row.try_get::<f32, _>(name)?.to_string(),
2175 SqlType::Double => row.try_get::<f64, _>(name)?.to_string(),
2176 SqlType::Boolean => if row.try_get::<bool, _>(name)? {
2177 "true"
2178 } else {
2179 "false"
2180 }
2181 .to_string(),
2182 SqlType::Text => row.try_get::<String, _>(name)?,
2183 SqlType::Date => row.try_get::<NaiveDate, _>(name)?.to_string(),
2184 SqlType::Time => row.try_get::<NaiveTime, _>(name)?.to_string(),
2185 SqlType::Timestamptz => row.try_get::<DateTime<Utc>, _>(name)?.to_rfc3339(),
2186 SqlType::Uuid => row.try_get::<Uuid, _>(name)?.to_string(),
2187 SqlType::Json => row.try_get::<Value, _>(name)?.to_string(),
2188 SqlType::Array(_) => panic_array_unsupported(&col.name),
2189 SqlType::Inet
2190 | SqlType::Cidr
2191 | SqlType::MacAddr
2192 | SqlType::Xml
2193 | SqlType::Ltree
2194 | SqlType::Bit
2195 | SqlType::FullText => panic_pg_only_unsupported(&col.name),
2196 SqlType::ForeignKey => match fk_target_pk_sql_type(col) {
2197 Some(SqlType::Text) => row.try_get::<String, _>(name)?,
2198 Some(SqlType::Uuid) => row.try_get::<Uuid, _>(name)?.to_string(),
2199 _ => row.try_get::<i64, _>(name)?.to_string(),
2200 },
2201 SqlType::Bytes => hex_encode(&row.try_get::<Vec<u8>, _>(name)?),
2202 SqlType::Decimal => panic_pg_only_unsupported(&col.name),
2203 })
2204}
2205
2206pub fn decode_pg_to_string(
2217 row: &sqlx::postgres::PgRow,
2218 col: &Column,
2219) -> Result<String, sqlx::Error> {
2220 use chrono::{DateTime, NaiveDate, NaiveTime, Utc};
2221 use serde_json::Value;
2222 use uuid::Uuid;
2223
2224 let name = col.name.as_str();
2225 if col.nullable {
2226 return Ok(match col.ty {
2227 SqlType::SmallInt => row
2228 .try_get::<Option<i16>, _>(name)?
2229 .map_or(String::new(), |v| v.to_string()),
2230 SqlType::Integer => row
2231 .try_get::<Option<i32>, _>(name)?
2232 .map_or(String::new(), |v| v.to_string()),
2233 SqlType::BigInt => row
2234 .try_get::<Option<i64>, _>(name)?
2235 .map_or(String::new(), |v| v.to_string()),
2236 SqlType::Real => row
2237 .try_get::<Option<f32>, _>(name)?
2238 .map_or(String::new(), |v| v.to_string()),
2239 SqlType::Double => row
2240 .try_get::<Option<f64>, _>(name)?
2241 .map_or(String::new(), |v| v.to_string()),
2242 SqlType::Boolean => row
2243 .try_get::<Option<bool>, _>(name)?
2244 .map_or(String::new(), |v| {
2245 if v { "true" } else { "false" }.to_string()
2246 }),
2247 SqlType::Text => row.try_get::<Option<String>, _>(name)?.unwrap_or_default(),
2248 SqlType::Date => row
2249 .try_get::<Option<NaiveDate>, _>(name)?
2250 .map_or(String::new(), |v| v.to_string()),
2251 SqlType::Time => row
2252 .try_get::<Option<NaiveTime>, _>(name)?
2253 .map_or(String::new(), |v| v.to_string()),
2254 SqlType::Timestamptz => row
2255 .try_get::<Option<DateTime<Utc>>, _>(name)?
2256 .map_or(String::new(), |v| v.to_rfc3339()),
2257 SqlType::Uuid => row
2258 .try_get::<Option<Uuid>, _>(name)?
2259 .map_or(String::new(), |v| v.to_string()),
2260 SqlType::Json => row
2261 .try_get::<Option<Value>, _>(name)?
2262 .map_or(String::new(), |v| v.to_string()),
2263 SqlType::Array(_)
2269 | SqlType::Inet
2270 | SqlType::Cidr
2271 | SqlType::MacAddr
2272 | SqlType::Xml
2273 | SqlType::Ltree
2274 | SqlType::Bit
2275 | SqlType::FullText => row
2276 .try_get::<Option<String>, _>(name)
2277 .ok()
2278 .flatten()
2279 .unwrap_or_default(),
2280 SqlType::ForeignKey => match fk_target_pk_sql_type(col) {
2283 Some(SqlType::Text) => row.try_get::<Option<String>, _>(name)?.unwrap_or_default(),
2284 Some(SqlType::Uuid) => row
2285 .try_get::<Option<Uuid>, _>(name)?
2286 .map_or(String::new(), |v| v.to_string()),
2287 _ => row
2288 .try_get::<Option<i64>, _>(name)?
2289 .map_or(String::new(), |v| v.to_string()),
2290 },
2291 SqlType::Bytes => row
2292 .try_get::<Option<Vec<u8>>, _>(name)?
2293 .map_or(String::new(), |b| hex_encode(&b)),
2294 SqlType::Decimal => panic_pg_only_unsupported(&col.name),
2295 });
2296 }
2297 Ok(match col.ty {
2298 SqlType::SmallInt => row.try_get::<i16, _>(name)?.to_string(),
2299 SqlType::Integer => row.try_get::<i32, _>(name)?.to_string(),
2300 SqlType::BigInt => row.try_get::<i64, _>(name)?.to_string(),
2301 SqlType::Real => row.try_get::<f32, _>(name)?.to_string(),
2302 SqlType::Double => row.try_get::<f64, _>(name)?.to_string(),
2303 SqlType::Boolean => if row.try_get::<bool, _>(name)? {
2304 "true"
2305 } else {
2306 "false"
2307 }
2308 .to_string(),
2309 SqlType::Text => row.try_get::<String, _>(name)?,
2310 SqlType::Date => row.try_get::<NaiveDate, _>(name)?.to_string(),
2311 SqlType::Time => row.try_get::<NaiveTime, _>(name)?.to_string(),
2312 SqlType::Timestamptz => row.try_get::<DateTime<Utc>, _>(name)?.to_rfc3339(),
2313 SqlType::Uuid => row.try_get::<Uuid, _>(name)?.to_string(),
2314 SqlType::Json => row.try_get::<Value, _>(name)?.to_string(),
2315 SqlType::Array(_)
2317 | SqlType::Inet
2318 | SqlType::Cidr
2319 | SqlType::MacAddr
2320 | SqlType::Xml
2321 | SqlType::Ltree
2322 | SqlType::Bit
2323 | SqlType::FullText => row.try_get::<String, _>(name).unwrap_or_default(),
2324 SqlType::ForeignKey => match fk_target_pk_sql_type(col) {
2325 Some(SqlType::Text) => row.try_get::<String, _>(name)?,
2326 Some(SqlType::Uuid) => row.try_get::<Uuid, _>(name)?.to_string(),
2327 _ => row.try_get::<i64, _>(name)?.to_string(),
2328 },
2329 SqlType::Bytes => hex_encode(&row.try_get::<Vec<u8>, _>(name)?),
2330 SqlType::Decimal => panic_pg_only_unsupported(&col.name),
2331 })
2332}
2333
2334pub fn decode_to_json_aliased(
2347 row: &sqlx::sqlite::SqliteRow,
2348 col: &Column,
2349 alias: &str,
2350) -> Result<serde_json::Value, sqlx::Error> {
2351 let mut aliased = col.clone();
2352 aliased.name = alias.to_string();
2353 decode_to_json(row, &aliased)
2354}
2355
2356pub fn decode_pg_to_json_aliased(
2358 row: &sqlx::postgres::PgRow,
2359 col: &Column,
2360 alias: &str,
2361) -> Result<serde_json::Value, sqlx::Error> {
2362 let mut aliased = col.clone();
2363 aliased.name = alias.to_string();
2364 decode_pg_to_json(row, &aliased)
2365}
2366
2367fn fk_target_pk_sql_type(col: &Column) -> Option<SqlType> {
2391 if !matches!(col.ty, SqlType::ForeignKey) {
2392 return None;
2393 }
2394 let target_table = col.fk_target.as_deref()?;
2395 crate::migrate::pk_meta_for_table(target_table).map(|(_, ty)| ty)
2396}
2397
2398pub fn decode_to_json(
2399 row: &sqlx::sqlite::SqliteRow,
2400 col: &Column,
2401) -> Result<serde_json::Value, sqlx::Error> {
2402 use chrono::{DateTime, NaiveDate, NaiveTime, Utc};
2403 use serde_json::Value;
2404 use uuid::Uuid;
2405
2406 let name = col.name.as_str();
2407 if col.nullable {
2408 return Ok(match col.ty {
2409 SqlType::SmallInt | SqlType::Integer => row
2410 .try_get::<Option<i32>, _>(name)?
2411 .map_or(Value::Null, Value::from),
2412 SqlType::BigInt => row
2413 .try_get::<Option<i64>, _>(name)?
2414 .map_or(Value::Null, Value::from),
2415 SqlType::Real => row
2416 .try_get::<Option<f32>, _>(name)?
2417 .map_or(Value::Null, |v| Value::from(v as f64)),
2418 SqlType::Double => row
2419 .try_get::<Option<f64>, _>(name)?
2420 .map_or(Value::Null, Value::from),
2421 SqlType::Boolean => row
2422 .try_get::<Option<bool>, _>(name)?
2423 .map_or(Value::Null, Value::from),
2424 SqlType::Text => row
2425 .try_get::<Option<String>, _>(name)?
2426 .map_or(Value::Null, Value::from),
2427 SqlType::Date => row
2428 .try_get::<Option<NaiveDate>, _>(name)?
2429 .map_or(Value::Null, |v| Value::from(v.to_string())),
2430 SqlType::Time => row
2431 .try_get::<Option<NaiveTime>, _>(name)?
2432 .map_or(Value::Null, |v| Value::from(v.to_string())),
2433 SqlType::Timestamptz => row
2434 .try_get::<Option<DateTime<Utc>>, _>(name)?
2435 .map_or(Value::Null, |v| Value::from(v.to_rfc3339())),
2436 SqlType::Uuid => row
2437 .try_get::<Option<Uuid>, _>(name)?
2438 .map_or(Value::Null, |v| Value::from(v.to_string())),
2439 SqlType::Json => row
2440 .try_get::<Option<Value>, _>(name)?
2441 .unwrap_or(Value::Null),
2442 SqlType::Array(_) => panic_array_unsupported(&col.name),
2443 SqlType::Inet
2444 | SqlType::Cidr
2445 | SqlType::MacAddr
2446 | SqlType::Xml
2447 | SqlType::Ltree
2448 | SqlType::Bit
2449 | SqlType::FullText => panic_pg_only_unsupported(&col.name),
2450 SqlType::ForeignKey => match fk_target_pk_sql_type(col) {
2455 Some(SqlType::Text) => row
2456 .try_get::<Option<String>, _>(name)?
2457 .map_or(Value::Null, Value::from),
2458 Some(SqlType::Uuid) => row
2459 .try_get::<Option<Uuid>, _>(name)?
2460 .map_or(Value::Null, |v| Value::from(v.to_string())),
2461 _ => row
2462 .try_get::<Option<i64>, _>(name)?
2463 .map_or(Value::Null, Value::from),
2464 },
2465 SqlType::Bytes => row
2466 .try_get::<Option<Vec<u8>>, _>(name)?
2467 .map_or(Value::Null, |b| bytes_to_json(&b)),
2468 SqlType::Decimal => panic_pg_only_unsupported(&col.name),
2469 });
2470 }
2471 Ok(match col.ty {
2472 SqlType::SmallInt | SqlType::Integer => Value::from(row.try_get::<i32, _>(name)?),
2473 SqlType::BigInt => Value::from(row.try_get::<i64, _>(name)?),
2474 SqlType::Real => Value::from(row.try_get::<f32, _>(name)? as f64),
2475 SqlType::Double => Value::from(row.try_get::<f64, _>(name)?),
2476 SqlType::Boolean => Value::from(row.try_get::<bool, _>(name)?),
2477 SqlType::Text => Value::from(row.try_get::<String, _>(name)?),
2478 SqlType::Date => Value::from(row.try_get::<NaiveDate, _>(name)?.to_string()),
2479 SqlType::Time => Value::from(row.try_get::<NaiveTime, _>(name)?.to_string()),
2480 SqlType::Timestamptz => Value::from(row.try_get::<DateTime<Utc>, _>(name)?.to_rfc3339()),
2481 SqlType::Uuid => Value::from(row.try_get::<Uuid, _>(name)?.to_string()),
2482 SqlType::Json => row.try_get::<Value, _>(name)?,
2483 SqlType::Array(_) => panic_array_unsupported(&col.name),
2484 SqlType::Inet
2485 | SqlType::Cidr
2486 | SqlType::MacAddr
2487 | SqlType::Xml
2488 | SqlType::Ltree
2489 | SqlType::Bit
2490 | SqlType::FullText => panic_pg_only_unsupported(&col.name),
2491 SqlType::ForeignKey => match fk_target_pk_sql_type(col) {
2494 Some(SqlType::Text) => Value::from(row.try_get::<String, _>(name)?),
2495 Some(SqlType::Uuid) => Value::from(row.try_get::<Uuid, _>(name)?.to_string()),
2496 _ => Value::from(row.try_get::<i64, _>(name)?),
2497 },
2498 SqlType::Bytes => bytes_to_json(&row.try_get::<Vec<u8>, _>(name)?),
2499 SqlType::Decimal => panic_pg_only_unsupported(&col.name),
2500 })
2501}
2502
2503pub fn decode_pg_to_json(
2507 row: &sqlx::postgres::PgRow,
2508 col: &Column,
2509) -> Result<serde_json::Value, sqlx::Error> {
2510 use chrono::{DateTime, NaiveDate, NaiveTime, Utc};
2511 use serde_json::Value;
2512 use uuid::Uuid;
2513
2514 let name = col.name.as_str();
2515 if col.nullable {
2516 return Ok(match col.ty {
2517 SqlType::SmallInt => row
2518 .try_get::<Option<i16>, _>(name)?
2519 .map_or(Value::Null, Value::from),
2520 SqlType::Integer => row
2521 .try_get::<Option<i32>, _>(name)?
2522 .map_or(Value::Null, Value::from),
2523 SqlType::BigInt => row
2524 .try_get::<Option<i64>, _>(name)?
2525 .map_or(Value::Null, Value::from),
2526 SqlType::Real => row
2527 .try_get::<Option<f32>, _>(name)?
2528 .map_or(Value::Null, |v| Value::from(v as f64)),
2529 SqlType::Double => row
2530 .try_get::<Option<f64>, _>(name)?
2531 .map_or(Value::Null, Value::from),
2532 SqlType::Boolean => row
2533 .try_get::<Option<bool>, _>(name)?
2534 .map_or(Value::Null, Value::from),
2535 SqlType::Text => row
2536 .try_get::<Option<String>, _>(name)?
2537 .map_or(Value::Null, Value::from),
2538 SqlType::Date => row
2539 .try_get::<Option<NaiveDate>, _>(name)?
2540 .map_or(Value::Null, |v| Value::from(v.to_string())),
2541 SqlType::Time => row
2542 .try_get::<Option<NaiveTime>, _>(name)?
2543 .map_or(Value::Null, |v| Value::from(v.to_string())),
2544 SqlType::Timestamptz => row
2545 .try_get::<Option<DateTime<Utc>>, _>(name)?
2546 .map_or(Value::Null, |v| Value::from(v.to_rfc3339())),
2547 SqlType::Uuid => row
2548 .try_get::<Option<Uuid>, _>(name)?
2549 .map_or(Value::Null, |v| Value::from(v.to_string())),
2550 SqlType::Json => row
2551 .try_get::<Option<Value>, _>(name)?
2552 .unwrap_or(Value::Null),
2553 SqlType::Array(_)
2554 | SqlType::Inet
2555 | SqlType::Cidr
2556 | SqlType::MacAddr
2557 | SqlType::Xml
2558 | SqlType::Ltree
2559 | SqlType::Bit
2560 | SqlType::FullText => row
2561 .try_get::<Option<String>, _>(name)
2562 .ok()
2563 .flatten()
2564 .map_or(Value::Null, Value::from),
2565 SqlType::ForeignKey => match fk_target_pk_sql_type(col) {
2568 Some(SqlType::Text) => row
2569 .try_get::<Option<String>, _>(name)?
2570 .map_or(Value::Null, Value::from),
2571 Some(SqlType::Uuid) => row
2572 .try_get::<Option<Uuid>, _>(name)?
2573 .map_or(Value::Null, |v| Value::from(v.to_string())),
2574 _ => row
2575 .try_get::<Option<i64>, _>(name)?
2576 .map_or(Value::Null, Value::from),
2577 },
2578 SqlType::Bytes => row
2579 .try_get::<Option<Vec<u8>>, _>(name)?
2580 .map_or(Value::Null, |b| bytes_to_json(&b)),
2581 SqlType::Decimal => panic_pg_only_unsupported(&col.name),
2582 });
2583 }
2584 Ok(match col.ty {
2585 SqlType::SmallInt => Value::from(row.try_get::<i16, _>(name)?),
2586 SqlType::Integer => Value::from(row.try_get::<i32, _>(name)?),
2587 SqlType::BigInt => Value::from(row.try_get::<i64, _>(name)?),
2588 SqlType::Real => Value::from(row.try_get::<f32, _>(name)? as f64),
2589 SqlType::Double => Value::from(row.try_get::<f64, _>(name)?),
2590 SqlType::Boolean => Value::from(row.try_get::<bool, _>(name)?),
2591 SqlType::Text => Value::from(row.try_get::<String, _>(name)?),
2592 SqlType::Date => Value::from(row.try_get::<NaiveDate, _>(name)?.to_string()),
2593 SqlType::Time => Value::from(row.try_get::<NaiveTime, _>(name)?.to_string()),
2594 SqlType::Timestamptz => Value::from(row.try_get::<DateTime<Utc>, _>(name)?.to_rfc3339()),
2595 SqlType::Uuid => Value::from(row.try_get::<Uuid, _>(name)?.to_string()),
2596 SqlType::Json => row.try_get::<Value, _>(name)?,
2597 SqlType::Array(_)
2598 | SqlType::Inet
2599 | SqlType::Cidr
2600 | SqlType::MacAddr
2601 | SqlType::Xml
2602 | SqlType::Ltree
2603 | SqlType::Bit
2604 | SqlType::FullText => row
2605 .try_get::<String, _>(name)
2606 .map(Value::from)
2607 .unwrap_or(Value::Null),
2608 SqlType::ForeignKey => match fk_target_pk_sql_type(col) {
2611 Some(SqlType::Text) => Value::from(row.try_get::<String, _>(name)?),
2612 Some(SqlType::Uuid) => Value::from(row.try_get::<Uuid, _>(name)?.to_string()),
2613 _ => Value::from(row.try_get::<i64, _>(name)?),
2614 },
2615 SqlType::Bytes => bytes_to_json(&row.try_get::<Vec<u8>, _>(name)?),
2616 SqlType::Decimal => panic_pg_only_unsupported(&col.name),
2617 })
2618}
2619
2620fn normalize_str(col: &Column, s: &str) -> String {
2626 let trimmed = if col.trim { s.trim() } else { s };
2627 if col.lowercase {
2628 trimmed.to_lowercase()
2629 } else {
2630 trimmed.to_string()
2631 }
2632}
2633
2634fn normalize_json_for_col(col: &Column, v: &serde_json::Value) -> Option<serde_json::Value> {
2640 if !(col.trim || col.lowercase) {
2641 return None;
2642 }
2643 let s = v.as_str()?;
2644 Some(serde_json::Value::String(normalize_str(col, s)))
2645}
2646
2647fn form_str_to_sea_value(col: &Column, raw: &str) -> Result<SeaValue, WriteError> {
2655 let normalized = normalize_str(col, raw);
2659 let raw = normalized.as_str();
2660 if raw.is_empty() {
2661 if col.ty == SqlType::Boolean {
2662 return Ok(SeaValue::Bool(Some(false)));
2664 }
2665 if col.nullable {
2666 return Ok(null_for(col.ty));
2667 }
2668 return Err(WriteError::RequiredFieldMissing {
2669 field: col.name.clone(),
2670 });
2671 }
2672 if matches!(col.ty, SqlType::Json | SqlType::Array(_)) {
2684 let parsed: serde_json::Value =
2685 serde_json::from_str(raw).map_err(|e| WriteError::Validator {
2686 field: col.name.clone(),
2687 message: format!("Not valid JSON: {e}"),
2688 })?;
2689 return json_to_sea_value(col.ty, &parsed, col.nullable, &col.name, None);
2690 }
2691 if matches!(col.ty, SqlType::ForeignKey) {
2692 return match fk_target_pk_sql_type(col) {
2693 Some(SqlType::Text) => Ok(SeaValue::String(Some(Box::new(raw.to_string())))),
2694 Some(SqlType::Uuid) => uuid::Uuid::parse_str(raw)
2695 .map(|v| SeaValue::Uuid(Some(Box::new(v))))
2696 .map_err(|_| WriteError::TypeMismatch {
2697 field: col.name.clone(),
2698 expected: SqlType::Uuid,
2699 got: raw.to_string(),
2700 }),
2701 _ => raw
2702 .parse::<i64>()
2703 .map(|v| SeaValue::BigInt(Some(v)))
2704 .map_err(|_| WriteError::TypeMismatch {
2705 field: col.name.clone(),
2706 expected: SqlType::BigInt,
2707 got: raw.to_string(),
2708 }),
2709 };
2710 }
2711 let json = serde_json::Value::String(raw.to_string());
2712 if let Some(sealed) = crate::orm::write::seal_masked_json(col, &json)? {
2715 return json_to_sea_value(col.ty, &sealed, col.nullable, &col.name, None);
2716 }
2717 json_to_sea_value(col.ty, &json, col.nullable, &col.name, None)
2718}
2719
2720fn hex_encode(bytes: &[u8]) -> String {
2724 let mut out = String::with_capacity(bytes.len() * 2);
2725 for b in bytes {
2726 out.push_str(&format!("{b:02x}"));
2727 }
2728 out
2729}
2730
2731fn bytes_to_json(bytes: &[u8]) -> serde_json::Value {
2734 serde_json::Value::Array(bytes.iter().map(|b| serde_json::Value::from(*b)).collect())
2735}
2736
2737fn panic_array_unsupported(column: &str) -> ! {
2738 panic!(
2739 "DynQuerySet: column `{column}` is a Postgres-only Array; the \
2740 field/backend system check should have failed boot."
2741 )
2742}
2743
2744fn panic_pg_only_unsupported(column: &str) -> ! {
2745 panic!(
2746 "DynQuerySet: column `{column}` is a Postgres-only network type \
2747 (Inet/Cidr/MacAddr); the field/backend system check should \
2748 have failed boot."
2749 )
2750}
2751
2752fn classify_or_sqlx(
2758 e: sqlx::Error,
2759 body: &serde_json::Map<String, serde_json::Value>,
2760) -> crate::orm::write::WriteError {
2761 if let Some(classified) = crate::orm::validation::classify_sql_error(&e, body) {
2762 return classified;
2763 }
2764 crate::orm::write::WriteError::Sqlx(e)
2765}
2766
2767fn validate_numeric_bounds(
2768 col: &Column,
2769 json: &serde_json::Value,
2770) -> Result<(), crate::orm::write::WriteError> {
2771 let Some(n) = json.as_f64() else {
2772 return Ok(());
2773 };
2774 if let Some(min) = col.min {
2775 if n < min as f64 {
2776 return Err(crate::orm::write::WriteError::Validator {
2777 field: col.name.clone(),
2778 message: format!("must be >= {min} (got {n})."),
2779 });
2780 }
2781 }
2782 if let Some(max) = col.max {
2783 if n > max as f64 {
2784 return Err(crate::orm::write::WriteError::Validator {
2785 field: col.name.clone(),
2786 message: format!("must be <= {max} (got {n})."),
2787 });
2788 }
2789 }
2790 Ok(())
2791}
2792
2793fn json_pk_to_sea(v: &serde_json::Value) -> Option<sea_query::Value> {
2799 match v {
2800 serde_json::Value::Number(n) => n.as_i64().map(|i| sea_query::Value::BigInt(Some(i))),
2801 serde_json::Value::String(s) => Some(sea_query::Value::String(Some(Box::new(s.clone())))),
2802 _ => None,
2803 }
2804}
2805
2806fn normalize_sr_token(name: &str) -> String {
2823 name.replace("__", ".")
2824}
2825
2826fn validate_sr_chain(root_meta: &crate::migrate::ModelMeta, chain: &str) -> Option<Vec<String>> {
2836 let hops: Vec<&str> = chain.split('.').filter(|s| !s.is_empty()).collect();
2837 if hops.is_empty() {
2838 return None;
2839 }
2840 let registered = crate::migrate::registered_models();
2841 let mut targets: Vec<String> = Vec::with_capacity(hops.len());
2842 let mut current_table: String = root_meta.table.clone();
2843 let mut current_meta: Option<crate::migrate::ModelMeta> = None;
2844 for hop in &hops {
2845 let meta_ref: &crate::migrate::ModelMeta =
2846 if current_table == root_meta.table && current_meta.is_none() {
2847 root_meta
2848 } else {
2849 current_meta = registered
2850 .iter()
2851 .find(|m| m.table == current_table)
2852 .cloned();
2853 current_meta.as_ref()?
2854 };
2855 let col = meta_ref.fields.iter().find(|c| &c.name == hop)?;
2856 let target = col.fk_target.clone()?;
2857 targets.push(target.clone());
2858 current_table = target;
2859 }
2860 Some(targets)
2861}
2862
2863async fn hydrate_select_related_into(
2880 meta: &crate::migrate::ModelMeta,
2881 sr_fields: &[String],
2882 rows: &mut [serde_json::Map<String, serde_json::Value>],
2883) -> Result<(), sqlx::Error> {
2884 let pool = resolve_pool_dyn(meta, crate::db::RouteOp::Read);
2885 for chain in sr_fields {
2886 let hops: Vec<&str> = chain.split('.').filter(|s| !s.is_empty()).collect();
2887 if hops.is_empty() {
2888 continue;
2889 }
2890 let Some(targets) = validate_sr_chain(meta, chain) else {
2891 continue;
2896 };
2897
2898 let registered = crate::migrate::registered_models();
2907 let hop_target_pk: Vec<(String, SqlType)> = targets
2908 .iter()
2909 .filter_map(|t| {
2910 registered
2911 .iter()
2912 .find(|m| &m.table == t)
2913 .and_then(|m| m.pk_column().map(|c| (c.name.clone(), c.ty)))
2914 })
2915 .collect();
2916 if hop_target_pk.len() != hops.len() {
2917 continue;
2921 }
2922 let hop_target_soft_delete: Vec<bool> = targets
2923 .iter()
2924 .map(|t| {
2925 registered
2926 .iter()
2927 .find(|m| &m.table == t)
2928 .is_some_and(|m| m.soft_delete)
2929 })
2930 .collect();
2931
2932 let first_field = hops[0];
2936 let mut ids: Vec<serde_json::Value> = Vec::with_capacity(rows.len());
2937 for row in rows.iter() {
2938 let Some(v) = row.get(first_field) else {
2939 continue;
2940 };
2941 if v.is_null() {
2942 continue;
2943 }
2944 ids.push(v.clone());
2945 }
2946 if ids.is_empty() {
2947 continue;
2948 }
2949 dedup_by_pk_key(&mut ids);
2950 let mut levels: Vec<Vec<serde_json::Value>> = Vec::with_capacity(hops.len());
2951 levels.push(
2952 crate::orm::queryset::hydration::fetch_related_as_json_by_pk(
2953 &targets[0],
2954 &hop_target_pk[0].0,
2955 hop_target_pk[0].1,
2956 hop_target_soft_delete[0],
2957 &ids,
2958 &pool,
2959 )
2960 .await?,
2961 );
2962
2963 for hop_idx in 1..hops.len() {
2964 let hop_field = hops[hop_idx];
2965 let hop_target = &targets[hop_idx];
2966 let prev_lvl = &levels[hop_idx - 1];
2967 let mut next_ids: Vec<serde_json::Value> = prev_lvl
2968 .iter()
2969 .filter_map(|r| {
2970 let v = r.as_object()?.get(hop_field)?;
2971 if v.is_null() { None } else { Some(v.clone()) }
2972 })
2973 .collect();
2974 if next_ids.is_empty() {
2975 break;
2980 }
2981 dedup_by_pk_key(&mut next_ids);
2982 levels.push(
2983 crate::orm::queryset::hydration::fetch_related_as_json_by_pk(
2984 hop_target,
2985 &hop_target_pk[hop_idx].0,
2986 hop_target_pk[hop_idx].1,
2987 hop_target_soft_delete[hop_idx],
2988 &next_ids,
2989 &pool,
2990 )
2991 .await?,
2992 );
2993 }
2994
2995 if levels.len() > 1 {
3000 for i in (0..levels.len() - 1).rev() {
3001 let next_pk_col = &hop_target_pk[i + 1].0;
3002 let next_by_pk: HashMap<String, serde_json::Value> = levels[i + 1]
3003 .iter()
3004 .filter_map(|obj| {
3005 let map = obj.as_object()?;
3006 let pk_val = map.get(next_pk_col.as_str())?;
3007 Some((pk_json_key(pk_val), obj.clone()))
3008 })
3009 .collect();
3010 let hop_field = hops[i + 1];
3011 for row in levels[i].iter_mut() {
3012 let Some(map) = row.as_object_mut() else {
3013 continue;
3014 };
3015 let Some(fk_val) = map.get(hop_field) else {
3016 continue;
3017 };
3018 if fk_val.is_null() {
3019 continue;
3020 }
3021 let key = pk_json_key(fk_val);
3022 if let Some(next_json) = next_by_pk.get(&key) {
3023 map.insert(hop_field.to_string(), next_json.clone());
3024 }
3025 }
3026 }
3027 }
3028
3029 let first_pk_col = &hop_target_pk[0].0;
3036 let first_by_pk: HashMap<String, serde_json::Value> = levels
3037 .into_iter()
3038 .next()
3039 .unwrap_or_default()
3040 .into_iter()
3041 .filter_map(|obj| {
3042 let map = obj.as_object()?;
3043 let pk_val = map.get(first_pk_col.as_str())?;
3044 Some((pk_json_key(pk_val), obj.clone()))
3045 })
3046 .collect();
3047 for row in rows.iter_mut() {
3048 let Some(fk_val) = row.get(first_field) else {
3049 continue;
3050 };
3051 if fk_val.is_null() {
3052 continue;
3053 }
3054 let key = pk_json_key(fk_val);
3055 if let Some(resolved) = first_by_pk.get(&key) {
3056 row.insert(first_field.to_string(), resolved.clone());
3057 }
3058 }
3059 }
3060 Ok(())
3061}
3062
3063fn dedup_by_pk_key(ids: &mut Vec<serde_json::Value>) {
3068 let mut seen: std::collections::HashSet<String> = std::collections::HashSet::new();
3069 ids.retain(|v| seen.insert(pk_json_key(v)));
3070}
3071
3072async fn hydrate_m2m_batched(
3086 meta: &crate::migrate::ModelMeta,
3087 pk_name: &str,
3088 rows: &mut [serde_json::Map<String, serde_json::Value>],
3089) -> Result<(), sqlx::Error> {
3090 if meta.m2m_relations.is_empty() || rows.is_empty() {
3091 return Ok(());
3092 }
3093
3094 for row in rows.iter_mut() {
3099 for rel in &meta.m2m_relations {
3100 row.insert(rel.field_name.clone(), serde_json::Value::Array(Vec::new()));
3101 }
3102 }
3103
3104 let mut parent_sea_vals: Vec<sea_query::Value> = Vec::with_capacity(rows.len());
3108 let mut seen_keys: std::collections::HashSet<String> = std::collections::HashSet::new();
3109 for row in rows.iter() {
3110 let Some(pk_json) = row.get(pk_name) else {
3111 continue;
3112 };
3113 let Some(sea_val) = json_pk_to_sea(pk_json) else {
3114 continue;
3115 };
3116 let key = pk_json_key(pk_json);
3117 if seen_keys.insert(key) {
3118 parent_sea_vals.push(sea_val);
3119 }
3120 }
3121 if parent_sea_vals.is_empty() {
3122 return Ok(());
3123 }
3124
3125 for rel in &meta.m2m_relations {
3126 let junction_table = format!("{}_{}", meta.table, rel.field_name);
3127 let mut sel = Query::select();
3128 sel.from(crate::db::router::schema_qualified_table(&junction_table));
3129 sel.column(Alias::new("parent_id"));
3130 sel.column(Alias::new("child_id"));
3131 sel.and_where(Expr::col(Alias::new("parent_id")).is_in(parent_sea_vals.clone()));
3132
3133 let mut children_by_parent: HashMap<String, Vec<serde_json::Value>> = HashMap::new();
3134 match resolve_pool_dyn(meta, crate::db::RouteOp::Read) {
3135 DbPool::Sqlite(pool) => {
3136 let (sql, values) = sel.build_sqlx(SqliteQueryBuilder);
3137 let db_rows = sqlx::query_with(&sql, values).fetch_all(&pool).await?;
3138 for r in &db_rows {
3139 let parent = read_junction_id_sqlite(r, "parent_id")?;
3140 let child = read_junction_id_sqlite(r, "child_id")?;
3141 children_by_parent
3142 .entry(pk_json_key(&parent))
3143 .or_default()
3144 .push(child);
3145 }
3146 }
3147 DbPool::Postgres(pool) => {
3148 let (sql, values) = sel.build_sqlx(PostgresQueryBuilder);
3149 let db_rows = sqlx::query_with(&sql, values).fetch_all(&pool).await?;
3150 for r in &db_rows {
3151 let parent = read_junction_id_pg(r, "parent_id")?;
3152 let child = read_junction_id_pg(r, "child_id")?;
3153 children_by_parent
3154 .entry(pk_json_key(&parent))
3155 .or_default()
3156 .push(child);
3157 }
3158 }
3159 }
3160
3161 for row in rows.iter_mut() {
3162 let Some(pk_json) = row.get(pk_name) else {
3163 continue;
3164 };
3165 let key = pk_json_key(pk_json);
3166 if let Some(children) = children_by_parent.remove(&key) {
3167 row.insert(rel.field_name.clone(), serde_json::Value::Array(children));
3168 }
3169 }
3170 }
3171 Ok(())
3172}
3173
3174fn pk_json_key(v: &serde_json::Value) -> String {
3180 match v {
3181 serde_json::Value::Number(n) => format!("n:{n}"),
3182 serde_json::Value::String(s) => format!("s:{s}"),
3183 other => format!("o:{other}"),
3184 }
3185}
3186
3187fn read_junction_id_sqlite(
3192 row: &sqlx::sqlite::SqliteRow,
3193 col: &str,
3194) -> Result<serde_json::Value, sqlx::Error> {
3195 if let Ok(i) = row.try_get::<i64, _>(col) {
3196 return Ok(serde_json::Value::Number(i.into()));
3197 }
3198 let s = row.try_get::<String, _>(col)?;
3199 Ok(serde_json::Value::String(s))
3200}
3201
3202fn read_junction_id_pg(
3203 row: &sqlx::postgres::PgRow,
3204 col: &str,
3205) -> Result<serde_json::Value, sqlx::Error> {
3206 if let Ok(i) = row.try_get::<i64, _>(col) {
3207 return Ok(serde_json::Value::Number(i.into()));
3208 }
3209 let s = row.try_get::<String, _>(col)?;
3210 Ok(serde_json::Value::String(s))
3211}
3212
3213async fn collect_parent_pks(
3220 meta: &crate::migrate::ModelMeta,
3221 pk_col: &crate::migrate::Column,
3222 where_clauses: &[Condition],
3223) -> Result<Vec<serde_json::Value>, crate::orm::write::WriteError> {
3224 let mut sel = Query::select();
3225 sel.from(crate::db::router::schema_qualified_table(&meta.table));
3226 sel.column(Alias::new(&pk_col.name));
3227 for cond in where_clauses {
3228 sel.cond_where(cond.clone());
3229 }
3230 match resolve_pool_dyn(meta, crate::db::RouteOp::Read) {
3231 DbPool::Sqlite(pool) => {
3232 let (sql, values) = sel.build_sqlx(SqliteQueryBuilder);
3233 let rows = sqlx::query_with(&sql, values).fetch_all(&pool).await?;
3234 rows.iter()
3235 .map(|row| decode_to_json(row, pk_col))
3236 .collect::<Result<Vec<_>, _>>()
3237 .map_err(crate::orm::write::WriteError::Sqlx)
3238 }
3239 DbPool::Postgres(pool) => {
3240 let (sql, values) = sel.build_sqlx(PostgresQueryBuilder);
3241 let rows = sqlx::query_with(&sql, values).fetch_all(&pool).await?;
3242 rows.iter()
3243 .map(|row| decode_pg_to_json(row, pk_col))
3244 .collect::<Result<Vec<_>, _>>()
3245 .map_err(crate::orm::write::WriteError::Sqlx)
3246 }
3247 }
3248}
3249
3250async fn collect_parent_pks_in_tx(
3254 meta: &crate::migrate::ModelMeta,
3255 pk_col: &crate::migrate::Column,
3256 where_clauses: &[Condition],
3257 tx: &mut crate::db::Transaction,
3258) -> Result<Vec<serde_json::Value>, crate::orm::write::WriteError> {
3259 let mut sel = Query::select();
3260 sel.from(crate::db::router::schema_qualified_table(&meta.table));
3261 sel.column(Alias::new(&pk_col.name));
3262 for cond in where_clauses {
3263 sel.cond_where(cond.clone());
3264 }
3265 match tx.backend_name() {
3266 "sqlite" => {
3267 let (sql, values) = sel.build_sqlx(SqliteQueryBuilder);
3268 let inner = tx.as_sqlite_mut().expect("sqlite backend_name");
3269 let rows = sqlx::query_with(&sql, values)
3270 .fetch_all(&mut **inner)
3271 .await?;
3272 rows.iter()
3273 .map(|row| decode_to_json(row, pk_col))
3274 .collect::<Result<Vec<_>, _>>()
3275 .map_err(crate::orm::write::WriteError::Sqlx)
3276 }
3277 _ => {
3278 let (sql, values) = sel.build_sqlx(PostgresQueryBuilder);
3279 let inner = tx.as_pg_mut().expect("postgres backend_name");
3280 let rows = sqlx::query_with(&sql, values)
3281 .fetch_all(&mut **inner)
3282 .await?;
3283 rows.iter()
3284 .map(|row| decode_pg_to_json(row, pk_col))
3285 .collect::<Result<Vec<_>, _>>()
3286 .map_err(crate::orm::write::WriteError::Sqlx)
3287 }
3288 }
3289}
3290
3291fn is_unauthorized_privileged(col: &crate::migrate::Column, allow_privileged: &[String]) -> bool {
3310 col.privileged && !allow_privileged.iter().any(|a| a == &col.name)
3311}
3312
3313fn normalise_insert_body(
3314 meta: &crate::migrate::ModelMeta,
3315 body: &serde_json::Map<String, serde_json::Value>,
3316 allow_privileged: &[String],
3317) -> Option<serde_json::Map<String, serde_json::Value>> {
3318 let needs_owned = meta.fields.iter().any(|c| {
3319 c.noform || c.slug_from.is_some() || is_unauthorized_privileged(c, allow_privileged)
3320 });
3321 if !needs_owned {
3322 return None;
3323 }
3324 let mut owned = body.clone();
3325 for col in &meta.fields {
3326 if col.noform || is_unauthorized_privileged(col, allow_privileged) {
3327 owned.remove(&col.name);
3328 }
3329 }
3330 crate::orm::write::apply_slug_from(&meta.fields, &mut owned, false);
3331 Some(owned)
3332}
3333
3334fn normalise_update_body(
3339 meta: &crate::migrate::ModelMeta,
3340 body: &serde_json::Map<String, serde_json::Value>,
3341 allow_privileged: &[String],
3342) -> Option<serde_json::Map<String, serde_json::Value>> {
3343 let needs_owned = meta.fields.iter().any(|c| {
3344 c.noform || c.slug_from.is_some() || is_unauthorized_privileged(c, allow_privileged)
3345 });
3346 if !needs_owned {
3347 return None;
3348 }
3349 let mut owned = body.clone();
3350 for col in &meta.fields {
3351 if col.noform || is_unauthorized_privileged(col, allow_privileged) {
3352 owned.remove(&col.name);
3353 }
3354 }
3355 crate::orm::write::apply_slug_from(&meta.fields, &mut owned, true);
3356 Some(owned)
3357}
3358
3359struct InsertPlan {
3361 q: sea_query::InsertStatement,
3362 pk_name: String,
3363 pk_ty: SqlType,
3364}
3365
3366fn build_insert_plan(
3375 meta: &crate::migrate::ModelMeta,
3376 body: &serde_json::Map<String, serde_json::Value>,
3377) -> Result<InsertPlan, crate::orm::write::WriteError> {
3378 use crate::orm::write::{WriteError, is_default_pk};
3379
3380 let mut cols: Vec<&str> = Vec::new();
3381 let mut values: Vec<SeaValue> = Vec::new();
3382 for col in &meta.fields {
3383 if col.primary_key {
3384 let supplied = body.get(&col.name);
3385 let is_sentinel = match supplied {
3386 None | Some(serde_json::Value::Null) => true,
3387 Some(v) => is_default_pk(col.ty, v),
3388 };
3389 if matches!(
3390 col.ty,
3391 SqlType::Integer | SqlType::BigInt | SqlType::SmallInt
3392 ) && is_sentinel
3393 {
3394 continue;
3395 }
3396 }
3397 let Some(json) = body.get(&col.name) else {
3398 if col.auto_now_add || col.auto_now {
3399 let now_value = crate::orm::write::now_for_column(col.ty);
3400 cols.push(&col.name);
3401 values.push(now_value);
3402 continue;
3403 }
3404 continue;
3405 };
3406 if json.is_null() {
3407 continue;
3408 }
3409 validate_numeric_bounds(col, json)?;
3410 if let (Some(fmt), Some(s)) = (col.text_format.as_deref(), json.as_str()) {
3411 if let Err(e) = crate::orm::validators::validate_text_format(fmt, s) {
3412 return Err(WriteError::Validator {
3413 field: col.name.clone(),
3414 message: e.to_string(),
3415 });
3416 }
3417 }
3418 let normalized_json = normalize_json_for_col(col, json);
3421 let json = normalized_json.as_ref().unwrap_or(json);
3422 let sealed = crate::orm::write::seal_masked_json(col, json)?;
3424 let sea_value = crate::orm::write::json_to_sea_value(
3425 col.ty,
3426 sealed.as_ref().unwrap_or(json),
3427 col.nullable,
3428 &col.name,
3429 fk_target_pk_sql_type(col),
3430 )?;
3431 cols.push(&col.name);
3432 values.push(sea_value);
3433 }
3434
3435 let pk_col = meta.fields.iter().find(|c| c.primary_key).ok_or_else(|| {
3436 WriteError::Sqlx(sqlx::Error::Protocol(
3437 "insert_json: model has no PK".to_string(),
3438 ))
3439 })?;
3440 let pk_name = pk_col.name.clone();
3441 let pk_ty = pk_col.ty;
3442
3443 let mut q = Query::insert();
3444 q.into_table(crate::db::router::schema_qualified_table(&meta.table));
3445 q.columns(cols.iter().map(|c| Alias::new(*c)).collect::<Vec<_>>());
3446 let exprs: Vec<sea_query::SimpleExpr> = values.into_iter().map(Into::into).collect();
3447 q.values_panic(exprs);
3448
3449 Ok(InsertPlan { q, pk_name, pk_ty })
3450}
3451
3452async fn write_m2m_junctions_in_tx(
3456 meta: &crate::migrate::ModelMeta,
3457 parent_pk_json: Option<&serde_json::Value>,
3458 body: &serde_json::Map<String, serde_json::Value>,
3459 tx: &mut crate::db::Transaction,
3460) -> Result<(), crate::orm::write::WriteError> {
3461 if meta.m2m_relations.is_empty() {
3462 return Ok(());
3463 }
3464 let Some(parent_pk_value) = parent_pk_json.and_then(json_pk_to_sea) else {
3465 return Ok(());
3466 };
3467 for rel in &meta.m2m_relations {
3468 let Some(value) = body.get(&rel.field_name) else {
3469 continue;
3470 };
3471 let Some(items) = value.as_array() else {
3472 continue;
3473 };
3474 let mut child_ids: Vec<sea_query::Value> = Vec::with_capacity(items.len());
3475 for item in items {
3476 if item.is_null() {
3477 continue;
3478 }
3479 if let Some(v) = json_pk_to_sea(item) {
3480 child_ids.push(v);
3481 }
3482 }
3483 let junction_table = format!("{}_{}", meta.table, rel.field_name);
3484 crate::orm::m2m::set_junction_dynamic_in_tx(
3485 &junction_table,
3486 parent_pk_value.clone(),
3487 child_ids,
3488 tx,
3489 )
3490 .await
3491 .map_err(crate::orm::write::WriteError::Sqlx)?;
3492 }
3493 Ok(())
3494}
3495
3496async fn hydrate_m2m_into_tx(
3501 meta: &crate::migrate::ModelMeta,
3502 parent_pk_json: Option<&serde_json::Value>,
3503 out: &mut serde_json::Map<String, serde_json::Value>,
3504 tx: &mut crate::db::Transaction,
3505) -> Result<(), sqlx::Error> {
3506 if meta.m2m_relations.is_empty() {
3507 return Ok(());
3508 }
3509 let Some(parent_pk_value) = parent_pk_json.and_then(json_pk_to_sea) else {
3510 return Ok(());
3511 };
3512 for rel in &meta.m2m_relations {
3513 let junction_table = format!("{}_{}", meta.table, rel.field_name);
3514 let mut sel = Query::select();
3515 sel.from(crate::db::router::schema_qualified_table(&junction_table));
3516 sel.column(Alias::new("child_id"));
3517 sel.and_where(Expr::col(Alias::new("parent_id")).eq(parent_pk_value.clone()));
3518 let children: Vec<serde_json::Value> = match tx.backend_name() {
3519 "sqlite" => {
3520 let inner = tx.as_sqlite_mut().expect("sqlite backend_name");
3521 let (sql, values) = sel.build_sqlx(SqliteQueryBuilder);
3522 let rows = sqlx::query_with(&sql, values)
3523 .fetch_all(&mut **inner)
3524 .await?;
3525 rows.iter()
3526 .map(|r| {
3527 r.try_get::<i64, _>("child_id")
3528 .map(|i| serde_json::Value::Number(i.into()))
3529 .or_else(|_| {
3530 r.try_get::<String, _>("child_id")
3531 .map(serde_json::Value::String)
3532 })
3533 })
3534 .collect::<Result<Vec<_>, _>>()?
3535 }
3536 _ => {
3537 let inner = tx.as_pg_mut().expect("postgres backend_name");
3538 let (sql, values) = sel.build_sqlx(PostgresQueryBuilder);
3539 let rows = sqlx::query_with(&sql, values)
3540 .fetch_all(&mut **inner)
3541 .await?;
3542 rows.iter()
3543 .map(|r| {
3544 r.try_get::<i64, _>("child_id")
3545 .map(|i| serde_json::Value::Number(i.into()))
3546 .or_else(|_| {
3547 r.try_get::<String, _>("child_id")
3548 .map(serde_json::Value::String)
3549 })
3550 })
3551 .collect::<Result<Vec<_>, _>>()?
3552 }
3553 };
3554 out.insert(rel.field_name.clone(), serde_json::Value::Array(children));
3555 }
3556 Ok(())
3557}
3558
3559fn coerce_csv_cell(ty: SqlType, nullable: bool, raw: &str) -> serde_json::Value {
3576 use serde_json::Value;
3577 if raw.is_empty() && nullable {
3578 return Value::Null;
3579 }
3580 match ty {
3581 SqlType::SmallInt | SqlType::Integer | SqlType::BigInt | SqlType::ForeignKey => raw
3582 .parse::<i64>()
3583 .map(Value::from)
3584 .unwrap_or_else(|_| Value::String(raw.to_string())),
3585 SqlType::Real | SqlType::Double => raw
3586 .parse::<f64>()
3587 .ok()
3588 .and_then(serde_json::Number::from_f64)
3589 .map(Value::Number)
3590 .unwrap_or_else(|| Value::String(raw.to_string())),
3591 SqlType::Boolean => match raw.trim().to_ascii_lowercase().as_str() {
3592 "true" | "1" | "t" | "yes" | "y" => Value::Bool(true),
3593 "false" | "0" | "f" | "no" | "n" => Value::Bool(false),
3594 _ => Value::String(raw.to_string()),
3595 },
3596 SqlType::Json => {
3597 serde_json::from_str(raw).unwrap_or_else(|_| Value::String(raw.to_string()))
3598 }
3599 _ => Value::String(raw.to_string()),
3600 }
3601}
3602
3603#[derive(Debug, Default)]
3610pub struct CsvImportReport {
3611 pub inserted: usize,
3612 pub errors: Vec<(usize, String)>,
3613}
3614
3615pub async fn import_table_rows(
3626 meta: &ModelMeta,
3627 headers: &[String],
3628 rows: &[Vec<String>],
3629) -> CsvImportReport {
3630 let col_for: HashMap<&str, &Column> =
3631 meta.fields.iter().map(|c| (c.name.as_str(), c)).collect();
3632
3633 let mut report = CsvImportReport::default();
3634 for (i, row) in rows.iter().enumerate() {
3635 let mut obj = serde_json::Map::new();
3636 for (header, cell) in headers.iter().zip(row.iter()) {
3637 if let Some(col) = col_for.get(header.as_str()) {
3638 obj.insert(header.clone(), coerce_csv_cell(col.ty, col.nullable, cell));
3639 }
3640 }
3641 match DynQuerySet::for_meta(meta).insert_json(&obj).await {
3642 Ok(_) => report.inserted += 1,
3643 Err(e) => report.errors.push((i + 2, e.to_string())),
3644 }
3645 }
3646 report
3647}
3648
3649#[cfg(test)]
3650mod tests {
3651 use super::form_str_to_sea_value;
3652 use crate::migrate::Column;
3653 use crate::orm::{FkAction, SqlType};
3654 use sea_query::Value as SeaValue;
3655
3656 fn col(name: &str, ty: SqlType, nullable: bool) -> Column {
3657 Column {
3658 name: name.to_string(),
3659 ty,
3660 primary_key: false,
3661 nullable,
3662 fk_target: None,
3663 noform: false,
3664 privileged: false,
3665 db_constraint: true,
3666 noedit: false,
3667 is_string_repr: false,
3668 max_length: 0,
3669 choices: Vec::new(),
3670 choice_labels: Vec::new(),
3671 default: String::new(),
3672 is_multichoice: false,
3673 unique: false,
3674 on_delete: FkAction::NoAction,
3675 on_update: FkAction::NoAction,
3676 index: false,
3677 auto_now_add: false,
3678 auto_now: false,
3679 trim: false,
3680 lowercase: false,
3681 case_insensitive: false,
3682 help: String::new(),
3683 example: String::new(),
3684 widget: None,
3685 supported_backends: Vec::new(),
3686 min: None,
3687 max: None,
3688 text_format: None,
3689 slug_from: None,
3690 }
3691 }
3692
3693 #[test]
3694 fn form_fk_numeric_string_binds_as_bigint() {
3695 let mut plugin = col("plugin", SqlType::ForeignKey, false);
3696 plugin.fk_target = Some("plugin".to_string());
3697
3698 let value = form_str_to_sea_value(&plugin, "1").expect("coerce FK id");
3699
3700 assert_eq!(
3701 value,
3702 SeaValue::BigInt(Some(1)),
3703 "integer-backed FK form values must bind as bigint, not text"
3704 );
3705 }
3706
3707 #[test]
3708 fn nullable_form_fk_blank_binds_as_null_bigint() {
3709 let mut parent = col("parent", SqlType::ForeignKey, true);
3710 parent.fk_target = Some("plugin_comment".to_string());
3711
3712 let value = form_str_to_sea_value(&parent, "").expect("blank nullable FK");
3713
3714 assert_eq!(
3715 value,
3716 SeaValue::BigInt(None),
3717 "blank nullable integer-backed FK should bind SQL NULL"
3718 );
3719 }
3720}