/
mcmare
/
RustAPI
Обзор
Документация
Войти
/
mcmare
/
RustAPI
Код
Вики
Пакеты
0
Релизы
0
Аналитика
Безопасность
main
engine/src/db/crud.rs
422 строки
17 KB
mcmare
Add multi-DB CRUD engine: pooled connections, auto table creation, create/read/list/update/delete
18 июл 2026, 16:21
18 июл 2026, 16:21
c40e4bc
Код
Авторство
О чём код?
use rustapi_core::ApiError; use serde_json::Value; use std::collections::HashMap; use super::DbPool; use crate::config::FieldType; fn map_db_err(e: sqlx::Error) -> ApiError { if let sqlx::Error::Database(db_err) = &e { if db_err.is_unique_violation() { return ApiError::Conflict("unique constraint violated".to_string()); } } ApiError::Internal(anyhow::anyhow!(e)) } /// Converts a raw path/query-param string into a JSON value matching the target field /// type, so it can go through the same typed bind path as request-body values. fn param_to_value(ty: &FieldType, raw: &str) -> Result<Value, ApiError> { Ok(match ty { FieldType::Number => { let n = raw .parse::<f64>() .map_err(|_| ApiError::Validation(format!("invalid number '{raw}'")))?; Value::Number( serde_json::Number::from_f64(n) .ok_or_else(|| ApiError::Validation(format!("invalid number '{raw}'")))?, ) } FieldType::Bool => { let b = raw .parse::<bool>() .map_err(|_| ApiError::Validation(format!("invalid bool '{raw}'")))?; Value::Bool(b) } _ => Value::String(raw.to_string()), }) } /// Resolves values to INSERT for a `create`: client-supplied `input` fields, plus /// server-generated fields present in `output` but absent from `input` — a UUID for the /// primary key, `now()` for any other datetime column. Anything else stays unset (DB /// default/NULL). fn compute_insert_values( primary_key: &str, input: &HashMap<String, FieldType>, output: &HashMap<String, FieldType>, body: &Value, ) -> Result<HashMap<String, (FieldType, Value)>, ApiError> { let body_obj = body .as_object() .ok_or_else(|| ApiError::Validation("body must be a JSON object".to_string()))?; let mut values = HashMap::new(); for (name, ty) in input { let v = body_obj .get(name) .ok_or_else(|| ApiError::Validation(format!("missing required field '{name}'")))?; values.insert(name.clone(), (ty.clone(), v.clone())); } for (name, ty) in output { if values.contains_key(name) { continue; } if name == primary_key && *ty == FieldType::Uuid { values.insert(name.clone(), (ty.clone(), Value::String(uuid::Uuid::new_v4().to_string()))); } else if *ty == FieldType::Datetime { values.insert(name.clone(), (ty.clone(), Value::String(chrono::Utc::now().to_rfc3339()))); } } Ok(values) } macro_rules! impl_backend { ($mod_name:ident, $db:ty, $pool:ty, $row:ty, $qopen:literal, $qclose:literal) => { mod $mod_name { use super::*; use sqlx::{QueryBuilder, Row}; fn quote_ident(ident: &str) -> String { format!("{}{}{}", $qopen, ident, $qclose) } fn push_bind(builder: &mut QueryBuilder<'_, $db>, ty: &FieldType, value: &Value) -> Result<(), ApiError> { match ty { FieldType::String => { let v = value .as_str() .ok_or_else(|| ApiError::Validation("expected a string".to_string()))? .to_string(); builder.push_bind(v); } FieldType::Number => { let v = value .as_f64() .ok_or_else(|| ApiError::Validation("expected a number".to_string()))?; builder.push_bind(v); } FieldType::Bool => { let v = value .as_bool() .ok_or_else(|| ApiError::Validation("expected a bool".to_string()))?; builder.push_bind(v); } FieldType::Uuid => { let s = value .as_str() .ok_or_else(|| ApiError::Validation("expected a uuid string".to_string()))?; let v = uuid::Uuid::parse_str(s) .map_err(|e| ApiError::Validation(format!("invalid uuid: {e}")))?; builder.push_bind(v); } FieldType::Datetime => { let s = value .as_str() .ok_or_else(|| ApiError::Validation("expected an RFC3339 datetime string".to_string()))?; let v = chrono::DateTime::parse_from_rfc3339(s) .map_err(|e| ApiError::Validation(format!("invalid datetime: {e}")))? .with_timezone(&chrono::Utc); builder.push_bind(v); } FieldType::Object | FieldType::Array(_) => { builder.push_bind(value.clone()); } } Ok(()) } fn row_to_json(row: &$row, name: &str, ty: &FieldType) -> Result<Value, ApiError> { let err = |e: sqlx::Error| ApiError::Internal(anyhow::anyhow!("column '{name}': {e}")); Ok(match ty { FieldType::String => row .try_get::<Option<String>, _>(name) .map_err(err)? .map(Value::String) .unwrap_or(Value::Null), FieldType::Number => row .try_get::<Option<f64>, _>(name) .map_err(err)? .and_then(|n| serde_json::Number::from_f64(n).map(Value::Number)) .unwrap_or(Value::Null), FieldType::Bool => row .try_get::<Option<bool>, _>(name) .map_err(err)? .map(Value::Bool) .unwrap_or(Value::Null), FieldType::Uuid => row .try_get::<Option<uuid::Uuid>, _>(name) .map_err(err)? .map(|u| Value::String(u.to_string())) .unwrap_or(Value::Null), FieldType::Datetime => row .try_get::<Option<chrono::DateTime<chrono::Utc>>, _>(name) .map_err(err)? .map(|dt| Value::String(dt.to_rfc3339())) .unwrap_or(Value::Null), FieldType::Object | FieldType::Array(_) => { row.try_get::<Option<Value>, _>(name).map_err(err)?.unwrap_or(Value::Null) } }) } pub(super) async fn insert( pool: &$pool, table: &str, values: &HashMap<String, (FieldType, Value)>, ) -> Result<(), ApiError> { let cols: Vec<&String> = values.keys().collect(); let mut qb: QueryBuilder<$db> = QueryBuilder::new(format!("INSERT INTO {} (", quote_ident(table))); for (i, c) in cols.iter().enumerate() { if i > 0 { qb.push(", "); } qb.push(quote_ident(c)); } qb.push(") VALUES ("); for (i, c) in cols.iter().enumerate() { if i > 0 { qb.push(", "); } let (ty, val) = &values[*c]; push_bind(&mut qb, ty, val)?; } qb.push(")"); qb.build().execute(pool).await.map_err(map_db_err)?; Ok(()) } pub(super) async fn select_by_id( pool: &$pool, table: &str, primary_key: &str, id_ty: &FieldType, id: &Value, output: &HashMap<String, FieldType>, ) -> Result<Option<Value>, ApiError> { let mut qb: QueryBuilder<$db> = QueryBuilder::new(format!( "SELECT * FROM {} WHERE {} = ", quote_ident(table), quote_ident(primary_key) )); push_bind(&mut qb, id_ty, id)?; let row = qb.build().fetch_optional(pool).await.map_err(map_db_err)?; match row { None => Ok(None), Some(row) => { let mut obj = serde_json::Map::new(); for (name, ty) in output { obj.insert(name.clone(), row_to_json(&row, name, ty)?); } Ok(Some(Value::Object(obj))) } } } pub(super) async fn select_list( pool: &$pool, table: &str, output: &HashMap<String, FieldType>, limit: i64, offset: i64, filters: &[(String, FieldType, Value)], ) -> Result<Vec<Value>, ApiError> { let mut qb: QueryBuilder<$db> = QueryBuilder::new(format!("SELECT * FROM {}", quote_ident(table))); if !filters.is_empty() { qb.push(" WHERE "); for (i, (name, ty, val)) in filters.iter().enumerate() { if i > 0 { qb.push(" AND "); } qb.push(quote_ident(name)); qb.push(" = "); push_bind(&mut qb, ty, val)?; } } qb.push(" LIMIT "); qb.push_bind(limit); qb.push(" OFFSET "); qb.push_bind(offset); let rows = qb.build().fetch_all(pool).await.map_err(map_db_err)?; let mut result = Vec::with_capacity(rows.len()); for row in &rows { let mut obj = serde_json::Map::new(); for (name, ty) in output { obj.insert(name.clone(), row_to_json(row, name, ty)?); } result.push(Value::Object(obj)); } Ok(result) } pub(super) async fn update_by_id( pool: &$pool, table: &str, primary_key: &str, id_ty: &FieldType, id: &Value, values: &HashMap<String, (FieldType, Value)>, ) -> Result<bool, ApiError> { let mut qb: QueryBuilder<$db> = QueryBuilder::new(format!("UPDATE {} SET ", quote_ident(table))); for (i, (name, (ty, val))) in values.iter().enumerate() { if i > 0 { qb.push(", "); } qb.push(quote_ident(name)); qb.push(" = "); push_bind(&mut qb, ty, val)?; } qb.push(" WHERE "); qb.push(quote_ident(primary_key)); qb.push(" = "); push_bind(&mut qb, id_ty, id)?; let result = qb.build().execute(pool).await.map_err(map_db_err)?; Ok(result.rows_affected() > 0) } pub(super) async fn delete_by_id( pool: &$pool, table: &str, primary_key: &str, id_ty: &FieldType, id: &Value, ) -> Result<bool, ApiError> { let mut qb: QueryBuilder<$db> = QueryBuilder::new(format!("DELETE FROM {} WHERE ", quote_ident(table))); qb.push(quote_ident(primary_key)); qb.push(" = "); push_bind(&mut qb, id_ty, id)?; let result = qb.build().execute(pool).await.map_err(map_db_err)?; Ok(result.rows_affected() > 0) } } }; } impl_backend!(pg, sqlx::Postgres, sqlx::PgPool, sqlx::postgres::PgRow, "\"", "\""); impl_backend!(mysql, sqlx::MySql, sqlx::MySqlPool, sqlx::mysql::MySqlRow, "`", "`"); impl_backend!(sqlite, sqlx::Sqlite, sqlx::SqlitePool, sqlx::sqlite::SqliteRow, "\"", "\""); pub async fn create( pool: &DbPool, table: &str, primary_key: &str, input: &HashMap<String, FieldType>, output: &HashMap<String, FieldType>, body: &Value, ) -> Result<Value, ApiError> { let values = compute_insert_values(primary_key, input, output, body)?; if values.is_empty() { return Err(ApiError::Validation("no fields to insert".to_string())); } match pool { DbPool::Postgres(p) => pg::insert(p, table, &values).await?, DbPool::MySql(p) => mysql::insert(p, table, &values).await?, DbPool::Sqlite(p) => sqlite::insert(p, table, &values).await?, } let mut obj = serde_json::Map::new(); for name in output.keys() { let v = values.get(name).map(|(_, v)| v.clone()).unwrap_or(Value::Null); obj.insert(name.clone(), v); } Ok(Value::Object(obj)) } pub async fn read( pool: &DbPool, table: &str, primary_key: &str, id_raw: &str, output: &HashMap<String, FieldType>, ) -> Result<Value, ApiError> { let id_ty = output.get(primary_key).cloned().unwrap_or(FieldType::Uuid); let id = param_to_value(&id_ty, id_raw)?; let row = match pool { DbPool::Postgres(p) => pg::select_by_id(p, table, primary_key, &id_ty, &id, output).await?, DbPool::MySql(p) => mysql::select_by_id(p, table, primary_key, &id_ty, &id, output).await?, DbPool::Sqlite(p) => sqlite::select_by_id(p, table, primary_key, &id_ty, &id, output).await?, }; row.ok_or(ApiError::NotFound) } pub async fn list( pool: &DbPool, table: &str, output: &HashMap<String, FieldType>, limit: i64, offset: i64, filters_raw: &HashMap<String, String>, ) -> Result<Value, ApiError> { let mut filters = Vec::with_capacity(filters_raw.len()); for (name, raw) in filters_raw { let ty = output .get(name) .ok_or_else(|| ApiError::Validation(format!("unknown filter field '{name}'")))?; let val = param_to_value(ty, raw)?; filters.push((name.clone(), ty.clone(), val)); } let rows = match pool { DbPool::Postgres(p) => pg::select_list(p, table, output, limit, offset, &filters).await?, DbPool::MySql(p) => mysql::select_list(p, table, output, limit, offset, &filters).await?, DbPool::Sqlite(p) => sqlite::select_list(p, table, output, limit, offset, &filters).await?, }; Ok(Value::Array(rows)) } pub async fn update( pool: &DbPool, table: &str, primary_key: &str, id_raw: &str, input: &HashMap<String, FieldType>, output: &HashMap<String, FieldType>, body: &Value, ) -> Result<Value, ApiError> { let id_ty = output.get(primary_key).cloned().unwrap_or(FieldType::Uuid); let id = param_to_value(&id_ty, id_raw)?; let body_obj = body .as_object() .ok_or_else(|| ApiError::Validation("body must be a JSON object".to_string()))?; let mut values = HashMap::new(); for (name, ty) in input { if let Some(v) = body_obj.get(name) { values.insert(name.clone(), (ty.clone(), v.clone())); } } if values.is_empty() { return Err(ApiError::Validation("no updatable fields present in body".to_string())); } let updated = match pool { DbPool::Postgres(p) => pg::update_by_id(p, table, primary_key, &id_ty, &id, &values).await?, DbPool::MySql(p) => mysql::update_by_id(p, table, primary_key, &id_ty, &id, &values).await?, DbPool::Sqlite(p) => sqlite::update_by_id(p, table, primary_key, &id_ty, &id, &values).await?, }; if !updated { return Err(ApiError::NotFound); } read(pool, table, primary_key, id_raw, output).await } pub async fn delete( pool: &DbPool, table: &str, primary_key: &str, id_raw: &str, output: &HashMap<String, FieldType>, ) -> Result<(), ApiError> { let id_ty = output.get(primary_key).cloned().unwrap_or(FieldType::Uuid); let id = param_to_value(&id_ty, id_raw)?; let deleted = match pool { DbPool::Postgres(p) => pg::delete_by_id(p, table, primary_key, &id_ty, &id).await?, DbPool::MySql(p) => mysql::delete_by_id(p, table, primary_key, &id_ty, &id).await?, DbPool::Sqlite(p) => sqlite::delete_by_id(p, table, primary_key, &id_ty, &id).await?, }; if !deleted { return Err(ApiError::NotFound); } Ok(()) }