/
NovAl
/
rust_rules
Обзор
Документация
Войти
/
NovAl
/
rust_rules
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
master
bundle.txt
1 492 строки
52 KB
Novikov Aleksandr
Первый коммит: проект без target и бинарников
21 июл 2026, 22:09
21 июл 2026, 22:09
f6008f1
Код
Авторство
О чём код?
// === crates/engine/src/dispatcher.rs === use std::collections::{HashMap, HashSet}; use std::sync::Arc; use tokio::sync::RwLock; use tokio::sync::mpsc; use rumqttc::{AsyncClient, QoS}; pub type Event = (String, String); // (topic, payload) // НОВАЯ ВЕРСИЯ: для сценариев через каналы pub struct ChannelDispatcher { subscribers: Arc<RwLock<HashMap<String, Vec<mpsc::UnboundedSender<Event>>>>>, alias_map: Arc<HashMap<String, String>>, mqtt_client: AsyncClient, all_topics: Arc<HashSet<String>>, command_suffix: String, all_subscribers: Arc<RwLock<Vec<mpsc::UnboundedSender<Event>>>>, } impl ChannelDispatcher { pub fn new( aliases: HashMap<String, String>, mqtt_client: AsyncClient, all_topics: HashSet<String>, command_suffix: String ) -> Self { Self { subscribers: Arc::new(RwLock::new(HashMap::new())), alias_map: Arc::new(aliases), mqtt_client, all_topics: Arc::new(all_topics), command_suffix, all_subscribers: Arc::new(RwLock::new(Vec::new())), } } /// Подписка на топик — возвращает Receiver pub async fn subscribe_raw(&self, topic: String) -> mpsc::UnboundedReceiver<Event> { let (tx, rx) = mpsc::unbounded_channel(); let mut subs = self.subscribers.write().await; subs.entry(topic).or_insert_with(Vec::new).push(tx); rx } /// Подписка на топик (поддерживает алиасы) /// Возвращает Receiver для чтения событий pub async fn subscribe(&self, key: &str) -> mpsc::UnboundedReceiver<Event> { let real_topic = self.resolve(key); self.subscribe_raw(real_topic).await } /// Отправка события всем подписчикам pub async fn dispatch_raw(&self, topic: String, payload: String) { // создаем переменную которая содержит кортеж топик и значение топика let event: Event = (topic.clone(), payload); //поле subscribers содержит HashMap таблицу где ключ это топик а значение вектор писателей (tx) тех кто подписанна этот топик // например скрипт, но сейчас мы поняли что нельзя передавать в rune rx // senders содержит блок кода (что бы как можно быстрее снять блокировку) в котром переменная subs держит блокирову на четение //другие то же могут читать но никто не может писать // subs.get получаем из таблицы вектор писателей для данного топика // в итоге senders это вектор писателей тех кто подписан на данный топика // ТОЧЕЧНАЯ РАССЫЛКА СОДЕРЖИТ ПИСАТЕЛЕЙ ДЛЯ КОНКРЕТНОГО ТОПИКА (ПИСАТЕЛЬ ОТПРАВЛЯЕТ ЧИТАТЕЛЮ В ДРУГОЙ ПОТОК Event) let senders = { let subs = self.subscribers.read().await; subs.get(&topic).cloned() }; // проверяем является ли senders элементом Some перечисления Option //если да то запускаем цикл по senders и отправляем каждому подписчику копию event if let Some(senders) = senders { for tx in senders { let _ = tx.send(event.clone()); } } //поле all_subscribers содержит вектор писателей (tx), здесь уже нет таблицы с топиками потому что //нам нужно что бы писатели отправили своим читателям (rx) сообщения от любого топика //all_subs держит блокировку на четение // ОБЩАЯ РАССЫЛКА СОДЕРЖИТ ПИСАТЕЛЕЙ ЧИТАТЕЛЯМ КОТОРЫХ НУЖНЫ СООБЩЕНИЯ ОТ ВСЕХ ТОПИКОВ (ПИСАТЕЛЬ ОТПРАВЛЯЕТ ЧИТАТЕЛЮ В ДРУГОЙ ПОТОК Event с любым топиком) let all_subs = self.all_subscribers.read().await; //цикл пробегает по вектору содержащимуся в all_subscribers и отправляет каждому копю event for tx in all_subs.iter() { let _ = tx.send(event.clone()); } } /// Очистка мертвых каналов pub async fn cleanup(&self) { let mut subs = self.subscribers.write().await; for senders in subs.values_mut() { senders.retain(|tx| !tx.is_closed()); } subs.retain(|_, senders| !senders.is_empty()); } pub fn resolve(&self, key: &str) -> String { self.alias_map .get(key) .cloned() .unwrap_or_else(|| key.to_string()) } /// Отправить команду (через MQTT) pub async fn set_raw(&self, key: &str, payload: String) { let real_topic = key.to_string(); // 1. Внутренняя рассылка (подписчикам в системе) self.dispatch_raw(real_topic.clone(), payload.clone()).await; // 2. Отправка MQTT команды устройству let _ = self.mqtt_client .publish(real_topic, QoS::AtLeastOnce, false, payload) .await; } /// Отправить команду (через MQTT Безопасная отправка (только топики из конфига)) pub async fn set(&self, key: &str, payload: String) { let real_topic = self.resolve(key); if !self.all_topics.contains(&real_topic) { eprintln!("⚠️ Ошибка: топик '{}' не найден в конфиге", real_topic); return; } let command_topic = format!("{}{}", real_topic, self.command_suffix); // 3. Логируем, что и куда отправляем (полезно для отладки). println!("📤 [CMD] '{}' -> '{}': '{}'", key, command_topic, payload); self.set_raw(&command_topic, payload).await; } // при вызове этого метода создается mpsc канал с неограниченным буфером //переменная all_subs блокирует поле all_subscribers на запись, пока блокировка стоит никто не может ни читать ни изменять поле all_subscribers // all_subs.push(tx); добавляет в вектор нового писателя // в конце метол возвращает читателя pub async fn subscribe_all(&self) -> mpsc::UnboundedReceiver<Event> { let (tx, rx) = mpsc::unbounded_channel(); let mut all_subs = self.all_subscribers.write().await; all_subs.push(tx); rx } /// Доступ к списку всех топиков (для валидации) pub fn all_topics(&self) -> &HashSet<String> { &self.all_topics } /// Доступ к суффиксу команд pub fn command_suffix(&self) -> &str { &self.command_suffix } } impl Clone for ChannelDispatcher { fn clone(&self) -> Self { Self { subscribers: self.subscribers.clone(), alias_map: self.alias_map.clone(), mqtt_client: self.mqtt_client.clone(), all_topics: self.all_topics.clone(), command_suffix: self.command_suffix.clone(), all_subscribers: self.all_subscribers.clone(), } } } impl Default for ChannelDispatcher { fn default() -> Self { // Для default нужен фейковый клиент let (mqtt_client, _) = AsyncClient::new( rumqttc::MqttOptions::new("default", "localhost", 1883), 10 ); Self::new(HashMap::new(), mqtt_client, HashSet::new(), "".to_string()) } } // === crates/engine/src/lib.rs === use rumqttc::{AsyncClient, Event, EventLoop, Incoming, MqttOptions, QoS}; use std::time::Duration; use rust_rules_config::BrokerConfig; mod dispatcher; pub use dispatcher::{ChannelDispatcher, Event as SmartHomeEvent}; #[derive(Debug, Clone)] pub struct MqttEvent { pub topic: String, pub payload: String, } pub struct MqttEngine { client: AsyncClient, eventloop: EventLoop, dispatcher: Option<ChannelDispatcher>, } impl MqttEngine { pub fn new(settings: &BrokerConfig) -> Self { let mut mqttoptions = MqttOptions::new("rust_rules", &settings.host, settings.port); mqttoptions.set_keep_alive(Duration::from_secs(5)); let (client, eventloop) = AsyncClient::new(mqttoptions, 10); Self { client, eventloop, dispatcher: None, } } pub async fn subscribe(&self, topic: &str) { if let Err(e) = self.client .subscribe(topic, QoS::AtLeastOnce) .await { println!("Ошибка подписки на топик {}: {:?}", topic, e); } } pub async fn poll(&mut self) { match self.eventloop.poll().await { Ok(Event::Incoming(Incoming::Publish(p))) => { let payload = String::from_utf8_lossy(&p.payload).to_string(); if let Some(dispatcher) = &self.dispatcher { dispatcher.dispatch_raw(p.topic, payload).await; } }, Ok(_) => {}, Err(e) => eprintln!("MQTT ошибка {}", e), } } pub async fn run(mut self) { loop { self.poll().await; } } pub fn client(&self) -> &AsyncClient { &self.client } pub fn set_dispatcher(&mut self, dispatcher: ChannelDispatcher){ self.dispatcher = Some(dispatcher); } } // === crates/rust_rules/src/main.rs === mod state; mod mqtt_worker; use rust_rules_config::{Config, TopicTypes}; use engine::{MqttEngine, ChannelDispatcher}; use state::State; use std::collections::HashSet; use std::path::PathBuf; use mqtt_worker::{MqttWorker, WorkerCommand}; #[tokio::main] async fn main() { println!("Запуск системы автоматизации..."); // ============================================ // 1. ЗАГРУЗКА КОНФИГА // ============================================ let config: Config = Config::load_config("rust_rules_config.toml"); println!("Конфиг успешно загружен!"); // ============================================ // 2. СОЗДАНИЕ STATE (пустой, честный) // ============================================ let state = State::new(); println!("✅ State создан (пустой, без дефолтных значений)"); // ============================================ // 3. СОЗДАНИЕ DISPATCHER И MQTT ENGINE // ============================================ let aliases = config.aliases(); let mut mqtt_engine = MqttEngine::new(&config.broker); let mqtt_client = mqtt_engine.client().clone(); let topic_types = TopicTypes::from_config(&config); let all_topics_set: HashSet<String> = topic_types.all_topics().into_iter().collect(); let command_suffix = config.broker.command_suffix.clone(); let dispatcher = ChannelDispatcher::new(aliases, mqtt_client, all_topics_set, command_suffix); mqtt_engine.set_dispatcher(dispatcher.clone()); println!("✅ Dispatcher создан"); // ============================================ // 4. ПОДПИСКА НА ВСЕ ТОПИКИ (MQTT) // ============================================ let all_topics = config.all_topics(); println!("📡 Всего топиков для подписки: {}", all_topics.len()); for topic in &all_topics { mqtt_engine.subscribe(topic).await; println!(" ✅ Подписались: {}", topic); } // ============================================ // 5. СОЗДАНИЕ КАНАЛОВ ДЛЯ MQTT WORKER // ============================================ let (worker_tx, worker_rx) = tokio::sync::mpsc::unbounded_channel(); let (mqtt_out_tx, mut mqtt_out_rx) = tokio::sync::mpsc::unbounded_channel(); // ============================================ // 6. МОСТ: DISPATCHER → STATE & WORKER // (создаём ДО запуска двигателя — чтобы не потерять retained!) // ============================================ let state_for_mqtt = state.clone(); let mut rx = dispatcher.subscribe_all().await; let worker_tx_for_bridge = worker_tx.clone(); tokio::spawn(async move { while let Some((topic, payload)) = rx.recv().await { // А) Обновляем внутренний стейт let type_name = topic_types.get(&topic); match type_name { "bool" => { let value = payload == "1" || payload == "true"; state_for_mqtt.mqtt.set_bool(&topic, value).await; } "u8" | "i8" | "u16" | "i16" | "u32" | "i32" => { let value = payload.parse().unwrap_or(0); state_for_mqtt.mqtt.set_i32(&topic, value).await; } "f32" | "f64" => { let value = payload.parse().unwrap_or(0.0); state_for_mqtt.mqtt.set_f64(&topic, value).await; } _ => { state_for_mqtt.mqtt.set_string(&topic, payload.clone()).await; } } // Б) Прокидываем событие в Rune-воркер let _ = worker_tx_for_bridge.send(WorkerCommand::TriggerEvent { topic: topic, payload: payload, }); } }); println!("✅ Мост Dispatcher → State & MqttWorker запущен"); // ============================================ // 7. ЗАПУСК MQTT WORKER (изолированный актёр) // ============================================ let worker = MqttWorker::new( worker_rx, mqtt_out_tx, state.clone(), // ← State для доступа из скриптов dispatcher.clone(), // ← Dispatcher для resolve() алиасов ); tokio::spawn(async move { worker.run().await; }); println!("✅ MqttWorker запущен в изолированном потоке"); // ============================================ // 8. МОСТ ОБРАТНОЙ СВЯЗИ: Worker → Dispatcher // ============================================ let dispatcher_for_out = dispatcher.clone(); tokio::spawn(async move { while let Some(msg) = mqtt_out_rx.recv().await { dispatcher_for_out.set(&msg.topic, msg.payload).await; } }); println!("✅ Мост MqttWorker → Dispatcher (Обратная связь) запущен"); // ============================================ // 9. ЗАГРУЗКА СКРИПТОВ Rune // ============================================ let scripts_dir = PathBuf::from(&config.scripts.directory); worker_tx.send(WorkerCommand::LoadScriptsFromDir { dir_path: scripts_dir }) .expect("Не удалось отправить команду LoadScriptsFromDir в воркер"); println!("📜 Команда загрузки скриптов отправлена воркеру"); // ============================================ // 10. ЗАПУСК MQTT ДВИГАТЕЛЯ // (только теперь — всё готово!) // ============================================ tokio::spawn(async move { mqtt_engine.run().await; }); println!("✅ MQTT двигатель запущен"); // ============================================ // 11. ИМИТАЦИЯ ДЛЯ ТЕСТИРОВАНИЯ (ВРЕМЕННО) // ============================================ let worker_tx_test = worker_tx.clone(); tokio::spawn(async move { tokio::time::sleep(tokio::time::Duration::from_secs(1)).await; println!("\n--- СТАРТ ТЕСТОВЫХ СОБЫТИЙ ---"); let _ = worker_tx_test.send(WorkerCommand::TriggerEvent { topic: "/devices/wb-msw-v4_139/controls/Current Motion".to_string(), payload: "1".to_string(), }); tokio::time::sleep(tokio::time::Duration::from_millis(500)).await; let _ = worker_tx_test.send(WorkerCommand::TriggerEvent { topic: "/devices/wb-msw-v4_139/controls/Temperature".to_string(), payload: "26.4".to_string(), }); tokio::time::sleep(tokio::time::Duration::from_millis(500)).await; let _ = worker_tx_test.send(WorkerCommand::TriggerEvent { topic: "/devices/wb-msw-v4_139/controls/Current Motion".to_string(), payload: "0".to_string(), }); }); // Ожидаем сигнала закрытия (Ctrl+C) tokio::signal::ctrl_c().await.unwrap(); println!("Завершение работы..."); } // === crates/rust_rules/src/mqtt_worker.rs === use std::collections::HashMap; use std::path::{Path, PathBuf}; use std::sync::Arc; use tokio::sync::mpsc; use rune::{Unit, Source, Sources}; use rune_macros::{FromValue, Any}; use crate::state::State; use engine::ChannelDispatcher; /// Структура правила из Rune #[derive(FromValue)] struct RuleManifest { when: Vec<String>, then: rune::runtime::Function, } /// Умный контейнер события для UX мечты в скриптах #[derive(Any, Clone)] pub struct MqttEvent { #[rune(get)] pub topic: String, #[rune(get)] pub payload: String, mqtt_sender: mpsc::UnboundedSender<OutgoingMqttMessage>, // ДОБАВЛЯЕМ State и Dispatcher для доступа из скриптов state: State, dispatcher: ChannelDispatcher, script_name: String, // имя текущего скрипта (для var-доступа) } impl MqttEvent { pub fn new( topic: String, payload: String, mqtt_sender: mpsc::UnboundedSender<OutgoingMqttMessage>, state: State, dispatcher: ChannelDispatcher, script_name: String, ) -> Self { Self { topic, payload, mqtt_sender, state, dispatcher, script_name, } } pub fn as_bool(&self) -> bool { self.payload == "1" || self.payload == "true" || self.payload == "ON" } pub fn as_int(&self) -> i64 { self.payload.parse::<i64>().unwrap_or(0) } pub fn as_float(&self) -> f64 { self.payload.parse::<f64>().unwrap_or(0.0) } pub fn publish(&self, topic: String, payload: String) { let _ = self.mqtt_sender.send(OutgoingMqttMessage { topic, payload, }); } // ======================================================== // НОВЫЕ МЕТОДЫ ДЛЯ ДОСТУПА К STATE ИЗ RUNE-СКРИПТОВ // ======================================================== /// Прочитать bool значение из MQTT-состояния (с алиасом!) pub fn mqtt_get_bool(&self, key: String) -> Option<bool> { let topic = self.dispatcher.resolve(&key); println!("🔍 [Rune] mqtt_get_bool: '{}' -> '{}'", key, topic); self.state.mqtt.smart_get_bool(&topic) } /// Прочитать f64 значение из MQTT-состояния pub fn mqtt_get_f64(&self, key: String) -> Option<f64> { let topic = self.dispatcher.resolve(&key); println!("🔍 [Rune] mqtt_get_f64: '{}' -> '{}'", key, topic); self.state.mqtt.smart_get_f64(&topic) } /// Прочитать i32 значение из MQTT-состояния pub fn mqtt_get_i32(&self, key: String) -> Option<i32> { let topic = self.dispatcher.resolve(&key); println!("🔍 [Rune] mqtt_get_i32: '{}' -> '{}'", key, topic); self.state.mqtt.smart_get_i32(&topic) } /// Прочитать String значение из MQTT-состояния pub fn mqtt_get_string(&self, key: String) -> Option<String> { let topic = self.dispatcher.resolve(&key); println!("🔍 [Rune] mqtt_get_string: '{}' -> '{}'", key, topic); self.state.mqtt.smart_get_string(&topic) } /// Записать bool переменную ТЕКУЩЕГО скрипта pub fn var_set_bool(&self, name: String, value: bool) { if let Some(table) = self.state.var.get_sync(&self.script_name) { if let Ok(mut guard) = table.try_lock() { guard.set_bool(&name, value); println!("📝 [Rune] var_set_bool: {}.{} = {}", self.script_name, name, value); } } } /// Прочитать bool переменную ТЕКУЩЕГО скрипта pub fn var_get_bool(&self, name: String) -> Option<bool> { let table = self.state.var.get_sync(&self.script_name)?; let guard = table.try_lock().ok()?; guard.get_bool(&name) } /// Записать i32 переменную pub fn var_set_i32(&self, name: String, value: i32) { if let Some(table) = self.state.var.get_sync(&self.script_name) { if let Ok(mut guard) = table.try_lock() { guard.set_i32(&name, value); println!("📝 [Rune] var_set_i32: {}.{} = {}", self.script_name, name, value); } } } /// Прочитать i32 переменную pub fn var_get_i32(&self, name: String) -> Option<i32> { let table = self.state.var.get_sync(&self.script_name)?; let guard = table.try_lock().ok()?; guard.get_i32(&name) } /// Записать f64 переменную pub fn var_set_f64(&self, name: String, value: f64) { if let Some(table) = self.state.var.get_sync(&self.script_name) { if let Ok(mut guard) = table.try_lock() { guard.set_f64(&name, value); println!("📝 [Rune] var_set_f64: {}.{} = {}", self.script_name, name, value); } } } /// Прочитать f64 переменную pub fn var_get_f64(&self, name: String) -> Option<f64> { let table = self.state.var.get_sync(&self.script_name)?; let guard = table.try_lock().ok()?; guard.get_f64(&name) } /// Записать String переменную pub fn var_set_string(&self, name: String, value: String) { if let Some(table) = self.state.var.get_sync(&self.script_name) { if let Ok(mut guard) = table.try_lock() { guard.set_string(&name, value.clone()); println!("📝 [Rune] var_set_string: {}.{} = {}", self.script_name, name, value); } } } /// Прочитать String переменную pub fn var_get_string(&self, name: String) -> Option<String> { let table = self.state.var.get_sync(&self.script_name)?; let guard = table.try_lock().ok()?; guard.get_string(&name) } } /// Команда на отправку в сеть MQTT #[derive(Debug)] pub struct OutgoingMqttMessage { pub topic: String, pub payload: String, } /// Информация о зарегистрированном правиле #[derive(Clone)] struct CompiledRule { func_hash: rune::Hash, unit: Arc<Unit>, script_name: String, } /// Наш изолированный Актёр pub struct MqttWorker { receiver: mpsc::UnboundedReceiver<WorkerCommand>, subscriptions: HashMap<String, Vec<CompiledRule>>, runtime: Option<Arc<rune::runtime::RuntimeContext>>, mqtt_sender: mpsc::UnboundedSender<OutgoingMqttMessage>, state: State, dispatcher: ChannelDispatcher, } #[derive(Debug)] pub enum WorkerCommand { LoadScriptsFromDir { dir_path: PathBuf }, TriggerEvent { topic: String, payload: String }, } impl MqttWorker { pub fn new( receiver: mpsc::UnboundedReceiver<WorkerCommand>, mqtt_sender: mpsc::UnboundedSender<OutgoingMqttMessage>, state: State, dispatcher: ChannelDispatcher, ) -> Self { Self { receiver, subscriptions: HashMap::new(), runtime: None, mqtt_sender, state, dispatcher, } } /// Создаёт кастомный модуль Rune и регистрирует ВСЕ методы fn create_custom_module(&self) -> Result<rune::Module, rune::ContextError> { let mut module = rune::Module::new(); // Регистрируем тип события и ВСЕ его методы module.ty::<MqttEvent>()?; // Старые методы module.associated_function("as_bool", MqttEvent::as_bool)?; module.associated_function("as_int", MqttEvent::as_int)?; module.associated_function("as_float", MqttEvent::as_float)?; module.associated_function("publish", MqttEvent::publish)?; // НОВЫЕ методы для чтения MQTT-состояния module.associated_function("mqtt_get_bool", MqttEvent::mqtt_get_bool)?; module.associated_function("mqtt_get_f64", MqttEvent::mqtt_get_f64)?; module.associated_function("mqtt_get_i32", MqttEvent::mqtt_get_i32)?; module.associated_function("mqtt_get_string", MqttEvent::mqtt_get_string)?; // НОВЫЕ методы для работы с переменными скрипта module.associated_function("var_set_bool", MqttEvent::var_set_bool)?; module.associated_function("var_get_bool", MqttEvent::var_get_bool)?; module.associated_function("var_set_i32", MqttEvent::var_set_i32)?; module.associated_function("var_get_i32", MqttEvent::var_get_i32)?; module.associated_function("var_set_f64", MqttEvent::var_set_f64)?; module.associated_function("var_get_f64", MqttEvent::var_get_f64)?; module.associated_function("var_set_string", MqttEvent::var_set_string)?; module.associated_function("var_get_string", MqttEvent::var_get_string)?; Ok(module) } /// Сканируем папку, собираем .rn файлы и компилируем fn compile_from_directory(&mut self, dir: &Path) -> Result<(), Box<dyn std::error::Error>> { println!("MqttWorker: Сканируем папку со скриптами: {:?}", dir); if !dir.exists() { return Err(format!("Папка со скриптами не найдена: {:?}", dir).into()); } self.subscriptions.clear(); let mut context = rune_modules::default_context()?; context.install(self.create_custom_module()?)?; let arc_runtime = Arc::new(context.runtime()?); self.runtime = Some(arc_runtime.clone()); for entry in std::fs::read_dir(dir)? { let entry = entry?; let file_path = entry.path(); if file_path.is_file() && file_path.extension().and_then(|s| s.to_str()) == Some("rn") { let file_name = file_path.file_name().unwrap().to_str().unwrap(); let script_name = file_name.to_string(); println!(" -> Компиляция файла: {}", file_name); // Регистрируем таблицу переменных для этого скрипта self.state.var.register_sync(&script_name); println!(" 📝 Зарегистрирована таблица переменных для '{}'", script_name); let code = std::fs::read_to_string(&file_path)?; let mut sources = Sources::new(); sources.insert(Source::memory(code)?)?; let mut diagnostics = rune::Diagnostics::new(); let unit = match rune::prepare(&mut sources) .with_context(&context) .with_diagnostics(&mut diagnostics) .build() { Ok(u) => Arc::new(u), Err(_e) => { eprintln!("❌ Ошибка компиляции в файле {}:", file_name); let mut writer = rune::termcolor::StandardStream::stderr( rune::termcolor::ColorChoice::Always ); let _ = diagnostics.emit(&mut writer, &sources); continue; } }; let mut vm = rune::Vm::new(arc_runtime.clone(), unit.clone()); let main_hash = rune::Hash::type_hash(["main"]); let raw_value: rune::runtime::Value = vm.call(main_hash, ())?; let manifests: Vec<RuleManifest> = rune::from_value(raw_value)?; for manifest in manifests { let func_hash = manifest.then.type_hash(); let rule = CompiledRule { func_hash, unit: unit.clone(), script_name: script_name.clone(), }; for topic in manifest.when { let resolved_topic = self.dispatcher.resolve(&topic); if resolved_topic != topic { println!(" 🔄 Алиас: '{}' -> '{}'", topic, resolved_topic); } self.subscriptions .entry(resolved_topic) .or_insert_with(Vec::new) .push(rule.clone()); } } } } println!("MqttWorker: Сборка завершена. Подписок на топики: {}", self.subscriptions.len()); Ok(()) } pub async fn run(mut self) { println!("MqttWorker: запущен и ожидает команды..."); while let Some(command) = self.receiver.recv().await { match command { WorkerCommand::LoadScriptsFromDir { dir_path } => { if let Err(e) = self.compile_from_directory(&dir_path) { eprintln!("MqttWorker: ❌ Ошибка сборки директории скриптов: {}", e); } } WorkerCommand::TriggerEvent { topic, payload } => { if let Some(rule_list) = self.subscriptions.get(&topic) { println!( "MqttWorker: найдены подписки ({}) для топика {}", rule_list.len(), topic ); for rule in rule_list { let event = MqttEvent::new( topic.clone(), payload.clone(), self.mqtt_sender.clone(), self.state.clone(), self.dispatcher.clone(), rule.script_name.clone(), ); if let Err(e) = self.execute_stateless( rule.func_hash, rule.unit.clone(), event.clone() ) { eprintln!("MqttWorker: ❌ Ошибка выполнения правила: {}", e); } } } } } } println!("MqttWorker: канал закрыт, работа завершена."); } /// 100% Stateless выполнение на стеке потока воркера fn execute_stateless( &self, func_hash: rune::Hash, unit: Arc<Unit>, event: MqttEvent ) -> Result<(), Box<dyn std::error::Error>> { if let Some(runtime) = &self.runtime { let mut vm = rune::Vm::new(runtime.clone(), unit); vm.call(func_hash, (event,))?; } Ok(()) } } // === crates/rust_rules/src/state.rs === use std::collections::HashMap; use std::sync::Arc; use tokio::sync::{RwLock, Mutex}; // ============================================================ // ГЛАВНОЕ ХРАНИЛИЩЕ СОСТОЯНИЯ // ============================================================ #[derive(Clone)] pub struct State { pub mqtt: MqttTables, pub var: ScriptVarManager, } impl State { pub fn new() -> Self { Self { mqtt: MqttTables::new(), var: ScriptVarManager::new(), } } } // ============================================================ // MQTT-ТАБЛИЦЫ (состояние физических устройств) // ============================================================ #[derive(Clone)] pub struct MqttTables { pub i32: Arc<RwLock<HashMap<String, i32>>>, pub f64: Arc<RwLock<HashMap<String, f64>>>, pub bool: Arc<RwLock<HashMap<String, bool>>>, pub str: Arc<RwLock<HashMap<String, String>>>, } impl MqttTables { pub fn new() -> Self { Self { i32: Arc::new(RwLock::new(HashMap::new())), f64: Arc::new(RwLock::new(HashMap::new())), bool: Arc::new(RwLock::new(HashMap::new())), str: Arc::new(RwLock::new(HashMap::new())), } } // ======================================================== // Асинхронные методы (для MQTT-моста и Web-интерфейса) // ======================================================== pub async fn get_i32(&self, topic: &str) -> Option<i32> { let map = self.i32.read().await; map.get(topic).copied() } pub async fn get_f64(&self, topic: &str) -> Option<f64> { let map = self.f64.read().await; map.get(topic).copied() } pub async fn get_bool(&self, topic: &str) -> Option<bool> { let map = self.bool.read().await; map.get(topic).copied() } pub async fn get_string(&self, topic: &str) -> Option<String> { let map = self.str.read().await; map.get(topic).cloned() } pub async fn set_i32(&self, topic: &str, value: i32) { let mut map = self.i32.write().await; map.insert(topic.to_string(), value); } pub async fn set_f64(&self, topic: &str, value: f64) { let mut map = self.f64.write().await; map.insert(topic.to_string(), value); } pub async fn set_bool(&self, topic: &str, value: bool) { let mut map = self.bool.write().await; map.insert(topic.to_string(), value); } pub async fn set_string(&self, topic: &str, value: String) { let mut map = self.str.write().await; map.insert(topic.to_string(), value); } // ======================================================== // СИНХРОННЫЕ МЕТОДЫ ДЛЯ RUNE (умное чтение с fallback) // ======================================================== /// Умное чтение bool: быстрые попытки → блокирующее чтение pub fn smart_get_bool(&self, topic: &str) -> Option<bool> { for attempt in 1..=10 { if let Ok(map) = self.bool.try_read() { return map.get(topic).copied(); } std::thread::yield_now(); if attempt == 5 { eprintln!( "⚠️ smart_get_bool: попытка {} для '{}' — writer держит блокировку", attempt, topic ); } } eprintln!( "🔴 smart_get_bool: 10 попыток не удались, блокируем поток для '{}'", topic ); let map = self.bool.blocking_read(); map.get(topic).copied() } /// Умное чтение i32 pub fn smart_get_i32(&self, topic: &str) -> Option<i32> { for attempt in 1..=10 { if let Ok(map) = self.i32.try_read() { return map.get(topic).copied(); } std::thread::yield_now(); if attempt == 5 { eprintln!( "⚠️ smart_get_i32: попытка {} для '{}' — writer держит блокировку", attempt, topic ); } } eprintln!( "🔴 smart_get_i32: 10 попыток не удались, блокируем поток для '{}'", topic ); let map = self.i32.blocking_read(); map.get(topic).copied() } /// Умное чтение f64 pub fn smart_get_f64(&self, topic: &str) -> Option<f64> { for attempt in 1..=10 { if let Ok(map) = self.f64.try_read() { return map.get(topic).copied(); } std::thread::yield_now(); if attempt == 5 { eprintln!( "⚠️ smart_get_f64: попытка {} для '{}' — writer держит блокировку", attempt, topic ); } } eprintln!( "🔴 smart_get_f64: 10 попыток не удались, блокируем поток для '{}'", topic ); let map = self.f64.blocking_read(); map.get(topic).copied() } /// Умное чтение String pub fn smart_get_string(&self, topic: &str) -> Option<String> { for attempt in 1..=10 { if let Ok(map) = self.str.try_read() { return map.get(topic).cloned(); } std::thread::yield_now(); if attempt == 5 { eprintln!( "⚠️ smart_get_string: попытка {} для '{}' — writer держит блокировку", attempt, topic ); } } eprintln!( "🔴 smart_get_string: 10 попыток не удались, блокируем поток для '{}'", topic ); let map = self.str.blocking_read(); map.get(topic).cloned() } // Для отладки pub async fn snapshot(&self) -> String { let i32_map = self.i32.read().await; let f64_map = self.f64.read().await; let bool_map = self.bool.read().await; let str_map = self.str.read().await; format!( "MqttTables:\n i32: {:?}\n f64: {:?}\n bool: {:?}\n str: {:?}", *i32_map, *f64_map, *bool_map, *str_map ) } } // ============================================================ // МЕНЕДЖЕР ПЕРЕМЕННЫХ СКРИПТОВ // ============================================================ #[derive(Clone)] pub struct ScriptVarManager { scripts: Arc<RwLock<HashMap<String, Arc<Mutex<VarTables>>>>>, } impl ScriptVarManager { pub fn new() -> Self { Self { scripts: Arc::new(RwLock::new(HashMap::new())), } } /// Зарегистрировать новый скрипт при загрузке (асинхронный) pub async fn register(&self, script_name: &str) -> Arc<Mutex<VarTables>> { let table = Arc::new(Mutex::new(VarTables::new())); let mut scripts = self.scripts.write().await; scripts.insert(script_name.to_string(), table.clone()); println!( "📝 ScriptVarManager: зарегистрирована таблица для '{}'", script_name ); table } /// СИНХРОННАЯ регистрация скрипта (для вызова из того же потока, без tokio) pub fn register_sync(&self, script_name: &str) -> Arc<Mutex<VarTables>> { let table = Arc::new(Mutex::new(VarTables::new())); // Пытаемся получить write-блокировку (с fallback) let mut scripts = loop { if let Ok(guard) = self.scripts.try_write() { break guard; } // Короткая пауза перед повторной попыткой std::thread::yield_now(); }; scripts.insert(script_name.to_string(), table.clone()); println!( "📝 ScriptVarManager: зарегистрирована таблица для '{}'", script_name ); table } /// Удалить скрипт (асинхронный) pub async fn unregister(&self, script_name: &str) { let mut scripts = self.scripts.write().await; if scripts.remove(script_name).is_some() { println!( "🗑️ ScriptVarManager: удалена таблица для '{}'", script_name ); } } /// Получить таблицу переменных скрипта (асинхронный) pub async fn get(&self, script_name: &str) -> Option<Arc<Mutex<VarTables>>> { let scripts = self.scripts.read().await; scripts.get(script_name).cloned() } /// СИНХРОННАЯ версия get() для Rune-функций pub fn get_sync(&self, script_name: &str) -> Option<Arc<Mutex<VarTables>>> { for attempt in 1..=10 { if let Ok(scripts) = self.scripts.try_read() { return scripts.get(script_name).cloned(); } std::thread::yield_now(); if attempt == 5 { eprintln!( "⚠️ ScriptVarManager::get_sync: попытка {} для '{}' — writer держит блокировку", attempt, script_name ); } } eprintln!( "🔴 ScriptVarManager::get_sync: 10 попыток не удались, блокируем поток для '{}'", script_name ); let scripts = self.scripts.blocking_read(); scripts.get(script_name).cloned() } /// Очистить ВСЕ таблицы (асинхронный) pub async fn clear_all(&self) { let mut scripts = self.scripts.write().await; let count = scripts.len(); scripts.clear(); println!("🗑️ ScriptVarManager: удалены все таблицы (было {})", count); } /// Для Web-интерфейса: получить снапшоты всех таблиц pub async fn snapshot_all(&self) -> HashMap<String, String> { let scripts = self.scripts.read().await; let mut result = HashMap::new(); for (name, table_arc) in scripts.iter() { let guard = table_arc.lock().await; result.insert(name.clone(), guard.snapshot()); } result } } // ============================================================ // ТАБЛИЦА ПЕРЕМЕННЫХ ОДНОГО СКРИПТА // ============================================================ #[derive(Clone, Debug)] pub struct VarTables { pub i32: HashMap<String, i32>, pub f64: HashMap<String, f64>, pub bool: HashMap<String, bool>, pub str: HashMap<String, String>, } impl VarTables { pub fn new() -> Self { Self { i32: HashMap::new(), f64: HashMap::new(), bool: HashMap::new(), str: HashMap::new(), } } // Синхронные методы для чтения pub fn get_i32(&self, name: &str) -> Option<i32> { self.i32.get(name).copied() } pub fn get_f64(&self, name: &str) -> Option<f64> { self.f64.get(name).copied() } pub fn get_bool(&self, name: &str) -> Option<bool> { self.bool.get(name).copied() } pub fn get_string(&self, name: &str) -> Option<String> { self.str.get(name).cloned() } // Синхронные методы для записи pub fn set_i32(&mut self, name: &str, value: i32) { self.i32.insert(name.to_string(), value); } pub fn set_f64(&mut self, name: &str, value: f64) { self.f64.insert(name.to_string(), value); } pub fn set_bool(&mut self, name: &str, value: bool) { self.bool.insert(name.to_string(), value); } pub fn set_string(&mut self, name: &str, value: String) { self.str.insert(name.to_string(), value); } // Для отладки pub fn snapshot(&self) -> String { format!( "VarTables:\n i32: {:?}\n f64: {:?}\n bool: {:?}\n str: {:?}", self.i32, self.f64, self.bool, self.str ) } } // ============================================================ // ТЕСТЫ // ============================================================ #[cfg(test)] mod tests { use super::*; #[tokio::test] async fn test_script_var_manager() { let manager = ScriptVarManager::new(); let kitchen_table = manager.register("kitchen.rn").await; { let mut guard = kitchen_table.lock().await; guard.set_bool("lights_on", true); guard.set_i32("brightness", 75); } { let guard = kitchen_table.lock().await; assert_eq!(guard.get_bool("lights_on"), Some(true)); assert_eq!(guard.get_i32("brightness"), Some(75)); } manager.unregister("kitchen.rn").await; assert!(manager.get("kitchen.rn").await.is_none()); } #[tokio::test] async fn test_independent_scripts() { let manager = ScriptVarManager::new(); let t1 = manager.register("script_a.rn").await; let t2 = manager.register("script_b.rn").await; { let mut guard = t1.lock().await; guard.set_i32("x", 10); } { let mut guard = t2.lock().await; guard.set_i32("y", 20); } let g1 = t1.lock().await; let g2 = t2.lock().await; assert_eq!(g1.get_i32("x"), Some(10)); assert_eq!(g1.get_i32("y"), None); assert_eq!(g2.get_i32("y"), Some(20)); assert_eq!(g2.get_i32("x"), None); } #[tokio::test] async fn test_state_clone() { let state = State::new(); let state2 = state.clone(); state.mqtt.set_bool("/devices/test", true).await; assert_eq!(state2.mqtt.get_bool("/devices/test").await, Some(true)); } #[test] fn test_smart_get_sync() { let tables = MqttTables::new(); // Синхронно пишем через блокирующий write (для теста) { let mut map = tables.bool.blocking_write(); map.insert("/devices/motion".to_string(), true); } // Синхронно читаем через smart_get let value = tables.smart_get_bool("/devices/motion"); assert_eq!(value, Some(true)); } } // === crates/rust_rules_config/src/lib.rs === //cd /d/Rust/projects/rune_rules/crates/wb_rust_config use std::collections::HashMap; use serde::Deserialize; #[derive(Debug, Deserialize)] pub struct Config { pub broker: BrokerConfig, pub devices: HashMap<String, HashMap<String, Control>>, pub scripts: ScriptsConfig, pub r#virtual: Option<HashMap<String, HashMap<String, VirtualControl>>>, } impl Config { pub fn load_config(path: &str) -> Self { let content = std::fs::read_to_string(path) .expect(&format!("Не удалось найти файл конфигурации: {}", path)); toml::from_str(&content) .expect("Ошибка в формате TOML") } pub fn all_topics(&self) -> Vec<String> { let mut topics = Vec::new(); for (_, controls) in &self.devices { for (_, info) in controls { topics.push(info.topic.clone()); } } topics } pub fn aliases(&self) -> HashMap<String, String> { let mut aliases = HashMap::new(); // Авто-алиасы из реальных устройств for (device_name, controls) in &self.devices { for (control_name, control) in controls { let alias = format!("{}.{}", device_name, control_name); aliases.insert(alias, control.topic.clone()); } } // Авто-алиасы из виртуальных устройств if let Some(virtual_devices) = &self.r#virtual { for (device_name, controls) in virtual_devices { for control_name in controls.keys() { let alias = format!("{}.{}", device_name, control_name); let topic = format!("/devices/{}/controls/{}", device_name, control_name); aliases.insert(alias, topic); } } } aliases } // /// Получить тип топика по его полному пути // pub fn get_type(&self, topic: &str) -> Option<String> { // } } #[derive(Debug, Deserialize)] pub struct ScriptsConfig { pub directory: String, // путь к папке со скриптами } #[derive(Debug, Deserialize)] pub struct BrokerConfig { pub host: String, pub port: u16, pub command_suffix: String, } #[derive(Debug, Deserialize)] pub struct Control { pub topic: String, pub r#type: String, } pub struct TopicTypes { map: HashMap<String, String>, // topic → type } impl TopicTypes { /// Создать из Config pub fn from_config(config: &Config) -> Self { let mut map = HashMap::new(); for (_, controls) in &config.devices { for (_, info) in controls { map.insert(info.topic.clone(), info.r#type.clone()); } } Self {map} } /// Получить тип топика. Если топик не найден — логирует и возвращает "string" pub fn get(&self, topic: &str) -> &str { match self.map.get(topic) { Some(t) => t.as_str(), None => { eprintln!("⚠️ Топик '{}' не найден в конфиге, тип = string", topic); "string" } } } /// Получить все топики pub fn all_topics(&self) -> Vec<String> { self.map.keys().cloned().collect() } } #[derive(Debug, Deserialize, Clone)] pub struct VirtualControl { /// Тип топика Wiren Board: "switch", "range", "value", "text", "rgb", "pushbutton" pub topic_type: String, /// Значение по умолчанию pub value: Option<serde_json::Value>, /// Тип данных: "bool", "i32", "f64", "string" #[serde(rename = "type")] pub value_type: String, // Опциональные мета-поля pub min: Option<f64>, pub max: Option<f64>, pub precision: Option<u32>, pub readonly: Option<bool>, #[serde(default)] pub force_default: bool, #[serde(default)] pub lazy_init: bool, #[serde(default)] pub retain: bool, } // === crates/scenarios/src/lib.rs === // mod motion_light; // pub use motion_light::register_motion_light_scenario; // === crates/scenarios/src/motion_light.rs === // === crates/scripting/src/lib.rs === pub fn add(left: u64, right: u64) -> u64 { left + right } #[cfg(test)] mod tests { use super::*; #[test] fn it_works() { let result = add(2, 2); assert_eq!(result, 4); } }