Files
astrololo/services/data/app/providers/excel_provider.py
T
gitea caf4fd80d1 feat(astroklient): pule plików per konto i izolacja od produkcji (PRE-29)
Demo ma być rozdawane szeroko i różnym osobom, więc pierwsza wersja — jedno konto
na produkcyjnej warstwie danych — nie nadawała się do użycia: każdy dostawałby
dostęp do oryginalnych baz, a wgrania jednego klienta widzieliby wszyscy.

IZOLACJA OD PRODUKCJI. Warstwa danych i logiczna demo są osobne (manifesty w repo
deploy). Osobna musi być TEŻ LOGICZNA, bo zna ona jeden adres warstwy danych —
demo korzystające z produkcyjnej logiki i tak trafiłoby na produkcyjne bazy.

PULE PER KONTO w warstwie danych. Zapytanie i lista plików niosą nazwę puli;
puste = cały udział, czyli produkcja działa dokładnie jak dotąd i o pulach nic
nie wie. Nazwa puli przechodzi przez sito dopuszczające wyłącznie znaki bezpieczne
w nazwie katalogu — „../..” albo ukośnik wyprowadziłyby zapytanie wprost do cudzych
baz, więc sito ZAMIENIA podejrzane znaki zamiast ufać, że nikt ich nie poda.

PULA MUSI BYĆ W KLUCZU CACHE ZAPYTAŃ. Bez tego wynik policzony dla jednego konta
trafiłby z cache do drugiego — cicha wymiana treści baz między klientami,
niewidoczna w logach i nie do wykrycia z zewnątrz. Osobny test tego pilnuje.

PULA WYNIKA Z LOGINU, nigdy z żądania. Klient warstwy logicznej jest budowany
per żądanie i związany z pulą zalogowanej osoby; gdyby nazwa przychodziła
z formularza, wystarczyłoby podstawić cudzy login. Test wysyła `tenant`, `user`
i `login` w polach formularza i sprawdza, że nie mają na nią wpływu.

Pulę wstrzykujemy w INSTANCJĘ klienta, nie w sygnatury metod. Argumentem trzeba
by ją przeprowadzić przez protokół DataSource i build_report — kod, który o kontach
nie ma prawa nic wiedzieć — a każde nowe wywołanie byłoby okazją, żeby o nią
zapomnieć i sięgnąć nie tam.

Konta demo to lista `login:sekret` (DEMO_USERS), bo jedno wspólne konto oznaczałoby
wspólną pulę. Format i skrypt haseł te same, co w głównej aplikacji.

Pula klienta to JEDEN KATALOG, więc przejście na pełną wersję nie oznacza utraty
wgrań — procedurę importu opisuje runbook w repo deploy.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-26 11:57:05 +02:00

190 lines
8.0 KiB
Python

"""ExcelDataProvider — wyszukiwanie w setkach plików .xlsx z 4-poziomowym cache.
Ścieżka zapytania (od najszybszej):
1) QueryCache (in-memory) -> gotowy wynik
2) InvertedIndex (SQLite) -> które pliki w ogóle otwierać (zamiast skanu setek)
3) FrameCache (Parquet) -> wczytanie pliku bez parsowania .xlsx
4) SchemaCache (SQLite) -> bez ponownego wykrywania nagłówka/układu kolumn
...dopiero gdy wszystko spudłuje, czytamy .xlsx i wypełniamy cache.
Cała ta złożoność jest UKRYTA za interfejsem DataProvider.
"""
from __future__ import annotations
import time
from pathlib import Path
import pandas as pd
from app.cache.fingerprint import fingerprint
from app.cache.frame_cache import FrameCache
from app.cache.index import InvertedIndex
from app.cache.query_cache import QueryCache
from app.cache.schema_cache import SchemaCache
from app.config import Settings
from app.excel.header_detect import detect_header_row
from app.excel.layout import build_column_mapping
from app.models import HealthInfo, SearchQuery, SearchResult
from app.providers.base import DataProvider
class ExcelDataProvider(DataProvider):
name = "excel"
def __init__(self, settings: Settings) -> None:
self.s = settings
self.schema = SchemaCache(settings.cache_dir)
self.frames = FrameCache(settings.cache_dir)
self.index = InvertedIndex(settings.cache_dir)
self.queries = QueryCache(settings.query_cache_size, settings.query_cache_ttl)
# ---- ładowanie pojedynczego arkusza z pełnym cache ----
def _load_frame(self, path: str, sheet: str | int = 0) -> pd.DataFrame:
fp = fingerprint(path)
sheet_key = str(sheet)
cached = self.frames.get(fp, sheet_key) # poziom 2: Parquet
if cached is not None:
return cached
raw = pd.read_excel(path, sheet_name=sheet, header=None, dtype=object)
meta = self.schema.get(fp, sheet_key) # poziom 1: schemat
if meta is None:
header_row = detect_header_row(raw, self.s.header_scan_rows)
header_cells = [str(c) for c in raw.iloc[header_row].tolist()]
mapping = build_column_mapping(header_cells)
self.schema.put(fp, sheet_key, header_row, mapping)
else:
header_row, mapping = meta
header_cells = [str(c) for c in raw.iloc[header_row].tolist()]
data = raw.iloc[header_row + 1 :].copy()
data.columns = header_cells
data = data.dropna(how="all")
inverse = {orig: canon for canon, orig in mapping.items()}
data = data.rename(columns=inverse).reset_index(drop=True)
self.frames.put(fp, sheet_key, data) # zapisz Parquet na przyszłość
return data
# ---- budowa odwróconego indeksu (warmup / po zmianie pliku) ----
def _ensure_indexed(self, path: str) -> None:
fp = fingerprint(path)
if self.index.file_fingerprint(path) == fp:
return # aktualny
frame = self._load_frame(path)
rows: list[tuple[str, str, str]] = []
for key in self.s.indexed_keys:
if key in frame.columns:
for v in frame[key].dropna().astype(str).unique():
rows.append((key, v, "0"))
self.index.reindex_file(path, fp, rows)
def warmup(self) -> None:
for path in self._excel_files():
try:
self._ensure_indexed(path)
except Exception as e: # jeden uszkodzony plik nie może zablokować startu
print(f"[data] pominięto plik przy indeksowaniu: {path}{e}")
def _excel_files(self) -> list[str]:
base = Path(self.s.excel_dir)
return [str(p) for p in sorted(base.glob("**/*.xlsx")) if not p.name.startswith("~$")]
def _enabled_files(self, paths: list[str], tenant: str = "") -> list[str]:
"""Bazy biorące udział w wyszukiwaniu.
Źródłem prawdy jest REJESTR PLIKÓW (DAN-27) — stan klikany z ekranu,
trwały na udziale. Zmienna DISABLED_BASES z DAN-15 zostaje jako awaryjne
wyłączenie z konfiguracji: gdy jest ustawiona, odsiewa DODATKOWO. Nie
odwrotnie — inaczej ktoś z dostępem do ekranu mógłby włączyć bazę
wyłączoną świadomie na poziomie wdrożenia.
"""
# `bases` MUSI być zaimportowane tutaj — modułowego importu nie ma,
# a przepisując tę funkcję pod rejestr usunąłem lokalny. Efekt: NameError
# przy KAŻDYM wyszukiwaniu, czyli 500 z warstwy danych.
from app import bases, files
usable = set(files.usable_paths(files.tenant_root(self.s.excel_dir, tenant)))
out = [p for p in paths if p in usable]
entries = bases.disabled_entries()
if entries:
out = [p for p in out if bases.is_enabled(p, self.s.excel_dir, entries)]
return out
def list_bases(self) -> list[dict]:
"""Bazy dostępne na udziale + metaopis + stan włączenia (DAN-15/PRE-09)."""
from app import bases
from app import files
return files.registry(self.s.excel_dir, for_admin=True)
# ---- publiczne API ----
def search(self, query: SearchQuery) -> SearchResult:
t0 = time.perf_counter()
# Lista wyłączonych baz wchodzi do klucza cache: bez tego zmiana ustawień
# oddawałaby wynik sprzed zmiany, czyli treść bazy uznanej za wyłączoną.
from app import bases
disabled = ",".join(bases.disabled_entries())
# PULA MUSI BYĆ W KLUCZU. Bez niej wynik policzony dla jednego konta
# trafiłby z cache do drugiego — czyli cicha wymiana treści baz między
# kontami, niewidoczna w logach i nie do wykrycia z zewnątrz.
cache_key = (f"{query.key}|{query.value}|{query.exact}|{query.limit}"
f"|{query.fields}|{disabled}|{query.tenant}")
hit = self.queries.get(cache_key) # poziom 3: wynik zapytania
if hit is not None:
hit = hit.model_copy(update={"cache": "hit", "elapsed_ms": _ms(t0)})
return hit
candidates = self.index.lookup(query.key, query.value, query.exact)
if not candidates:
# brak w indeksie (np. klucz nieindeksowany) -> przeszukaj wszystkie pliki
candidates = [(p, "0") for p in self._excel_files()]
# bazy wyłączone globalnie (DAN-15) pomijamy niezależnie od źródła kandydatów
allowed = set(self._enabled_files([p for p, _ in candidates], query.tenant))
candidates = [(p, s) for p, s in candidates if p in allowed]
rows: list[dict] = []
for path, _sheet in candidates:
frame = self._load_frame(path)
if query.key not in frame.columns:
continue
col = frame[query.key].astype(str)
if query.exact:
mask = col.str.lower() == query.value.lower()
else:
# regex=False: wartości sygnifikatorów zawierają znaki [ + itd.,
# które są metaznakami regex — szukamy dosłownie.
mask = col.str.lower().str.contains(query.value.lower(), na=False, regex=False)
matched = frame[mask]
if query.fields:
keep = [c for c in query.fields if c in matched.columns]
matched = matched[keep]
rows.extend(matched.to_dict(orient="records"))
if len(rows) >= query.limit:
break
result = SearchResult(
rows=rows[: query.limit],
total=len(rows),
elapsed_ms=_ms(t0),
cache="miss",
provider=self.name,
)
self.queries.put(cache_key, result)
return result
def health(self) -> HealthInfo:
return HealthInfo(
provider=self.name,
indexed_files=self.index.count_files(),
details={"excel_dir": str(self.s.excel_dir), "files_on_disk": len(self._excel_files())},
)
def _ms(t0: float) -> float:
return round((time.perf_counter() - t0) * 1000, 2)