diff --git a/docs/PROVIDER_RPC_SPEC.md b/docs/PROVIDER_RPC_SPEC.md index 2324f6f..d9bc69b 100644 --- a/docs/PROVIDER_RPC_SPEC.md +++ b/docs/PROVIDER_RPC_SPEC.md @@ -99,10 +99,21 @@ Tout se branche dans le `when (method)` existant de - **retour** : `result_b64` = base64 du profil (toutes colonnes de `/profiles`, dont `system_prompt_override` multiligne). `ok=false, error=not_found` si l'id n'existe pas. +### 2.4 `rag_dump_json` +- **arg** : `{"source": ""}`. +- **retour** : `result_b64` = base64 de `{"source": "...", "text": ""}`. +- **pourquoi** : `ragDocCursor` (lignes ~209-213) renvoie la colonne `text` = la **fiche + RAG complète** (titres, `:`, listes, sauts de ligne). La relire via + `query /rag/{source}` la corrompt exactement comme un prompt. Symétrique de + `rag_upsert_json` (§3.4) — sans cette méthode, l'authoring RAG serait écriture-seule. + `ok=false, error=not_found` si la source n'existe pas. + > Les autres lectures (`/state`, `/turns`, `/crashes`, `/models`, `/voices`, -> `/rag_status`, `/updates`, `/dist_config`, liste des sessions `/conversations`) -> **restent en `query`** : leurs valeurs sont numériques ou des chaînes courtes sans -> `:`/saut de ligne. Pas de dump nécessaire. +> `/rag_status`, `/updates`, `/dist_config`, l'**index** RAG `/rag`, liste des sessions +> `/conversations`) **restent en `query`** : leurs valeurs sont numériques ou des chaînes +> courtes sans `:`/saut de ligne. Seul le **texte** d'une fiche (`/rag/{source}`) exige +> le dump §2.4 ; son index (`/rag` : `source, char_count, chunk_count, updated_at`) reste +> en `query`. --- @@ -129,6 +140,10 @@ Toutes diffusent le broadcast de reload approprié (déjà définis : `notes`, `is_default`). `id` absent → création. - **effet** : réutiliser `upsertProfile(ContentValues)` en construisant le `ContentValues` depuis le JSON. Broadcast `RELOAD_PROFILES`. +- **PIN** (sémantique déjà gérée par `upsertProfile`, l. 493-508) — passer le PIN **en + clair** dans la clé `pin`, jamais le hash : `pin` **absent** = PIN inchangé ; `pin=""` + = efface le PIN ; `pin="1234"` = définit (le provider hashe en SHA-256). Le client + central ne connaît jamais `pin_hash` (cf. `profile_dump_json` qui rend `has_pin`). - **retour** : `ok`, et `result_b64` = base64 `{"id": ""}`. ### 3.4 `rag_upsert_json` @@ -183,15 +198,24 @@ Forme **imbriquée** utilisée par `cfg_dump_json` (complète) et `cfg_apply_jso (partielle). Noms de clés = colonnes du provider en `snake_case`. ### 5.1 Objet racine +> ⚠️ Forme = **miroir exact de `runtime.json`** (ce que `ConfigStore.toJson()` +> sérialise, vérifié l. 204-231). Les clés RAG sont **à plat** +> (`rag_enabled`/`rag_threshold`/`rag_top_k`), **PAS** imbriquées sous `"rag"` : +> c'est ainsi que `toJson()`/`fromJson()` les lisent et écrivent, donc +> `cfg_dump_json`/`cfg_apply_json` réutilisent ces helpers **sans couche de +> mapping**. Seuls `speaker`/`thinker`/`presets` sont imbriqués (idem `toJson`). + ```json { - "cascade_enabled": true, + "cascade_enabled": false, "tts_enabled": true, - "stt_engine": "whisper", - "llm_engine": "engine", + "stt_engine": "prod", + "llm_engine": "prod", "tts_engine": "cosyvoice", + "rag_enabled": false, + "rag_threshold": 0.82, + "rag_top_k": 3, "debug_enabled": false, - "rag": { "enabled": false, "threshold": 0.82, "top_k": 3 }, "speaker": { /* §5.2 */ }, "thinker": { /* §5.2 */ }, "presets": [ /* §5.3 */ ] @@ -262,7 +286,7 @@ Pour piloter une flotte de tablettes potentiellement sur des versions différent "supported_calls": [ "update_check","update_install", "cfg_dump_json","cfg_apply_json","presets_dump_json","presets_apply_json", - "profile_dump_json","profile_upsert_json","rag_upsert_json", + "profile_dump_json","profile_upsert_json","rag_upsert_json","rag_dump_json", "conversations_export","export_purge","capabilities","rag_sync","voices_reload" ] } @@ -287,7 +311,7 @@ Pour piloter une flotte de tablettes potentiellement sur des versions différent - [ ] Constantes des nouveaux noms `call()` dans le `companion object`. - [ ] Helpers `ok/err/decode/encode` (§1.2). -- [ ] `cfg_dump_json`, `presets_dump_json`, `profile_dump_json` (§2). +- [ ] `cfg_dump_json`, `presets_dump_json`, `profile_dump_json`, `rag_dump_json` (§2). - [ ] `cfg_apply_json`, `presets_apply_json`, `profile_upsert_json`, `rag_upsert_json` (§3) — **réutiliser** `updateConfig/mergeModel/upsertProfile/insert` existants. - [ ] `conversations_export` (+ `export_purge`) (§4) — décider chiffré vs plaintext-MVP, le **documenter**. - [ ] Gate token (§6) en plomberie désactivée. diff --git a/kazeia_central/__main__.py b/kazeia_central/__main__.py new file mode 100644 index 0000000..f08ea56 --- /dev/null +++ b/kazeia_central/__main__.py @@ -0,0 +1,38 @@ +"""Point d'entrée : `python -m kazeia_central`. + +Hôte/port configurables par variables d'environnement (local-first par défaut) : + KAZEIA_HOST défaut 127.0.0.1 — mettre 0.0.0.0 pour exposer sur le réseau + KAZEIA_PORT défaut 8000 + KAZEIA_USER / KAZEIA_PASS — si définis, active l'auth HTTP Basic (UI + API) + +⚠️ Exposer sur 0.0.0.0 sans KAZEIA_USER/KAZEIA_PASS = PII de santé + pilotage flotte +accessibles à tout le réseau. À éviter (cf. §8 RGPD). +""" + +from __future__ import annotations + +import os + + +def main() -> None: + import uvicorn + + host = os.environ.get("KAZEIA_HOST", "127.0.0.1") + port = int(os.environ.get("KAZEIA_PORT", "8000")) + has_auth = bool(os.environ.get("KAZEIA_USER") and os.environ.get("KAZEIA_PASS")) + exposed = host not in ("127.0.0.1", "localhost", "::1") + + if exposed and not has_auth: + print("\033[33m" + + f"⚠️ Kazeia-central écoute sur {host}:{port} SANS authentification.\n" + " Données PII santé + pilotage flotte exposés sur le réseau.\n" + " → définis KAZEIA_USER et KAZEIA_PASS pour activer l'auth Basic." + + "\033[0m", flush=True) + elif exposed: + print(f"Kazeia-central sur {host}:{port} — auth Basic active.", flush=True) + + uvicorn.run("kazeia_central.api.app:app", host=host, port=port) + + +if __name__ == "__main__": + main() diff --git a/kazeia_central/adb/client.py b/kazeia_central/adb/client.py index 3084513..2e41e7b 100644 --- a/kazeia_central/adb/client.py +++ b/kazeia_central/adb/client.py @@ -16,6 +16,7 @@ import base64 import json import re import subprocess +import time from dataclasses import dataclass, field from typing import Any @@ -24,6 +25,17 @@ PROVIDER_URI = f"content://{AUTHORITY}" ContentRow = dict[str, Any] +# États possibles d'une tablette dans `adb devices -l` (2e colonne). Sert à +# distinguer une vraie ligne device des lignes de bruit (`* daemon ... *`, +# l'en-tête "List of devices attached", lignes vides). +_DEVICE_STATES = frozenset({ + "device", "offline", "unauthorized", "authorizing", "bootloader", + "recovery", "sideload", "connecting", "host", "no", # "no permissions" +}) + +# Délai avant le retry de réveil quand le provider est endormi/throttlé (§3.5). +_WAKE_RETRY_DELAY = 0.4 + class AdbError(RuntimeError): """Échec d'une commande adb (code de retour non nul ou adb introuvable).""" @@ -85,11 +97,14 @@ class Adb: def devices(self) -> list[AdbDevice]: out = self._run(["devices", "-l"]).stdout.decode("utf-8", "replace") devices: list[AdbDevice] = [] - for line in out.splitlines()[1:]: # saute l'en-tête "List of devices attached" + for line in out.splitlines(): + # Ne pas se fier à la position de l'en-tête : adb émet des lignes + # `* daemon not running... *` AVANT "List of devices attached". + # Une vraie ligne device a ` [props...]`. line = line.strip() - if not line: - continue parts = line.split() + if len(parts) < 2 or parts[1] not in _DEVICE_STATES: + continue serial, state = parts[0], parts[1] props = {} for tok in parts[2:]: @@ -104,18 +119,38 @@ class Adb: # ---- provider : lecture (query) --------------------------------------- def content_query(self, path: str, *, serial: str | None = None, - where: str | None = None) -> list[ContentRow]: + where: str | None = None, expect_rows: bool = False, + _retry: bool = True) -> list[ContentRow]: """`content query` sur content://com.kazeia.provider/. ⚠️ Le format de sortie d'`adb content query` (`Row: N k=v, k=v`) est FRAGILE : il casse si une valeur contient `, ` ou un saut de ligne. N'utiliser que pour les endpoints à valeurs simples. Pour le texte riche, attendre les méthodes call() base64-JSON (PROVIDER_RPC_SPEC.md). + + Réveil/retry (§3.5) : le provider peut être endormi/throttlé et alors la + commande échoue ou rend 0 ligne. La 1ʳᵉ requête sert de réveil ; on + re-tente UNE fois. `expect_rows=True` pour les endpoints mono-ligne + (/state, /rag_status, /updates, /dist_config) où 0 ligne = endormi, pas + un résultat légitime. Les endpoints-liste laissent `expect_rows=False` + (0 ligne est une réponse valide → pas de retry inutile). """ cmd = f"content query --uri {PROVIDER_URI}/{path}" if where is not None: cmd += f' --where "{where}"' - return parse_content_query(self.shell(cmd, serial=serial)) + try: + rows = parse_content_query(self.shell(cmd, serial=serial)) + except AdbError: + if not _retry: + raise + time.sleep(_WAKE_RETRY_DELAY) + return self.content_query(path, serial=serial, where=where, + expect_rows=expect_rows, _retry=False) + if expect_rows and not rows and _retry: + time.sleep(_WAKE_RETRY_DELAY) + return self.content_query(path, serial=serial, where=where, + expect_rows=expect_rows, _retry=False) + return rows # ---- provider : RPC (call) — pour les méthodes base64-JSON futures ----- def content_call(self, method: str, *, arg_json: Any | None = None, diff --git a/kazeia_central/api/app.py b/kazeia_central/api/app.py index 92900e7..89dc818 100644 --- a/kazeia_central/api/app.py +++ b/kazeia_central/api/app.py @@ -10,15 +10,50 @@ Lancer : `uvicorn kazeia_central.api.app:app --reload` from __future__ import annotations +import base64 +import os +import secrets +from pathlib import Path + from fastapi import FastAPI, HTTPException +from fastapi.staticfiles import StaticFiles +from starlette.responses import Response from ..adb import Adb, AdbError from ..provider import ProviderClient +_WEB_DIR = Path(__file__).resolve().parent.parent / "web" + + +def _install_basic_auth(app: FastAPI) -> None: + """Auth HTTP Basic optionnelle, active SEULEMENT si KAZEIA_USER+KAZEIA_PASS + sont définis. Couvre tout (UI statique incluse) via middleware. Pensée pour + sécuriser une exposition réseau (0.0.0.0) — sans elle, l'app reste ouverte.""" + user = os.environ.get("KAZEIA_USER") + pwd = os.environ.get("KAZEIA_PASS") + if not (user and pwd): + return + + @app.middleware("http") + async def _basic_auth(request, call_next): + hdr = request.headers.get("authorization", "") + ok = False + if hdr.startswith("Basic "): + try: + u, _, p = base64.b64decode(hdr[6:]).decode("utf-8").partition(":") + ok = secrets.compare_digest(u, user) and secrets.compare_digest(p, pwd) + except Exception: + ok = False + if not ok: + return Response("Authentification requise", status_code=401, + headers={"WWW-Authenticate": 'Basic realm="Kazeia-central"'}) + return await call_next(request) + def create_app(adb: Adb | None = None) -> FastAPI: adb = adb or Adb() app = FastAPI(title="Kazeia-central", version="0.0.1") + _install_basic_auth(app) def client(serial: str) -> ProviderClient: return ProviderClient(adb, serial) @@ -79,6 +114,11 @@ def create_app(adb: Adb | None = None) -> FastAPI: def sessions(serial: str, profile_id: str | None = None): return [s.model_dump() for s in _guard(lambda: client(serial).sessions(profile_id))] + # UI web locale servie à la racine — APRÈS les routes /api pour ne pas les + # masquer. `html=True` → sert index.html sur "/". + if _WEB_DIR.is_dir(): + app.mount("/", StaticFiles(directory=str(_WEB_DIR), html=True), name="web") + return app diff --git a/kazeia_central/provider/client.py b/kazeia_central/provider/client.py index bf59323..21cc93f 100644 --- a/kazeia_central/provider/client.py +++ b/kazeia_central/provider/client.py @@ -9,7 +9,7 @@ NotImplemented explicites pour cadrer l'intégration future. from __future__ import annotations -from ..adb import Adb +from ..adb import Adb, AdbError from . import models as m @@ -21,9 +21,17 @@ class ProviderClient: def _q(self, path: str, **kw) -> list[dict]: return self.adb.content_query(path, serial=self.serial, **kw) + def _one(self, path: str, **kw) -> dict: + """Endpoint mono-ligne : réveille le provider si besoin (§3.5) et lève + une erreur explicite plutôt qu'un IndexError si toujours vide.""" + rows = self._q(path, expect_rows=True, **kw) + if not rows: + raise AdbError(f"{path}: provider sans réponse (endormi ?) après réveil") + return rows[0] + # ---- lecture (endpoints actuels, valeurs simples) ---------------------- def state(self) -> m.StateRow: - return m.StateRow(**self._q("state")[0]) + return m.StateRow(**self._one("state")) def turns(self) -> list[m.TurnRow]: return [m.TurnRow(**r) for r in self._q("turns")] @@ -38,7 +46,7 @@ class ProviderClient: return [m.VoiceRow(**r) for r in self._q("voices")] def rag_status(self) -> m.RagStatus: - return m.RagStatus(**self._q("rag_status")[0]) + return m.RagStatus(**self._one("rag_status")) def rag_index(self) -> list[m.RagDocRow]: return [m.RagDocRow(**r) for r in self._q("rag")] @@ -48,10 +56,10 @@ class ProviderClient: return [m.RagHit(**r) for r in self._q("rag_query", where=text)] def updates(self) -> m.UpdateStatus: - return m.UpdateStatus(**self._q("updates")[0]) + return m.UpdateStatus(**self._one("updates")) def dist_config(self) -> m.DistConfig: - return m.DistConfig(**self._q("dist_config")[0]) + return m.DistConfig(**self._one("dist_config")) def sessions(self, profile_id: str | None = None) -> list[m.SessionRow]: path = f"conversations/profile/{profile_id}" if profile_id else "conversations" @@ -61,7 +69,7 @@ class ProviderClient: return [m.ProfileRow(**r) for r in self._q("profiles")] def active_profile_id(self) -> str | None: - rows = self._q("profiles/active") + rows = self._q("profiles/active", expect_rows=True) return rows[0].get("active_profile_id") if rows else None # ---- MAJ OTA (call() déjà disponibles côté app) ------------------------ diff --git a/kazeia_central/provider/models.py b/kazeia_central/provider/models.py index 30e5cbf..d29c482 100644 --- a/kazeia_central/provider/models.py +++ b/kazeia_central/provider/models.py @@ -21,10 +21,14 @@ class _Base(BaseModel): class StateRow(_Base): - pid: int - uptime_seconds: int - pss_mb: int - ion_mb: int + # Tous optionnels : sur une flotte hétérogène (§8 handshake), une tablette + # plus ancienne peut omettre une colonne. Mieux vaut une cellule vide qu'un + # endpoint /state mort (extra="ignore" tolère déjà les colonnes en PLUS ; + # ceci tolère les colonnes en MOINS). + pid: int | None = None + uptime_seconds: int | None = None + pss_mb: int | None = None + ion_mb: int | None = None # Colonnes présentes dans le code mais non documentées au §4.1 du CLAUDE.md # (relevées par le dev) : anon_pages_mb: int | None = None @@ -105,7 +109,7 @@ class RagHit(_Base): class UpdateStatus(_Base): - phase: str + phase: str | None = None app_update_available: bool = False app_version_name: str | None = None app_size: int | None = None diff --git a/kazeia_central/store/__init__.py b/kazeia_central/store/__init__.py new file mode 100644 index 0000000..b0a77c1 --- /dev/null +++ b/kazeia_central/store/__init__.py @@ -0,0 +1,6 @@ +"""Store local chiffré (archive clinique, mapping flotte, audit).""" + +from .crypto import BadPassword, Vault +from .db import Store, StoreLocked + +__all__ = ["Store", "StoreLocked", "Vault", "BadPassword"] diff --git a/kazeia_central/store/crypto.py b/kazeia_central/store/crypto.py new file mode 100644 index 0000000..68c70cc --- /dev/null +++ b/kazeia_central/store/crypto.py @@ -0,0 +1,105 @@ +"""Chiffrement au niveau champ pour le store local de Kazeia-central. + +Choix (cf. décision 2026-06-18) : libsodium via **PyNaCl**, pas SQLCipher. +`sqlcipher3-binary` n'a de wheels que pour Linux x86_64 → impose une compilation +sur Mac/Windows, ce qui casse « facilement installable ». PyNaCl publie des wheels +`abi3` Linux/macOS/Windows (un seul couvre tout CPython ≥3.8, dont 3.14), et est +**déjà** une dépendance (déchiffrement des exports `crypto_box_seal`, §4 de la spec). + +Modèle : la clé maître est dérivée du **mot de passe opérateur** (Argon2id) + un sel +aléatoire persisté. Chaque valeur sensible (texte des tours, labels patients) est +scellée indépendamment par `SecretBox` (XSalsa20-Poly1305) → seules les lignes lues +sont déchiffrées (granularité voulue, sans base entièrement chiffrée au repos). +Les métadonnées requêtables (timestamps, ids) restent en clair dans la base ; la +protection au repos de CES colonnes repose sur le chiffrement disque OS (§8 CLAUDE.md). +""" + +from __future__ import annotations + +from dataclasses import dataclass + +import nacl.pwhash +import nacl.secret +import nacl.utils +from nacl.exceptions import CryptoError + +# Paramètres KDF par défaut. MODERATE = compromis sûr/réactif pour un déverrouillage +# unique au lancement (quelques centaines de ms). Persistés en base → un changement +# futur de coût n'invalide pas les bases existantes. +_DEFAULT_OPS = nacl.pwhash.argon2id.OPSLIMIT_MODERATE +_DEFAULT_MEM = nacl.pwhash.argon2id.MEMLIMIT_MODERATE +SALT_BYTES = nacl.pwhash.argon2id.SALTBYTES # 16 +KEY_BYTES = nacl.secret.SecretBox.KEY_SIZE # 32 + +# Sentinelle chiffrée à la création, re-déchiffrée à chaque unlock pour valider le +# mot de passe sans jamais le stocker. +_VERIFIER_PLAINTEXT = b"kazeia-central-vault-v1" + + +class BadPassword(Exception): + """Mot de passe opérateur incorrect (échec de déchiffrement du vérificateur).""" + + +@dataclass(frozen=True) +class KdfParams: + salt: bytes + ops: int = _DEFAULT_OPS + mem: int = _DEFAULT_MEM + + @staticmethod + def fresh(ops: int = _DEFAULT_OPS, mem: int = _DEFAULT_MEM) -> "KdfParams": + return KdfParams(salt=nacl.utils.random(SALT_BYTES), ops=ops, mem=mem) + + +def derive_key(password: str, p: KdfParams) -> bytes: + return nacl.pwhash.argon2id.kdf( + KEY_BYTES, password.encode("utf-8"), p.salt, + opslimit=p.ops, memlimit=p.mem, + ) + + +class Vault: + """Détient la clé maître dérivée. Scelle/ouvre des valeurs individuelles. + + Ne jamais persister la clé ; instancier via `Vault.unlock(...)` au lancement, + garder en mémoire le temps de la session opérateur. + """ + + def __init__(self, key: bytes) -> None: + self._box = nacl.secret.SecretBox(key) + + # ---- cycle de vie ----------------------------------------------------- + @classmethod + def setup(cls, password: str, ops: int = _DEFAULT_OPS, + mem: int = _DEFAULT_MEM) -> tuple["Vault", KdfParams, bytes]: + """Première initialisation : dérive la clé, renvoie le vault, les params + KDF à persister et le blob vérificateur à stocker. `ops`/`mem` réglables + (défaut prod MODERATE ; tests peuvent baisser pour la vitesse).""" + params = KdfParams.fresh(ops, mem) + v = cls(derive_key(password, params)) + return v, params, v.seal_bytes(_VERIFIER_PLAINTEXT) + + @classmethod + def unlock(cls, password: str, params: KdfParams, verifier: bytes) -> "Vault": + """Ré-ouverture : dérive la clé et valide contre le vérificateur stocké.""" + v = cls(derive_key(password, params)) + try: + if v.open_bytes(verifier) != _VERIFIER_PLAINTEXT: + raise BadPassword() + except CryptoError as e: + raise BadPassword() from e + return v + + # ---- scellage --------------------------------------------------------- + def seal_bytes(self, data: bytes) -> bytes: + # SecretBox.encrypt préfixe un nonce aléatoire → blob auto-suffisant. + return bytes(self._box.encrypt(data)) + + def open_bytes(self, blob: bytes) -> bytes: + return self._box.decrypt(blob) + + def seal(self, text: str | None) -> bytes | None: + return None if text is None else self.seal_bytes(text.encode("utf-8")) + + def open(self, blob: bytes | None) -> str | None: + return None if blob is None else self.open_bytes(blob).decode("utf-8") diff --git a/kazeia_central/store/db.py b/kazeia_central/store/db.py new file mode 100644 index 0000000..da0018e --- /dev/null +++ b/kazeia_central/store/db.py @@ -0,0 +1,244 @@ +"""Store local chiffré de Kazeia-central. + +`sqlite3` (stdlib, zéro install) + chiffrement au niveau champ (`crypto.Vault`, +PyNaCl). Contient l'archive clinique rapatriée (tours de conversation, PII de santé +chiffrée), le mapping `serial → label patient` (§9), et la piste d'audit (§8). + +Posture : +- Texte des tours et labels patients → **chiffrés** (colonnes `*_enc BLOB`). +- Métadonnées (serial, profile_id, session_id, timestamps) → **clair**, pour rester + requêtables (dashboard, collecte incrémentale, fan-out). Protégées au repos par le + chiffrement disque OS (§8). +- Tout accès d'export est journalisé (`audit`). + +Le store s'ouvre verrouillé ; `unlock(password)` dérive la clé et la garde en mémoire +pour la session. Sans unlock, seules les métadonnées sont lisibles, jamais le contenu. +""" + +from __future__ import annotations + +import json +import sqlite3 +from pathlib import Path +from typing import Any, Iterable + +from .crypto import KdfParams, Vault + +_SCHEMA = """ +CREATE TABLE IF NOT EXISTS meta ( + key TEXT PRIMARY KEY, + value BLOB +); +CREATE TABLE IF NOT EXISTS devices ( + serial TEXT PRIMARY KEY, + label_enc BLOB, -- label patient, chiffré + created_at INTEGER NOT NULL, + last_seen_at INTEGER +); +CREATE TABLE IF NOT EXISTS turns ( + serial TEXT NOT NULL, + profile_id TEXT, + session_id TEXT NOT NULL, + turn_id TEXT NOT NULL, -- `id` du tour côté provider + ts INTEGER, -- timestamp du tour (clair → tri/incrémental) + role TEXT, -- PATIENT | KAZEIA + text_enc BLOB, -- contenu, chiffré (PII) + ttft_ms INTEGER, + total_ms INTEGER, + voice_used TEXT, + model_used TEXT, + archived_at INTEGER NOT NULL, + PRIMARY KEY (serial, session_id, turn_id) -- dédoublonne les re-pulls +); +CREATE INDEX IF NOT EXISTS idx_turns_session ON turns(serial, session_id, ts); +CREATE INDEX IF NOT EXISTS idx_turns_profile ON turns(serial, profile_id, ts); +CREATE TABLE IF NOT EXISTS audit ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + ts INTEGER NOT NULL, + actor TEXT, + action TEXT NOT NULL, + target TEXT, + detail TEXT -- métadonnée d'audit (pas de PII) → clair +); +""" + + +class StoreLocked(Exception): + """Opération sur du contenu chiffré alors que le store n'est pas déverrouillé.""" + + +class Store: + def __init__(self, path: str | Path, *, kdf_ops: int | None = None, + kdf_mem: int | None = None) -> None: + # kdf_ops/kdf_mem : coût Argon2id à l'initialisation (1ᵉʳ unlock). None = + # défaut prod MODERATE. Les tests les baissent pour la vitesse ; ils sont + # persistés, donc les ré-ouvertures réutilisent les mêmes. + self.path = str(path) + self._kdf_ops = kdf_ops + self._kdf_mem = kdf_mem + self._db = sqlite3.connect(self.path) + self._db.row_factory = sqlite3.Row + self._db.execute("PRAGMA journal_mode=WAL") + self._db.execute("PRAGMA foreign_keys=ON") + self._db.executescript(_SCHEMA) + self._db.commit() + self._vault: Vault | None = None + + def close(self) -> None: + self._db.close() + + def __enter__(self) -> "Store": + return self + + def __exit__(self, *exc: object) -> None: + self.close() + + # ---- meta ------------------------------------------------------------- + def _get_meta(self, key: str) -> bytes | None: + row = self._db.execute("SELECT value FROM meta WHERE key=?", (key,)).fetchone() + return row["value"] if row else None + + def _set_meta(self, key: str, value: bytes) -> None: + self._db.execute( + "INSERT INTO meta(key, value) VALUES(?, ?) " + "ON CONFLICT(key) DO UPDATE SET value=excluded.value", + (key, value), + ) + self._db.commit() + + @property + def is_initialized(self) -> bool: + return self._get_meta("kdf_salt") is not None + + @property + def is_unlocked(self) -> bool: + return self._vault is not None + + # ---- déverrouillage --------------------------------------------------- + def unlock(self, password: str) -> None: + """Initialise le coffre au 1ᵉʳ appel, sinon valide le mot de passe. + Lève `crypto.BadPassword` si incorrect.""" + if not self.is_initialized: + from .crypto import _DEFAULT_MEM, _DEFAULT_OPS + vault, params, verifier = Vault.setup( + password, + ops=self._kdf_ops if self._kdf_ops is not None else _DEFAULT_OPS, + mem=self._kdf_mem if self._kdf_mem is not None else _DEFAULT_MEM, + ) + self._set_meta("kdf_salt", params.salt) + self._set_meta("kdf_ops", str(params.ops).encode()) + self._set_meta("kdf_mem", str(params.mem).encode()) + self._set_meta("verifier", verifier) + self._vault = vault + return + params = KdfParams( + salt=self._get_meta("kdf_salt"), + ops=int(self._get_meta("kdf_ops")), + mem=int(self._get_meta("kdf_mem")), + ) + self._vault = Vault.unlock(password, params, self._get_meta("verifier")) + + def _require_vault(self) -> Vault: + if self._vault is None: + raise StoreLocked("store verrouillé : appeler unlock(password) d'abord") + return self._vault + + # ---- mapping serial → label patient (§9) ------------------------------ + def set_device_label(self, serial: str, label: str, *, now: int) -> None: + v = self._require_vault() + self._db.execute( + "INSERT INTO devices(serial, label_enc, created_at, last_seen_at) " + "VALUES(?,?,?,?) ON CONFLICT(serial) DO UPDATE SET " + "label_enc=excluded.label_enc, last_seen_at=excluded.last_seen_at", + (serial, v.seal(label), now, now), + ) + self._db.commit() + + def device_label(self, serial: str) -> str | None: + v = self._require_vault() + row = self._db.execute( + "SELECT label_enc FROM devices WHERE serial=?", (serial,) + ).fetchone() + return v.open(row["label_enc"]) if row else None + + def devices(self) -> list[dict[str, Any]]: + v = self._require_vault() + out = [] + for r in self._db.execute("SELECT * FROM devices ORDER BY created_at"): + out.append({ + "serial": r["serial"], + "label": v.open(r["label_enc"]), + "created_at": r["created_at"], + "last_seen_at": r["last_seen_at"], + }) + return out + + # ---- archive des tours (collecte) ------------------------------------- + def archive_turns(self, serial: str, turns: Iterable[dict[str, Any]], *, now: int) -> int: + """Insère/ignore des tours (idempotent par (serial, session_id, turn_id)). + Chaque `text` est chiffré. Renvoie le nombre de NOUVEAUX tours archivés.""" + v = self._require_vault() + before = self._db.total_changes + self._db.executemany( + "INSERT OR IGNORE INTO turns(serial, profile_id, session_id, turn_id, ts, " + "role, text_enc, ttft_ms, total_ms, voice_used, model_used, archived_at) " + "VALUES(:serial,:profile_id,:session_id,:turn_id,:ts,:role,:text_enc," + ":ttft_ms,:total_ms,:voice_used,:model_used,:archived_at)", + [{ + "serial": serial, + "profile_id": t.get("profile_id"), + "session_id": t["session_id"], + "turn_id": str(t.get("id", t.get("turn_id"))), + "ts": t.get("timestamp", t.get("ts")), + "role": t.get("role"), + "text_enc": v.seal(t.get("text")), + "ttft_ms": t.get("ttft_ms"), + "total_ms": t.get("total_ms"), + "voice_used": t.get("voice_used"), + "model_used": t.get("model_used"), + "archived_at": now, + } for t in turns], + ) + self._db.commit() + return self._db.total_changes - before + + def session_turns(self, serial: str, session_id: str) -> list[dict[str, Any]]: + v = self._require_vault() + rows = self._db.execute( + "SELECT * FROM turns WHERE serial=? AND session_id=? ORDER BY ts, turn_id", + (serial, session_id), + ).fetchall() + out = [] + for r in rows: + d = dict(r) + d["text"] = v.open(d.pop("text_enc")) + out.append(d) + return out + + def last_turn_ts(self, serial: str, session_id: str | None = None) -> int | None: + """Dernier timestamp archivé → borne `since` pour la collecte incrémentale + (autonomie). Pas besoin du vault (métadonnée en clair).""" + if session_id is None: + row = self._db.execute( + "SELECT MAX(ts) m FROM turns WHERE serial=?", (serial,) + ).fetchone() + else: + row = self._db.execute( + "SELECT MAX(ts) m FROM turns WHERE serial=? AND session_id=?", + (serial, session_id), + ).fetchone() + return row["m"] if row else None + + # ---- audit (§8) ------------------------------------------------------- + def audit(self, action: str, *, actor: str | None = None, target: str | None = None, + detail: dict[str, Any] | None = None, now: int) -> None: + self._db.execute( + "INSERT INTO audit(ts, actor, action, target, detail) VALUES(?,?,?,?,?)", + (now, actor, action, target, json.dumps(detail) if detail else None), + ) + self._db.commit() + + def audit_log(self, limit: int = 200) -> list[dict[str, Any]]: + return [dict(r) for r in self._db.execute( + "SELECT * FROM audit ORDER BY id DESC LIMIT ?", (limit,) + )] diff --git a/kazeia_central/web/index.html b/kazeia_central/web/index.html new file mode 100644 index 0000000..d02efa0 --- /dev/null +++ b/kazeia_central/web/index.html @@ -0,0 +1,145 @@ + + + + + +Kazeia-central — console de flotte + + + +
+

Kazeia-central

+ console de flotte + +
+
+ +
← sélectionne une tablette
+
+ + + + diff --git a/tests/test_adb_parse.py b/tests/test_adb_parse.py index 4f102ec..c5932b1 100644 --- a/tests/test_adb_parse.py +++ b/tests/test_adb_parse.py @@ -1,5 +1,7 @@ +import subprocess + from kazeia_central.adb import parse_content_query -from kazeia_central.adb.client import parse_call_bundle +from kazeia_central.adb.client import Adb, AdbError, parse_call_bundle from kazeia_central.provider import models as m @@ -30,3 +32,60 @@ def test_bool_coercion(): def test_call_bundle(): b = parse_call_bundle("Result: Bundle[{ok=true, count=3}]") assert b["ok"] == "true" and b["count"] == "3" + + +def test_devices_skips_daemon_noise(): + """adb émet des lignes `* daemon ... *` AVANT l'en-tête — ne pas les prendre + pour des tablettes (régression : `splitlines()[1:]` en fabriquait).""" + out = ( + b"* daemon not running; starting now at tcp:5037 *\n" + b"* daemon started successfully *\n" + b"List of devices attached\n" + b"ABC123 device usb:1-1 product:OnePlus model:Pad3 transport_id:5\n\n" + ) + a = Adb() + a._run = lambda *x, **k: subprocess.CompletedProcess(x, 0, stdout=out, stderr=b"") + devs = a.devices() + assert [d.serial for d in devs] == ["ABC123"] + assert devs[0].model == "Pad3" and devs[0].ready + + +def test_query_wakes_provider_on_empty(monkeypatch=None): + """Endpoint mono-ligne vide = provider endormi (§3.5) → 1 retry de réveil.""" + calls = {"n": 0} + + def shell(cmd, serial=None, timeout=None): + calls["n"] += 1 + return "" if calls["n"] == 1 else "Row: 0 pid=1, uptime_seconds=2, pss_mb=3, ion_mb=4" + + a = Adb() + a.shell = shell + rows = a.content_query("state", expect_rows=True) + assert calls["n"] == 2 and rows # réveil puis succès + + +def test_query_no_retry_on_empty_list(): + """Endpoint-liste vide = réponse valide → pas de retry inutile.""" + calls = {"n": 0} + + def shell(cmd, serial=None, timeout=None): + calls["n"] += 1 + return "No result found." + + a = Adb() + a.shell = shell + assert a.content_query("crashes") == [] and calls["n"] == 1 + + +def test_query_retries_on_adb_error(): + calls = {"n": 0} + + def shell(cmd, serial=None, timeout=None): + calls["n"] += 1 + if calls["n"] == 1: + raise AdbError("provider down") + return "Row: 0 pid=9, uptime_seconds=1, pss_mb=1, ion_mb=1" + + a = Adb() + a.shell = shell + assert a.content_query("state", expect_rows=True) and calls["n"] == 2 diff --git a/tests/test_store.py b/tests/test_store.py new file mode 100644 index 0000000..6e394c6 --- /dev/null +++ b/tests/test_store.py @@ -0,0 +1,101 @@ +"""Tests du store chiffré. KDF en coût MIN (INTERACTIVE) pour la vitesse.""" + +import nacl.pwhash + +from kazeia_central.store import BadPassword, Store, StoreLocked + +_OPS = nacl.pwhash.argon2id.OPSLIMIT_MIN +_MEM = nacl.pwhash.argon2id.MEMLIMIT_MIN + +NOW = 1_700_000_000_000 + + +def _store(tmp_path): + return Store(tmp_path / "central.db", kdf_ops=_OPS, kdf_mem=_MEM) + + +def _turns(): + return [ + {"id": "t1", "profile_id": "p1", "session_id": "s1", "timestamp": NOW + 1, + "role": "PATIENT", "text": "j'ai mal dormi", "total_ms": 10}, + {"id": "t2", "profile_id": "p1", "session_id": "s1", "timestamp": NOW + 2, + "role": "KAZEIA", "text": "tu te sens fatiguée ?", "ttft_ms": 5}, + ] + + +def test_setup_then_reopen_with_password(tmp_path): + s = _store(tmp_path) + assert not s.is_initialized and not s.is_unlocked + s.unlock("hunter2") + assert s.is_initialized and s.is_unlocked + s.set_device_label("ABC123", "Mme D.", now=NOW) + s.close() + + # Ré-ouverture : bon mot de passe → OK et déchiffre + s2 = Store(tmp_path / "central.db") # défauts persistés (INTERACTIVE) + s2.unlock("hunter2") + assert s2.device_label("ABC123") == "Mme D." + s2.close() + + +def test_wrong_password_rejected(tmp_path): + s = _store(tmp_path) + s.unlock("correct horse") + s.close() + s2 = Store(tmp_path / "central.db") + try: + s2.unlock("wrong") + assert False, "aurait dû lever BadPassword" + except BadPassword: + pass + finally: + s2.close() + + +def test_locked_guard(tmp_path): + s = _store(tmp_path) + try: + s.devices() + assert False, "aurait dû lever StoreLocked" + except StoreLocked: + pass + finally: + s.close() + + +def test_archive_decrypt_and_idempotent(tmp_path): + s = _store(tmp_path) + s.unlock("pw") + n = s.archive_turns("ABC123", _turns(), now=NOW) + assert n == 2 + # re-pull du même historique → 0 nouveau (dédoublonnage par clé) + assert s.archive_turns("ABC123", _turns(), now=NOW) == 0 + rows = s.session_turns("ABC123", "s1") + assert [r["text"] for r in rows] == ["j'ai mal dormi", "tu te sens fatiguée ?"] + assert s.last_turn_ts("ABC123") == NOW + 2 + assert s.last_turn_ts("ABC123", "s1") == NOW + 2 + s.close() + + +def test_content_encrypted_at_rest(tmp_path): + s = _store(tmp_path) + s.unlock("pw") + s.archive_turns("ABC123", _turns(), now=NOW) + s.set_device_label("ABC123", "Mme D.", now=NOW) + s.close() + # Lecture BRUTE du fichier : aucun texte PII ni label en clair + raw = (tmp_path / "central.db").read_bytes() + assert b"j'ai mal dormi" not in raw + assert b"Mme D." not in raw + # ... mais les métadonnées requêtables, oui (par construction) + assert b"s1" in raw + + +def test_audit_log(tmp_path): + s = _store(tmp_path) + s.unlock("pw") + s.audit("export", actor="op1", target="ABC123/s1", + detail={"format": "ndjson", "turns": 2}, now=NOW) + log = s.audit_log() + assert log[0]["action"] == "export" and log[0]["actor"] == "op1" + s.close()