meraproject/services/user-reader/app/main.py

641 lines
22 KiB
Python
Raw Normal View History

import os
import re
import traceback
from datetime import date, datetime
from decimal import Decimal
from typing import Annotated, Any
import pymysql
from fastapi import Depends, FastAPI, Header, HTTPException, Query, Security
from fastapi.middleware.cors import CORSMiddleware
from fastapi.responses import FileResponse
from fastapi.security import APIKeyHeader
from pymysql.cursors import DictCursor
from app.emp_departments import enrich_employee_items, load_emp_departments
from app.emp_staffing import enrich_employee_staffing, load_staffing_dict
from app.pagination_helpers import DEFAULT_FETCH_ALL_MAX
from app.emp_schema import (
EMP_EXCLUDE_SUFFIXES,
EMP_FIELD_SUFFIXES,
EMP_TABLE_CANONICAL,
PHP_REFERENCE,
)
MYSQL = {
"host": os.environ.get("MYSQL_HOST", "db"),
"port": int(os.environ.get("MYSQL_PORT", "3306")),
"user": os.environ.get("MYSQL_USER", "root"),
"password": os.environ.get("MYSQL_PASSWORD", ""),
"database": os.environ.get("MYSQL_DATABASE", "j7508239_tracker"),
"charset": "utf8mb4",
"cursorclass": DictCursor,
}
app = FastAPI(title="Mera user reader", version="1.0.0")
app.add_middleware(
CORSMiddleware,
allow_origins=["*"],
allow_methods=["*"],
allow_headers=["*"],
)
_api_key_header = APIKeyHeader(name="X-Api-Key", auto_error=False)
def require_api_key(
x_api_key: Annotated[str | None, Security(_api_key_header)] = None,
authorization: Annotated[str | None, Header()] = None,
) -> None:
"""Если задан USER_READER_API_KEY — доступ только с ключом (заголовок или Bearer)."""
expected = os.environ.get("USER_READER_API_KEY", "").strip()
if not expected:
return
got = (x_api_key or "").strip()
if not got and authorization:
auth = authorization.strip()
if auth.lower().startswith("bearer "):
got = auth[7:].strip()
if got != expected:
raise HTTPException(
status_code=401,
detail="Требуется ключ: заголовок X-Api-Key или Authorization: Bearer <ключ>",
headers={"WWW-Authenticate": "Bearer"},
)
def get_conn():
return pymysql.connect(**MYSQL)
def _json_cell(v: Any) -> Any:
if v is None:
return None
if isinstance(v, datetime):
return v.isoformat(sep=" ", timespec="seconds")
if isinstance(v, date):
return v.isoformat()
if isinstance(v, Decimal):
return float(v)
if isinstance(v, bytes):
return v.decode("utf-8", errors="replace")
return v
def _is_secret_col(name: str) -> bool:
n = name.lower()
return n.endswith("_pass") or "password" in n or "secret" in n
def json_rows(rows: list[dict]) -> list[dict]:
out = []
for row in rows:
clean = {k: v for k, v in row.items() if not _is_secret_col(k)}
out.append({k: _json_cell(v) for k, v in clean.items()})
return out
def _db_name(cur) -> str:
cur.execute("SELECT DATABASE() AS d")
return cur.fetchone()["d"]
def _list_tables(cur, db: str) -> list[str]:
cur.execute(
"""
SELECT TABLE_NAME AS n FROM information_schema.TABLES
WHERE TABLE_SCHEMA = %s ORDER BY TABLE_NAME
""",
(db,),
)
return [r["n"] for r in cur.fetchall()]
def _quote_ident(name: str) -> str:
return "`" + name.replace("`", "``") + "`"
def resolve_emp_table(cur, db: str) -> str:
"""Физическое имя таблицы (регистр как в MySQL)."""
env = os.environ.get("EMP_TABLE", "").strip().strip("`")
if env:
cur.execute(
"""
SELECT COUNT(*) AS c FROM information_schema.TABLES
WHERE TABLE_SCHEMA = %s AND TABLE_NAME = %s
""",
(db, env),
)
if cur.fetchone()["c"]:
return env
raise RuntimeError(
f"EMP_TABLE={env!r} не найдена в БД {db!r}. "
f"Таблицы: {', '.join(_list_tables(cur, db)[:40])}"
)
want = EMP_TABLE_CANONICAL.lower()
cur.execute(
"""
SELECT TABLE_NAME AS n
FROM information_schema.TABLES
WHERE TABLE_SCHEMA = %s
AND LOWER(TABLE_NAME) = %s
LIMIT 1
""",
(db, want),
)
row = cur.fetchone()
if row:
return row["n"]
cur.execute(
"""
SELECT TABLE_NAME AS n
FROM information_schema.TABLES
WHERE TABLE_SCHEMA = %s
AND LOWER(TABLE_NAME) REGEXP 'merakomis.*emp'
AND LOWER(TABLE_NAME) NOT LIKE '%%children%%'
AND LOWER(TABLE_NAME) NOT LIKE '%%group%%'
ORDER BY CHAR_LENGTH(TABLE_NAME)
LIMIT 1
""",
(db,),
)
row = cur.fetchone()
if row:
return row["n"]
fb = (
os.environ.get("USER_READER_FALLBACK_TABLE", "")
or os.environ.get("USER_TABLE", "")
).strip().strip("`")
if fb:
cur.execute(
"""
SELECT TABLE_NAME AS n FROM information_schema.TABLES
WHERE TABLE_SCHEMA = %s AND TABLE_NAME = %s
LIMIT 1
""",
(db, fb),
)
row = cur.fetchone()
if row:
return row["n"]
raise RuntimeError(
f"USER_TABLE / USER_READER_FALLBACK_TABLE={fb!r} не найдена в БД {db!r}."
)
all_t = _list_tables(cur, db)
hint = ", ".join([t for t in all_t if re.search(r"emp|user|staff|profile", t, re.I)][:25])
raise RuntimeError(
f"Нет таблицы {EMP_TABLE_CANONICAL!r} (см. {PHP_REFERENCE}). "
f"Импортируйте дамп Merakomis или задайте окружение, например: "
f"USER_TABLE=account (аккаунты CMS, module/core/user/account/model.php). "
f"Таблица achiprogressuser — это прогресс достижений (module/achi/progress/user/model.php), "
f"она может быть пустой. "
f"База {db!r}, похожие по имени: {hint or ''}"
)
def _table_columns(cur, db: str, table: str) -> list[str]:
cur.execute(
"""
SELECT COLUMN_NAME AS c
FROM information_schema.COLUMNS
WHERE TABLE_SCHEMA = %s AND TABLE_NAME = %s
ORDER BY ORDINAL_POSITION
""",
(db, table),
)
return [r["c"] for r in cur.fetchall()]
def _column_lookup(cols: list[str]) -> dict[str, str]:
return {c.lower(): c for c in cols}
def table_uses_merakomis_emp_schema(cols: list[str]) -> bool:
p = EMP_TABLE_CANONICAL.lower() + "_"
return any(c.lower().startswith(p) for c in cols)
def build_generic_select(
cols: list[str],
) -> tuple[str, str | None, str, list[str]]:
"""Произвольная таблица пользователей: все колонки кроме секретных."""
selected: list[str] = []
parts: list[str] = []
for c in cols:
if _is_secret_col(c):
continue
selected.append(c)
parts.append(_quote_ident(c))
if not parts:
raise RuntimeError("Нет ни одной не-секретной колонки для выборки.")
removed_db: str | None = None
for c in cols:
cl = c.lower()
if cl == "removed" or cl.endswith("_removed"):
removed_db = c
break
order_db: str | None = None
for prefer in ("name", "login", "title", "email", "id"):
for c in cols:
cl = c.lower()
if cl == prefer or cl.endswith("_" + prefer):
order_db = c
break
if order_db:
break
if not order_db:
order_db = cols[0]
return ", ".join(parts), removed_db, order_db, selected
def build_emp_select(
cols: list[str],
) -> tuple[str, str | None, str, list[str], list[str]]:
lut = _column_lookup(cols)
prefix = EMP_TABLE_CANONICAL
parts: list[str] = []
aliases: list[str] = []
skipped: list[str] = []
removed_db: str | None = None
order_db: str | None = None
for suf in EMP_FIELD_SUFFIXES:
if suf in EMP_EXCLUDE_SUFFIXES:
continue
logical_full = f"{prefix}_{suf}".lower()
db_col = lut.get(logical_full)
if not db_col:
skipped.append(suf)
continue
parts.append(f"{_quote_ident(db_col)} AS {_quote_ident(suf)}")
aliases.append(suf)
if suf == "removed":
removed_db = db_col
if suf == "name":
order_db = db_col
if not parts:
raise RuntimeError(
"В таблице нет ни одной колонки из списка Emp в emp_schema.py "
f"(префикс {prefix!r}). Колонки БД: {cols[:30]}"
)
if not order_db:
first_suffix = aliases[0]
order_db = lut[f"{prefix}_{first_suffix}".lower()]
select_sql = ", ".join(parts)
return select_sql, removed_db, order_db, aliases, skipped
def _delta_timestamp_column_name(cols: list[str], merakomis: bool) -> str:
"""Имя колонки БД для сравнения 'изменено после' (unix time или сопоставимое целое)."""
c = _delta_column_db_name_optional(cols, merakomis)
if c:
return c
if merakomis:
raise HTTPException(
status_code=400,
detail="Для Merakomis Emp нужна колонка updated (tMerakomisEmp_updated).",
)
raise HTTPException(
status_code=400,
detail="Для delta нужна колонка *_updated, updated, *_date или date.",
)
def _delta_column_db_name_optional(cols: list[str], merakomis: bool) -> str | None:
lut = _column_lookup(cols)
if merakomis:
return lut.get(f"{EMP_TABLE_CANONICAL.lower()}_updated")
for c in cols:
cl = c.lower()
if cl.endswith("_updated") or cl == "updated":
return c
for c in cols:
cl = c.lower()
if cl.endswith("_date") or cl == "date":
return c
return None
def _delta_json_field_name(cols: list[str], merakomis: bool) -> str | None:
"""Имя поля в JSON ответа (алиас или физическое имя колонки)."""
if merakomis:
return "updated" if _delta_column_db_name_optional(cols, True) else None
dbn = _delta_column_db_name_optional(cols, False)
return dbn
def _row_ts_value(v: Any) -> int:
if v is None:
return 0
if isinstance(v, datetime):
return int(v.timestamp())
if isinstance(v, date):
return int(datetime(v.year, v.month, v.day).timestamp())
if isinstance(v, Decimal):
return int(v)
return int(v)
@app.get("/")
def index():
return FileResponse(
os.path.join(os.path.dirname(__file__), "..", "static", "index.html")
)
@app.get("/api/health")
def health() -> dict[str, Any]:
try:
with get_conn() as c:
with c.cursor() as cur:
cur.execute("SELECT 1 AS ok")
row = cur.fetchone()
db = _db_name(cur)
return {"ok": True, "db": bool(row and row.get("ok") == 1), "database": db}
except pymysql.Error as e:
return {"ok": False, "db": False, "error": str(e)}
@app.get("/api/meta")
def meta(_auth: Annotated[None, Depends(require_api_key)]) -> dict[str, Any]:
try:
with get_conn() as c:
with c.cursor() as cur:
db = _db_name(cur)
table = resolve_emp_table(cur, db)
cols = _table_columns(cur, db, table)
if table_uses_merakomis_emp_schema(cols):
_, removed_db, order_db, aliases, skipped = build_emp_select(
cols
)
mode = "merakomis_emp"
else:
_, removed_db, order_db, aliases = build_generic_select(cols)
skipped = []
mode = "generic_user_table"
delta_field = _delta_json_field_name(cols, mode == "merakomis_emp")
return {
"php_model": PHP_REFERENCE,
"logical_table": EMP_TABLE_CANONICAL,
"database": db,
"physical_table": table,
"schema_mode": mode,
"delta_field": delta_field,
"columns_in_db": len(cols),
"selected_fields": aliases,
"skipped_unknown_in_db": skipped,
"where_removed": removed_db,
"order_by": order_db,
}
except pymysql.Error as e:
raise HTTPException(status_code=500, detail=str(e)) from e
except Exception as e:
raise HTTPException(status_code=500, detail=str(e)) from e
@app.get("/api/employees")
def employees(
_auth: Annotated[None, Depends(require_api_key)],
limit: int = Query(100, ge=1, le=500),
offset: int = Query(0, ge=0),
fetch_all: bool = Query(False, description="Вернуть всех сотрудников без пагинации"),
) -> dict[str, Any]:
try:
with get_conn() as c:
with c.cursor() as cur:
db = _db_name(cur)
table = resolve_emp_table(cur, db)
cols = _table_columns(cur, db, table)
if table_uses_merakomis_emp_schema(cols):
select_sql, removed_db, order_db, _, skipped = build_emp_select(
cols
)
else:
select_sql, removed_db, order_db, _sel = build_generic_select(
cols
)
skipped = []
tq = _quote_ident(table)
oq = _quote_ident(order_db)
# Фильтр "не удалён" только для Merakomis Emp; в произвольных таблицах
# колонка *removed* часто NULL или с другой семантикой — иначе получается 0 строк.
where_sql = ""
if table_uses_merakomis_emp_schema(cols) and removed_db:
where_sql = f" WHERE {_quote_ident(removed_db)} = 0 "
sql_base = (
f"SELECT {select_sql} FROM {tq} {where_sql}"
f" ORDER BY {oq}"
)
cnt_sql = f"SELECT COUNT(*) AS n FROM {tq} {where_sql}"
cur.execute(cnt_sql)
total = int(cur.fetchone()["n"])
if fetch_all:
if total > DEFAULT_FETCH_ALL_MAX:
raise HTTPException(
status_code=400,
detail={
"code": "too_many_rows",
"message": (
f"Слишком много строк ({total}); "
f"максимум {DEFAULT_FETCH_ALL_MAX}"
),
"total": total,
},
)
cur.execute(sql_base)
rows = cur.fetchall()
resp_limit, resp_offset = total, 0
else:
cur.execute(f"{sql_base} LIMIT %s OFFSET %s", (limit, offset))
rows = cur.fetchall()
resp_limit, resp_offset = limit, offset
items = json_rows(rows)
if table_uses_merakomis_emp_schema(cols):
dept_map = load_emp_departments(cur, db)
staffing_dict = load_staffing_dict(cur, db)
enrich_employee_items(items, dept_map)
enrich_employee_staffing(items, staffing_dict)
return {
"php_model": PHP_REFERENCE,
"logical_table": EMP_TABLE_CANONICAL,
"physical_table": table,
"table": table,
"skipped_unknown_in_db": skipped,
"total": total,
"limit": resp_limit,
"offset": resp_offset,
"fetch_all": fetch_all,
"count": len(items),
"items": items,
}
except pymysql.Error as e:
raise HTTPException(status_code=500, detail=str(e)) from e
except Exception as e:
dbg = os.environ.get("DEBUG", "")
msg = str(e)
if dbg == "1":
msg = f"{msg}\n{traceback.format_exc()}"
raise HTTPException(status_code=500, detail=msg) from e
@app.get("/api/employees/delta")
def employees_delta(
_auth: Annotated[None, Depends(require_api_key)],
since_updated: int = Query(0, ge=0, description="Unix time: строки строго новее этого значения"),
limit: int = Query(500, ge=1, le=500),
) -> dict[str, Any]:
"""Инкрементальная выборка по колонке времени (Merakomis: updated; иначе *_updated / date)."""
try:
with get_conn() as c:
with c.cursor() as cur:
db = _db_name(cur)
table = resolve_emp_table(cur, db)
cols = _table_columns(cur, db, table)
merakomis = table_uses_merakomis_emp_schema(cols)
delta_col = _delta_timestamp_column_name(cols, merakomis)
if merakomis:
select_sql, removed_db, order_db, _, skipped = build_emp_select(
cols
)
else:
select_sql, removed_db, order_db, _sel = build_generic_select(
cols
)
skipped = []
tq = _quote_ident(table)
dq = _quote_ident(delta_col)
conds: list[str] = [f"{dq} > %s"]
params: list[Any] = [since_updated]
if merakomis and removed_db:
conds.insert(0, f"{_quote_ident(removed_db)} = 0")
where_sql = " WHERE " + " AND ".join(conds)
lut = _column_lookup(cols)
order_parts = [f"{dq} ASC"]
if merakomis:
id_col = lut.get(f"{EMP_TABLE_CANONICAL.lower()}_id")
else:
id_col = None
for c in cols:
if c.lower() == "id" or c.lower().endswith("_id"):
id_col = c
break
if id_col:
order_parts.append(f"{_quote_ident(id_col)} ASC")
sql = (
f"SELECT {select_sql} FROM {tq}{where_sql}"
f" ORDER BY {', '.join(order_parts)} LIMIT %s"
)
cur.execute(sql, (*params, limit))
raw_rows = cur.fetchall()
row_key = "updated" if merakomis else delta_col
max_ts = since_updated
for row in raw_rows:
v = row.get(row_key)
if v is not None:
max_ts = max(max_ts, _row_ts_value(v))
items = json_rows(raw_rows)
if merakomis:
dept_map = load_emp_departments(cur, db)
staffing_dict = load_staffing_dict(cur, db)
enrich_employee_items(items, dept_map)
enrich_employee_staffing(items, staffing_dict)
return {
"php_model": PHP_REFERENCE,
"logical_table": EMP_TABLE_CANONICAL,
"physical_table": table,
"table": table,
"skipped_unknown_in_db": skipped,
"since_updated": since_updated,
"max_updated": max_ts,
"delta_column": row_key,
"schema_mode": "merakomis_emp" if merakomis else "generic_user_table",
"limit": limit,
"count": len(raw_rows),
"items": items,
}
except HTTPException:
raise
except pymysql.Error as e:
raise HTTPException(status_code=500, detail=str(e)) from e
except Exception as e:
dbg = os.environ.get("DEBUG", "")
msg = str(e)
if dbg == "1":
msg = f"{msg}\n{traceback.format_exc()}"
raise HTTPException(status_code=500, detail=msg) from e
@app.get("/labor")
def labor_page():
return FileResponse(
os.path.join(os.path.dirname(__file__), "..", "static", "labor.html")
)
@app.get("/summary")
def summary_page():
return FileResponse(
os.path.join(os.path.dirname(__file__), "..", "static", "summary.html")
)
@app.get("/project-report")
def project_report_page():
return FileResponse(
os.path.join(
os.path.dirname(__file__), "..", "static", "project-report.html"
)
)
@app.get("/time-entry")
def time_entry_page():
return FileResponse(
os.path.join(os.path.dirname(__file__), "..", "static", "time-entry.html")
)
@app.get("/project-member")
def project_member_page():
return FileResponse(
os.path.join(os.path.dirname(__file__), "..", "static", "project-member.html")
)
from app.batch_api import router as batch_router
from app.labor import router as labor_router
from app.labor_calendar import router as labor_calendar_router
from app.labor_write import router as labor_write_router
from app.project_members_read import router as project_members_read_router
from app.project_members_write import router as project_members_write_router
app.include_router(labor_router)
app.include_router(batch_router)
app.include_router(labor_write_router)
app.include_router(labor_calendar_router)
app.include_router(project_members_read_router)
app.include_router(project_members_write_router)