feat: store chiffré, dashboard web flotte, accès réseau + durcissements

- store/ : archive clinique chiffrée (sqlite3 + PyNaCl SecretBox/Argon2id),
  mapping serial→label patient, audit. Choix sqlite+PyNaCl plutôt que SQLCipher
  (wheels Linux-only) pour rester facilement installable Mac/Win.
- web/ + montage StaticFiles : dashboard flotte en lecture (état, profils,
  sessions, RAG, updates, crashes), HTML+JS vanilla sans build.
- __main__.py : `python -m kazeia_central`, host/port via env + auth HTTP Basic
  optionnelle (KAZEIA_USER/PASS) pour exposition réseau (0.0.0.0).
- adb : devices() robuste au bruit daemon, réveil/retry du provider (§3.5).
- models : champs non-clés optionnels (résilience flotte hétérogène).
- spec : RAG à plat (miroir ConfigStore.toJson), ajout rag_dump_json, note PIN.
- tests : 14/14 (parsing adb + store chiffré).

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
alf 2026-06-18 22:41:55 +02:00
parent 311236cecc
commit 465ddd4e7a
12 changed files with 835 additions and 26 deletions

View File

@ -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": "<id>"}`.
- **retour** : `result_b64` = base64 de `{"source": "...", "text": "<fiche multiligne>"}`.
- **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": "<id résultant>"}`.
### 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.

View File

@ -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()

View File

@ -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 `<serial> <state> [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/<path>.
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) 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,

View File

@ -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

View File

@ -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) ------------------------

View File

@ -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

View File

@ -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"]

View File

@ -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")

244
kazeia_central/store/db.py Normal file
View File

@ -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,)
)]

View File

@ -0,0 +1,145 @@
<!DOCTYPE html>
<html lang="fr">
<head>
<meta charset="utf-8">
<meta name="viewport" content="width=device-width, initial-scale=1">
<title>Kazeia-central — console de flotte</title>
<style>
:root {
--bg:#0f1115; --panel:#171a21; --panel2:#1e222b; --line:#2a2f3a;
--txt:#e6e8ec; --muted:#8b93a3; --accent:#7aa2f7; --ok:#7ee787;
--warn:#e3b341; --err:#f7768e;
}
* { box-sizing:border-box; }
body { margin:0; font:14px/1.5 system-ui,Segoe UI,Roboto,sans-serif;
background:var(--bg); color:var(--txt); }
header { padding:14px 20px; border-bottom:1px solid var(--line);
display:flex; align-items:center; gap:12px; }
header h1 { font-size:16px; margin:0; font-weight:600; }
header .sub { color:var(--muted); font-size:12px; }
header button { margin-left:auto; }
.layout { display:flex; height:calc(100vh - 53px); }
aside { width:300px; border-right:1px solid var(--line); overflow:auto; padding:10px; }
main { flex:1; overflow:auto; padding:18px 22px; }
.dev { padding:10px 12px; border:1px solid var(--line); border-radius:8px;
margin-bottom:8px; cursor:pointer; background:var(--panel); }
.dev:hover { border-color:var(--accent); }
.dev.sel { border-color:var(--accent); background:var(--panel2); }
.dev .serial { font-family:ui-monospace,monospace; font-weight:600; }
.dev .meta { color:var(--muted); font-size:12px; }
.badge { display:inline-block; padding:1px 7px; border-radius:10px; font-size:11px;
border:1px solid var(--line); }
.badge.device { color:var(--ok); border-color:var(--ok); }
.badge.offline,.badge.unauthorized { color:var(--err); border-color:var(--err); }
button { background:var(--panel2); color:var(--txt); border:1px solid var(--line);
padding:6px 12px; border-radius:7px; cursor:pointer; font-size:13px; }
button:hover { border-color:var(--accent); }
.grid { display:grid; grid-template-columns:repeat(auto-fill,minmax(260px,1fr)); gap:14px; }
.card { background:var(--panel); border:1px solid var(--line); border-radius:10px;
padding:14px 16px; }
.card h3 { margin:0 0 10px; font-size:13px; color:var(--muted); text-transform:uppercase;
letter-spacing:.04em; font-weight:600; }
.kv { display:flex; justify-content:space-between; gap:10px; padding:3px 0;
border-bottom:1px dashed var(--line); }
.kv:last-child { border-bottom:0; }
.kv .k { color:var(--muted); } .kv .v { font-family:ui-monospace,monospace; }
table { width:100%; border-collapse:collapse; font-size:13px; }
th,td { text-align:left; padding:6px 8px; border-bottom:1px solid var(--line); }
th { color:var(--muted); font-weight:600; font-size:12px; }
.full { grid-column:1/-1; }
.err { color:var(--err); font-size:12px; }
.muted { color:var(--muted); }
.dot { width:7px; height:7px; border-radius:50%; display:inline-block; margin-right:6px; }
.dot.on { background:var(--ok); } .dot.off { background:var(--err); }
#placeholder { color:var(--muted); margin-top:40px; text-align:center; }
</style>
</head>
<body>
<header>
<h1>Kazeia-central</h1>
<span class="sub" id="sub">console de flotte</span>
<button onclick="loadDevices()">↻ Rafraîchir le parc</button>
</header>
<div class="layout">
<aside><div id="devices" class="muted">chargement…</div></aside>
<main><div id="detail"><div id="placeholder">← sélectionne une tablette</div></div></main>
</div>
<script>
const $ = (id) => document.getElementById(id);
let selected = null;
async function api(path) {
const r = await fetch("/api" + path);
if (!r.ok) {
let detail = r.status;
try { detail = (await r.json()).detail || detail; } catch {}
throw new Error(detail);
}
return r.json();
}
async function loadDevices() {
const box = $("devices");
try {
const devs = await api("/devices");
$("sub").textContent = devs.length + " tablette(s) découverte(s)";
if (!devs.length) { box.innerHTML = '<div class="muted">aucune tablette branchée</div>'; return; }
box.innerHTML = "";
for (const d of devs) {
const el = document.createElement("div");
el.className = "dev" + (d.serial === selected ? " sel" : "");
el.innerHTML = `<div class="serial">${d.serial}</div>
<div class="meta">${d.props?.model || "?"} · <span class="badge ${d.state}">${d.state}</span></div>`;
el.onclick = () => selectDevice(d.serial);
box.appendChild(el);
}
} catch (e) {
box.innerHTML = `<div class="err">erreur adb : ${e.message}</div>`;
}
}
function card(title, body, cls="") { return `<div class="card ${cls}"><h3>${title}</h3>${body}</div>`; }
function kv(k, v) { return `<div class="kv"><span class="k">${k}</span><span class="v">${v ?? "—"}</span></div>`; }
function dot(on) { return `<span class="dot ${on ? "on":"off"}"></span>`; }
async function panel(serial, path, render) {
try { return render(await api(`/devices/${serial}/${path}`)); }
catch (e) { return `<div class="err">${path} : ${e.message}</div>`; }
}
async function selectDevice(serial) {
selected = serial;
loadDevices();
const d = $("detail");
d.innerHTML = `<div class="muted">interrogation de <b>${serial}</b></div>`;
const [state, updates, rag, profiles, convs, crashes] = await Promise.all([
panel(serial, "state", s => card("État runtime",
kv("pid", s.pid) + kv("uptime (s)", s.uptime_seconds) + kv("PSS (Mo)", s.pss_mb) +
kv("ION (Mo)", s.ion_mb) + kv("RAM dispo (Mo)", s.mem_available_mb))),
panel(serial, "updates", u => card("Mises à jour",
kv("phase", u.phase) + kv("app MAJ dispo", u.app_update_available) +
kv("version", u.app_version_name) + kv("dernier check", u.last_check))),
panel(serial, "rag/status", r => card("RAG",
`<div class="kv"><span class="k">état</span><span class="v">${dot(r.ready)}${r.ready?"prêt":"non prêt"}</span></div>` +
kv("modèle", r.model) + kv("dim", r.dim) + kv("docs", r.doc_count) + kv("chunks", r.chunk_count))),
panel(serial, "profiles", ps => card("Profils ("+ps.length+")",
`<table><tr><th>id</th><th>nom</th><th>voix</th><th>PIN</th><th>tours</th></tr>` +
ps.map(p=>`<tr><td>${p.id}</td><td>${p.display_name??"—"}</td><td>${p.voice_id??"—"}</td><td>${p.has_pin?"🔒":"—"}</td><td>${p.turn_count??0}</td></tr>`).join("") +
`</table>`, "full")),
panel(serial, "conversations", cs => card("Conversations ("+cs.length+" sessions)",
cs.length ? `<table><tr><th>profil</th><th>session</th><th>tours</th><th>début</th></tr>` +
cs.map(c=>`<tr><td>${c.profile_id}</td><td>${c.session_id}</td><td>${c.turn_count??0}</td><td>${c.started_at?new Date(c.started_at).toLocaleString():"—"}</td></tr>`).join("") +
`</table>` : `<span class="muted">aucune session</span>`, "full")),
panel(serial, "crashes", cr => card("Crashes ("+cr.length+")",
cr.length ? cr.map(c=>kv(c.component||"?", c.message||"")).join("") : `<span class="muted">aucun</span>`)),
]);
d.innerHTML = `<div class="grid">${state}${updates}${rag}${profiles}${convs}${crashes}</div>`;
}
loadDevices();
</script>
</body>
</html>

View File

@ -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

101
tests/test_store.py Normal file
View File

@ -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()