Files
scud_ai/services/scud_etl/pipeline.py
T

99 lines
4.0 KiB
Python

"""
===============================================================================
FILE: services/scud_etl/pipeline.py
ROLE: Загрузка наилучших срезов СКУД и штата/отсутствий 1С с умным fallback-ом.
===============================================================================
"""
import os
import logging
from typing import Optional, Dict, Any, Tuple
import pandas as pd
from core.connection import get_connection
from core.database import load_scud_from_db_by_snapshot
from config import DATA_DIR
from services.data_loader import load_1c_data_smart, load_staff_data, load_absent_data
logger = logging.getLogger("SCUD_PIPELINE")
def load_best_snapshot_for_date(date_str: str, prefer_final_y: bool = False) -> Optional[pd.DataFrame]:
"""
Загружает наилучший срез СКУД за дату.
Если prefer_final_y=True — отдает предпочтение финишному Y (23:59:59).
"""
with get_connection() as conn:
cursor = conn.cursor()
target_snap_id = None
if prefer_final_y:
cursor.execute("""
SELECT snapshot_id
FROM scud_logs
WHERE log_date = ?
AND (snapshot_id LIKE 'Y%' OR snapshot_time LIKE '%23:59:59' OR snapshot_time LIKE '%22:00:00')
ORDER BY id DESC LIMIT 1
""", (date_str,))
row = cursor.fetchone()
if row:
target_snap_id = row[0]
if not target_snap_id:
cursor.execute("""
SELECT snapshot_id
FROM scud_logs
WHERE log_date = ?
ORDER BY id DESC LIMIT 1
""", (date_str,))
row = cursor.fetchone()
if row:
target_snap_id = row[0]
if not target_snap_id:
return None
return load_scud_from_db_by_snapshot(date_str, snapshot_param=target_snap_id)
def load_1c_files_for_date(date_str: str) -> Tuple[Optional[pd.DataFrame], Optional[pd.DataFrame]]:
"""
Загружает штат и отсутствия с каскадным fallback:
1. Синхронизирует свежие файлы с сетевой шары.
2. Загружает штат и отсутствия через load_1c_data_smart(..., use_db=True).
3. Если штат пуст — берет последний доступный срез штата из zup_staff в SQLite.
"""
clean_date = date_str.replace('_', '.')
try:
from services.share_copier import copy_1c_files_from_share
copy_1c_files_from_share()
except Exception as e:
logger.warning(f"[Pipeline] Ошибка копирования с шары: {e}")
df_staff, df_abs = load_1c_data_smart(clean_date, use_db=True)
# Fallback за штат: ако данашњи штат још увек није доступан, користи се претходни из базе
if df_staff is None or df_staff.empty:
with get_connection(row_factory=True) as conn:
cursor = conn.cursor()
cursor.execute("""
SELECT snapshot_date
FROM zup_staff
ORDER BY
SUBSTR(snapshot_date, 7, 4) DESC,
SUBSTR(snapshot_date, 4, 2) DESC,
SUBSTR(snapshot_date, 1, 2) DESC,
id DESC
LIMIT 1
""")
row = cursor.fetchone()
if row and row[0]:
fallback_date = row[0]
df_staff = pd.read_sql_query(
"SELECT fio as 'ФИО', fio_clean, department as 'Подразделение', position as 'Должность' FROM zup_staff WHERE snapshot_date = ?",
conn, params=(fallback_date,)
)
logger.info(f"[Pipeline] Для даты {clean_date} применен штат за {fallback_date} ({len(df_staff)} чел.)")
return df_staff, df_abs