/
AxeOn
/
ProjectAi
Обзор
Документация
Войти
/
AxeOn
/
ProjectAi
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
master
database.py
4 329 строк
204 KB
Дмитрий
refactor(backend): pivot surcharge logic from automated rules engine to manual hourly calculation records
07 май 2026, 09:01
07 май 2026, 09:01
4eeb141
Код
Авторство
О чём код?
# database.py """ Менеджер базы данных для Офис-Планировщика v5.0 УЛУЧШЕНИЯ: - Thread-safe без check_same_thread=False - Connection pooling для многопоточности - Валидация данных перед записью - CHECK constraints в схеме БД - Система миграций - Структурированное логирование ОБРАТНАЯ СОВМЕСТИМОСТЬ: Все методы из старой версии сохранены и работают. """ import sqlite3 import threading import logging from pathlib import Path from datetime import datetime from typing import Optional, List, Dict, Any from contextlib import contextmanager from queue import Queue, Empty from models import Employee, ManagementTask, DeadlineCase, VacationRange, WorkingGroup, ManualSurcharge logger = logging.getLogger(__name__) class ValidationError(Exception): """Ошибка валидации данных""" pass class ConnectionPool: """ Пул соединений для thread-safe работы с SQLite ВАЖНО: SQLite требует, чтобы каждое соединение использовалось только в том потоке, где оно было создано. Поэтому мы создаем НОВОЕ соединение для каждого потока вместо переиспользования. """ def __init__(self, db_path: str): self.db_path = db_path self.lock = threading.Lock() # Словарь для хранения соединений по thread_id self._thread_connections = {} logger.debug(f"Connection pool initialized for {db_path}") def _create_connection(self) -> sqlite3.Connection: """Создать новое соединение""" conn = sqlite3.connect(self.db_path, timeout=30.0, isolation_level=None) conn.row_factory = sqlite3.Row # ВАЖНО: Используем режим DELETE вместо WAL для консистентности conn.execute("PRAGMA journal_mode=DELETE") conn.execute("PRAGMA synchronous=NORMAL") conn.execute("PRAGMA foreign_keys=ON") return conn @contextmanager def get_connection(self): """ Context manager для получения соединения. Каждый поток получает своё соединение, которое переиспользуется. Подготовительная секция (проверка / создание) защищена self.lock, чтобы исключить гонку между "проверить → удалить → создать". yield выполняется вне лока — другие потоки не блокируются во время работы с БД. """ thread_id = threading.get_ident() with self.lock: conn = self._thread_connections.get(thread_id) if conn is not None: try: conn.execute("SELECT 1") except sqlite3.Error: try: conn.close() except Exception: pass del self._thread_connections[thread_id] conn = None if conn is None: conn = self._create_connection() self._thread_connections[thread_id] = conn try: yield conn except Exception as e: logger.error(f"Connection error in thread {thread_id}: {e}") try: conn.close() except Exception: pass with self.lock: self._thread_connections.pop(thread_id, None) raise def close_thread_connection(self, thread_id): """Безопасно закрыть соединение конкретного потока""" with self.lock: if thread_id in self._thread_connections: try: self._thread_connections[thread_id].close() except Exception: pass del self._thread_connections[thread_id] def close_all(self): """Закрыть все соединения""" with self.lock: for conn in list(self._thread_connections.values()): try: conn.close() except Exception: pass self._thread_connections.clear() logger.info("All connections closed") class DBManager: """Менеджер базы данных с улучшенной архитектурой""" def __init__(self, db_path=None): if db_path: self.db_path = db_path else: import os base_dir = os.path.dirname(os.path.abspath(__file__)) self.db_path = os.path.join(base_dir, "planner_v4.db") logger.info(f"DBManager initializing: {self.db_path}") # Connection pool вместо одного соединения self.pool = ConnectionPool(str(self.db_path)) self.lock = threading.Lock() # Инициализация self.init_db() self.check_migrations() # === ВСПОМОГАТЕЛЬНЫЕ МЕТОДЫ === @staticmethod def _validate_date(date_str: str, field_name: str = "date") -> str: """Валидация даты в формате YYYY-MM-DD""" if not date_str: raise ValidationError(f"{field_name} cannot be empty") if not isinstance(date_str, str): raise ValidationError(f"{field_name} must be string, got {type(date_str)}") if len(date_str) != 10: raise ValidationError( f"{field_name} must be 10 chars (YYYY-MM-DD), got {len(date_str)}: '{date_str}'" ) try: datetime.strptime(date_str, "%Y-%m-%d") except ValueError as e: raise ValidationError(f"Invalid {field_name} '{date_str}': {e}") return date_str def init_db(self): """Создание схемы БД""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() # 1. Сотрудники cursor.execute(""" CREATE TABLE IF NOT EXISTS employees ( id INTEGER PRIMARY KEY AUTOINCREMENT, name TEXT UNIQUE NOT NULL, email TEXT, tg_id TEXT, department TEXT, phone TEXT, managers TEXT, is_manager INTEGER DEFAULT 0, is_temp INTEGER DEFAULT 0, birth_date TEXT, category TEXT, is_reso INTEGER DEFAULT 0, working_group_id INTEGER, base_salary REAL DEFAULT 0 ) """) # 2. Рабочие группы cursor.execute(""" CREATE TABLE IF NOT EXISTS working_groups ( id INTEGER PRIMARY KEY AUTOINCREMENT, name TEXT UNIQUE, tg_id TEXT, department TEXT ) """) # 3. Задачи (для модуля Загрузка) cursor.execute(""" CREATE TABLE IF NOT EXISTS tasks ( id INTEGER PRIMARY KEY AUTOINCREMENT, employee_id INTEGER, task_date TEXT, task_name TEXT, client_name TEXT, status TEXT DEFAULT 'New', priority TEXT DEFAULT 'Medium', parent_id INTEGER DEFAULT 0, deadline TEXT, comment TEXT, poll_id TEXT, FOREIGN KEY(employee_id) REFERENCES employees(id) ) """) # 3a. Задачи управления (для модуля "Текущие задачи" - планировщик) cursor.execute(""" CREATE TABLE IF NOT EXISTS management_tasks ( id INTEGER PRIMARY KEY AUTOINCREMENT, title TEXT NOT NULL, priority TEXT DEFAULT 'Medium', employee_id INTEGER, due_date TEXT NOT NULL, status TEXT DEFAULT 'New', parent_id INTEGER DEFAULT 0, comment TEXT, client_name TEXT, created_at TEXT DEFAULT CURRENT_TIMESTAMP, updated_at TEXT, poll_id TEXT, poll_sent_at TEXT, poll_answered_at TEXT, FOREIGN KEY(employee_id) REFERENCES employees(id) ) """) logger.info("Management tasks table initialized") # Миграция: добавляем новые колонки в существующие БД try: cursor.execute("PRAGMA table_info(management_tasks)") existing = [row[1] for row in cursor.fetchall()] for col, defn in [("poll_sent_at", "TEXT"), ("poll_answered_at", "TEXT")]: if col not in existing: cursor.execute(f"ALTER TABLE management_tasks ADD COLUMN {col} {defn}") logger.info(f"Added {col} column to management_tasks") except Exception as _e: logger.warning(f"Migration management_tasks: {_e}") # 4. Календарь сотрудников cursor.execute(""" CREATE TABLE IF NOT EXISTS employee_calendar ( employee_id INTEGER, date TEXT, is_working INTEGER, PRIMARY KEY(employee_id, date), FOREIGN KEY(employee_id) REFERENCES employees(id) ) """) # Добавляем колонку is_absent если её нет (для ручных отметок "Н") try: cursor.execute("PRAGMA table_info(employee_calendar)") columns = [row[1] for row in cursor.fetchall()] if 'is_absent' not in columns: cursor.execute(""" ALTER TABLE employee_calendar ADD COLUMN is_absent INTEGER DEFAULT 0 """) logger.info("Added is_absent column to employee_calendar") if 'status' not in columns: cursor.execute(""" ALTER TABLE employee_calendar ADD COLUMN status TEXT DEFAULT NULL """) logger.info("Added status column to employee_calendar") except Exception as e: logger.warning(f"Could not add column to employee_calendar: {e}") # 5. Производственный календарь cursor.execute(""" CREATE TABLE IF NOT EXISTS production_calendar_exceptions ( date TEXT PRIMARY KEY, is_working INTEGER ) """) # 6. Отпуска cursor.execute(""" CREATE TABLE IF NOT EXISTS vacations ( id INTEGER PRIMARY KEY AUTOINCREMENT, employee_id INTEGER, start_date TEXT, end_date TEXT, v_type TEXT, days INTEGER, type TEXT, FOREIGN KEY(employee_id) REFERENCES employees(id) ) """) # 7. Ограничения умного распределения cursor.execute(""" CREATE TABLE IF NOT EXISTS distribution_restrictions ( id INTEGER PRIMARY KEY AUTOINCREMENT, expert_name TEXT NOT NULL UNIQUE, blocked_services TEXT DEFAULT '[]', blocked_regions TEXT DEFAULT '[]', blocked_risks TEXT DEFAULT '[]', blocked_prefixes TEXT DEFAULT '[]', mode_services TEXT DEFAULT 'block', mode_regions TEXT DEFAULT 'block', mode_risks TEXT DEFAULT 'block', mode_prefixes TEXT DEFAULT 'block' ) """) # Миграция для существующих БД: добавляем новые колонки если их нет for _col, _def in ( ('mode_services', "'block'"), ('mode_regions', "'block'"), ('mode_risks', "'block'"), ('blocked_prefixes', "'[]'"), ('mode_prefixes', "'block'"), ): try: cursor.execute( f"ALTER TABLE distribution_restrictions " f"ADD COLUMN {_col} TEXT DEFAULT {_def}" ) except Exception: pass # колонка уже существует # 8. Лимиты распределения cursor.execute(""" CREATE TABLE IF NOT EXISTS distribution_limits ( id INTEGER PRIMARY KEY AUTOINCREMENT, expert_name TEXT NOT NULL UNIQUE, max_cases INTEGER NOT NULL ) """) # 8a. Матрица маршрутизации (одна строка на эксперта) # Каждое поле — JSON-объект вида {"type": "deny"|"only", "values": [...]} cursor.execute(""" CREATE TABLE IF NOT EXISTS distribution_rules ( id INTEGER PRIMARY KEY AUTOINCREMENT, expert_name TEXT NOT NULL, service_types TEXT DEFAULT NULL, regions TEXT DEFAULT NULL, risks TEXT DEFAULT NULL, prefixes TEXT DEFAULT NULL, is_active INTEGER DEFAULT 1 ) """) # 9. Цены cursor.execute(""" CREATE TABLE IF NOT EXISTS prices ( id INTEGER PRIMARY KEY AUTOINCREMENT, name TEXT NOT NULL, calc_rule_month TEXT, calc_rule_salary TEXT, price_high REAL DEFAULT 0, price_first REAL DEFAULT 0, price_second REAL DEFAULT 0, service_cost REAL DEFAULT 0, payment_type TEXT DEFAULT '0', tax_percent REAL DEFAULT 0, year TEXT ) """) cursor.execute("CREATE INDEX IF NOT EXISTS idx_prices_year ON prices(year)") cursor.execute("CREATE INDEX IF NOT EXISTS idx_prices_name ON prices(name)") # 8. Услуги cursor.execute(""" CREATE TABLE IF NOT EXISTS services ( id INTEGER PRIMARY KEY AUTOINCREMENT, name TEXT UNIQUE, is_base INTEGER DEFAULT 1 ) """) # 9. Корректировки зарплаты cursor.execute(""" CREATE TABLE IF NOT EXISTS salary_adjustments ( id INTEGER PRIMARY KEY AUTOINCREMENT, employee_id INTEGER, period TEXT, amount REAL, comment TEXT, FOREIGN KEY(employee_id) REFERENCES employees(id) ) """) # 9a. Доплаты зарплаты cursor.execute(""" CREATE TABLE IF NOT EXISTS salary_doplaty ( id INTEGER PRIMARY KEY AUTOINCREMENT, period TEXT NOT NULL, employee_id INTEGER, dept TEXT, emp_name TEXT, deal_number TEXT, service_name TEXT, base_cost REAL DEFAULT 0, hours REAL DEFAULT 0, hour_rate REAL DEFAULT 0, doplata REAL DEFAULT 0, comment TEXT, FOREIGN KEY(employee_id) REFERENCES employees(id) ) """) # 10. Настройки cursor.execute(""" CREATE TABLE IF NOT EXISTS settings ( key TEXT PRIMARY KEY, value TEXT ) """) # 11. Telegram статистика cursor.execute(""" CREATE TABLE IF NOT EXISTS tg_stats ( id INTEGER PRIMARY KEY AUTOINCREMENT, event_type TEXT, poll_id TEXT, created_at TEXT ) """) # 12. Настройки рассылок Telegram cursor.execute(""" CREATE TABLE IF NOT EXISTS tg_notification_rules ( id INTEGER PRIMARY KEY AUTOINCREMENT, rule_type TEXT NOT NULL, subject TEXT NOT NULL, recipients TEXT, is_enabled INTEGER DEFAULT 1, description TEXT, created_at TEXT ) """) # Предзаполнение дефолтными правилами (добавляем отсутствующие по ключу тип+тема) now_iso = datetime.now().isoformat() default_rules = [ ("Уведомление", "Отпуск (за 5 дней)", "Руководители сотрудника", "За 5 рабочих дней до начала отпуска — руководителям"), ("Уведомление", "Отпуск (за 1 день)", "Руководители + Группа", "За 1 рабочий день до отпуска — руководителям и рабочей группе"), ("Уведомление", "День рождения", "Руководители сотрудника", "За 1 рабочий день до дня рождения — руководителям"), ("Опрос", "Задача", "Исполнитель", "Опрос о статусе задачи при создании — исполнителю"), ("Уведомление", "Нет Базовых услуг", "Руководители сотрудника", "Сотрудники без закрытых базовых услуг за прошлый рабочий день"), ("Уведомление", "Итоги по отделам", "Руководители сотрудника", "Количество закрытых базовых заявок по отделам за прошлый рабочий день"), ("Повторная отправка", "Неотвеченные опросы", "Исполнитель", "Повторно отправлять опрос по задачам, на которые не был дан ответ"), ("Автоотправка по срокам", "За 3 рабочих дня", "Исполнитель", "Отправить опрос исполнителю за 3 рабочих дня до истечения срока задачи"), ("Автоотправка по срокам", "В день истечения срока", "Исполнитель", "Отправить опрос с предупреждением в день истечения срока задачи"), ("Автоотправка по срокам", "После истечения срока", "Исполнитель", "Отправить опрос с требованием связаться с руководством после истечения срока"), ] cursor.execute( "SELECT rule_type || '|' || subject FROM tg_notification_rules" ) existing_keys = {row[0] for row in cursor.fetchall()} for r in default_rules: key = f"{r[0]}|{r[1]}" if key not in existing_keys: cursor.execute(""" INSERT INTO tg_notification_rules (rule_type, subject, recipients, is_enabled, description, created_at) VALUES (?, ?, ?, 1, ?, ?) """, (r[0], r[1], r[2], r[3], now_iso)) existing_keys.add(key) # 13. Workload (Загрузка сотрудников - daily и monthly) cursor.execute(""" CREATE TABLE IF NOT EXISTS workload ( id INTEGER PRIMARY KEY AUTOINCREMENT, period TEXT, emp_name TEXT, service_name TEXT, count REAL, price REAL DEFAULT 0, total_sum REAL DEFAULT 0, created_at TEXT, source TEXT DEFAULT 'daily' ) """) cursor.execute(""" CREATE INDEX IF NOT EXISTS idx_workload_period ON workload(period) """) cursor.execute(""" CREATE INDEX IF NOT EXISTS idx_workload_source ON workload(source) """) cursor.execute(""" CREATE INDEX IF NOT EXISTS idx_workload_period_source ON workload(period, source) """) # 13. Контроль сроков дел cursor.execute(""" CREATE TABLE IF NOT EXISTS deadline_cases ( id INTEGER PRIMARY KEY AUTOINCREMENT, case_number TEXT NOT NULL, create_date TEXT, inspection_date TEXT, risk TEXT, expert_inspection TEXT, expert_calculation TEXT, payment_info TEXT, region TEXT, department TEXT, deadline INTEGER, comment TEXT, load_date TEXT NOT NULL, address TEXT, created_at TEXT DEFAULT CURRENT_TIMESTAMP ) """) cursor.execute(""" CREATE INDEX IF NOT EXISTS idx_deadline_case_number ON deadline_cases(case_number) """) cursor.execute(""" CREATE INDEX IF NOT EXISTS idx_deadline_load_date ON deadline_cases(load_date) """) cursor.execute(""" CREATE INDEX IF NOT EXISTS idx_deadline_expert ON deadline_cases(expert_calculation) """) # Миграция: добавляем address если её нет (для уже существующих БД) try: cursor.execute( "ALTER TABLE deadline_cases ADD COLUMN address TEXT" ) except Exception: pass # колонка уже существует logger.info("Deadline cases table initialized") # 14. История изменений комментариев cursor.execute(""" CREATE TABLE IF NOT EXISTS deadline_comment_history ( id INTEGER PRIMARY KEY AUTOINCREMENT, case_number TEXT NOT NULL, load_date TEXT NOT NULL, old_comment TEXT, new_comment TEXT, changed_at TEXT DEFAULT CURRENT_TIMESTAMP ) """) cursor.execute(""" CREATE INDEX IF NOT EXISTS idx_comment_history_case ON deadline_comment_history(case_number, load_date) """) logger.info("Comment history table initialized") # Контроль качества cursor.execute(""" CREATE TABLE IF NOT EXISTS quality_cases ( id INTEGER PRIMARY KEY AUTOINCREMENT, nomer_polisa TEXT NOT NULL, risk TEXT, razdel TEXT, data_zakr TEXT, mes_zakr INTEGER, god_zakr INTEGER, osmotrowshik TEXT, region_osm TEXT, narusheniya TEXT, naim_podr TEXT, opisanie TEXT, ispravlenie TEXT DEFAULT 'Нет', data_osm TEXT, kto_obnaruzhil TEXT, UNIQUE(nomer_polisa, data_osm) ) """) logger.info("Quality cases table initialized") # Миграция: заменяем UNIQUE(nomer_polisa) на UNIQUE(nomer_polisa, data_osm) cursor.execute( "SELECT sql FROM sqlite_master WHERE type='table' AND name='quality_cases'" ) _qc = cursor.fetchone() _qc_sql = (_qc[0] if _qc else '') or '' if 'nomer_polisa TEXT NOT NULL UNIQUE' in _qc_sql: logger.info( "quality_cases: выполняем миграцию UNIQUE " "nomer_polisa → (nomer_polisa, data_osm)" ) cursor.executescript(""" CREATE TABLE IF NOT EXISTS quality_cases_mig ( id INTEGER PRIMARY KEY AUTOINCREMENT, nomer_polisa TEXT NOT NULL, risk TEXT, razdel TEXT, data_zakr TEXT, mes_zakr INTEGER, god_zakr INTEGER, osmotrowshik TEXT, region_osm TEXT, narusheniya TEXT, naim_podr TEXT, opisanie TEXT, ispravlenie TEXT DEFAULT 'Нет', data_osm TEXT, kto_obnaruzhil TEXT, UNIQUE(nomer_polisa, data_osm) ); INSERT OR IGNORE INTO quality_cases_mig SELECT * FROM quality_cases; DROP TABLE quality_cases; ALTER TABLE quality_cases_mig RENAME TO quality_cases; """) logger.info("quality_cases: миграция завершена") # 16. Ручные доплаты (Manual Surcharges) cursor.execute(""" CREATE TABLE IF NOT EXISTS manual_surcharges ( id INTEGER PRIMARY KEY AUTOINCREMENT, expert_name TEXT NOT NULL, department TEXT, case_number TEXT, service_name TEXT, base_cost REAL DEFAULT 0, hourly_rate REAL DEFAULT 0, hours REAL DEFAULT 0, total_amount REAL DEFAULT 0, operation_type TEXT NOT NULL DEFAULT 'Доплата' CHECK(operation_type IN ('Доплата', 'Взаимовычет')), note TEXT, period_date TEXT NOT NULL ) """) logger.info("Manual surcharges table initialized") conn.commit() logger.info("Database schema initialized") def check_migrations(self): """Запуск миграций для обновления схемы""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() try: # Миграция: добавление отсутствующих полей в employees cursor.execute("PRAGMA table_info(employees)") cols = [row['name'] for row in cursor.fetchall()] if 'base_salary' not in cols: cursor.execute("ALTER TABLE employees ADD COLUMN base_salary REAL DEFAULT 0") logger.info("Added column: employees.base_salary") if 'working_group_id' not in cols: cursor.execute("ALTER TABLE employees ADD COLUMN working_group_id INTEGER") logger.info("Added column: employees.working_group_id") if 'is_fired' not in cols: cursor.execute("ALTER TABLE employees ADD COLUMN is_fired INTEGER DEFAULT 0") logger.info("Added column: employees.is_fired") if 'fired_date' not in cols: cursor.execute("ALTER TABLE employees ADD COLUMN fired_date TEXT") logger.info("Added column: employees.fired_date") if 'schedule_pattern' not in cols: cursor.execute( "ALTER TABLE employees ADD COLUMN schedule_pattern TEXT DEFAULT '5/2_mon_fri'" ) logger.info("Added column: employees.schedule_pattern") if 'svc_types' not in cols: cursor.execute( "ALTER TABLE employees ADD COLUMN svc_types TEXT DEFAULT '[]'" ) logger.info("Added column: employees.svc_types") # Миграция: переименование tasks.date -> tasks.task_date cursor.execute("PRAGMA table_info(tasks)") task_cols = [row['name'] for row in cursor.fetchall()] if 'date' in task_cols and 'task_date' not in task_cols: cursor.execute("ALTER TABLE tasks RENAME COLUMN date TO task_date") logger.info("Renamed column: tasks.date -> tasks.task_date") # Обновляем список колонок cursor.execute("PRAGMA table_info(tasks)") task_cols = [row['name'] for row in cursor.fetchall()] # Миграция: добавление обязательных полей в tasks required_fields = { 'employee_id': "INTEGER", 'task_date': "TEXT", 'task_name': "TEXT", 'client_name': "TEXT", 'status': "TEXT DEFAULT 'New'", 'priority': "TEXT DEFAULT 'Medium'", 'parent_id': "INTEGER DEFAULT 0", 'deadline': "TEXT", 'comment': "TEXT", 'poll_id': "TEXT" } for field, definition in required_fields.items(): if field not in task_cols: try: cursor.execute(f"ALTER TABLE tasks ADD COLUMN {field} {definition}") logger.info(f"Added column: tasks.{field}") except Exception as e: logger.warning(f"Could not add tasks.{field}: {e}") # Миграция: добавление tg_id в working_groups cursor.execute("PRAGMA table_info(working_groups)") wg_cols = [row['name'] for row in cursor.fetchall()] if 'tg_id' not in wg_cols: cursor.execute("ALTER TABLE working_groups ADD COLUMN tg_id TEXT") logger.info("Added column: working_groups.tg_id") # Миграция: добавление is_active в distribution_rules cursor.execute("PRAGMA table_info(distribution_rules)") dr_cols = [row['name'] for row in cursor.fetchall()] if 'is_active' not in dr_cols: cursor.execute( "ALTER TABLE distribution_rules ADD COLUMN is_active INTEGER DEFAULT 1" ) logger.info("Added column: distribution_rules.is_active") # Миграция: добавление ne_value в workload (фильтр по офису/НЭ) cursor.execute("PRAGMA table_info(workload)") wl_cols = [row['name'] for row in cursor.fetchall()] if 'ne_value' not in wl_cols: cursor.execute("ALTER TABLE workload ADD COLUMN ne_value TEXT DEFAULT NULL") logger.info("Added column: workload.ne_value") # Миграция: добавление services_list и candidate_expert в deadline_cases cursor.execute("PRAGMA table_info(deadline_cases)") dc_cols = [row['name'] for row in cursor.fetchall()] if 'services_list' not in dc_cols: cursor.execute("ALTER TABLE deadline_cases ADD COLUMN services_list TEXT") logger.info("Added column: deadline_cases.services_list") if 'candidate_expert' not in dc_cols: cursor.execute("ALTER TABLE deadline_cases ADD COLUMN candidate_expert TEXT") logger.info("Added column: deadline_cases.candidate_expert") # Миграция: признак «Работает» (производственная необходимость) в отпусках cursor.execute("PRAGMA table_info(vacations)") vac_cols = [row['name'] for row in cursor.fetchall()] if 'is_working' not in vac_cols: cursor.execute( "ALTER TABLE vacations ADD COLUMN is_working INTEGER DEFAULT 0" ) logger.info("Added column: vacations.is_working") # КРИТИЧЕСКАЯ МИГРАЦИЯ: Очистка некорректных дат cursor.execute(""" DELETE FROM employee_calendar WHERE length(date) != 10 OR date NOT LIKE '____-__-__' """) deleted = cursor.rowcount if deleted > 0: logger.warning(f"Deleted {deleted} invalid date records from employee_calendar") # Миграция: снять UNIQUE с expert_name в distribution_rules cursor.execute( "SELECT sql FROM sqlite_master WHERE type='table' AND name='distribution_rules'" ) _dr_row = cursor.fetchone() _dr_sql = (_dr_row['sql'] if _dr_row else '') or '' if 'UNIQUE' in _dr_sql.upper(): conn.executescript(""" BEGIN TRANSACTION; CREATE TABLE IF NOT EXISTS distribution_rules_new ( id INTEGER PRIMARY KEY AUTOINCREMENT, expert_name TEXT NOT NULL, service_types TEXT DEFAULT NULL, regions TEXT DEFAULT NULL, risks TEXT DEFAULT NULL, prefixes TEXT DEFAULT NULL, is_active INTEGER NOT NULL DEFAULT 1 ); INSERT INTO distribution_rules_new (id, expert_name, service_types, regions, risks, prefixes, is_active) SELECT id, expert_name, service_types, regions, risks, prefixes, COALESCE(is_active, 1) FROM distribution_rules; DROP TABLE distribution_rules; ALTER TABLE distribution_rules_new RENAME TO distribution_rules; COMMIT; """) logger.info("Migrated distribution_rules: removed UNIQUE constraint from expert_name") # Создание индексов для производительности indexes = [ "CREATE INDEX IF NOT EXISTS idx_tasks_employee ON tasks(employee_id)", "CREATE INDEX IF NOT EXISTS idx_tasks_date ON tasks(task_date)", "CREATE INDEX IF NOT EXISTS idx_vacations_employee ON vacations(employee_id)", ] for idx_sql in indexes: try: cursor.execute(idx_sql) except Exception: pass # Индекс уже существует conn.commit() logger.info("Migrations completed successfully") except Exception as e: logger.error(f"Migration error: {e}", exc_info=True) conn.rollback() # === СОТРУДНИКИ === def get_employees(self, department=None, only_managers=False, fired_only=False): """Получить список сотрудников. fired_only=True — только уволенные; False (умолч.) — только активные.""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() sql = "SELECT * FROM employees WHERE 1=1" params = [] if fired_only: sql += " AND is_fired=1" else: sql += " AND (is_fired IS NULL OR is_fired=0)" if department and department != "Все": sql += " AND department=?" params.append(department) if only_managers: sql += " AND is_manager=1" sql += " ORDER BY name" cursor.execute(sql, params) return [dict(row) for row in cursor.fetchall()] def get_employee_by_id(self, emp_id) -> Optional[Employee]: """Получить сотрудника по ID""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute("SELECT * FROM employees WHERE id=?", (emp_id,)) row = cursor.fetchone() return Employee.from_row(row) if row else None def get_employee_by_name(self, name) -> Optional[Employee]: """Получить сотрудника по имени""" if not name or not name.strip(): return None with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute("SELECT * FROM employees WHERE name=?", (name.strip(),)) row = cursor.fetchone() return Employee.from_row(row) if row else None def get_employee_id_by_name(self, name): """Получить ID сотрудника по имени""" emp = self.get_employee_by_name(name) return emp.id if emp else None def upsert_employee(self, data): """Добавить или обновить сотрудника""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() emp_id = data.get('id') if emp_id: # UPDATE cursor.execute(""" UPDATE employees SET name=?, email=?, tg_id=?, department=?, phone=?, managers=?, is_manager=?, is_temp=?, birth_date=?, category=?, is_reso=?, working_group_id=?, base_salary=?, is_fired=?, fired_date=?, schedule_pattern=?, svc_types=? WHERE id=? """, ( data['name'], data.get('email'), data.get('tg_id'), data.get('department'), data.get('phone'), data.get('managers'), data.get('is_manager', 0), data.get('is_temp', 0), data.get('birth_date'), data.get('category'), data.get('is_reso', 0), data.get('working_group_id'), data.get('base_salary', 0), data.get('is_fired', 0), data.get('fired_date'), data.get('schedule_pattern', '5/2_mon_fri'), data.get('svc_types', '[]'), emp_id )) else: # INSERT cursor.execute(""" INSERT INTO employees ( name, email, tg_id, department, phone, managers, is_manager, is_temp, birth_date, category, is_reso, working_group_id, base_salary, is_fired, fired_date, schedule_pattern, svc_types ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) """, ( data['name'], data.get('email'), data.get('tg_id'), data.get('department'), data.get('phone'), data.get('managers'), data.get('is_manager', 0), data.get('is_temp', 0), data.get('birth_date'), data.get('category'), data.get('is_reso', 0), data.get('working_group_id'), data.get('base_salary', 0), data.get('is_fired', 0), data.get('fired_date'), data.get('schedule_pattern', '5/2_mon_fri'), data.get('svc_types', '[]') )) emp_id = cursor.lastrowid conn.commit() return emp_id def delete_employee(self, emp_id): """Удалить сотрудника""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute("DELETE FROM employees WHERE id=?", (emp_id,)) conn.commit() # === РАБОЧИЕ ГРУППЫ === def get_working_groups(self) -> List[WorkingGroup]: """Получить список рабочих групп""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute("SELECT * FROM working_groups ORDER BY name") return [WorkingGroup.from_row(row) for row in cursor.fetchall()] def get_working_group_tg_id(self, group_id): """Получить Telegram ID рабочей группы""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute("SELECT tg_id FROM working_groups WHERE id=?", (group_id,)) row = cursor.fetchone() return row['tg_id'] if row else None def upsert_working_group(self, *args, **kwargs): """ Добавить или обновить рабочую группу Поддерживает два варианта вызова: 1. upsert_working_group(data_dict) - новый стиль 2. upsert_working_group(id, name, tg_id, department) - старый стиль """ # Определяем какой стиль вызова используется if len(args) == 1 and isinstance(args[0], dict): # Новый стиль: один словарь data = args[0] elif len(args) >= 2: # Старый стиль: отдельные параметры (id, name, tg_id, department) group_id = args[0] if args[0] else None name = args[1] if len(args) > 1 else None tg_id = args[2] if len(args) > 2 else None department = args[3] if len(args) > 3 else None data = { 'id': group_id, 'name': name, 'tg_id': tg_id, 'department': department } else: raise ValueError("Invalid arguments for upsert_working_group") with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() group_id = data.get('id') if group_id: cursor.execute(""" UPDATE working_groups SET name=?, tg_id=?, department=? WHERE id=? """, (data['name'], data.get('tg_id'), data.get('department'), group_id)) else: cursor.execute(""" INSERT INTO working_groups (name, tg_id, department) VALUES (?, ?, ?) """, (data['name'], data.get('tg_id'), data.get('department'))) group_id = cursor.lastrowid conn.commit() return group_id def delete_working_group(self, group_id): """Удалить рабочую группу""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute("DELETE FROM working_groups WHERE id=?", (group_id,)) conn.commit() # === ЗАДАЧИ === def get_all_tasks(self, filters=None): """Получить все задачи управления (для планировщика)""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() sql = """ SELECT mt.*, e.name as assignee, mt.title as title, mt.due_date as due_date, COALESCE(mt.priority, 'Medium') as priority, COALESCE(mt.status, 'New') as status, COALESCE(mt.comment, '') as comment, mt.id as id, mt.parent_id as parent_id FROM management_tasks mt LEFT JOIN employees e ON mt.employee_id = e.id WHERE 1=1 """ params = [] if filters: if filters.get('assignee') and filters['assignee'] != "Все": sql += " AND e.name = ?" params.append(filters['assignee']) if filters.get('status') and filters['status'] != "Все": sql += " AND mt.status = ?" params.append(filters['status']) if filters.get('priority') and filters['priority'] != "Все": sql += " AND mt.priority = ?" params.append(filters['priority']) if filters.get('search'): sql += " AND (mt.title LIKE ? OR mt.client_name LIKE ? OR mt.comment LIKE ?)" search_param = f"%{filters['search']}%" params.extend([search_param, search_param, search_param]) sql += " ORDER BY mt.due_date, mt.id" cursor.execute(sql, params) tasks = [dict(row) for row in cursor.fetchall()] # Добавляем поля created_at и updated_at если их нет for task in tasks: if 'created_at' not in task: task['created_at'] = task.get('due_date', '') if 'updated_at' not in task: task['updated_at'] = None return tasks def get_task_by_id(self, task_id): """Получить задачу управления по ID""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute(""" SELECT mt.*, e.name as assignee FROM management_tasks mt LEFT JOIN employees e ON mt.employee_id = e.id WHERE mt.id = ? """, (task_id,)) row = cursor.fetchone() return dict(row) if row else None def add_task(self, title, priority, assignee, due_date, status="New", parent_id=0, comment="", client_name=""): """Добавить задачу управления""" # Валидация даты try: due_date = self._validate_date(due_date, "due_date") except ValidationError as e: logger.error(f"Invalid due_date: {e}") raise with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() # Получаем employee_id по имени в той же транзакции emp_id = None if assignee: cursor.execute("SELECT id FROM employees WHERE name=?", (assignee,)) emp_row = cursor.fetchone() emp_id = emp_row['id'] if emp_row else None # Вставляем задачу в management_tasks cursor.execute(""" INSERT INTO management_tasks ( title, priority, employee_id, due_date, status, parent_id, comment, client_name, created_at ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, datetime('now')) """, (title, priority, emp_id, due_date, status, parent_id or 0, comment or "", client_name or "")) conn.commit() task_id = cursor.lastrowid logger.info(f"Задача управления добавлена: ID={task_id}, title='{title}'") return task_id def update_task(self, task_id, title, priority, assignee, due_date, status, comment): """Обновить задачу управления""" # Валидация даты try: due_date = self._validate_date(due_date, "due_date") except ValidationError as e: logger.error(f"Invalid due_date: {e}") raise with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() # Получаем employee_id по имени в той же транзакции emp_id = None if assignee: cursor.execute("SELECT id FROM employees WHERE name=?", (assignee,)) emp_row = cursor.fetchone() emp_id = emp_row['id'] if emp_row else None cursor.execute(""" UPDATE management_tasks SET title=?, priority=?, employee_id=?, due_date=?, status=?, comment=?, updated_at=datetime('now') WHERE id=? """, (title, priority, emp_id, due_date, status, comment or "", task_id)) conn.commit() logger.info(f"Задача управления обновлена: ID={task_id}") def update_task_status(self, task_id, status): """Обновить статус задачи""" with self.lock: thread_id = threading.get_ident() with self.pool.get_connection() as conn: cursor = conn.cursor() # ИСПРАВЛЕНО: Обновляем management_tasks вместо tasks! cursor.execute("UPDATE management_tasks SET status=?, updated_at=CURRENT_TIMESTAMP WHERE id=?", (status, task_id)) conn.commit() logger.info(f"Updated task {task_id} status to '{status}'") # Закрываем соединение этого потока, чтобы следующий запрос # получил актуальные данные (сброс кэша соединения) self.pool.close_thread_connection(thread_id) logger.debug(f"Closed connection for thread {thread_id} after status update") def update_task_comment(self, task_id, comment): """Обновить комментарий задачи""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute("UPDATE management_tasks SET comment=?, updated_at=CURRENT_TIMESTAMP WHERE id=?", (comment, task_id)) conn.commit() def update_task_date(self, task_id, new_date): """Обновить дату задачи""" try: new_date = self._validate_date(new_date, "new_date") except ValidationError as e: logger.error(f"Invalid new_date: {e}") raise with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute("UPDATE management_tasks SET due_date=?, updated_at=CURRENT_TIMESTAMP WHERE id=?", (new_date, task_id)) conn.commit() def update_task_parent(self, task_id, parent_id): """Обновить родителя задачи""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute("UPDATE management_tasks SET parent_id=?, updated_at=CURRENT_TIMESTAMP WHERE id=?", (parent_id, task_id)) conn.commit() def update_task_poll(self, task_id, poll_id): """Обновить poll_id и время отправки опроса""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute( "UPDATE management_tasks SET poll_id=?, poll_sent_at=CURRENT_TIMESTAMP, updated_at=CURRENT_TIMESTAMP WHERE id=?", (poll_id, task_id)) conn.commit() def update_task_answered_at(self, task_id): """Зафиксировать дату/время получения ответа на опрос""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute( "UPDATE management_tasks SET poll_answered_at=CURRENT_TIMESTAMP WHERE id=?", (task_id,)) conn.commit() def find_task_by_poll(self, poll_id) -> Optional[ManagementTask]: """Найти задачу по poll_id""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute("SELECT * FROM management_tasks WHERE poll_id=?", (poll_id,)) row = cursor.fetchone() return ManagementTask.from_row(row) if row else None def delete_task(self, task_id): """Удалить задачу управления""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute("DELETE FROM management_tasks WHERE id=?", (task_id,)) conn.commit() logger.info(f"Задача управления удалена: ID={task_id}") def get_children(self, parent_id): """Получить подзадачи управления""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute(""" SELECT mt.*, e.name as assignee FROM management_tasks mt LEFT JOIN employees e ON mt.employee_id = e.id WHERE mt.parent_id = ? """, (parent_id,)) return [dict(row) for row in cursor.fetchall()] def get_tasks_by_date(self, date_str): """Получить все задачи за указанную дату""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() try: cursor.execute(""" SELECT t.*, e.name as employee_name FROM tasks t LEFT JOIN employees e ON t.employee_id = e.id WHERE t.task_date = ? ORDER BY e.name, t.task_name """, (date_str,)) return [dict(row) for row in cursor.fetchall()] except Exception as e: logger.error(f"Error getting tasks by date: {e}") return [] def get_employee_id_by_name(self, name): """Получить ID сотрудника по имени""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() try: cursor.execute("SELECT id FROM employees WHERE name = ?", (name,)) row = cursor.fetchone() return row[0] if row else None except Exception as e: logger.error(f"Error getting employee ID: {e}") return None def add_task_simple(self, employee_id, task_date, task_name, client_name="", status="New", comment=""): """Добавить задачу (упрощенный метод для модуля задач)""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() try: cursor.execute(""" INSERT INTO tasks (employee_id, task_date, task_name, client_name, status, comment) VALUES (?, ?, ?, ?, ?, ?) """, (employee_id, task_date, task_name, client_name, status, comment)) conn.commit() return cursor.lastrowid except Exception as e: logger.error(f"Error adding task: {e}") conn.rollback() raise def update_task_simple(self, task_id, employee_id, task_date, task_name, client_name="", status="New", comment=""): """Обновить задачу (упрощенный метод для модуля задач)""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() try: cursor.execute(""" UPDATE tasks SET employee_id=?, task_date=?, task_name=?, client_name=?, status=?, comment=? WHERE id=? """, (employee_id, task_date, task_name, client_name, status, comment, task_id)) conn.commit() except Exception as e: logger.error(f"Error updating task: {e}") conn.rollback() raise # === КАЛЕНДАРЬ === def get_calendar_exception(self, date_str): """Получить исключение производственного календаря""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute( "SELECT * FROM production_calendar_exceptions WHERE date=?", (date_str,) ) row = cursor.fetchone() return dict(row) if row else None def set_calendar_exception(self, date_str, is_working): """Установить исключение производственного календаря""" try: date_str = self._validate_date(date_str, "date") except ValidationError as e: logger.error(f"Invalid date for calendar exception: {e}") raise with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute(""" INSERT OR REPLACE INTO production_calendar_exceptions (date, is_working) VALUES (?, ?) """, (date_str, is_working)) conn.commit() def delete_calendar_exception(self, date_str): """Удалить исключение производственного календаря""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute("DELETE FROM production_calendar_exceptions WHERE date=?", (date_str,)) conn.commit() def get_employee_calendar(self, emp_id, year=None, month=None): """Получить календарь сотрудника""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() sql = "SELECT * FROM employee_calendar WHERE employee_id=?" params = [emp_id] if year and month: sql += " AND date LIKE ?" params.append(f"{year}-{month:02d}%") elif year: sql += " AND date LIKE ?" params.append(f"{year}%") try: cursor.execute(sql, params) return [dict(row) for row in cursor.fetchall()] except Exception as e: logger.error(f"Error getting employee calendar: {e}") return [] def get_employee_calendar_range(self, emp_id, d_start, d_end): """Получить календарь сотрудника за диапазон дат""" s_start = str(d_start) s_end = str(d_end) with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() try: cursor.execute(""" SELECT * FROM employee_calendar WHERE employee_id=? AND date BETWEEN ? AND ? """, (emp_id, s_start, s_end)) return [dict(row) for row in cursor.fetchall()] except Exception as e: logger.error(f"Error getting employee calendar range: {e}") return [] def get_all_employee_calendars_range(self, d_start, d_end): """Получить данные employee_calendar для всех сотрудников за диапазон дат. Возвращает {emp_id: [dict, ...]}""" s_start = str(d_start) s_end = str(d_end) with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() try: cursor.execute(""" SELECT employee_id, date, is_working, is_absent, status FROM employee_calendar WHERE date BETWEEN ? AND ? """, (s_start, s_end)) result = {} for row in cursor.fetchall(): eid = row['employee_id'] result.setdefault(eid, []).append(dict(row)) return result except Exception as e: logger.error(f"Error getting all employee calendars range: {e}") return {} def get_days_with_tasks(self, emp_id, d_start, d_end): """ Получить дни с БАЗОВЫМИ задачами для сотрудника за период Проверяет наличие выполненных базовых услуг в указанный период. """ s_start = str(d_start) s_end = str(d_end) with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() try: # Получаем имя сотрудника по ID cursor.execute("SELECT name FROM employees WHERE id = ?", (emp_id,)) emp_row = cursor.fetchone() if not emp_row: return [] emp_name = emp_row[0] # Получаем список базовых услуг cursor.execute("SELECT name FROM services WHERE is_basic = 1") basic_services = {row[0].strip().lower() for row in cursor.fetchall()} if not basic_services: # Если нет базовых услуг в справочнике - считаем все услуги базовыми basic_services = None # Проверяем наличие таблицы workload cursor.execute(""" SELECT name FROM sqlite_master WHERE type='table' AND name='workload' """) if cursor.fetchone(): # Ищем в workload # Структура: period, emp_name, service_name, count, created_at if basic_services: # Фильтруем только базовые услуги # Используем >= и < вместо BETWEEN для корректной работы с датами разных форматов cursor.execute(""" SELECT DISTINCT created_at FROM workload WHERE emp_name = ? AND created_at >= ? AND created_at < ? """, (emp_name, s_start, s_end)) # Проверяем каждую дату на наличие базовых услуг all_dates = [row[0] for row in cursor.fetchall()] result_dates = [] for date in all_dates: cursor.execute(""" SELECT service_name, count FROM workload WHERE emp_name = ? AND created_at = ? """, (emp_name, date)) has_basic = False for svc_row in cursor.fetchall(): svc_name = svc_row[0].strip().lower() count = float(svc_row[1]) if svc_row[1] else 0 if svc_name in basic_services and count > 0: has_basic = True break if has_basic: # Обрезаем время, оставляем только дату YYYY-MM-DD date_only = date[:10] if len(date) >= 10 else date result_dates.append(date_only) return result_dates else: # Нет базовых услуг - считаем любую работу # Используем >= и < вместо BETWEEN для корректной работы с датами разных форматов cursor.execute(""" SELECT DISTINCT created_at FROM workload WHERE emp_name = ? AND created_at >= ? AND created_at < ? """, (emp_name, s_start, s_end)) # Обрезаем время, оставляем только даты return [row[0][:10] if len(row[0]) >= 10 else row[0] for row in cursor.fetchall()] else: # Используем таблицу tasks cursor.execute(""" SELECT DISTINCT task_date FROM tasks WHERE employee_id=? AND task_date BETWEEN ? AND ? """, (emp_id, s_start, s_end)) return [row[0] for row in cursor.fetchall()] except Exception as e: logger.error(f"Error getting days with tasks: {e}") return [] def get_all_days_with_tasks_range(self, d_start, d_end): """Получить дни с базовыми задачами для всех сотрудников за диапазон дат. Возвращает {emp_id: set(date_str)}""" s_start = str(d_start) s_end = str(d_end) with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() try: cursor.execute("SELECT name FROM services WHERE is_basic = 1") basic_services = {row[0].strip().lower() for row in cursor.fetchall()} if not basic_services: basic_services = None cursor.execute( "SELECT name FROM sqlite_master WHERE type='table' AND name='workload'" ) if cursor.fetchone(): if basic_services: cursor.execute(""" SELECT DISTINCT e.id AS emp_id, SUBSTR(w.created_at, 1, 10) AS date_only FROM workload w JOIN employees e ON e.name = w.emp_name WHERE w.created_at >= ? AND w.created_at < ? AND EXISTS ( SELECT 1 FROM services s WHERE s.is_basic = 1 AND LOWER(TRIM(s.name)) = LOWER(TRIM(w.service_name)) AND CAST(w.count AS REAL) > 0 ) """, (s_start, s_end)) else: cursor.execute(""" SELECT DISTINCT e.id AS emp_id, SUBSTR(w.created_at, 1, 10) AS date_only FROM workload w JOIN employees e ON e.name = w.emp_name WHERE w.created_at >= ? AND w.created_at < ? """, (s_start, s_end)) else: cursor.execute(""" SELECT DISTINCT employee_id AS emp_id, task_date AS date_only FROM tasks WHERE task_date BETWEEN ? AND ? """, (s_start, s_end)) result = {} for row in cursor.fetchall(): eid = row['emp_id'] result.setdefault(eid, set()).add(row['date_only']) return result except Exception as e: logger.error(f"Error getting all days with tasks range: {e}") return {} def get_employee_day_override(self, emp_id, date_str): """Получить переопределение рабочего дня сотрудника""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute( "SELECT is_working FROM employee_calendar WHERE employee_id=? AND date=?", (emp_id, date_str) ) row = cursor.fetchone() return row['is_working'] if row else None def set_employee_day_override(self, emp_id, date_str, is_working): """Установить переопределение рабочего дня сотрудника""" # КРИТИЧЕСКАЯ ВАЛИДАЦИЯ ДАТЫ try: date_str = self._validate_date(date_str, "date") except ValidationError as e: logger.error(f"Invalid date for employee calendar: {e}") raise with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() # Проверяем есть ли запись cursor.execute( "SELECT is_absent FROM employee_calendar WHERE employee_id=? AND date=?", (emp_id, date_str) ) row = cursor.fetchone() if row: # Обновляем существующую запись, сохраняя is_absent, сбрасываем статус cursor.execute(""" UPDATE employee_calendar SET is_working=?, status=NULL WHERE employee_id=? AND date=? """, (is_working, emp_id, date_str)) else: # Создаём новую запись cursor.execute(""" INSERT INTO employee_calendar (employee_id, date, is_working, is_absent, status) VALUES (?, ?, ?, 0, NULL) """, (emp_id, date_str, is_working)) conn.commit() def get_employee_absent(self, emp_id, date_str): """Получить отметку об отсутствии сотрудника ('Н')""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute( "SELECT is_absent FROM employee_calendar WHERE employee_id=? AND date=?", (emp_id, date_str) ) row = cursor.fetchone() return row['is_absent'] if row else 0 def set_employee_absent(self, emp_id, date_str, is_absent): """Установить отметку об отсутствии сотрудника ('Н')""" # КРИТИЧЕСКАЯ ВАЛИДАЦИЯ ДАТЫ try: date_str = self._validate_date(date_str, "date") except ValidationError as e: logger.error(f"Invalid date for employee calendar: {e}") raise with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() # Проверяем есть ли запись cursor.execute( "SELECT is_working FROM employee_calendar WHERE employee_id=? AND date=?", (emp_id, date_str) ) row = cursor.fetchone() if row: # Обновляем существующую запись cursor.execute(""" UPDATE employee_calendar SET is_absent=? WHERE employee_id=? AND date=? """, (is_absent, emp_id, date_str)) else: # Создаём новую запись (день рабочий по умолчанию, но человек отсутствует) cursor.execute(""" INSERT INTO employee_calendar (employee_id, date, is_working, is_absent) VALUES (?, ?, 1, ?) """, (emp_id, date_str, is_absent)) conn.commit() def delete_employee_day_override(self, emp_id, date_str): """Удалить переопределение дня для сотрудника (сброс к календарю)""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() try: cursor.execute(""" DELETE FROM employee_calendar WHERE employee_id=? AND date=? """, (emp_id, date_str)) conn.commit() logger.info(f"Deleted day override for employee {emp_id}, date {date_str}") except Exception as e: logger.error(f"Error deleting day override: {e}") conn.rollback() def get_employee_manual_absent(self, emp_id, date_str): """Проверить ручную отметку отсутствия (Н)""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() try: cursor.execute(""" SELECT is_absent FROM employee_calendar WHERE employee_id=? AND date=? """, (emp_id, date_str)) row = cursor.fetchone() return row[0] if row and 'is_absent' in row.keys() else 0 except Exception as e: logger.error(f"Error getting manual absent: {e}") return 0 def set_employee_manual_absent(self, emp_id, date_str, is_absent=1): """Установить ручную отметку отсутствия (Н)""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() try: # Используем INSERT OR REPLACE, сбрасываем статус К/Б при смене cursor.execute(""" INSERT OR REPLACE INTO employee_calendar (employee_id, date, is_working, is_absent, status) VALUES (?, ?, COALESCE((SELECT is_working FROM employee_calendar WHERE employee_id=? AND date=?), 1), ?, NULL) """, (emp_id, date_str, emp_id, date_str, is_absent)) conn.commit() logger.info(f"Set manual absent for employee {emp_id}, date {date_str}, is_absent={is_absent}") except Exception as e: logger.error(f"Error setting manual absent: {e}") conn.rollback() def set_employee_status(self, emp_id: int, date_str: str, status: str) -> None: """Store a custom day status ('К' or 'Б') — marks day working + clears absent.""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() try: cursor.execute(""" INSERT OR REPLACE INTO employee_calendar (employee_id, date, is_working, is_absent, status) VALUES (?, ?, 1, 0, ?) """, (emp_id, date_str, status)) conn.commit() except Exception as e: logger.error(f"Error setting employee status: {e}") conn.rollback() def set_calendar_exception(self, date_str, is_working): """Установить исключение в производственном календаре (праздник/перенос)""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() try: # Создаем таблицу если её нет cursor.execute(""" CREATE TABLE IF NOT EXISTS production_calendar ( date TEXT PRIMARY KEY, is_working INTEGER DEFAULT 1, description TEXT ) """) cursor.execute(""" INSERT OR REPLACE INTO production_calendar (date, is_working) VALUES (?, ?) """, (date_str, is_working)) conn.commit() logger.info(f"Set calendar exception: {date_str} = {is_working}") except Exception as e: logger.error(f"Error setting calendar exception: {e}") conn.rollback() # === ОТПУСКА === def get_vacations(self, department=None, year=None): """Получить список отпусков""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() sql = """ SELECT v.*, e.name as emp_name, e.name, e.department FROM vacations v JOIN employees e ON v.employee_id = e.id WHERE 1=1 """ params = [] if department and department != "Все": sql += " AND e.department = ?" params.append(department) if year: sql += " AND (v.start_date LIKE ? OR v.end_date LIKE ?)" params.extend([f"{year}%", f"{year}%"]) sql += " ORDER BY v.start_date" cursor.execute(sql, params) return [dict(row) for row in cursor.fetchall()] def get_employee_vacations_range(self, emp_id, d_start, d_end): """Получить отпуска сотрудника в диапазоне дат. Записи с is_working=1 не возвращаются, так как сотрудник фактически работает и не считается «в отпуске». """ s_start = str(d_start) s_end = str(d_end) with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() try: cursor.execute(""" SELECT start_date as start, end_date as end, v_type FROM vacations WHERE employee_id=? AND NOT (end_date < ? OR start_date > ?) AND (is_working IS NULL OR is_working = 0) """, (emp_id, s_start, s_end)) return [VacationRange.from_row(row) for row in cursor.fetchall()] except Exception as e: logger.error(f"Error getting vacation range: {e}") return [] def get_all_employee_vacations_range(self, d_start, d_end): """Получить отпуска всех сотрудников в диапазоне дат. Возвращает {emp_id: [{'start':..,'end':..,'v_type':...}, ...]} Сотрудники с is_working=1 («Работает») НЕ включаются: они продолжают быть доступны для умного распределения несмотря на запись об отпуске. """ s_start = str(d_start) s_end = str(d_end) with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() try: cursor.execute(""" SELECT employee_id, start_date AS start, end_date AS end, v_type FROM vacations WHERE NOT (end_date < ? OR start_date > ?) AND (is_working IS NULL OR is_working = 0) """, (s_start, s_end)) result: Dict[int, List[VacationRange]] = {} for row in cursor.fetchall(): eid = row['employee_id'] result.setdefault(eid, []).append(VacationRange.from_row(row)) return result except Exception as e: logger.error(f"Error getting all employee vacations range: {e}") return {} def add_vacation(self, emp_id, start_date, end_date, v_type, days, is_working=False): """Добавить отпуск. is_working=True означает «отпуск числится, но сотрудник работает» (производственная необходимость). Такой сотрудник не исключается из умного распределения. """ try: start_date = self._validate_date(start_date, "start_date") end_date = self._validate_date(end_date, "end_date") except ValidationError as e: logger.error(f"Invalid vacation dates: {e}") raise with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute(""" INSERT INTO vacations (employee_id, start_date, end_date, v_type, days, is_working) VALUES (?, ?, ?, ?, ?, ?) """, (emp_id, start_date, end_date, v_type, days, 1 if is_working else 0)) conn.commit() return cursor.lastrowid def update_vacation(self, vac_id, emp_id, start_date, end_date, v_type, days=None, is_working=False): """Обновить отпуск. is_working=True — «отпуск числится, но сотрудник работает». """ try: start_date = self._validate_date(start_date, "start_date") end_date = self._validate_date(end_date, "end_date") except ValidationError as e: logger.error(f"Invalid vacation dates: {e}") raise if days is None: d_start = datetime.strptime(start_date, "%Y-%m-%d") d_end = datetime.strptime(end_date, "%Y-%m-%d") days = (d_end - d_start).days + 1 with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute(""" UPDATE vacations SET employee_id=?, start_date=?, end_date=?, v_type=?, days=?, is_working=? WHERE id=? """, (emp_id, start_date, end_date, v_type, days, 1 if is_working else 0, vac_id)) conn.commit() # ── Ограничения умного распределения ────────────────────────────────────── def get_distribution_restrictions(self) -> list: """Вернуть все ограничения в виде списка dict. Каждый dict содержит: id, expert_name, blocked_services/regions/risks/prefixes (list), mode_services/regions/risks/prefixes ('block'|'only') """ import json as _json with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute( "SELECT id, expert_name, blocked_services, blocked_regions, " "blocked_risks, blocked_prefixes, " "mode_services, mode_regions, mode_risks, mode_prefixes " "FROM distribution_restrictions ORDER BY expert_name" ) result = [] for row in cursor.fetchall(): result.append({ 'id': row['id'], 'expert_name': row['expert_name'], 'blocked_services': _json.loads(row['blocked_services'] or '[]'), 'blocked_regions': _json.loads(row['blocked_regions'] or '[]'), 'blocked_risks': _json.loads(row['blocked_risks'] or '[]'), 'blocked_prefixes': _json.loads(row['blocked_prefixes'] or '[]'), 'mode_services': row['mode_services'] or 'block', 'mode_regions': row['mode_regions'] or 'block', 'mode_risks': row['mode_risks'] or 'block', 'mode_prefixes': row['mode_prefixes'] or 'block', }) return result def get_restriction_by_id(self, restriction_id: int) -> dict | None: """Вернуть одну запись об ограничении по id.""" import json as _json with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute( "SELECT id, expert_name, blocked_services, blocked_regions, " "blocked_risks, blocked_prefixes, " "mode_services, mode_regions, mode_risks, mode_prefixes " "FROM distribution_restrictions WHERE id = ?", (restriction_id,) ) row = cursor.fetchone() if not row: return None return { 'id': row['id'], 'expert_name': row['expert_name'], 'blocked_services': _json.loads(row['blocked_services'] or '[]'), 'blocked_regions': _json.loads(row['blocked_regions'] or '[]'), 'blocked_risks': _json.loads(row['blocked_risks'] or '[]'), 'blocked_prefixes': _json.loads(row['blocked_prefixes'] or '[]'), 'mode_services': row['mode_services'] or 'block', 'mode_regions': row['mode_regions'] or 'block', 'mode_risks': row['mode_risks'] or 'block', 'mode_prefixes': row['mode_prefixes'] or 'block', } def save_restriction(self, expert_name: str, blocked_services: list, blocked_regions: list, blocked_risks: list, blocked_prefixes: list = None, mode_services: str = 'block', mode_regions: str = 'block', mode_risks: str = 'block', mode_prefixes: str = 'block') -> None: """Создать или обновить ограничение для эксперта (UPSERT по expert_name).""" import json as _json with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute(""" INSERT INTO distribution_restrictions (expert_name, blocked_services, blocked_regions, blocked_risks, blocked_prefixes, mode_services, mode_regions, mode_risks, mode_prefixes) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?) ON CONFLICT(expert_name) DO UPDATE SET blocked_services = excluded.blocked_services, blocked_regions = excluded.blocked_regions, blocked_risks = excluded.blocked_risks, blocked_prefixes = excluded.blocked_prefixes, mode_services = excluded.mode_services, mode_regions = excluded.mode_regions, mode_risks = excluded.mode_risks, mode_prefixes = excluded.mode_prefixes """, ( expert_name, _json.dumps(blocked_services, ensure_ascii=False), _json.dumps(blocked_regions, ensure_ascii=False), _json.dumps(blocked_risks, ensure_ascii=False), _json.dumps(blocked_prefixes or [], ensure_ascii=False), mode_services, mode_regions, mode_risks, mode_prefixes, )) conn.commit() def get_distinct_prefixes(self) -> list: """Список уникальных нормализованных префиксов из deadline_cases. Использует ту же логику нормализации, что и SmartDistributor._get_prefix(). """ import re as _re FOCUS = {'КВ', 'ДЧ', 'ИН'} prefixes: set = set() with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute( "SELECT DISTINCT case_number FROM deadline_cases " "WHERE case_number IS NOT NULL LIMIT 20000" ) for row in cursor.fetchall(): cn = (row['case_number'] or '').strip().upper() m = _re.match(r'^([А-ЯЁA-Z]+)', cn) if not m: prefixes.add('Без букв') else: pfx = m.group(1) if pfx in FOCUS: prefixes.add(pfx) else: prefixes.add(f'Другое({pfx})') return sorted(prefixes) def delete_restriction(self, restriction_id: int) -> None: """Удалить ограничение по id.""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute( "DELETE FROM distribution_restrictions WHERE id = ?", (restriction_id,) ) conn.commit() # ── Матрица маршрутизации (distribution_rules) ──────────────────────────── def get_distribution_rules(self) -> list: """Вернуть все правила маршрутизации. Каждый элемент — dict с полями: id, expert_name, service_types, regions, risks, prefixes (каждое — dict вида {"type":"deny"|"only","values":[...]} или None), is_active (bool). """ import json as _json def _parse(v): try: return _json.loads(v) if v else None except Exception: return None with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute( "SELECT id, expert_name, service_types, regions, risks, prefixes, is_active " "FROM distribution_rules ORDER BY expert_name" ) return [ { 'id': row['id'], 'expert_name': row['expert_name'], 'service_types': _parse(row['service_types']), 'regions': _parse(row['regions']), 'risks': _parse(row['risks']), 'prefixes': _parse(row['prefixes']), 'is_active': bool(row['is_active']), } for row in cursor.fetchall() ] def save_distribution_rule(self, expert_name: str, service_types=None, regions=None, risks=None, prefixes=None, is_active: bool = True, rule_id: int | None = None) -> int: """Создать или обновить правило маршрутизации. Если rule_id задан — UPDATE по id, иначе INSERT. Возвращает id строки.""" import json as _json def _enc(v): return _json.dumps(v, ensure_ascii=False) if v is not None else None with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() if rule_id is not None: cursor.execute(""" UPDATE distribution_rules SET expert_name = ?, service_types = ?, regions = ?, risks = ?, prefixes = ?, is_active = ? WHERE id = ? """, (expert_name, _enc(service_types), _enc(regions), _enc(risks), _enc(prefixes), 1 if is_active else 0, rule_id)) else: cursor.execute(""" INSERT INTO distribution_rules (expert_name, service_types, regions, risks, prefixes, is_active) VALUES (?, ?, ?, ?, ?, ?) """, (expert_name, _enc(service_types), _enc(regions), _enc(risks), _enc(prefixes), 1 if is_active else 0)) conn.commit() return cursor.lastrowid or rule_id def toggle_distribution_rule(self, rule_id: int) -> bool: """Переключить is_active правила. Возвращает новое значение.""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute( "UPDATE distribution_rules SET is_active = 1 - is_active WHERE id = ?", (rule_id,) ) conn.commit() cursor.execute( "SELECT is_active FROM distribution_rules WHERE id = ?", (rule_id,) ) row = cursor.fetchone() return bool(row['is_active']) if row else False def delete_distribution_rule(self, rule_id: int) -> None: """Удалить правило маршрутизации по id.""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute( "DELETE FROM distribution_rules WHERE id = ?", (rule_id,) ) conn.commit() # ── Manual Surcharges CRUD ─────────────────────────────────────────── def get_manual_surcharges(self, year: int, month: int) -> list[ManualSurcharge]: """Вернуть все ручные доплаты за указанный период (year-month).""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() prefix = f"{year:04d}-{month:02d}" cursor.execute( "SELECT * FROM manual_surcharges WHERE period_date LIKE ? ORDER BY id", (prefix + '%',) ) return [ManualSurcharge.from_row(row) for row in cursor.fetchall()] def save_manual_surcharge(self, rec: ManualSurcharge) -> int: """Создать или обновить запись ручной доплаты. Возвращает id.""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() if rec.id and rec.id > 0: cursor.execute(""" UPDATE manual_surcharges SET expert_name = ?, department = ?, case_number = ?, service_name = ?, base_cost = ?, hourly_rate = ?, hours = ?, total_amount = ?, operation_type = ?, note = ?, period_date = ? WHERE id = ? """, ( rec.expert_name, rec.department, rec.case_number, rec.service_name, rec.base_cost, rec.hourly_rate, rec.hours, rec.total_amount, rec.operation_type, rec.note, rec.period_date, rec.id )) conn.commit() return rec.id else: cursor.execute(""" INSERT INTO manual_surcharges (expert_name, department, case_number, service_name, base_cost, hourly_rate, hours, total_amount, operation_type, note, period_date) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) """, ( rec.expert_name, rec.department, rec.case_number, rec.service_name, rec.base_cost, rec.hourly_rate, rec.hours, rec.total_amount, rec.operation_type, rec.note, rec.period_date )) conn.commit() return cursor.lastrowid def delete_manual_surcharge(self, rec_id: int) -> None: """Удалить запись ручной доплаты по id.""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute("DELETE FROM manual_surcharges WHERE id = ?", (rec_id,)) conn.commit() def get_distinct_regions(self) -> list: """Список уникальных регионов из deadline_cases (для настройки ограничений).""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute(""" SELECT DISTINCT region FROM deadline_cases WHERE region IS NOT NULL AND trim(region) != '' ORDER BY region """) return [row['region'] for row in cursor.fetchall()] def get_distinct_risks(self) -> list: """Список уникальных рисков из deadline_cases (для настройки ограничений).""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute(""" SELECT DISTINCT risk FROM deadline_cases WHERE risk IS NOT NULL AND trim(risk) != '' ORDER BY risk """) return [row['risk'] for row in cursor.fetchall()] def get_distinct_prefixes_normalized(self) -> list: """Уникальные нормализованные префиксы из case_number (Без букв, КВ, ДЧ, ...).""" import re as _re _focus = frozenset({'КВ', 'ДЧ', 'ИН'}) with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute( "SELECT DISTINCT substr(case_number, 1, 6) AS pfx " "FROM deadline_cases " "WHERE case_number IS NOT NULL AND trim(case_number) != ''" ) seen: set = set() for row in cursor.fetchall(): s = (row['pfx'] or '').strip() if not s: continue m = _re.match(r'^([А-ЯA-ZЁ]{1,4})\d', s) if m: p = m.group(1).upper() seen.add(p if p in _focus else f'Другое({p})') elif _re.match(r'^\d', s): seen.add('Без букв') return sorted(seen) # ── Лимиты распределения ───────────────────────────────────────────────── def get_all_distribution_limits(self) -> dict: """Вернуть все лимиты: {expert_name: max_cases}.""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute( "SELECT expert_name, max_cases FROM distribution_limits" ) return {row['expert_name']: row['max_cases'] for row in cursor.fetchall()} def save_distribution_limits_bulk(self, limits: dict) -> None: """Сохранить лимиты оптом {expert_name: max_cases}. Полностью заменяет предыдущий набор. Чтобы сбросить все лимиты, передайте пустой словарь. """ with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute("DELETE FROM distribution_limits") for expert, max_cases in limits.items(): if max_cases and int(max_cases) > 0: cursor.execute( "INSERT INTO distribution_limits " "(expert_name, max_cases) VALUES (?, ?)", (expert, int(max_cases)) ) conn.commit() def get_limits_with_workload(self) -> list: """Return all active employees with in-work case count and their limit. Each item: {id, expert_name, in_work, max_limit (int|None)}. in_work = cases assigned to expert on the latest load_date. """ with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() # Latest load date cursor.execute("SELECT MAX(load_date) AS ld FROM deadline_cases") ld_row = cursor.fetchone() latest_load = ld_row['ld'] if ld_row else None # Active employees cursor.execute( "SELECT id, name FROM employees " "WHERE (is_fired IS NULL OR is_fired = 0) ORDER BY name" ) employees = [(r['id'], r['name']) for r in cursor.fetchall()] # Cases in work per expert on latest load in_work_map: dict = {} if latest_load: cursor.execute( "SELECT expert_calculation, COUNT(*) AS cnt " "FROM deadline_cases " "WHERE load_date = ? " " AND expert_calculation IS NOT NULL " " AND lower(trim(expert_calculation)) NOT IN ('', 'nan', 'none', 'nat') " "GROUP BY expert_calculation", (latest_load,) ) for r in cursor.fetchall(): exp = (r['expert_calculation'] or '').strip() if exp: in_work_map[exp] = r['cnt'] # Existing limits cursor.execute("SELECT expert_name, max_cases FROM distribution_limits") limits_map = {r['expert_name']: r['max_cases'] for r in cursor.fetchall()} return [ { 'id': emp_id, 'expert_name': emp_name, 'in_work': in_work_map.get(emp_name, 0), 'max_limit': limits_map.get(emp_name), # None = no limit } for emp_id, emp_name in employees ] def set_distribution_limit(self, expert_name: str, max_limit) -> None: """Set or clear the per-expert distribution limit. max_limit=None (or empty string) removes the limit row entirely. """ with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() if max_limit is None or max_limit == '': cursor.execute( "DELETE FROM distribution_limits WHERE expert_name = ?", (expert_name,) ) else: cursor.execute(""" INSERT INTO distribution_limits (expert_name, max_cases) VALUES (?, ?) ON CONFLICT(expert_name) DO UPDATE SET max_cases = excluded.max_cases """, (expert_name, int(max_limit))) conn.commit() def get_expert_main_loads(self) -> dict: """Текущая загрузка экспертов по последней дате без учёта follow-up заявок. Follow-up заявки (Доп калька / Доп Смета (имущ) / Доп Калька (имущ)) не входят в лимит, поэтому не учитываются при подсчёте. Возвращает {expert_name: count}. """ FOLLOWUP = ('Доп калька', 'Доп Смета (имущ)', 'Доп Калька (имущ)') with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute( "SELECT MAX(load_date) as latest FROM deadline_cases" ) row = cursor.fetchone() if not row or not row['latest']: return {} latest = row['latest'] cursor.execute( "SELECT expert_calculation, payment_info " "FROM deadline_cases " "WHERE load_date = ? " " AND expert_calculation IS NOT NULL " " AND trim(expert_calculation) NOT IN ('', 'nan', 'None', 'NaT')", (latest,) ) loads: dict = {} for r in cursor.fetchall(): exp = (r['expert_calculation'] or '').strip() if not exp or exp.lower() in ('nan', 'none', 'nat'): continue pinfo = r['payment_info'] or '' if any(svc in pinfo for svc in FOLLOWUP): continue loads[exp] = loads.get(exp, 0) + 1 return loads def get_case_last_expert(self, case_number: str) -> str | None: """Найти последнего эксперта, работавшего с указанным делом. Используется для follow-up услуг (Доп калька, Доп Смета (имущ), Доп Калька (имущ)): такие заявки должны назначаться тому же эксперту, который ранее выполнял базовую услугу по этому делу. Возвращает имя эксперта (строку) или None, если история не найдена. """ with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute(""" SELECT expert_calculation FROM deadline_cases WHERE case_number = ? AND expert_calculation IS NOT NULL AND trim(expert_calculation) != '' AND lower(trim(expert_calculation)) NOT IN ('nan', 'none', 'nat') ORDER BY load_date DESC LIMIT 1 """, (case_number,)) row = cursor.fetchone() if row: val = (row['expert_calculation'] or '').strip() return val if val else None return None def get_address_expert(self, address: str) -> str | None: """Найти эксперта, который последним работал с заявкой по указанному адресу. Поиск ведётся по всей истории БД (все даты загрузки). Используется точное совпадение адреса (без учёта регистра и пробелов). Возвращает имя эксперта или None, если совпадений нет. """ if not address or not address.strip(): return None norm = address.strip().lower() with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute(""" SELECT expert_calculation FROM deadline_cases WHERE lower(trim(address)) = ? AND expert_calculation IS NOT NULL AND trim(expert_calculation) != '' AND lower(trim(expert_calculation)) NOT IN ('nan', 'none', 'nat') ORDER BY load_date DESC, id DESC LIMIT 1 """, (norm,)) row = cursor.fetchone() if row: val = (row['expert_calculation'] or '').strip() return val if val else None return None def get_vacation_by_id(self, vac_id): """Получить запись об отпуске по id.""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute("SELECT * FROM vacations WHERE id=?", (vac_id,)) row = cursor.fetchone() return dict(row) if row else None def delete_vacation(self, vac_id): """Удалить отпуск""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute("DELETE FROM vacations WHERE id=?", (vac_id,)) conn.commit() def delete_vacations_in_range(self, emp_id: int, start_date: str, end_date: str) -> int: """Delete all vacation records that overlap [start_date, end_date] for emp_id.""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() # Overlap condition: vac.start <= range.end AND vac.end >= range.start cursor.execute(""" DELETE FROM vacations WHERE employee_id = ? AND start_date <= ? AND end_date >= ? """, (emp_id, end_date, start_date)) conn.commit() return cursor.rowcount def get_department_overlaps(self, start_date=None, end_date=None, department=None, exclude_id=None): """ Получить пересечения отпусков по отделу Возвращает отпуска других сотрудников из того же отдела, которые пересекаются с указанным периодом. Args: start_date: начало периода (YYYY-MM-DD) end_date: конец периода (YYYY-MM-DD) department: название отдела exclude_id: ID сотрудника, которого нужно ИСКЛЮЧИТЬ из результатов """ with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() sql = """ SELECT v.*, e.name as emp_name, e.name, e.department FROM vacations v JOIN employees e ON v.employee_id = e.id WHERE 1=1 """ params = [] # Фильтр по отделу if department and department != "Все": sql += " AND e.department = ?" params.append(department) # ВАЖНО: Исключаем самого сотрудника if exclude_id is not None: sql += " AND v.employee_id != ?" params.append(exclude_id) # Фильтр по пересечению дат # Отпуск пересекается если: # (v.start_date <= end_date) AND (v.end_date >= start_date) if start_date and end_date: sql += " AND v.start_date <= ? AND v.end_date >= ?" params.append(end_date) params.append(start_date) sql += " ORDER BY v.start_date" try: cursor.execute(sql, params) return [dict(row) for row in cursor.fetchall()] except Exception as e: logger.error(f"Error getting department overlaps: {e}") return [] # === УСЛУГИ === def get_services(self): """Получить список услуг""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute("SELECT * FROM services ORDER BY name") return [dict(row) for row in cursor.fetchall()] def add_service(self, name, is_base=1): """Добавить услугу""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() try: cursor.execute("INSERT INTO services (name, is_base) VALUES (?, ?)", (name, is_base)) conn.commit() return cursor.lastrowid except sqlite3.IntegrityError: logger.warning(f"Service '{name}' already exists") return None def delete_service(self, service_id): """Удалить услугу""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute("DELETE FROM services WHERE id=?", (service_id,)) conn.commit() def get_aliases(self): """ Получить алиасы услуг для нормализации названий Возвращает словарь: {вариант_названия: основное_название} """ # Базовые алиасы для часто встречающихся вариантов aliases = { # Осмотры 'осмотр': 'Осмотр', 'осм': 'Осмотр', # Оценка 'оценка': 'Оценка', 'оцен': 'Оценка', # Акт 'акт': 'Акт (имущ)', 'акт имущ': 'Акт (имущ)', # Калька 'калька': 'Калька (имущ)', 'калька имущ': 'Калька (имущ)', # Выезд 'выезд': 'Выезд', # Километраж 'километраж': 'Выезд (километраж)', 'км': 'Выезд (километраж)', # ПСО ЮЛ 'псо юл': 'ПСО ЮЛ', 'псо': 'ПСО ЮЛ', } # Получаем услуги из БД и добавляем их services = self.get_services() for service in services: name = service['name'] # Добавляем саму услугу (точное соответствие) aliases[name.lower()] = name aliases[name] = name return aliases # === ЗАРПЛАТА === def get_prices(self, year=None): """Получить цены""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() if year: cursor.execute("SELECT * FROM prices WHERE year=?", (year,)) else: cursor.execute("SELECT * FROM prices") return [dict(row) for row in cursor.fetchall()] def clear_prices(self, year=None): """Удалить все цены (опционально за конкретный год)""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() if year: cursor.execute("DELETE FROM prices WHERE year=?", (year,)) logger.info(f"Cleared prices for year {year}") else: cursor.execute("DELETE FROM prices") logger.info("Cleared all prices") conn.commit() def add_price_entry(self, name, calc_rule_month, calc_rule_salary, price_high, price_first, price_second, service_cost, payment_type, tax_percent, year=None): """ Добавить запись в справочник цен Args: name: название услуги calc_rule_month: правило расчёта для месячной загрузки calc_rule_salary: правило расчёта для зарплаты price_high: цена для высшей категории price_first: цена для первой категории price_second: цена для второй категории service_cost: стоимость услуги (выручка) payment_type: тип оплаты (0 или 1) tax_percent: процент налога year: год (опционально) """ with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute(""" INSERT INTO prices (name, calc_rule_month, calc_rule_salary, price_high, price_first, price_second, service_cost, payment_type, tax_percent, year) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?) """, (name, calc_rule_month, calc_rule_salary, price_high, price_first, price_second, service_cost, payment_type, tax_percent, year)) conn.commit() logger.info(f"Added price entry: {name} for year {year}") def add_price_entries_bulk(self, entries): """ Массовая вставка цен (оптимизировано) Args: entries: список кортежей (name, calc_rule_month, calc_rule_salary, price_high, price_first, price_second, service_cost, payment_type, tax_percent, year) """ if not entries: return with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.executemany(""" INSERT INTO prices (name, calc_rule_month, calc_rule_salary, price_high, price_first, price_second, service_cost, payment_type, tax_percent, year) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?) """, entries) conn.commit() logger.info(f"Bulk added {len(entries)} price entries") def set_price(self, year, category, price): """Установить цену""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute("SELECT id FROM prices WHERE year=? AND category=?", (year, category)) row = cursor.fetchone() if row: cursor.execute("UPDATE prices SET price=? WHERE id=?", (price, row['id'])) else: cursor.execute("INSERT INTO prices (year, category, price) VALUES (?, ?, ?)", (year, category, price)) conn.commit() def get_adjustments(self, period): """Получить корректировки зарплаты""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute("SELECT * FROM salary_adjustments WHERE period=?", (period,)) return [dict(row) for row in cursor.fetchall()] def add_adjustment(self, emp_id, period, amount, comment): """Добавить корректировку зарплаты""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute(""" INSERT INTO salary_adjustments (employee_id, period, amount, comment) VALUES (?, ?, ?, ?) """, (emp_id, period, amount, comment)) conn.commit() return cursor.lastrowid def delete_adjustment(self, adj_id): """Удалить корректировку зарплаты""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute("DELETE FROM salary_adjustments WHERE id=?", (adj_id,)) conn.commit() # === ДОПЛАТЫ ЗАРПЛАТЫ === def get_doplaty(self, period): """Получить записи доплат для периода""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute( "SELECT * FROM salary_doplaty WHERE period=? ORDER BY dept, emp_name", (period,) ) return [dict(row) for row in cursor.fetchall()] def upsert_doplata(self, rec): """Сохранить или обновить запись доплаты. rec — dict с полями таблицы.""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() if rec.get('id'): cursor.execute(""" UPDATE salary_doplaty SET period=?, employee_id=?, dept=?, emp_name=?, deal_number=?, service_name=?, base_cost=?, hours=?, hour_rate=?, doplata=?, comment=? WHERE id=? """, ( rec['period'], rec.get('employee_id'), rec.get('dept', ''), rec.get('emp_name', ''), rec.get('deal_number', ''), rec.get('service_name', ''), rec.get('base_cost', 0), rec.get('hours', 0), rec.get('hour_rate', 0), rec.get('doplata', 0), rec.get('comment', ''), rec['id'] )) else: cursor.execute(""" INSERT INTO salary_doplaty (period, employee_id, dept, emp_name, deal_number, service_name, base_cost, hours, hour_rate, doplata, comment) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) """, ( rec['period'], rec.get('employee_id'), rec.get('dept', ''), rec.get('emp_name', ''), rec.get('deal_number', ''), rec.get('service_name', ''), rec.get('base_cost', 0), rec.get('hours', 0), rec.get('hour_rate', 0), rec.get('doplata', 0), rec.get('comment', '') )) conn.commit() return cursor.lastrowid def delete_doplata(self, rec_id): """Удалить запись доплаты""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute("DELETE FROM salary_doplaty WHERE id=?", (rec_id,)) conn.commit() # === СТАТИСТИКА ЗАГРУЗКИ === def get_workload_by_year(self, year): """Получить загрузку по году""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() try: # Проверяем наличие таблицы workload cursor.execute(""" SELECT name FROM sqlite_master WHERE type='table' AND name='workload' """) if cursor.fetchone(): # Проверяем структуру таблицы workload cursor.execute("PRAGMA table_info(workload)") columns = [row[1] for row in cursor.fetchall()] # Используем правильную колонку для фильтрации по году if 'created_at' in columns: # Структура с created_at # Используем диапазон дат - работает для всех форматов cursor.execute(""" SELECT *, service_name as val, emp_name, created_at, source, count FROM workload WHERE created_at >= ? AND created_at < ? """, (f"{year}-01-01", f"{int(year)+1}-01-01")) results = [dict(row) for row in cursor.fetchall()] logger.info(f"get_workload_by_year({year}): found {len(results)} records from workload") return results elif 'date' in columns: cursor.execute(""" SELECT *, service_name as val FROM workload WHERE date LIKE ? """, (f"{year}%",)) elif 'task_date' in columns: cursor.execute(""" SELECT *, val, emp_name FROM workload WHERE task_date LIKE ? """, (f"{year}%",)) elif 'period' in columns: # Период в формате YYYY-MM cursor.execute(""" SELECT *, service_name as val, emp_name, created_at, source, count FROM workload WHERE period LIKE ? """, (f"{year}%",)) else: # Неизвестная структура - используем tasks logger.warning(f"Unknown workload structure, using tasks") cursor.execute(""" SELECT t.*, t.task_name as val, 'daily' as source, e.name as emp_name, e.name FROM tasks t LEFT JOIN employees e ON t.employee_id = e.id WHERE t.task_date LIKE ? """, (f"{year}%",)) else: # Используем таблицу tasks cursor.execute(""" SELECT t.*, t.task_name as val, 'daily' as source, e.name as emp_name, e.name FROM tasks t LEFT JOIN employees e ON t.employee_id = e.id WHERE t.task_date LIKE ? """, (f"{year}%",)) results = [dict(row) for row in cursor.fetchall()] logger.info(f"get_workload_by_year({year}): found {len(results)} records from tasks") return results except Exception as e: logger.error(f"Error getting workload by year: {e}") return [] # === НАСТРОЙКИ === def get_setting(self, key, default=""): """Получить настройку""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute("SELECT value FROM settings WHERE key=?", (key,)) row = cursor.fetchone() return row[0] if row else default def set_setting(self, key, value): """Установить настройку""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute("INSERT OR REPLACE INTO settings (key, value) VALUES (?, ?)", (key, value)) conn.commit() # === TELEGRAM СТАТИСТИКА === def log_tg_event(self, event_type, poll_id=None): """Логировать Telegram событие""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() try: cursor.execute(""" INSERT INTO tg_stats (event_type, poll_id, created_at) VALUES (?, ?, ?) """, (event_type, poll_id, datetime.now().isoformat())) conn.commit() except Exception as e: logger.error(f"Error logging TG event: {e}") def get_tg_stats_today(self): """Получить статистику Telegram за сегодня: (опросов, ответов, сообщений)""" today = datetime.now().strftime("%Y-%m-%d") with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute(""" SELECT COUNT(*) FROM tg_stats WHERE event_type='poll_sent' AND created_at LIKE ? """, (f"{today}%",)) sent = cursor.fetchone()[0] cursor.execute(""" SELECT COUNT(*) FROM tg_stats WHERE event_type='poll_response' AND created_at LIKE ? """, (f"{today}%",)) answered = cursor.fetchone()[0] cursor.execute(""" SELECT COUNT(*) FROM tg_stats WHERE event_type='message_sent' AND created_at LIKE ? """, (f"{today}%",)) messages = cursor.fetchone()[0] return sent, answered, messages # === МЕТОДЫ ДЛЯ НАСТРОЕК РАССЫЛОК === def get_tg_rules(self): """Получить все правила рассылок""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute(""" SELECT id, rule_type, subject, recipients, is_enabled, description FROM tg_notification_rules ORDER BY rule_type, subject """) return [dict(row) for row in cursor.fetchall()] def add_tg_rule(self, rule_type, subject, recipients, description=""): """Добавить правило рассылки""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute(""" INSERT INTO tg_notification_rules (rule_type, subject, recipients, is_enabled, description, created_at) VALUES (?, ?, ?, 1, ?, ?) """, (rule_type, subject, recipients, description, datetime.now().isoformat())) conn.commit() return cursor.lastrowid def update_tg_rule(self, rule_id, rule_type, subject, recipients, description=""): """Изменить правило рассылки""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute(""" UPDATE tg_notification_rules SET rule_type=?, subject=?, recipients=?, description=? WHERE id=? """, (rule_type, subject, recipients, description, rule_id)) conn.commit() def toggle_tg_rule(self, rule_id, is_enabled): """Включить/выключить правило рассылки""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute(""" UPDATE tg_notification_rules SET is_enabled=? WHERE id=? """, (1 if is_enabled else 0, rule_id)) conn.commit() def delete_tg_rule(self, rule_id): """Удалить правило рассылки""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute("DELETE FROM tg_notification_rules WHERE id=?", (rule_id,)) conn.commit() def get_unanswered_poll_tasks(self): """Задачи с отправленным опросом, не получившим ответа до сегодняшнего дня. Используется для повторной отправки на следующий рабочий день.""" today = datetime.now().strftime('%Y-%m-%d') with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute(""" SELECT mt.*, e.name as assignee FROM management_tasks mt LEFT JOIN employees e ON mt.employee_id = e.id WHERE mt.poll_id IS NOT NULL AND mt.poll_answered_at IS NULL AND DATE(mt.poll_sent_at) < ? AND mt.status NOT LIKE '%Выполнено%' ORDER BY mt.due_date """, (today,)) return [dict(row) for row in cursor.fetchall()] def get_tasks_for_deadline_check(self): """Задачи в статусе 'В работе' с заданным сроком. Используется для автоматической отправки опросов по срокам.""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute(""" SELECT mt.*, e.name as assignee FROM management_tasks mt LEFT JOIN employees e ON mt.employee_id = e.id WHERE mt.status = 'В работе' AND mt.due_date IS NOT NULL AND mt.due_date != '' ORDER BY mt.due_date """) return [dict(row) for row in cursor.fetchall()] # === МЕТОДЫ ДЛЯ КАЛЕНДАРЯ === def get_production_calendar_range(self, d_start, d_end): """Получить производственный календарь за диапазон дат (словарь date->is_working)""" s_start = str(d_start) s_end = str(d_end) with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() try: # Пробуем новую таблицу production_calendar cursor.execute(""" SELECT name FROM sqlite_master WHERE type='table' AND name='production_calendar' """) if cursor.fetchone(): cursor.execute(""" SELECT date, is_working FROM production_calendar WHERE date BETWEEN ? AND ? """, (s_start, s_end)) return {row[0]: row[1] for row in cursor.fetchall()} # Пробуем старую таблицу production_calendar_exceptions cursor.execute(""" SELECT name FROM sqlite_master WHERE type='table' AND name='production_calendar_exceptions' """) if cursor.fetchone(): cursor.execute(""" SELECT date, is_working FROM production_calendar_exceptions WHERE date BETWEEN ? AND ? """, (s_start, s_end)) return {row[0]: row[1] for row in cursor.fetchall()} # Таблицы нет - возвращаем пустой словарь return {} except Exception as e: logger.error(f"Error getting production calendar: {e}") return {} def get_loaded_dates_in_period(self, d_start, d_end): """Получить загруженные даты в периоде (для календаря)""" s_start = str(d_start) s_end = str(d_end) with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() try: # Проверяем наличие таблицы workload cursor.execute(""" SELECT name FROM sqlite_master WHERE type='table' AND name='workload' """) if cursor.fetchone(): # Проверяем структуру таблицы workload cursor.execute("PRAGMA table_info(workload)") columns = [row[1] for row in cursor.fetchall()] # Используем правильную колонку в зависимости от структуры if 'created_at' in columns: # Используем >= и < вместо BETWEEN для корректной работы с разными форматами cursor.execute(""" SELECT DISTINCT created_at as date FROM workload WHERE created_at >= ? AND created_at < ? """, (s_start, s_end)) # Обрезаем время, оставляем только даты YYYY-MM-DD return [row[0][:10] if len(row[0]) >= 10 else row[0] for row in cursor.fetchall()] elif 'date' in columns: cursor.execute(""" SELECT DISTINCT date FROM workload WHERE date >= ? AND date < ? """, (s_start, s_end)) return [row[0][:10] if len(row[0]) >= 10 else row[0] for row in cursor.fetchall()] elif 'task_date' in columns: cursor.execute(""" SELECT DISTINCT task_date as date FROM workload WHERE task_date >= ? AND task_date < ? """, (s_start, s_end)) return [row[0][:10] if len(row[0]) >= 10 else row[0] for row in cursor.fetchall()] else: # Нет подходящей колонки - используем tasks cursor.execute(""" SELECT DISTINCT task_date FROM tasks WHERE task_date >= ? AND task_date < ? """, (s_start, s_end)) return [row[0][:10] if len(row[0]) >= 10 else row[0] for row in cursor.fetchall()] else: # Используем таблицу tasks cursor.execute(""" SELECT DISTINCT task_date FROM tasks WHERE task_date BETWEEN ? AND ? """, (s_start, s_end)) return [row[0] for row in cursor.fetchall()] except Exception as e: logger.error(f"Error getting loaded dates: {e}") return [] # === МЕТОДЫ ДЛЯ ЗАГРУЗКИ === def get_workload_by_day(self, iso_date, source=None): """Получить загрузку по дню""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() try: # Проверяем наличие таблицы workload cursor.execute(""" SELECT name FROM sqlite_master WHERE type='table' AND name='workload' """) if cursor.fetchone(): # Проверяем структуру таблицы cursor.execute("PRAGMA table_info(workload)") columns = [row[1] for row in cursor.fetchall()] # Используем правильную колонку if 'created_at' in columns: # Структура: period, emp_name, service_name, count, created_at # DATE() нормализует форматы '2026-02-06' и '2026-02-06T12:00:00' sql = """ SELECT *, service_name as val, count FROM workload WHERE DATE(created_at) = ? """ params = [iso_date] elif 'date' in columns: sql = "SELECT * FROM workload WHERE date=?" params = [iso_date] elif 'task_date' in columns: sql = "SELECT * FROM workload WHERE task_date=?" params = [iso_date] else: # Таблица workload есть, но структура неизвестна - используем tasks logger.warning(f"Unknown workload table structure, columns: {columns}") sql = """ SELECT t.*, t.task_name as val, 'daily' as source, e.name as emp_name, e.name FROM tasks t LEFT JOIN employees e ON t.employee_id = e.id WHERE t.task_date = ? """ params = [iso_date] if source and 'workload' in sql: sql += " AND source=?" params.append(source) cursor.execute(sql, params) results = [dict(row) for row in cursor.fetchall()] logger.info(f"get_workload_by_day({iso_date}): found {len(results)} records") return results else: # Используем tasks sql = """ SELECT t.*, t.task_name as val, 'daily' as source, e.name as emp_name, e.name FROM tasks t LEFT JOIN employees e ON t.employee_id = e.id WHERE t.task_date = ? """ cursor.execute(sql, (iso_date,)) return [dict(row) for row in cursor.fetchall()] except Exception as e: logger.error(f"Error getting workload by day: {e}") return [] def get_workload_by_period(self, period, source=None): """Получить загрузку по периоду (месяц)""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() # Проверяем наличие таблицы workload cursor.execute(""" SELECT name FROM sqlite_master WHERE type='table' AND name='workload' """) if cursor.fetchone(): sql = "SELECT * FROM workload WHERE period=?" params = [period] if source: sql += " AND source=?" params.append(source) cursor.execute(sql, params) return [dict(row) for row in cursor.fetchall()] else: # Иначе используем tasks - группируем по месяцам sql = """ SELECT substr(t.task_date, 1, 7) as period, e.name as emp_name, COUNT(*) as task_count FROM tasks t LEFT JOIN employees e ON t.employee_id = e.id WHERE substr(t.task_date, 1, 7) = ? GROUP BY e.name """ cursor.execute(sql, (period,)) return [dict(row) for row in cursor.fetchall()] # === МЕТОДЫ ДЛЯ ЗАГРУЗКИ (WORKLOAD) === def get_workload_daily_for_month(self, year_month): """Все ежедневные записи за месяц (source='daily', period='YYYY-MM'). Дата конкретного дня берётся из created_at. Возвращает list[dict]: emp_name, service_name, count, work_date (YYYY-MM-DD), ne_value.""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() try: cursor.execute(""" SELECT emp_name, service_name, CAST(count AS INTEGER) AS count, DATE(created_at) AS work_date, ne_value FROM workload WHERE source = 'daily' AND period = ? """, (year_month,)) return [dict(r) for r in cursor.fetchall()] except Exception as e: logger.error(f"get_workload_daily_for_month({year_month}): {e}") return [] def delete_workload_day(self, iso_date, source='daily'): """Удалить записи загрузки за день""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() try: # Проверяем наличие таблицы workload cursor.execute(""" SELECT name FROM sqlite_master WHERE type='table' AND name='workload' """) if cursor.fetchone(): # Проверяем структуру таблицы workload cursor.execute("PRAGMA table_info(workload)") columns = [row[1] for row in cursor.fetchall()] # Удаляем используя правильную колонку if 'task_date' in columns: cursor.execute(""" DELETE FROM workload WHERE task_date = ? AND source = ? """, (iso_date, source)) elif 'date' in columns: cursor.execute(""" DELETE FROM workload WHERE date = ? AND source = ? """, (iso_date, source)) elif 'created_at' in columns: # Структура с created_at и period cursor.execute(""" DELETE FROM workload WHERE created_at = ? AND source = ? """, (iso_date, source)) else: # Неизвестная структура - удаляем из tasks logger.warning(f"Unknown workload structure, columns: {columns}") cursor.execute(""" DELETE FROM tasks WHERE task_date = ? """, (iso_date,)) else: # Удаляем из tasks (если workload нет) cursor.execute(""" DELETE FROM tasks WHERE task_date = ? """, (iso_date,)) rows_deleted = cursor.rowcount conn.commit() logger.info(f"Deleted {rows_deleted} workload records for date {iso_date}, source {source}") except Exception as e: logger.error(f"Error deleting workload: {e}") conn.rollback() def delete_workload_period(self, period, source='monthly'): """Удалить записи загрузки за период (месяц)""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() try: # Проверяем наличие таблицы workload cursor.execute(""" SELECT name FROM sqlite_master WHERE type='table' AND name='workload' """) if cursor.fetchone(): # Удаляем по периоду и источнику cursor.execute(""" DELETE FROM workload WHERE period = ? AND source = ? """, (period, source)) rows_deleted = cursor.rowcount conn.commit() logger.info(f"Deleted {rows_deleted} workload records for period {period}, source {source}") else: logger.warning(f"Table 'workload' does not exist") except Exception as e: logger.error(f"Error deleting workload period: {e}") conn.rollback() def save_workload_record(self, employee_id_or_name, iso_date, val, client_name, price=0, total_sum=0, created_at=None, source='daily', ne_value=None): """ Сохранить запись загрузки Args: employee_id_or_name: ID сотрудника (int) или его ФИО (str) iso_date: дата в формате YYYY-MM-DD (или период YYYY-MM для monthly) val: название услуги client_name: количество (передается как строка или число) price: цена (опционально) total_sum: общая сумма (опционально) created_at: дата создания записи (опционально) source: источник ('daily' или 'monthly') """ with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() try: count = float(client_name) if client_name else 1.0 # Определяем имя сотрудника if isinstance(employee_id_or_name, int): # Передан ID - ищем имя cursor.execute("SELECT name FROM employees WHERE id = ?", (employee_id_or_name,)) row = cursor.fetchone() if not row: logger.error(f"Employee with ID {employee_id_or_name} not found") return emp_name = row[0] else: # Передано ФИО - используем как есть emp_name = str(employee_id_or_name) # Период для workload if source == 'monthly': # Для месячных данных iso_date уже в формате YYYY-MM period = iso_date if len(iso_date) == 7 else iso_date[:7] else: # Для ежедневных - извлекаем месяц из даты period = iso_date[:7] # '2026-02-12' -> '2026-02' # Дата создания if not created_at: created_at = iso_date # Проверяем наличие таблицы workload cursor.execute(""" SELECT name FROM sqlite_master WHERE type='table' AND name='workload' """) if cursor.fetchone(): # Таблица workload существует # Структура: period, emp_name, service_name, count, price, total_sum, created_at, source, ne_value cursor.execute(""" INSERT INTO workload (period, emp_name, service_name, count, price, total_sum, created_at, source, ne_value) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?) """, (period, emp_name, val, count, price, total_sum, created_at, source, ne_value)) logger.info(f"Saved to workload: {emp_name}, {val}, {count}, period={period}, date={created_at}, source={source}") else: # Таблица workload не существует - используем tasks # Для tasks нужен ID if isinstance(employee_id_or_name, int): emp_id = employee_id_or_name else: # Ищем ID по ФИО cursor.execute("SELECT id FROM employees WHERE name = ?", (emp_name,)) row = cursor.fetchone() if not row: logger.warning(f"Employee '{emp_name}' not found in employees table, skipping") return emp_id = row[0] cursor.execute(""" INSERT INTO tasks (employee_id, task_date, task_name) VALUES (?, ?, ?) """, (emp_id, iso_date, val)) logger.info(f"Saved to tasks: {emp_name}, {val}, date={iso_date}") conn.commit() except Exception as e: logger.error(f"Error saving workload record: {e}") conn.rollback() # === КОНТРОЛЬ СРОКОВ ДАТЕЛ === def delete_deadline_cases_by_date(self, load_date): """Удалить все дела для указанной даты загрузки (перед повторной загрузкой)""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute( "DELETE FROM deadline_cases WHERE load_date = ?", (load_date,) ) conn.commit() return cursor.rowcount def get_deadline_upload_dates(self) -> list: """Return all distinct load_date values present in deadline_cases, sorted desc.""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute( "SELECT DISTINCT load_date FROM deadline_cases ORDER BY load_date DESC" ) return [row[0] for row in cursor.fetchall()] def delete_deadline_upload(self, load_date: str) -> int: """Delete all data for a specific upload date (cases + comment history). Both tables are wiped in a single transaction. Returns the number of deadline_cases rows deleted. """ with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute( "DELETE FROM deadline_cases WHERE load_date = ?", (load_date,) ) deleted = cursor.rowcount try: cursor.execute( "DELETE FROM deadline_comment_history WHERE load_date = ?", (load_date,) ) except Exception: pass # table may not exist in older DB schemas conn.commit() return deleted def save_deadline_case(self, data): """ Сохранить дело для контроля сроков Args: data: dict с полями case_number, create_date, inspection_date, risk, expert_inspection, expert_calculation, payment_info, region, department, deadline, comment, load_date """ with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() # Проверяем существует ли уже это дело на эту дату cursor.execute(""" SELECT id, deadline FROM deadline_cases WHERE case_number=? AND load_date=? """, (data['case_number'], data['load_date'])) existing = cursor.fetchone() if existing: # При дублях в файле сохраняем максимальный срок. # Остальные поля обновляем только если они заполнены в новой строке. existing_deadline = existing['deadline'] or 0 new_deadline = data.get('deadline') or 0 merged_deadline = max(existing_deadline, new_deadline) def _val(key): v = data.get(key) return v if v and str(v).strip() not in ('', 'nan') else None cursor.execute(""" UPDATE deadline_cases SET create_date=COALESCE(?, create_date), inspection_date=COALESCE(?, inspection_date), risk=COALESCE(?, risk), expert_inspection=COALESCE(?, expert_inspection), expert_calculation=COALESCE(?, expert_calculation), payment_info=COALESCE(?, payment_info), region=COALESCE(?, region), department=COALESCE(?, department), deadline=?, comment=COALESCE(?, comment), services_list=COALESCE(?, services_list), address=COALESCE(?, address) WHERE id=? """, ( _val('create_date'), _val('inspection_date'), _val('risk'), _val('expert_inspection'), _val('expert_calculation'), _val('payment_info'), _val('region'), _val('department'), merged_deadline, _val('comment'), _val('services_list'), _val('address'), existing['id'] )) else: # Вставляем новую запись cursor.execute(""" INSERT INTO deadline_cases ( case_number, create_date, inspection_date, risk, expert_inspection, expert_calculation, payment_info, region, department, deadline, comment, services_list, address, load_date ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) """, ( data['case_number'], data.get('create_date'), data.get('inspection_date'), data.get('risk'), data.get('expert_inspection'), data.get('expert_calculation'), data.get('payment_info'), data.get('region'), data.get('department'), data.get('deadline'), data.get('comment'), data.get('services_list'), data.get('address'), data['load_date'] )) conn.commit() def get_deadline_comments_latest(self): """Вернуть {case_number: comment} — последний непустой комментарий для каждого дела.""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute(""" SELECT dc.case_number, dc.comment FROM deadline_cases dc WHERE dc.comment IS NOT NULL AND dc.comment != '' AND dc.load_date = ( SELECT MAX(load_date) FROM deadline_cases WHERE case_number = dc.case_number AND comment IS NOT NULL AND comment != '' ) """) return {r['case_number']: r['comment'] for r in cursor.fetchall()} def update_deadline_comment(self, case_number, load_date, comment): """Обновить пользовательский комментарий по конкретному делу (с историей).""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute( "SELECT comment FROM deadline_cases WHERE case_number=? AND load_date=?", (case_number, load_date) ) row = cursor.fetchone() old_comment = row['comment'] if row else None cursor.execute(""" UPDATE deadline_cases SET comment=? WHERE case_number=? AND load_date=? """, (comment, case_number, load_date)) cursor.execute(""" INSERT INTO deadline_comment_history (case_number, load_date, old_comment, new_comment) VALUES (?, ?, ?, ?) """, (case_number, load_date, old_comment, comment)) conn.commit() def update_deadline_comment_all_dates(self, case_number, comment): """Обновить комментарий по делу во всех загрузках (импорт из Excel). История изменения фиксируется для самой последней загрузки. """ with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() # Получаем последнюю дату и старый комментарий для истории cursor.execute(""" SELECT load_date, comment FROM deadline_cases WHERE case_number=? ORDER BY load_date DESC LIMIT 1 """, (case_number,)) row = cursor.fetchone() if not row: return latest_date = row['load_date'] old_comment = row['comment'] # Обновляем во всех загрузках cursor.execute( "UPDATE deadline_cases SET comment=? WHERE case_number=?", (comment, case_number) ) # Историю пишем один раз — для последней загрузки cursor.execute(""" INSERT INTO deadline_comment_history (case_number, load_date, old_comment, new_comment) VALUES (?, ?, ?, ?) """, (case_number, latest_date, old_comment, comment)) conn.commit() def update_deadline_case_expert(self, case_number, load_date, expert_name): """Назначить эксперта на расчёт в конкретном деле.""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute( "UPDATE deadline_cases SET expert_calculation=? WHERE case_number=? AND load_date=?", (expert_name, case_number, load_date) ) conn.commit() def update_deadline_case_candidate(self, case_number, load_date, candidate_name): """Записать кандидата на расчёт (для тестирования умного распределения).""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute( "UPDATE deadline_cases SET candidate_expert=? WHERE case_number=? AND load_date=?", (candidate_name, case_number, load_date) ) conn.commit() def clear_deadline_candidates(self, load_date): """Очистить всех кандидатов для указанной даты загрузки.""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute( "UPDATE deadline_cases SET candidate_expert=NULL WHERE load_date=?", (load_date,) ) conn.commit() def get_comment_history(self, case_number, load_date=None): """История изменений комментария по номеру дела.""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() if load_date: cursor.execute(""" SELECT * FROM deadline_comment_history WHERE case_number=? AND load_date=? ORDER BY changed_at DESC """, (case_number, load_date)) else: cursor.execute(""" SELECT * FROM deadline_comment_history WHERE case_number=? ORDER BY changed_at DESC """, (case_number,)) return [dict(r) for r in cursor.fetchall()] def get_deadline_trend(self, threshold=3): """Динамика числа дел со сроком >= threshold по датам загрузки.""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute(""" SELECT load_date, COUNT(*) AS total, SUM(CASE WHEN deadline >= ? THEN 1 ELSE 0 END) AS over_threshold FROM deadline_cases GROUP BY load_date ORDER BY load_date """, (threshold,)) return [dict(r) for r in cursor.fetchall()] def get_act_imush_totals(self, period=None): """Суммы Акт (имущ) по Москве и регионам за указанный (или последний) период.""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() if period is None: cursor.execute( "SELECT MAX(period) FROM workload WHERE service_name='Акт (имущ)'" ) row = cursor.fetchone() period = row[0] if row and row[0] else None if not period: return {'period': None, 'moscow': 0, 'regions': 0} cursor.execute(""" SELECT CASE WHEN ne_value='Москва' THEN 'moscow' ELSE 'regions' END AS grp, SUM(count) AS total FROM workload WHERE service_name='Акт (имущ)' AND period=? AND ne_value IS NOT NULL AND ne_value NOT IN ('nan','') GROUP BY grp """, (period,)) result = {'period': period, 'moscow': 0, 'regions': 0} for r in cursor.fetchall(): result[r[0]] = int(r[1] or 0) return result def get_act_imush_history(self, num_months=6): """История Акт (имущ) по месяцам (ASC). Возвращает list[dict] с полями period, moscow, regions.""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute(""" SELECT period, SUM(CASE WHEN ne_value='Москва' THEN count ELSE 0 END) AS moscow, SUM(CASE WHEN ne_value IS NOT NULL AND ne_value NOT IN ('nan','','Москва') THEN count ELSE 0 END) AS regions FROM workload WHERE service_name='Акт (имущ)' GROUP BY period ORDER BY period DESC LIMIT ? """, (num_months,)) rows = [{'period': r[0], 'moscow': int(r[1] or 0), 'regions': int(r[2] or 0)} for r in cursor.fetchall()] return list(reversed(rows)) # ASC (старые → новые) def get_deadline_cases(self, load_date=None, expert=None): """ Получить дела для контроля сроков Args: load_date: дата загрузки (если None - последняя) expert: фильтр по эксперту на расчёт Returns: list of dict """ with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() # Определяем дату загрузки if load_date is None: cursor.execute(""" SELECT MAX(load_date) as last_date FROM deadline_cases """) result = cursor.fetchone() load_date = result['last_date'] if result else None if not load_date: return [] # Запрос sql = """ SELECT * FROM deadline_cases WHERE load_date = ? """ params = [load_date] if expert: sql += " AND expert_calculation = ?" params.append(expert) sql += " ORDER BY deadline DESC, case_number" cursor.execute(sql, params) return [DeadlineCase.from_row(row) for row in cursor.fetchall()] def get_deadline_cases_by_period(self, start_date: str, end_date: str) -> list[DeadlineCase]: """ Получить закрытые дела (с назначенным экспертом) за период дат загрузки. Args: start_date: начало периода (YYYY-MM-DD) end_date: конец периода (YYYY-MM-DD) Returns: list of DeadlineCase """ with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute(""" SELECT * FROM deadline_cases WHERE load_date >= ? AND load_date <= ? AND expert_calculation IS NOT NULL AND expert_calculation != '' ORDER BY expert_calculation, case_number """, (start_date, end_date)) return [DeadlineCase.from_row(row) for row in cursor.fetchall()] def get_case_work_periods(self, case_numbers): """ Для каждого номера дела вернуть дату первого и последнего появления в БД. Returns: dict {case_number: {'first': 'YYYY-MM-DD', 'last': 'YYYY-MM-DD'}} """ if not case_numbers: return {} with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() placeholders = ','.join('?' * len(case_numbers)) cursor.execute(f""" SELECT case_number, MIN(load_date) AS first_date, MAX(load_date) AS last_date FROM deadline_cases WHERE case_number IN ({placeholders}) GROUP BY case_number """, list(case_numbers)) return {row['case_number']: {'first': row['first_date'], 'last': row['last_date']} for row in cursor.fetchall()} def get_calc_days_per_case(self, load_date): """Вернуть {case_number: days} — сколько дней дело числится на расчёте. days = (load_date − первый load_date, когда expert_calculation был заполнен). Только для дел, у которых expert_calculation заполнен на указанный load_date. """ from datetime import datetime as _dt with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute(""" SELECT dc.case_number, MIN(dc2.load_date) AS first_calc_date FROM deadline_cases dc JOIN deadline_cases dc2 ON dc2.case_number = dc.case_number AND dc2.expert_calculation IS NOT NULL AND dc2.expert_calculation != '' AND dc2.load_date <= ? WHERE dc.load_date = ? AND dc.expert_calculation IS NOT NULL AND dc.expert_calculation != '' GROUP BY dc.case_number """, (load_date, load_date)) try: selected_dt = _dt.strptime(load_date, '%Y-%m-%d').date() except Exception: return {} result = {} for row in cursor.fetchall(): try: first_dt = _dt.strptime(row['first_calc_date'], '%Y-%m-%d').date() result[row['case_number']] = (selected_dt - first_dt).days except Exception: pass return result def get_calc_first_dates(self, load_date): """Вернуть {case_number: first_calc_date_str} — первую дату, когда текущий expert_calculation был назначен на дело (из загрузки load_date).""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute(""" SELECT dc.case_number, MIN(dc2.load_date) AS first_calc_date FROM deadline_cases dc JOIN deadline_cases dc2 ON dc2.case_number = dc.case_number AND TRIM(REPLACE(dc2.expert_calculation, char(160), ' ')) = TRIM(REPLACE(dc.expert_calculation, char(160), ' ')) AND dc2.load_date <= ? WHERE dc.load_date = ? AND dc.expert_calculation IS NOT NULL AND TRIM(REPLACE(dc.expert_calculation, char(160), ' ')) != '' AND TRIM(REPLACE(dc.expert_calculation, char(160), ' ')) != 'nan' GROUP BY dc.case_number """, (load_date, load_date)) return {r[0]: r[1] for r in cursor.fetchall()} def get_expert_closed_cases_since(self, expert_name, since_date, current_load_date): """Дела эксперта, которые присутствовали в загрузках начиная с since_date, но отсутствуют в current_load_date (т.е. были закрыты/выполнены), и у которых комментарий никогда не заполнялся ни в одной загрузке. Возвращает список номеров дел.""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute(""" SELECT DISTINCT case_number FROM deadline_cases WHERE TRIM(REPLACE(expert_calculation, char(160), ' ')) = ? AND load_date >= ? AND load_date < ? AND case_number NOT IN ( SELECT case_number FROM deadline_cases WHERE load_date = ? ) AND case_number NOT IN ( SELECT DISTINCT case_number FROM deadline_cases WHERE comment IS NOT NULL AND TRIM(comment) != '' ) ORDER BY case_number """, (expert_name, since_date, current_load_date, current_load_date)) return [r[0] for r in cursor.fetchall()] def get_expert_all_case_periods(self, expert_name, up_to_date): """ Для эксперта вернуть все дела, которые когда-либо у него были (до up_to_date), с датой первого и последнего появления. Логика "среднего срока": sum(last_date - first_date).days для каждого дела ──────────────────────────────────────────────── (up_to_date - MIN(first_date по всем делам)).days Returns: list of {'case_number': str, 'first_date': str, 'last_date': str} """ with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute(""" SELECT case_number, MIN(load_date) AS first_date, MAX(load_date) AS last_date FROM deadline_cases WHERE expert_calculation = ? AND load_date <= ? GROUP BY case_number """, (expert_name, up_to_date)) return [dict(row) for row in cursor.fetchall()] def get_cases_comment_timeline(self, case_numbers): """Вернуть временну́ю шкалу комментариев для списка дел. Returns: {case_number: [{'load_date': str, 'comment': str|None}, ...]} — список отсортирован по возрастанию load_date. """ if not case_numbers: return {} with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() placeholders = ','.join('?' * len(case_numbers)) cursor.execute(f""" SELECT case_number, load_date, comment FROM deadline_cases WHERE case_number IN ({placeholders}) ORDER BY case_number, load_date """, list(case_numbers)) result: dict = {} for row in cursor.fetchall(): cn = row['case_number'] result.setdefault(cn, []).append( {'load_date': row['load_date'], 'comment': row['comment']} ) return result def get_deadline_stats(self, load_date=None): """ Получить статистику по делам Returns: dict: { 'total': общее количество, 'unassigned': неназначенные (срок=3, эксперт пустой), 'by_expert': {expert: count}, 'by_region': {region: count} } """ cases = self.get_deadline_cases(load_date) stats = { 'total': len(cases), 'unassigned': 0, 'by_expert': {}, 'by_region': {} } for case in cases: # Неназначенные if case.deadline == 3 and not case.expert_calculation: stats['unassigned'] += 1 # По экспертам expert = (case.expert_calculation or '').strip() if expert and expert != 'nan': stats['by_expert'][expert] = stats['by_expert'].get(expert, 0) + 1 # По регионам region = (case.region or '').strip() if region: stats['by_region'][region] = stats['by_region'].get(region, 0) + 1 return stats # === КОНТРОЛЬ КАЧЕСТВА === def upsert_quality_cases(self, rows): """Добавить новые записи контроля качества (только INSERT, существующие не трогаем). Уникальность определяется составным ключом (nomer_polisa, data_osm). Если запись уже есть — пропускается (OR IGNORE). Возвращает количество реально добавленных строк. """ with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() inserted = 0 for r in rows: cursor.execute(""" INSERT OR IGNORE INTO quality_cases (nomer_polisa, risk, razdel, data_zakr, mes_zakr, god_zakr, osmotrowshik, region_osm, narusheniya, naim_podr, opisanie, ispravlenie, data_osm, kto_obnaruzhil) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?) """, ( r.get('nomer_polisa'), r.get('risk'), r.get('razdel'), r.get('data_zakr'), r.get('mes_zakr'), r.get('god_zakr'), r.get('osmotrowshik'), r.get('region_osm'), r.get('narusheniya'), r.get('naim_podr'), r.get('opisanie'), r.get('ispravlenie', 'Нет'), r.get('data_osm'), r.get('kto_obnaruzhil') )) if cursor.rowcount: inserted += 1 conn.commit() return inserted def get_quality_cases(self, osmotrowshik=None, naim_podr=None, ispravlenie=None, mes=None, god=None): """Получить записи контроля качества с опциональной фильтрацией""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() q = "SELECT * FROM quality_cases WHERE 1=1" params = [] if osmotrowshik: q += " AND TRIM(REPLACE(osmotrowshik, char(160), ' '))=?" params.append(osmotrowshik) if naim_podr: q += " AND naim_podr=?" params.append(naim_podr) if ispravlenie: q += " AND ispravlenie=?" params.append(ispravlenie) if mes: q += " AND mes_zakr=?" params.append(mes) if god: q += " AND god_zakr=?" params.append(god) q += " ORDER BY data_zakr DESC, nomer_polisa" cursor.execute(q, params) return [dict(r) for r in cursor.fetchall()] def delete_quality_storonny_akt(self): """Удалить из quality_cases все записи с 'Сторонний акт' в нарушениях""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute( "DELETE FROM quality_cases " "WHERE narusheniya LIKE '%Сторонний акт%'" ) deleted = cursor.rowcount conn.commit() return deleted def delete_quality_case(self, nomer_polisa, data_osm=None): """Удалить запись из quality_cases по составному ключу (nomer_polisa, data_osm).""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() if data_osm: cursor.execute( "DELETE FROM quality_cases WHERE nomer_polisa=? AND data_osm=?", (nomer_polisa, data_osm) ) else: cursor.execute( "DELETE FROM quality_cases WHERE nomer_polisa=?", (nomer_polisa,) ) conn.commit() def update_quality_ispravlenie(self, nomer_polisa, value, data_osm=None): """Обновить значение 'Исправление рассмотрено' по составному ключу (nomer_polisa, data_osm).""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() if data_osm: cursor.execute( "UPDATE quality_cases SET ispravlenie=? " "WHERE nomer_polisa=? AND data_osm=?", (value, nomer_polisa, data_osm) ) else: cursor.execute( "UPDATE quality_cases SET ispravlenie=? WHERE nomer_polisa=?", (value, nomer_polisa) ) conn.commit() def update_quality_data_osm(self, nomer_polisa, old_data_osm, new_data_osm): """Обновить дату осмотра в quality_cases. Ищет запись по (nomer_polisa, old_data_osm) и меняет data_osm на new_data_osm. Если new_data_osm=None — записывает NULL. """ with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() try: if old_data_osm: cursor.execute( "UPDATE quality_cases SET data_osm=? " "WHERE nomer_polisa=? AND data_osm=?", (new_data_osm, nomer_polisa, old_data_osm) ) else: cursor.execute( "UPDATE quality_cases SET data_osm=? " "WHERE nomer_polisa=? AND (data_osm IS NULL OR data_osm='')", (new_data_osm, nomer_polisa) ) conn.commit() except Exception as e: logger.error(f"update_quality_data_osm error: {e}") raise def get_quality_osmotrowshiki(self): """Список уникальных осмотровщиков из quality_cases""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute( "SELECT DISTINCT TRIM(REPLACE(osmotrowshik, char(160), ' ')) " "FROM quality_cases " "WHERE osmotrowshik IS NOT NULL AND osmotrowshik != '' " "ORDER BY TRIM(REPLACE(osmotrowshik, char(160), ' '))" ) return [r[0] for r in cursor.fetchall()] def get_quality_naim_podr(self): """Список уникальных наименований подразделений НЭ из quality_cases""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute( "SELECT DISTINCT naim_podr FROM quality_cases " "WHERE naim_podr IS NOT NULL AND naim_podr != '' ORDER BY naim_podr" ) return [r[0] for r in cursor.fetchall()] def get_quality_months(self): """Список уникальных месяцев из quality_cases, отсортированных (строки)""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute( "SELECT DISTINCT mes_zakr FROM quality_cases " "WHERE mes_zakr IS NOT NULL ORDER BY CAST(mes_zakr AS INTEGER)" ) return [str(r[0]) for r in cursor.fetchall()] def get_quality_years(self): """Список уникальных годов из quality_cases, отсортированных по убыванию (строки)""" with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute( "SELECT DISTINCT god_zakr FROM quality_cases " "WHERE god_zakr IS NOT NULL ORDER BY CAST(god_zakr AS INTEGER) DESC" ) return [str(r[0]) for r in cursor.fetchall()] # === SERVICE PROFILING === def backfill_employee_services(self, months_back: int = 6) -> dict: """Union-merge workload service names into employees.svc_types (additive only). Scans the last `months_back` months of workload records, groups unique service names per employee, and merges them into the existing svc_types JSON list — never removing services that were set manually. Returns: {'updated': <count of employees changed>, 'new_services': <total new entries added>} """ import json from datetime import date, timedelta from collections import defaultdict cutoff_period = (date.today() - timedelta(days=months_back * 30)).strftime('%Y-%m') def _norm(s: str) -> str: return ' '.join(str(s).strip().split()).upper() with self.lock: with self.pool.get_connection() as conn: cursor = conn.cursor() cursor.execute( "SELECT id, name, svc_types FROM employees " "WHERE is_fired IS NULL OR is_fired = 0" ) emp_by_norm = { _norm(r['name']): {'id': r['id'], 'svc_types': r['svc_types'] or '[]'} for r in cursor.fetchall() } cursor.execute( "SELECT emp_name, service_name FROM workload " "WHERE period >= ? AND service_name IS NOT NULL AND service_name != ''", (cutoff_period,) ) svc_by_emp: dict = defaultdict(set) for r in cursor.fetchall(): raw_name = (r['emp_name'] or '').strip() svc = (r['service_name'] or '').strip() if raw_name and svc: svc_by_emp[_norm(raw_name)].add(svc) updated = 0 new_total = 0 for norm_name, new_svcs in svc_by_emp.items(): emp = emp_by_norm.get(norm_name) if emp is None: continue try: existing: set = set(json.loads(emp['svc_types'])) except (json.JSONDecodeError, TypeError, ValueError): existing = set() merged = existing | new_svcs added = len(merged) - len(existing) if added > 0: cursor.execute( "UPDATE employees SET svc_types = ? WHERE id = ?", (json.dumps(sorted(merged), ensure_ascii=False), emp['id']) ) updated += 1 new_total += added conn.commit() return {'updated': updated, 'new_services': new_total} # === CLEANUP === def close(self): """Закрыть все соединения""" self.pool.close_all() logger.info("DBManager closed")