Files

296 lines
11 KiB
Python

# -*- coding: utf-8 -*-
"""RAG su export Irrigazione + Turni (P7)."""
from __future__ import annotations
import json
import logging
import os
from datetime import datetime, timezone
from typing import Any, Optional
from ..db import get_conn
from ..integrations import irrigazione, turni
from . import embeddings, qdrant_store
from .text_chunk import chunk_text as _chunk_text
LOGGER = logging.getLogger("loogle_mcp.apps_indexer")
INDEX_USER = os.environ.get("APPS_INDEX_USER", "daniele").lower()
def _int_env(name: str, default: int) -> int:
try:
return int(os.environ.get(name, str(default)))
except ValueError:
return default
def _enabled() -> bool:
flag = os.environ.get("APPS_INDEX_ENABLED", "yes").strip().lower()
return flag not in ("0", "false", "no", "off")
def _doc_id(source: str, record_id: str) -> int:
key = f"{source}:{record_id}"
return abs(hash(key)) % (2**31 - 1)
def _index_text(source: str, record_id: str, title: str, text: str) -> dict:
chunks = _chunk_text(_truncate(text))
if len(chunks) > 8:
chunks = chunks[:8]
if not chunks:
return {"source": source, "record_id": record_id, "chunks": 0, "skipped": True}
vectors = embeddings.embed_texts(chunks)
collection = qdrant_store.APPS_SHARED_COLLECTION
doc_id = _doc_id(source, record_id)
qdrant_store.delete_by_doc(collection, doc_id)
ids = []
payloads = []
for i, chunk in enumerate(chunks):
point_id = f"app-{source}-{record_id}-chunk-{i}"
ids.append(point_id)
payloads.append(
{
"doc_id": doc_id,
"source": source,
"record_id": record_id,
"chunk_index": i,
"title": title,
"text": chunk,
"owner": "family",
"visibility": "family",
}
)
qdrant_store.upsert_chunks(collection, ids, vectors, payloads)
get_conn().execute(
"INSERT INTO indexed_apps_records(source,record_id,title,owner,chunk_count,indexed_at)"
" VALUES (?,?,?,?,?,datetime('now'))"
" ON CONFLICT(source,record_id) DO UPDATE SET"
" title=excluded.title, chunk_count=excluded.chunk_count, indexed_at=datetime('now')",
(source, record_id, title, "family", len(chunks)),
)
get_conn().commit()
return {"source": source, "record_id": record_id, "chunks": len(chunks), "collection": collection}
def _format_irrigation_history_item(item: dict, idx: int) -> str:
parts = [f"Irrigazione storico #{idx}"]
for key in ("started_at", "ended_at", "zone", "zone_name", "duration_min", "volume_l", "mode", "note"):
if item.get(key) is not None:
parts.append(f"{key}: {item[key]}")
return "\n".join(parts)
def _format_irrigation_event(item: dict, idx: int) -> str:
parts = [f"Irrigazione evento #{idx}"]
for key in ("ts", "time", "type", "level", "message", "zone", "detail"):
if item.get(key) is not None:
parts.append(f"{key}: {item[key]}")
return "\n".join(parts)
def _format_turni_assignment(item: dict, idx: int) -> str:
parts = [f"Turno #{idx}"]
for key in (
"date", "startDate", "endDate", "doctorId", "doctorName", "doctor_name",
"slotId", "slotName", "slot_name", "uoc", "uocName", "shiftType", "notes",
):
if item.get(key) is not None:
parts.append(f"{key}: {item[key]}")
return "\n".join(parts)
def _truncate(text: str, limit: int = 6000) -> str:
if len(text) <= limit:
return text
return text[: limit - 20] + "\n… [truncated]"
def _summarize_irrigation_status(status: dict) -> str:
lines = ["Irrigazione — snapshot stato"]
for key in ("plan_mode", "program", "simulation", "hibernation", "draining"):
if key in status:
lines.append(f"{key}: {status[key]}")
ha = status.get("ha") or {}
lines.append(f"ha_connected: {ha.get('connected')}")
zones = status.get("zones") or []
lines.append(f"zone_count: {len(zones)}")
for z in zones[:12]:
if isinstance(z, dict):
lines.append(
f" - {z.get('name', z.get('id'))}: state={z.get('ha_state')} excluded={z.get('excluded')}"
)
analysis = status.get("analysis")
if isinstance(analysis, dict):
for k, v in list(analysis.items())[:8]:
lines.append(f"analysis.{k}: {v}")
elif analysis:
lines.append(f"analysis: {analysis}")
return "\n".join(lines)
def _summarize_zones(zones: Any) -> str:
items = zones if isinstance(zones, list) else _normalize_list(zones)
lines = [f"Irrigazione — zone ({len(items)})"]
for z in items[:20]:
if not isinstance(z, dict):
continue
lines.append(
f"- {z.get('name', z.get('id'))}: ha={z.get('ha_state')} "
f"rate_mmh={z.get('rate_mmh')} flow={z.get('zone_flow_lph')}"
)
return "\n".join(lines)
def _normalize_list(data: Any) -> list:
if isinstance(data, list):
return data
if isinstance(data, dict):
for key in ("items", "results", "history", "events", "assignments", "records"):
val = data.get(key)
if isinstance(val, list):
return val
return []
def index_irrigazione(*, username: Optional[str] = None, history_limit: Optional[int] = None, events_limit: Optional[int] = None) -> dict:
user = username or INDEX_USER
history_limit = history_limit if history_limit is not None else _int_env("APPS_INDEX_HISTORY_LIMIT", 25)
events_limit = events_limit if events_limit is not None else _int_env("APPS_INDEX_EVENTS_LIMIT", 40)
if not irrigazione.is_configured(user):
return {"source": "irrigazione", "skipped": True, "reason": "not configured"}
indexed = 0
errors = 0
try:
status = irrigazione.get_status(user)
status_text = _summarize_irrigation_status(status)
_index_text("irrigazione", "status-snapshot", "Irrigazione — stato attuale", status_text)
indexed += 1
except Exception as exc:
LOGGER.warning("Irrigazione status index failed: %s", exc)
errors += 1
try:
zones = irrigazione.get_zones(user)
zones_text = _summarize_zones(zones)
_index_text("irrigazione", "zones-snapshot", "Irrigazione — zone", zones_text)
indexed += 1
except Exception as exc:
LOGGER.warning("Irrigazione zones index failed: %s", exc)
errors += 1
try:
history = irrigazione.get_history(user, limit=history_limit)
for i, item in enumerate(_normalize_list(history)):
if not isinstance(item, dict):
continue
rid = str(item.get("id") or item.get("started_at") or i)
text = _format_irrigation_history_item(item, i)
_index_text("irrigazione", f"history-{rid}", f"Irrigazione storico {rid}", text)
indexed += 1
except Exception as exc:
LOGGER.warning("Irrigazione history index failed: %s", exc)
errors += 1
try:
events = irrigazione.get_events(user, limit=events_limit)
for i, item in enumerate(_normalize_list(events)):
if not isinstance(item, dict):
continue
rid = str(item.get("id") or item.get("ts") or item.get("time") or i)
text = _format_irrigation_event(item, i)
_index_text("irrigazione", f"event-{rid}", f"Irrigazione evento {rid}", text)
indexed += 1
except Exception as exc:
LOGGER.warning("Irrigazione events index failed: %s", exc)
errors += 1
try:
lavori = irrigazione.get_lavori_summary(user)
lavori_text = _truncate(json.dumps(lavori, ensure_ascii=False, indent=2), 4000)
_index_text("irrigazione", "lavori-summary", "Irrigazione — lavori manutenzione", lavori_text)
indexed += 1
except Exception as exc:
LOGGER.warning("Irrigazione lavori index failed: %s", exc)
errors += 1
return {"source": "irrigazione", "indexed": indexed, "errors": errors}
def index_turni(*, username: Optional[str] = None, assignments_limit: Optional[int] = None) -> dict:
user = username or INDEX_USER
assignments_limit = assignments_limit if assignments_limit is not None else _int_env("APPS_INDEX_ASSIGNMENTS_LIMIT", 80)
if not turni.is_configured(user):
return {"source": "turni", "skipped": True, "reason": "not configured"}
indexed = 0
errors = 0
try:
status = turni.get_status()
status_text = _truncate(json.dumps(status, ensure_ascii=False, indent=2), 2000)
_index_text("turni", "status-snapshot", "Turni — stato servizio", status_text)
indexed += 1
except Exception as exc:
LOGGER.warning("Turni status index failed: %s", exc)
errors += 1
try:
doctors = turni.list_doctors(user)
lines = ["Turni — medici"]
for d in (_normalize_list(doctors) if not isinstance(doctors, list) else doctors)[:40]:
if isinstance(d, dict):
lines.append(f"- {d.get('name', d.get('fullName'))} id={d.get('id')}")
_index_text("turni", "doctors-list", "Turni — elenco medici", "\n".join(lines))
indexed += 1
except Exception as exc:
LOGGER.warning("Turni doctors index failed: %s", exc)
errors += 1
try:
assignments = turni.get_shift_assignments(user, limit=assignments_limit)
for i, item in enumerate(_normalize_list(assignments)):
if not isinstance(item, dict):
continue
rid = str(item.get("id") or item.get("date") or i)
text = _format_turni_assignment(item, i)
_index_text("turni", f"assignment-{rid}", f"Turno {rid}", text)
indexed += 1
except Exception as exc:
LOGGER.warning("Turni assignments index failed: %s", exc)
errors += 1
return {"source": "turni", "indexed": indexed, "errors": errors}
def index_all(*, username: Optional[str] = None) -> dict:
if not _enabled():
return {"skipped": True, "reason": "APPS_INDEX_ENABLED=no"}
user = username or INDEX_USER
return {
"irrigazione": index_irrigazione(username=user),
"turni": index_turni(username=user),
"at": datetime.now(timezone.utc).isoformat(),
}
def search_apps_knowledge(query: str, limit: int = 8, source: Optional[str] = None) -> list:
vectors = embeddings.embed_texts([query])
flt = {"source": source} if source else None
hits = qdrant_store.search(
[qdrant_store.APPS_SHARED_COLLECTION],
vectors[0],
limit=limit,
visibility_filter=flt,
)
for hit in hits:
hit.setdefault("source_type", hit.get("source", "apps"))
return hits
def list_indexed_records(source: Optional[str] = None, limit: int = 40) -> list:
conn = get_conn()
if source:
rows = conn.execute(
"SELECT * FROM indexed_apps_records WHERE source=? ORDER BY indexed_at DESC LIMIT ?",
(source, limit),
).fetchall()
else:
rows = conn.execute(
"SELECT * FROM indexed_apps_records ORDER BY indexed_at DESC LIMIT ?",
(limit,),
).fetchall()
return [dict(r) for r in rows]