/
aristeh
/
gui
Обзор
Документация
Войти
/
aristeh
/
gui
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
Аналитика
Безопасность
master
src/db/mod.rs
830 строк
35 KB
aristeh
колонки
11 фев 2026, 14:12
11 фев 2026, 14:12
9f2c16d
Код
Авторство
О чём код?
use rumqttc::{Client, QoS}; use std::sync::mpsc; use std::sync::mpsc::{Receiver, Sender}; // use duckdb::params_from_iter; // use serde_json::Value; //use duckdb::arrow::array::RecordBatch; //use duckdb::arrow::util::pretty::print_batches; use duckdb::{Connection, Result as DbResult}; use serde_json; use std::error::Error; use std::fs; use std::path::Path; // use base64::{engine::general_purpose::STANDARD, Engine as _}; #[allow(dead_code)] pub struct DatabaseHandler { pub conn: Connection, pub mqtt_client: Client, } use crate::client; use crate::models; #[allow(dead_code)] impl DatabaseHandler { pub fn new(db_path: &str) -> Result<DatabaseHandler, Box<dyn Error>> { let conn = Connection::open(db_path)?; let (tx, rx): (Sender<models::Message>, Receiver<models::Message>) = mpsc::channel(); let mqtt_client = client::start(&String::from("db"), tx); mqtt_client.subscribe("db/query", QoS::AtLeastOnce).unwrap(); mqtt_client .subscribe("db/operation", QoS::AtLeastOnce) .unwrap(); // Start the message processing loop in a separate thread let conn_clone = conn.try_clone()?; let mqtt_client_clone = mqtt_client.clone(); // Clone the MQTT client to move into the thread std::thread::spawn(move || { while let Ok(message) = rx.recv() { match message { models::Message::QueryRequest(query_request) => { dbg!("db Получен запрос №: {}", &query_request.query_id); let is_select = query_request .sql .trim_start() .to_uppercase() .starts_with("SELECT"); let query_result = if is_select { match Self::execute_duckdb_select_static( &conn_clone, query_request.clone(), ) { Ok(result) => result, Err(e) => models::QueryResult { query_id: query_request.query_id.clone(), tip: query_request.tip.clone(), data: None, columns: None, row_count: 0, execution_time_ms: 0.0, error_message: Some(e.to_string()), }, } } else { match Self::execute_duckdb_modify_static( &conn_clone, query_request.clone(), ) { Ok(row_count) => models::QueryResult { query_id: query_request.query_id.clone(), tip: query_request.tip.clone(), data: None, columns: None, row_count, execution_time_ms: 0.0, error_message: None, }, Err(e) => models::QueryResult { query_id: query_request.query_id.clone(), tip: query_request.tip.clone(), data: None, columns: None, row_count: 0, execution_time_ms: 0.0, error_message: Some(e.to_string()), }, } }; // Send the result back via MQTT let response_topic = query_request.response_topic; let query_id = query_request.query_id.clone(); // Store query_id separately to avoid moving let result_message = query_result; let json_result = serde_json::to_string(&result_message).unwrap(); println!("db Отправка результата запроса в топик {}", &response_topic); if let Err(e) = mqtt_client_clone.publish( &response_topic, QoS::AtLeastOnce, false, json_result, ) { println!("db Ошибка отправки результата запроса {}: {}", query_id, e); } } _ => { println!("db Получено сообщение другого типа, пропускаем"); } } } }); let _k = models::TreeNode { id: String::from("root"), parent: None, name: String::from("root"), level: 0, is_leaf: false, }; Ok(Self { conn: conn, mqtt_client: mqtt_client, }) } /// Функция для изменения данных (INSERT, UPDATE, DELETE) /// Возвращает количество затронутых строк или ошибку #[allow(dead_code)] pub fn execute_duckdb_modify( &self, db_query: models::QueryRequest, ) -> Result<usize, duckdb::Error> { println!("45 {}", db_query.sql); match &db_query.data { Some(rows) => { let mut stmt = self.conn.prepare(&db_query.sql)?; let mut total = 0; for row in rows { let row_params = duckdb::params_from_iter(row.iter().map(|s: &String| s.as_str())); total += stmt.execute(row_params)?; } Ok(total) } None => { // Handle the case where parameters is Option<serde_json::Value> match &db_query.parameters { Some(serde_json::Value::Array(arr)) => { let params: Vec<String> = arr .iter() .map(|v| v.as_str().unwrap_or("").to_string()) .collect(); let params_from_iter = duckdb::params_from_iter(params.iter().map(|s| s.as_str())); self.conn.execute(&db_query.sql, params_from_iter).unwrap(); } _ => { // If parameters is null or not an array, execute without parameters self.conn.execute(&db_query.sql, []).unwrap(); } } Ok(0) } } } /// Static version of execute_duckdb_modify that takes Connection as parameter #[allow(dead_code)] pub fn execute_duckdb_modify_static( conn: &Connection, db_query: models::QueryRequest, ) -> Result<usize, duckdb::Error> { match &db_query.data { Some(rows) => { let mut stmt = conn.prepare(&db_query.sql)?; let mut total = 0; for row in rows { let row_params = duckdb::params_from_iter(row.iter().map(|s: &String| s.as_str())); total += stmt.execute(row_params)?; } Ok(total) } None => { // Handle the case where parameters is Option<serde_json::Value> match &db_query.parameters { Some(serde_json::Value::Array(arr)) => { let params: Vec<String> = arr .iter() .map(|v| v.as_str().unwrap_or("").to_string()) .collect(); let params_from_iter = duckdb::params_from_iter(params.iter().map(|s| s.as_str())); conn.execute(&db_query.sql, params_from_iter).unwrap(); } _ => { // If parameters is null or not an array, execute without parameters conn.execute(&db_query.sql, []).unwrap(); } } Ok(0) } } } // /// Static version that executes SELECT query with parameters and returns QueryResult // #[allow(dead_code)] // pub fn execute_duckdb_select_with_params( // conn: &Connection, // db_query: crate::models::QueryRequest, // ) -> DbResult<crate::models::QueryResult> { // let start_time = std::time::Instant::now(); // let mut stmt = conn.prepare(&db_query.sql)?; // let stmt_ptr = &stmt as *const duckdb::Statement; // let params: Vec<String> = match &db_query.parameters { // Some(serde_json::Value::Array(arr)) => arr // .iter() // .map(|v| v.as_str().unwrap_or("").to_string()) // .collect(), // _ => Vec::new(), // }; // let params_from_iter = duckdb::params_from_iter(params.iter().map(|s| s.as_str())); // let mut rows = stmt.query(params_from_iter)?; // let col_count = unsafe { (*stmt_ptr).column_count() }; // let mut data = Vec::new(); // while let Some(row) = rows.next()? { // let mut row_vec = Vec::new(); // for i in 0..col_count { // let val: Result<String, _> = row.get(i); // row_vec.push(val.unwrap_or_default()); // } // data.push(row_vec); // } // let execution_time_ms = start_time.elapsed().as_secs_f64() * 1000.0; // let row_count = data.len(); // Ok(crate::models::QueryResult { // query_id: db_query.query_id, // tip: db_query.tip, // data: Some(data), // columns: None, // row_count, // execution_time_ms, // error_message: None, // }) // } /// Static version of execute_duckdb_select that takes Connection as parameter #[allow(dead_code)] pub fn execute_duckdb_select_static( conn: &Connection, db_query: crate::models::QueryRequest, ) -> DbResult<crate::models::QueryResult> { let start_time = std::time::Instant::now(); let mut stmt = conn.prepare(&db_query.sql)?; let stmt_ptr = &stmt as *const duckdb::Statement; let results = match &db_query.parameters { Some(serde_json::Value::Array(arr)) => { let params: Vec<String> = arr .iter() .map(|v| v.as_str().unwrap_or("").to_string()) .collect(); let params_from_iter = duckdb::params_from_iter(params.iter().map(|s| s.as_str())); let mut rows = stmt.query(params_from_iter)?; let col_count = unsafe { (*stmt_ptr).column_count() }; let mut results = Vec::new(); while let Some(row) = rows.next()? { let mut row_vec = Vec::new(); for i in 0..col_count { let val: Result<String, _> = row.get(i); row_vec.push(val.unwrap_or_default()); } results.push(row_vec); } results } _ => { let mut rows = stmt.query([])?; let mut results = Vec::new(); let col_count = unsafe { (*stmt_ptr).column_count() }; while let Some(row) = rows.next()? { let mut row_vec = Vec::new(); for i in 0..col_count { let val: Result<String, _> = row.get(i); row_vec.push(val.unwrap_or_default()); } results.push(row_vec); } results } }; let col_count = unsafe { (*stmt_ptr).column_count() }; let mut columns = Vec::new(); for i in 0..col_count { let col_name = unsafe { (*stmt_ptr).column_name(i).unwrap() }; println!("Столбец {}: {}", i, col_name); columns.push(col_name.to_string()); } let execution_time_ms = start_time.elapsed().as_secs_f64() * 1000.0; let row_count = results.len(); Ok(crate::models::QueryResult { query_id: db_query.query_id, tip: db_query.tip, data: Some(results), columns: Some(columns), row_count, execution_time_ms: execution_time_ms, error_message: None, }) } #[allow(dead_code)] pub fn load_csvs_from_directory( &self, dir_path: &str, ) -> Result<usize, Box<dyn std::error::Error>> { let mut loaded_count = 0; let path = Path::new(dir_path); if !path.is_dir() { return Err(format!("Path '{}' is not a directory", dir_path).into()); } println!(" путь {}", &dir_path); if let Some(name) = path.file_name() { println!("{:?}", name); } for entry_result in fs::read_dir(path)? { let entry = entry_result?; let file_path = entry.path(); print!("{}", file_path.as_os_str().to_string_lossy()); if file_path.is_file() { if let Some(ext) = file_path.extension() { if ext == "csv" { if let Some(table_osstr) = file_path.file_stem() { let compare_fields = vec!["id"]; let mut table_name = table_osstr.to_string_lossy().to_string(); if table_name.len() > 2 { table_name = table_name.chars().skip(2).collect(); } let file_path_str = file_path.to_string_lossy(); if let Err(e) = self.make_copy_from_csv_query( &table_name, &file_path_str, &compare_fields, ) { eprintln!("Ошибка при загрузке CSV '{}': {}", file_path_str, e); } else { loaded_count += 1; } } } } } } Ok(loaded_count) } /// Формирует и выполняет SQL-запрос для добавления данных из CSV-файла в указанную таблицу fn make_copy_from_csv_query( &self, table: &str, csv_file: &str, compare_fields: &[&str], ) -> Result<(), duckdb::Error> { let temp_table = "temp"; // Создаём временную таблицу из CSV. let drop_sql = format!("DROP TABLE IF EXISTS \"{}\"", temp_table); self.conn.execute(&drop_sql, [])?; let csv_escaped = csv_file.replace("'", "''"); let create_sql = format!( "CREATE TABLE {temp} AS SELECT * FROM read_csv_auto('{path}')", //temp temp = temp_table, path = csv_escaped, ); // Выполнение запроса с обработкой ошибки if let Err(e) = self.conn.execute(&create_sql, []) { eprintln!( "Ошибка при создании временной таблицы из CSV файла '{}': {}", csv_file, e ); return Err(e); } // Генерируем часть WHERE для OR: любое из полей отличается let field_unequal_conditions = compare_fields .iter() .map(|field| { format!( "m.{field} <> {temp}.{field}", field = field, temp = temp_table ) }) .collect::<Vec<_>>() .join(" OR "); let insert_new_rows_sql = format!( "INSERT OR REPLACE INTO {table} SELECT * FROM {temp} WHERE NOT EXISTS ( SELECT 1 FROM {table} m WHERE m.id = {temp}.id ) OR (EXISTS ( SELECT 1 FROM {table} m WHERE m.id = {temp}.id AND ({unequal}) ))", table = table, temp = temp_table, unequal = field_unequal_conditions, ); if let Err(e) = self.conn.execute(&insert_new_rows_sql, []) { eprintln!( "Ошибка при добавлении записей из CSV файла '{}': {}", csv_file, e ); return Err(e); } Ok(()) } #[allow(dead_code)] pub fn load_metadata(&self, filepath: &str) { use std::fs::File; use std::io::{BufRead, BufReader}; if let Ok(file) = File::open(filepath) { let reader = BufReader::new(file); for line_result in reader.lines() { if let Ok(line) = line_result { match self.conn.execute(&line, []) { Ok(updated) => println!("{} Выполнили запрос: {}", updated, &line), Err(err) => println!("update failed: {}", err), } } } } } /// Sends a QueryResult message via MQTT to the client/bdresult topic pub fn send_query_result( &self, query_result: crate::models::QueryResult, ) -> Result<(), Box<dyn Error>> { let message = crate::models::Message::QueryResult(query_result); let json_result = serde_json::to_string(&message)?; self.mqtt_client .publish("client/bdresult", QoS::AtLeastOnce, false, json_result)?; Ok(()) } /// Выполняет рекурсивный запрос для иерархических данных в DuckDB. /// На вход подается имя таблицы, соединение к базе, фильтр закрытых узлов (id закрытых веток). /// Результат — Vec<TreeNode>, формат строки: " <name> (id: <id>, level <level>)". /// Сортировка: сначала родитель, затем его дети (по вложенности, по родителю). /// #[allow(dead_code)] pub fn fetch_hierarchical_list( &self, table_name: &str, closed_ids: &[String], filter: &str, ) -> DbResult<Vec<crate::models::TreeNode>> { // Рекурсивный CTE формирует path (цепочку id) для точного порядка вывода. // ORDER BY path — гарантирует порядок: родитель, потом его дети, и т.д. // добавим фильтр по наименованию, если filter не пустой let name_filter = if !filter.is_empty() { format!(" Where name LIKE '%{}%'", filter.replace('\'', "''")) } else { String::new() }; let sql = format!( " WITH RECURSIVE tree(id, parent, name, path, lvl, has_children) AS ( SELECT id, parent, name, CAST(id AS VARCHAR), 0, (SELECT COUNT(1) FROM {table} t2 WHERE t2.parent = {table}.id) > 0 AS has_children FROM {table} WHERE (parent IS NULL OR parent = '0' OR parent = '00000000-0000-0000-0000-000000000000' ) UNION ALL SELECT t.id, t.parent, t.name, tree.path || '/' || t.id, tree.lvl+1, (SELECT COUNT(1) FROM {table} t2 WHERE t2.parent = t.id) > 0 AS has_children FROM {table} t JOIN tree ON t.parent = tree.id ) SELECT id, parent, name, lvl, path, has_children FROM tree {name_filter} ORDER BY path, name ", table = table_name, name_filter = name_filter ); //WHERE ',' || tree.path || ',' NOT LIKE '%,{closed},%' let mut stmt = self.conn.prepare(&sql)?; let mut rows = stmt.query([])?; let mut result: Vec<crate::models::TreeNode> = Vec::new(); while let Some(row) = rows.next()? { let id: String = row.get(0)?; let _parent: Option<String> = row.get(1)?; let name: String = row.get(2)?; let level: i32 = row.get(3)?; let path_str: String = row.get(4)?; let is_leaf = row.get(5)?; // Проверяем, содержит ли path_str какой-либо из closed_ids — если да, этот элемент пропускаем let is_closed = closed_ids.iter().any(|cid| path_str.contains(cid)); if is_closed { // println!( // "{} path {} [hidden due to closed_id match] длина {}", // &id, // &path_str, // path_str.len() // ); if path_str.len() > 40 { continue; } } result.push(crate::models::TreeNode { id, parent: _parent, name, level: level as usize, is_leaf: is_leaf, // Вы можете определить is_leaf отдельно при необходимости }); } println!("Количество элементов в result 2: {}", result.len()); Ok(result) } } // /// Получить данные для формы по имени. // /// Возвращает map (Vec) из блоков данных для вывода в GUI. // /// - form_name: имя формы, по которому выбирается строка из forms // pub fn get_form_blocks_for_gui(&self, form_name: &str) -> DbResult<Vec<models::FormBlockData>> { // // 1. Найти форму по имени и получить её bindings как строку JSON // let mut stmt = self.conn.prepare("SELECT bindings FROM forms WHERE name = ? LIMIT 1")?; // let mut rows = stmt.query(params_from_iter([form_name]))?; // let bindings_json: String = match rows.next()? { // Some(row) => row.get(0)?, // None => return Ok(vec![]), // нет такой формы // }; // let mut result: Vec<models::FormBlockData> = Vec::new(); // // 2. Распарсить bindings как JSON-массив // let bindings_val: Value = serde_json::from_str(&bindings_json).expect("Invalid JSON"); // let blocks = bindings_val.as_array(); // //println!("{:?}", blocks); // if let Some(blocks_arr) = blocks { // for item in blocks_arr { // // Каждый item – массив: [запрос, массив-полей, параметры, объект с именем блока] // // Используем let binding для временного массива, чтобы избежать borrow error // let arr_owned: Vec<_>; // let arr: &[Value] = if let Some(arr_ref) = item.as_array() { // arr_ref // } else { // arr_owned = Vec::new(); // &arr_owned // }; // // if arr.len() < 4 { // // continue; // // } // let query = arr[0].as_str().unwrap_or(""); // let field_names_val = &arr[1]; // массив строк — имена полей результата // println!("field_names_val {:?}", field_names_val); // let params_val = &arr[2]; // параметры запроса, как объект // println!("params_val {:?}", params_val); // let block_name = arr[3] // .get("name") // .and_then(|v| v.as_str()) // .unwrap_or("block"); // // 3. Подготовить параметры запроса // // params_val – JSON-объект вида { "param": "value", ... } // let mut param_vec: Vec<String> = vec![]; // if let Some(obj) = params_val.as_object() { // for (_k, v) in obj.iter() { // // Поддерживаем только строковые и числовые значения // if let Some(s) = v.as_str() { // param_vec.push(s.to_string()); // } else if let Some(n) = v.as_i64() { // param_vec.push(n.to_string()); // } else if let Some(f) = v.as_f64() { // param_vec.push(f.to_string()); // } else { // param_vec.push(v.to_string()); // } // } // } // println!("param_vec {:?}", param_vec); // // 4. Выполнить запрос и забрать данные // let mut result_table: Vec<Vec<String>> = Vec::new(); // // первая строка — имена полей: // let field_names: Vec<String> = field_names_val // .as_array() // .unwrap_or(&vec![]) // .iter() // .filter_map(|v| v.as_str().map(|s| s.to_string())) // .collect(); // result_table.push(field_names.clone()); // println!("paraquerym_vec {:?}", query); // let mut stmt = self.conn.prepare(query)?; // let params = params_from_iter(param_vec.iter().map(|s| s.as_str())); // println!("params {:?}", params); // let mut rows = stmt.query(params)?; // while let Some(row) = rows.next()? { // let mut rowdata: Vec<String> = Vec::new(); // for i in 0..field_names.len() { // let value: Result<String, _> = row.get(i); // match value { // Ok(s) => rowdata.push(s), // Err(_) => rowdata.push("".to_string()), // } // } // result_table.push(rowdata); // } // result.push(models::FormBlockData{ // block_name: block_name.to_string(), // table: result_table, // }); // } // } // println!("result {:?} {:?}", &result[0].table, &result[0].block_name); // Ok(result) // } // // Функция для вставки BLOB данных в таблицу // fn insert_blob_data( // &self, // table: &str, // id: i32, // blob: Vec<u8>, // ) -> Result<usize, String> { // // Пример запроса: INSERT INTO images (id, data) VALUES (?, ?) // let query = format!("INSERT INTO {} (id, data) VALUES (?, ?)", table); // let params = vec![]; // let data = Some(vec![vec![ // id.to_string(), // // Кодируем blob в base64 строку (duckdb не поддерживает прямую вставку Vec<u8> через API) // STANDARD.encode(&blob), // ]]); // let db_query = models::QueryRequest { // query_id: String::from("df"), // sql: query, // parameters:params, // data:data, // timestamp: chrono::Utc::now().timestamp_millis(), // response_topic: String::from("db"), // }; // // Обрабатываем различие в типах Result: execute_duckdb_modify возвращает Result<_, duckdb::Error> // match self.execute_duckdb_modify(db_query) { // Ok(res) => Ok(res), // Err(e) => Err(format!("DuckDB modify error: {}", e)), // } // } // // Функция для извлечения BLOB данных из таблицы по id // fn get_blob_data( // &self, // table: &str, // id: i32, // ) -> Result<Option<Vec<u8>>, String> { // let query = format!("SELECT data FROM {} WHERE id = {}", table, id); // // Используем существующую функцию execute_duckdb_select, она возвращает Result<Vec<Vec<String>>, String> // match execute_duckdb_select(self.conn, DbQuery { // query, // params: vec![], // data: None, // }) { // Ok(rows) => { // if let Some(row) = rows.into_iter().next() { // // row[0] — строка base64 // let decoded = STANDARD // .decode(&row[0]) // .map_err(|e| format!("Base64 decode error: {}", e))?; // Ok(Some(decoded)) // } else { // Ok(None) // } // } // Err(e) => Err(format!("Select error: {}", e)), // } // } // /// Вставляет изображение из файла в базу данных, используя функцию insert_blob_data. // /// // /// # Аргументы // /// * `conn` - соединение с базой данных DuckDB. // /// * `table` - имя таблицы. // /// * `id` - идентификатор для вставки. // /// * `file_path` - путь к файлу изображения. // /// // /// # Возвращает // /// * `Result<usize, String>` — число вставленных строк или текст ошибки. // fn insert_image_from_file( // &self, // String, // id: i32, // file_path: &str, // ) -> Result<usize, String> { // use std::fs; // // Считываем все байты из файла // let blob = fs::read(file_path) // .map_err(|e| format!("Ошибка чтения файла: {}", e))?; // // Вставляем blob в базу данных // insert_blob_data(self.conn, table, id, blob) // } // /// Преобразует результат get_blob_data в строку, пригодную для передачи изображения в Slint. // /// В Slint изображение можно передавать как строку: "data:image/<format>;base64,<base64string>" // /// Обычно format = "png" или "jpeg", если вы это можете определить. // /// // /// # Аргументы // /// * `blob_result` - результат get_blob_data (Option<Vec<u8>>) // /// * `format` - расширение формата ("png", "jpeg" и т.д.) // /// // /// # Возвращает // /// * `Option<String>` — строка для передачи в Slint или None если blob отсутствует. // pub fn prepare_image_for_slint(blob_result: Option<Vec<u8>>, format: &str) -> Option<String> { // match blob_result { // Some(blob) => { // // Кодируем содержимое в base64 // let encoded = base64::engine::general_purpose::STANDARD.encode(&blob); // let slint_string = format!("data:image/{};base64,{}", format, encoded); // Some(slint_string) // } // None => None, // } // } // #[derive(Debug, Clone)] // pub struct Icon { // pub id: i64, // pub image_data: Vec<u8>, // pub width: u32, // pub height: u32, // } // /// Конвертирует иконку из базы данных в slint::Image // pub fn icon_to_slint_image(icon: &Icon) -> Result<Image, anyhow::Error> { // let buffer = SharedPixelBuffer::<Rgba8Pixel>::clone_from_slice( // &icon.image_data, // icon.width, // icon.height, // ); // Ok(Image::from_rgba8(buffer)) // } // } // use duckdb::{params, Connection, Result, ToSql}; // /// Универсальная функция для выполнения SELECT-запроса с параметрами. // /// Возвращает вектор имен (String) из первого столбца для примера. // fn query_data(conn: &Connection, sql: &str, params: &[&dyn ToSql]) -> Result<Vec<String>> { // let mut stmt = conn.prepare(sql)?; // // Выполняем запрос, передавая срез параметров // let mut rows = stmt.query(params)?; // let mut results = Vec::new(); // while let Some(row) = rows.next()? { // // В данном примере мы берем только первую колонку типа String // let name: String = row.get(0)?; // results.push(name); // } // Ok(results) // } // fn main() -> Result<()> { // // 1. Открываем соединение (в памяти) // let conn = Connection::open_in_memory()?; // // 2. Создаем таблицу для теста // conn.execute( // "CREATE TABLE users (id INTEGER, name TEXT, age INTEGER)", // [], // )?; // // 3. Вставляем данные // conn.execute( // "INSERT INTO users VALUES (1, 'Alice', 30), (2, 'Bob', 25), (3, 'Charlie', 35)", // [], // )?; // // --- ПРИМЕР ИСПОЛЬЗОВАНИЯ ФУНКЦИИ --- // let min_age = 28; // let name_pattern = "A%"; // // SQL запрос с плейсхолдерами ? // let sql = "SELECT name FROM users WHERE age > ? AND name LIKE ?"; // // Вызов функции с использованием макроса params! // let found_names = query_data( // &conn, // sql, // params![min_age, name_pattern] // )?; // println!("Найденные пользователи: {:?}", found_names); // Ok(()) // }