Перед тем, как исправлять выгрузку скуд с глупым Gemini-Lite

This commit is contained in:
2026-08-21 13:13:36 +03:00
parent 7bf1af7c8f
commit baf7c69395
28 changed files with 3064 additions and 1723 deletions
+46 -119
View File
@@ -3,137 +3,64 @@
FILE: services/scud_etl/pipeline.py
PROJECT: SCUD Orion AI (Unified Architecture)
MODULE: services / scud_etl
ROLE: Оркестратор этапов контроллинга (Загрузка -> Сверка -> Отчеты -> SQLite).
AI-CONTEXT-ANCHORS:
- ANCHOR[PIPELINE_RUN_CONTROLLING]: Главная функция выполнения ETL-конвейера.
ROLE: Оркестрация выборки снапшотов из SQLite и загрузки кадровых файлов 1С.
===============================================================================
"""
import os
from datetime import datetime, timedelta
import json
import logging
from typing import Optional
import pandas as pd
from config import DATE_TODAY, DATE_YESTERDAY, OUTPUT_DIR
from core.database import (
save_scud_to_db, save_staff_to_db, save_absences_to_db,
save_anomalies_to_db, load_scud_from_db_by_snapshot, get_latest_snapshot_time
)
from services.data_loader import load_1c_data_smart
from services.excel_exporter import generate_summary_excel, generate_detailed_excel
from services.ai_verifier import ai_verify_scud_against_staff, analyze_scud_mass_failure_ai
from services.text_reporter import generate_markdown_report
from services.feedback_loop import review_ai_decisions
from services.knowledge_base import load_knowledge_base
from core.connection import get_connection
from .merger import (
load_exceptions_config, apply_exceptions_from_json,
aggregate_scud_by_employee, filter_report_dataframe, load_static_reason_workers
)
from .anomaly_detector import detect_all_anomalies
logger = logging.getLogger("SCUD_PIPELINE")
# ANCHOR[PIPELINE_RUN_CONTROLLING]
def run_controlling_pipeline(snapshot_param: str = None, skip_export: bool = False, debug: bool = False, has_today_1c: bool = True) -> None:
"""Выполняет полный цикл сверки СКУД ⟷ 1С и сохранение результатов."""
kb_rules = load_knowledge_base().get("rules", [])
exceptions_cfg = load_exceptions_config()
static_reasons = load_static_reason_workers()
def load_best_snapshot_for_date(date_str: str, prefer_final_y: bool = True) -> pd.DataFrame:
"""
Извлекает срез за дату из SQLite. Для отчета за вчера строго ищет
финальный вечерний срез Y (с зафиксированными выходами за 22:00:00).
"""
with get_connection(row_factory=True) as conn:
cursor = conn.cursor()
if prefer_final_y:
cursor.execute("""
SELECT raw_data_json
FROM scud_snapshots
WHERE (snapshot_date = ? OR date_str = ?) AND (snapshot_id LIKE 'Y%' OR snapshot_time LIKE '22:00%')
ORDER BY id DESC LIMIT 1
""", (date_str, date_str))
row = cursor.fetchone()
if row and row["raw_data_json"]:
data = json.loads(row["raw_data_json"])
return pd.DataFrame(data)
print("[2/5] Загрузка данных из СКУД, 1С:ЗУП, реестра причин и исключений...")
cursor.execute("""
SELECT raw_data_json
FROM scud_snapshots
WHERE snapshot_date = ? OR date_str = ?
ORDER BY id DESC LIMIT 1
""", (date_str, date_str))
row = cursor.fetchone()
if row and row["raw_data_json"]:
data = json.loads(row["raw_data_json"])
return pd.DataFrame(data)
# 1. Определение дат целевого снапшота и предыдущей смены
if snapshot_param or skip_export:
snap_to_use = snapshot_param or get_latest_snapshot_time()
raw_scud_today_df = load_scud_from_db_by_snapshot(None, snapshot_param=snap_to_use)
if not raw_scud_today_df.empty and 'log_date' in raw_scud_today_df.columns:
target_date_str = str(raw_scud_today_df['log_date'].iloc[0])
else:
target_date_str = DATE_TODAY
else:
target_date_str = DATE_TODAY
raw_scud_today_df = load_scud_from_db_by_snapshot(target_date_str, snapshot_param=snapshot_param)
return pd.DataFrame()
dt_target = datetime.strptime(target_date_str, "%d.%m.%Y")
dt_yesterday = dt_target - timedelta(days=3 if dt_target.weekday() == 0 else 1)
yesterday_date_str = dt_yesterday.strftime("%d.%m.%Y")
print(f"[📸] СНАПШОТ ОПРЕДЕЛЕН: Целевая дата = {target_date_str}, Накануне = {yesterday_date_str}\n")
def load_1c_files_for_date(date_str: str) -> tuple[pd.DataFrame, pd.DataFrame]:
"""
Загружает файлы Штат_*.xlsx и Отсутствия_*.xlsx за указанную дату из data/1c/.
"""
formatted_date = date_str.replace(".", "_")
staff_file = f"data/1c/Штат_{formatted_date}.xlsx"
absences_file = f"data/1c/Отсутствия_{formatted_date}.xlsx"
raw_scud_yesterday_df = load_scud_from_db_by_snapshot(yesterday_date_str, snapshot_param=None)
df_staff_yesterday, df_absent_yesterday = load_1c_data_smart(yesterday_date_str, use_db=True)
df_staff_today, df_absent_today = load_1c_data_smart(target_date_str, use_db=True)
df_staff = pd.read_excel(staff_file) if os.path.exists(staff_file) else pd.DataFrame()
df_absences = pd.read_excel(absences_file) if os.path.exists(absences_file) else pd.DataFrame()
# 2. Этап 3: Обработка ВЧЕРА (Детальный отчет)
print(f"[3/5] Обработка и построение детального отчета за ВЧЕРА ({yesterday_date_str})...")
if df_staff_yesterday is not None:
save_staff_to_db(df_staff_yesterday, yesterday_date_str)
if df_absent_yesterday is not None:
save_absences_to_db(df_absent_yesterday, yesterday_date_str)
# Проверка опечаток ФИО через AI-аудитор
staff_fios_y_clean = df_staff_yesterday['fio_clean'].dropna().tolist() if df_staff_yesterday is not None else []
if not raw_scud_yesterday_df.empty:
raw_scud_yesterday_df['Пришел'] = raw_scud_yesterday_df['is_present'].astype(int) == 1 if 'is_present' in raw_scud_yesterday_df.columns else False
unrecog = raw_scud_yesterday_df[~raw_scud_yesterday_df['fio_clean'].isin(staff_fios_y_clean)]['fio_clean'].tolist()
fio_map = ai_verify_scud_against_staff(unrecog, staff_fios_y_clean)
if fio_map:
raw_scud_yesterday_df['fio_clean'] = raw_scud_yesterday_df['fio_clean'].apply(lambda x: fio_map[x]['staff_fio'] if x in fio_map else x)
raw_scud_yesterday_df = aggregate_scud_by_employee(raw_scud_yesterday_df)
merged_y = (df_staff_yesterday.copy() if df_staff_yesterday is not None else pd.DataFrame())
if not merged_y.empty:
if not raw_scud_yesterday_df.empty:
merged_y = merged_y.merge(
raw_scud_yesterday_df[['fio_clean', 'Пришел', 'Начало_дня', 'Первая_активность', 'Конец_дня', 'Находился_в_здании', 'anomaly_flag']],
on='fio_clean', how='left'
)
if df_absent_yesterday is not None and not df_absent_yesterday.empty:
merged_y = merged_y.merge(df_absent_yesterday[['fio_clean', 'Вид_отсутствия']], on='fio_clean', how='left')
merged_y['Пришел'] = merged_y['Пришел'].fillna(False) if 'Пришел' in merged_y.columns else False
if 'Сотрудник' not in merged_y.columns:
merged_y['Сотрудник'] = merged_y.get('ФИО', merged_y['fio_clean'])
merged_y = apply_exceptions_from_json(merged_y, exceptions_cfg)
scud_fios_y = set(raw_scud_yesterday_df['fio_clean'].dropna().tolist()) if not raw_scud_yesterday_df.empty else set()
anomalies_y = detect_all_anomalies(merged_y, static_reasons, kb_rules, scud_fios_set=scud_fios_y)
save_anomalies_to_db(anomalies_y, yesterday_date_str)
filtered_y = filter_report_dataframe(merged_y)
generate_detailed_excel(merged_df=filtered_y, date_str=yesterday_date_str)
save_scud_to_db(merged_y, yesterday_date_str, snapshot_time=f"{dt_yesterday.strftime('%Y-%m-%d')} 22:00:00", is_yesterday=True)
print(f"[✓] Детальный отчет за вчера сформирован и зафиксирован в SQLite за {yesterday_date_str}")
# 3. Этап 4 & 5: Обработка СЕГОДНЯ (Ежедневная сводка)
if has_today_1c and df_staff_today is not None and df_absent_today is not None:
print(f"\n[4/5] Обработка и построение Ежедневной сводки за {target_date_str}...")
save_staff_to_db(df_staff_today, target_date_str)
save_absences_to_db(df_absent_today, target_date_str)
raw_scud_today_df = aggregate_scud_by_employee(raw_scud_today_df)
merged_t = df_staff_today.copy()
if not raw_scud_today_df.empty:
merged_t = merged_t.merge(
raw_scud_today_df[['fio_clean', 'Пришел', 'Начало_дня', 'Первая_активность', 'Конец_дня', 'Находился_в_здании', 'anomaly_flag']],
on='fio_clean', how='left'
)
merged_t = merged_t.merge(df_absent_today[['fio_clean', 'Вид_отсутствия']], on='fio_clean', how='left')
merged_t['Пришел'] = merged_t['Пришел'].fillna(False) if 'Пришел' in merged_t.columns else False
if 'Сотрудник' not in merged_t.columns:
merged_t['Сотрудник'] = merged_t.get('ФИО', merged_t['fio_clean'])
merged_t = apply_exceptions_from_json(merged_t, exceptions_cfg)
scud_fios_t = set(raw_scud_today_df['fio_clean'].dropna().tolist()) if not raw_scud_today_df.empty else set()
anomalies_t = detect_all_anomalies(merged_t, static_reasons, kb_rules, scud_fios_set=scud_fios_t)
save_anomalies_to_db(anomalies_t, target_date_str)
print("[5/5] Сохранение Ежедневной сводки...")
generate_summary_excel(merged_df=merged_t, date_str=target_date_str)
print(f"[✓] Ежедневная сводка сохранена за {target_date_str}")
else:
print(f"\n[ℹ️] Формирование Ежедневной сводки за {target_date_str} ПРОПУЩЕНО.")
return df_staff, df_absences