feat: supersede di memorie false (qmem_correct, include_superseded, playbook sez. 11)

- gateway: supersedes_id validato (404/409), backlink superseded_by+superseded_at+supersede_reason sul vecchio record, indici keyword, search esclude i superseduti di default (include_superseded per lineage), risposta con lineage
- estensione: tool qmem_correct (memory_id o query), qmem_store con supersedes_id/supersede_reason, qmem_search con include_superseded
- README/playbook aggiornati, versione 1.1.0
This commit is contained in:
enne2
2026-08-13 12:41:14 +02:00
parent eadd2a756b
commit 1713e68b82
5 changed files with 513 additions and 12 deletions
+4 -2
View File
@@ -19,8 +19,9 @@ Oppure copia `extensions/index.ts` in `~/.pi/agent/extensions/pi-qmem/`.
| Tool | Descrizione |
|---|---|
| `qmem_store` | Salva un record di memoria (text, kind, agent_id, scope, project_id, source, expires_at) |
| `qmem_search` | Ricerca semantica su tutta la conoscenza condivisa (query, kind, project_id, scope, top_k) |
| `qmem_store` | Salva un record di memoria (text, kind, agent_id, scope, project_id, source, expires_at, supersedes_id, supersede_reason) |
| `qmem_search` | Ricerca semantica su tutta la conoscenza condivisa (query, kind, project_id, scope, top_k, include_superseded) |
| `qmem_correct` | Corregge una memoria falsa: crea un nuovo record che **supersede** il vecchio (che resta in archivio marcato superseded) |
## Comando
@@ -60,6 +61,7 @@ pi (estensione) ──HTTPS/VPN──▶ Memory Gateway (FastAPI) ──▶ Qdra
- Accesso condiviso: una chiave API con accesso completo in lettura/scrittura
- `agent_id` è solo metadata di provenienza, non isolamento
- Rate limit, audit log, cleanup automatico dei record scaduti
- **Correzioni**: i record sono immutabili; correggere una memoria falsa = nuovo record che supersede il vecchio (`qmem_correct` o `qmem_store` con `supersedes_id`). Il vecchio resta in archivio con `superseded_by`, escluso dalla ricerca di default (`include_superseded=true` per la lineage)
## Licenza
+355
View File
@@ -0,0 +1,355 @@
# Playbook Operativo — Memoria Centralizzata per Agenti AI (pi-qmem)
> Sessione: 11/08/2026 — da Graphiti a stack snello Qdrant + Gateway + BGE-M3.
> Questo playbook raccoglie procedure, comandi e best practice verificati operativamente.
---
## 1. Architettura e componenti
Stack snello per memoria condivisa di agenti AI, senza LLM in scrittura:
```
pi (estensione pi-qmem) ──HTTPS──▶ Nginx (enne2.net) ──VPN──▶ Memory Gateway (brain:8082) ──▶ Qdrant 1.19 (brain:6333)
└──▶ Ollama BGE-M3 (brain:11434, nativo)
```
| Componente | Dettaglio |
|---|---|
| Qdrant 1.19.0 | Container `memory-qdrant`, bind 127.0.0.1:6333/6334, API key admin + read-only, JWT RBAC |
| Memory Gateway | FastAPI `memory-gateway`, bind 10.8.0.3:8082 (solo VPN), chiave condivisa, rate limit, audit log |
| Ollama | Nativo sul server (non container), modello `bge-m3` (1024 dim, multilingue) |
| Estensione pi | `pi-qmem` (tool `qmem_store`/`qmem_search`, comando `/qmem:config`) |
| Reverse proxy | Nginx su enne2.net, `qmem.enne2.net` → HTTPS → 10.8.0.3:8082 |
Principi chiave:
- **Nessun LLM in scrittura**: l'agente salva record deliberati e strutturati
- **Accesso condiviso**: una chiave API con accesso completo (nessun isolamento per agente)
- **`agent_id` è solo provenienza**, non meccanismo di isolamento
- **Contenuto recuperato = evidenza non attendibile** (anti prompt-injection)
## 2. Deploy dello stack
### 2.1 Qdrant (docker-compose.yml)
```yaml
services:
qdrant:
image: qdrant/qdrant:v1.19.0
container_name: memory-qdrant
restart: unless-stopped
ports:
- "127.0.0.1:6333:6333"
- "127.0.0.1:6334:6334"
environment:
QDRANT__SERVICE__API_KEY: ${QDRANT_ADMIN_API_KEY}
QDRANT__SERVICE__READ_ONLY_API_KEY: ${QDRANT_READ_ONLY_API_KEY}
QDRANT__SERVICE__JWT_RBAC: "true"
volumes:
- qdrant_storage:/qdrant/storage
```
- Chiavi: `openssl rand -hex 32` in `.env` (0600), mai committate
- **Bind sempre su 127.0.0.1** (o interfaccia VPN), mai 0.0.0.0
- Pinnare la versione, mai `latest`
### 2.2 Gateway FastAPI
- Endpoint: `POST /v1/memories`, `POST /v1/memories:search`, `GET/DELETE /v1/memories/{id}`, `GET /v1/status`
- Auth: header `X-API-Key` (chiave condivisa), rate limit 120 req/min
- **Ollama da container**: usare `host.docker.internal` + `extra_hosts: ["host.docker.internal:host-gateway"]` (127.0.0.1 dentro il container ≠ host)
- Collection: `memories`, vettore 1024 dim (BGE-M3), indici keyword su `agent_id`, `project_id`, `scope`, `kind`
- `expires_at` salvato come **timestamp Unix** (i Range query Qdrant richiedono numeri, non stringhe ISO)
- Cleanup automatico record scaduti (loop orario)
### 2.3 Embedding
- `ollama pull bge-m3` (1.1GB, 1024 dim, 100+ lingue, italiano ottimo)
- **Mai mescolare embedding di modelli diversi nella stessa collection** — cambio modello = nuova collection + re-embed
## 3. Teardown di servizi legacy (Graphiti)
Procedura per rimuovere un servizio Docker Compose multi-container:
```bash
# 1. Backup dati (anche se sembra vuoto)
docker exec <container> redis-cli SAVE
cp -a <data-dir> /opt/backup/<nome>-$(date +%F)/
# 2. Individuare TUTTI i progetti compose (i container possono appartenere
# a progetti diversi con nomi diversi!)
docker inspect <container> --format '{{index .Config.Labels "com.docker.compose.project.working_dir"}}'
cd <working_dir> && docker compose down
# 3. Rimuovere volumi e directory
docker volume ls | grep <nome>
docker volume rm <volume>
rm -rf <dir>
# 4. Verificare porte libere
ss -tlnp | grep <porta>
```
Lezioni apprese:
- **Un container può appartenere a un progetto compose diverso** da quello atteso (es. `docker-graphiti-falkordb-1` era in `/opt/graphiti/mcp_server/docker/`, non in `/opt/graphiti/`)
- **Sempre backup prima del teardown**, anche se i dati sembrano irrilevanti
- **Pulire i cron morti** che referenziano servizi rimossi (`crontab -l | grep <nome>`)
## 4. Backup e disaster recovery
### 4.1 Backup Qdrant
```bash
# Snapshot full-storage (restore SOLO a startup con --storage-snapshot)
curl -X POST -H "api-key: $KEY" http://127.0.0.1:6333/snapshots
# Snapshot di collection (restore via API)
curl -X POST -H "api-key: $KEY" http://127.0.0.1:6333/collections/<name>/snapshots?wait=true
```
**Restore full-storage** (unico metodo per snapshot full):
```bash
docker run -v <snapshot-dir>:/snapshots:ro -v <vol>:/qdrant/storage \
qdrant/qdrant:v1.19.0 ./qdrant --storage-snapshot /snapshots/<file>.snapshot
```
**Restore collection** (via API):
```bash
curl -X PUT -H "api-key: $KEY" -H 'Content-Type: application/json' \
http://127.0.0.1:6333/collections/<nuova>/snapshots/recover?wait=true \
-d '{"location": "file:///qdrant/snapshots/<collection>/<file>.snapshot", "priority": "snapshot"}'
```
### 4.2 Backup rsync incrementale (best practice)
Root cause di backup sempre falliti: **virgolette letterali dentro una variabile shell**:
```bash
# ROTTO: le virgolette diventano parte del pattern → nessuna esclusione attiva
RSYNC_OPTS="-avz --exclude='.*' ..."
rsync $RSYNC_OPTS ...
# CORRETTO: array bash preserva le virgolette
RSYNC_OPTS=(-avz --delete --exclude='**/mail-state/' --exclude='**/mail-logs/')
rsync "${RSYNC_OPTS[@]}" ...
```
Best practice:
- **Array bash** per opzioni con pattern (mai stringa con virgolette)
- Pattern exclude: `**/percorso/` (doppio asterisco = qualsiasi profondità)
- **rc=23 (trasferimento parziale) = accettabile**: aggiornare comunque il symlink `latest` per mantenere la catena incrementale (`--link-dest`)
- Verifica incrementalità: `find <backup> -type f -links +1 | wc -l` (hardlink = file invariati)
- Escludere sempre: stato runtime container (`mail-state/`), log (`mail-logs/`), dir root-only (`ssh/`, `amavis`)
### 4.3 Destinazione backup su HD esterno
- Individuare il mount: `mount | grep -v tmpfs`, `df -h`, `/etc/fstab`
- Convenzione: `/home/enne2/archive/backups/<nome>/`
- Cron: `0 3 * * * /opt/memory/backup.sh` (snapshot Qdrant + rotazione 7 giorni)
## 5. Reverse proxy Nginx e HTTPS
### 5.1 Configurazione proxy
```nginx
server {
server_name qmem.enne2.net;
location / {
client_max_body_size 10m;
proxy_http_version 1.1;
proxy_set_header Host $host;
proxy_set_header X-Real-IP $remote_addr;
proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for;
proxy_set_header X-Forwarded-Proto $scheme;
proxy_pass http://10.8.0.3:8082;
}
}
```
- File in `/etc/nginx/conf.d/<dominio>.conf` (convenzione per-subdominio)
- **Verificare la connettività VPN prima** (`curl http://10.8.0.3:8082/v1/status` dal proxy server)
- Wildcard DNS `*.enne2.net` → nessun cambio DNS necessario
- `sudo nginx -t` prima di ogni reload
### 5.2 Certificato HTTPS
```bash
sudo certbot --nginx -d qmem.enne2.net --non-interactive --agree-tos --redirect
```
- Certbot aggiunge automaticamente: `listen 443 ssl`, certificato, redirect HTTP→HTTPS (301)
- Rinnovo automatico: `certbot renew` (cron/timer di sistema)
- Verifica: `echo | openssl s_client -connect <dom>:443 -servername <dom> | openssl x509 -noout -subject -dates`
## 6. Pulizia configurazioni Nginx (best practice)
Procedura per rimuovere un server block morto:
```bash
# 1. Individuare TUTTI i riferimenti al dominio
sudo grep -n '<dominio>' /etc/nginx/conf.d/*.conf
# 2. Identificare i confini del blocco (server { ... })
sudo cat /etc/nginx/conf.d/<file>.conf | sed -n '<start>,<end>p'
# 3. Rimuovere le righe del blocco
sudo sed -i '<start>,<end>d' /etc/nginx/conf.d/<file>.conf
# 4. Test e reload
sudo nginx -t && sudo systemctl reload nginx
# 5. Verificare che il dominio non risponda più
curl -s -o /dev/null -w '%{http_code}' https://<dominio>/ # atteso: 000/444
```
Best practice:
- **Mai rimuovere a mano i certificati**: usare `certbot delete --cert-name <dominio> --non-interactive`
- Dopo la rimozione del blocco, verificare che il catch-all (`server_name _; return 444;`) gestisca il dominio
- Controllare che il dominio rimosso non sia referenziato in altri file (redirect block, renewal config)
- Verificare che i domini attivi continuino a funzionare dopo il reload
## 7. Eliminazione certificati SSL orfani (best practice)
```bash
# 1. Elencare i certificati
sudo ls /etc/letsencrypt/live/
# 2. Identificare gli orfani (nessun server block li referenzia)
sudo grep -rn '<dominio>' /etc/nginx/conf.d/ # nessun match = orfano
# 3. Eliminare con certbot (pulisce live/, archive/, renewal/)
sudo certbot delete --cert-name <dominio> --non-interactive
# 4. Verificare
sudo ls /etc/letsencrypt/live/ | grep <dominio> # nessun output
```
Best practice:
- **Sempre `certbot delete`**, mai `rm -rf` manuale (certbot gestisce live/, archive/, renewal/ e i log)
- Un certificato è orfano se: nessun `server_name` lo referenzia E nessun renewal config attivo
- Dopo la cancellazione, il dominio non risponde più su HTTPS (connessione fallita = atteso)
## 8. Verifica servizi vector DB (best practice)
### 8.1 Health check
```bash
curl -s http://127.0.0.1:6333/healthz # "healthz check passed"
curl -s -H "api-key: $KEY" http://127.0.0.1:6333/ # versione
curl -s -H "api-key: $KEY" http://127.0.0.1:6333/collections # elenco
```
### 8.2 Verifica auth
```bash
curl -s -o /dev/null -w '%{http_code}' http://127.0.0.1:6333/collections # senza key → 401
```
### 8.3 Verifica dati e ricerca
```bash
# Conteggio punti
curl -s -H "api-key: $KEY" http://127.0.0.1:6333/collections/<name> | jq .result.points_count
# Ricerca di verifica
curl -s -X POST -H "api-key: $KEY" -H 'Content-Type: application/json' \
http://127.0.0.1:6333/collections/<name>/points/query \
-d '{"limit": 3, "with_payload": ["text"]}'
```
### 8.4 Test di restore (il test definitivo)
- **Collection snapshot**: recover in una collection di test → verificare `points_count` e contenuti → eliminare la collection di test
- **Full-storage snapshot**: container isolato con `--storage-snapshot` → verificare dati → rimuovere container e volume
- **Sempre su ambiente isolato**, mai sul dataset di produzione
### 8.5 Isolamento e sicurezza
- Test negativi: chiave errata → 401; scope non consentito → 403
- Verificare che il bind sia su 127.0.0.1 o interfaccia VPN (mai 0.0.0.0)
- Audit log: `docker logs <gateway> | grep '"action"'`
## 9. Estensione pi e pubblicazione
### 9.1 Struttura pi package
```json
{
"name": "pi-qmem",
"keywords": ["pi-package"],
"peerDependencies": {
"@earendil-works/pi-coding-agent": "*",
"typebox": "*"
},
"pi": { "extensions": ["./extensions"] }
}
```
- `typebox` e `@earendil-works/pi-coding-agent` in **peerDependencies** (pi li fornisce a runtime)
- Per test locale: package.json con le stesse dipendenze in `dependencies` + `npm install`
### 9.2 Pubblicazione su Gitea via tea
```bash
tea repo create --name <repo> --description "<desc>" [--private]
git remote add origin https://git.enne2.net/enne2/<repo>.git
git push -u origin main
pi install git:git.enne2.net/enne2/<repo>
```
- Verifica privacy: `curl -s -o /dev/null -w '%{http_code}' https://git.enne2.net/enne2/<repo>` → 404 = privato
- Aggiornamento: `git push` + `pi update git:git.enne2.net/enne2/<repo>`
### 9.3 Menu di configurazione estensione
- `/qmem:config` → menu TUI: URL gateway, API key, test connessione, mostra config
- Modalità CLI: `/qmem:config url <URL> | apikey <KEY> | test`
- Config in `~/.config/pi-qmem/config.json` (0600)
- `ctx.hasUI` per guardare i dialoghi interattivi (select/input) nei modi non-TUI
## 10. Troubleshooting
| Sintomo | Causa probabile | Fix |
|---|---|---|
| Gateway 500 su tutte le chiamate | Ollama non raggiungibile dal container | `host.docker.internal` + `extra_hosts` |
| `expires_at` cleanup error | Stringa ISO in Range query Qdrant | Salvare timestamp Unix |
| Backup rsync sempre FAILED | Virgolette letterali in variabile shell | Array bash per le opzioni |
| Restore snapshot "missing field params" | Snapshot full-storage usato con endpoint collection | Usare `--storage-snapshot` a startup |
| Porta occupata al deploy | Servizio esistente sulla porta | `ss -tlnp` per individuare, cambiare porta |
| Estensione non carica (typebox) | node_modules mancante | package.json + `npm install` nella dir estensione |
| 401 da Nginx (non dal gateway) | Basic auth globale o blocco sbagliato | `nginx -T` per vedere la config effettiva |
| 404 su supersede | `supersedes_id` inesistente | Verificare l'id con GET /v1/memories/{id} prima di supersedere |
| 409 su supersede | Record già superseduto | Correggere la versione attiva (cercare con include_superseded=false) |
---
## 11. Correzione di memorie false (supersede)
> Aggiunto 13/08/2026 — i record sono IMMUTABILI: correggere = nuovo record che supersede il vecchio.
### 11.1 Meccanica
- `POST /v1/memories` con `supersedes_id` → crea il nuovo record e marca il vecchio con `superseded_by` + `superseded_at` + `supersede_reason` (set_payload, in un unico handler)
- Errori: 404 se il target non esiste, 409 se è già superseduto (mai catene di supersede: correggere sempre la versione attiva)
- La ricerca di default esclude i superseduti (condizione `IsEmptyCondition` su `superseded_by`); `include_superseded=true` per la lineage
- Indici keyword su `supersedes_id` e `superseded_by` creati allo startup
### 11.2 Estensione pi
- `qmem_correct(memory_id | query, corrected_text, reason, ...)` → tool dedicato: cerca il record attivo (se serve) e lo supersede in una chiamata
- `qmem_store` accetta anche `supersedes_id` / `supersede_reason` (supersede esplicito)
- `qmem_search` accetta `include_superseded`; i risultati mostrano `supersedes_id`/`superseded_by`/`supersede_reason`
### 11.3 Verifica rapida
```bash
# crea record falso
OID=$(curl -s -X POST $B/memories -H "X-API-Key: $KEY" -H 'Content-Type: application/json' \
-d '{"text":"...falso...","kind":"fact"}' | jq -r .memory_id)
# supersede
curl -s -X POST $B/memories -H "X-API-Key: $KEY" -H 'Content-Type: application/json' \
-d "{\"text\":\"...corretto...\",\"supersedes_id\":\"$OID\",\"supersede_reason\":\"motivo\"}"
# verifiche: GET vecchio → superseded_by; search default esclude; include_superseded=true mostra
```
*Fine playbook — aggiornato 13/08/2026.*
+113 -2
View File
@@ -142,6 +142,8 @@ export default function qmemExtension(pi: ExtensionAPI) {
),
source: Type.Optional(Type.String({ description: "Origine del record (es. conversazione, file, ticket)." })),
expires_at: Type.Optional(Type.String({ description: "Scadenza ISO 8601 (es. 2026-09-01T00:00:00Z) per memoria volatile." })),
supersedes_id: Type.Optional(Type.String({ description: "UUID del record da supersedere (correzione): il nuovo record diventa la versione attiva, il vecchio resta in archivio marcato superseded." })),
supersede_reason: Type.Optional(Type.String({ description: "Motivo della correzione (visibile in audit e sul vecchio record)." })),
}),
async execute(toolCallId, params, signal, onUpdate, ctx) {
const cfg = loadConfig();
@@ -165,6 +167,8 @@ export default function qmemExtension(pi: ExtensionAPI) {
scope: p.scope ?? "agent",
source: p.source,
expires_at: p.expires_at,
supersedes_id: p.supersedes_id,
supersede_reason: p.supersede_reason,
},
signal,
);
@@ -175,7 +179,7 @@ export default function qmemExtension(pi: ExtensionAPI) {
};
}
return {
content: [{ type: "text", text: `Memoria salvata: ${data.memory_id} (${p.kind ?? "fact"}, scope ${p.scope ?? "agent"})` }],
content: [{ type: "text", text: `Memoria salvata: ${data.memory_id} (${p.kind ?? "fact"}, scope ${p.scope ?? "agent"})${p.supersedes_id ? ` — supersede ${p.supersedes_id}` : ""}` }],
details: { memory_id: data.memory_id, created_at: data.created_at },
};
},
@@ -207,6 +211,7 @@ export default function qmemExtension(pi: ExtensionAPI) {
description: "Filtra per scope di visibilità.",
}),
),
include_superseded: Type.Optional(Type.Boolean({ description: "Includi anche i record già superseduti/corretti (default: false)." })),
top_k: Type.Optional(Type.Integer({ description: "Numero massimo di risultati (default: 5, max 20)." })),
}),
async execute(toolCallId, params, signal, onUpdate, ctx) {
@@ -228,6 +233,7 @@ export default function qmemExtension(pi: ExtensionAPI) {
kind: p.kind,
project_id: p.project_id,
scope: p.scope,
include_superseded: p.include_superseded ?? false,
top_k: p.top_k ?? 5,
},
signal,
@@ -247,7 +253,7 @@ export default function qmemExtension(pi: ExtensionAPI) {
}
const lines = results.map(
(r: any, i: number) =>
`${i + 1}. [${r.kind}/${r.scope} score=${r.score}] ${r.text}\n (id: ${r.memory_id}, agente: ${r.agent_id ?? "?"}, creato: ${r.created_at ?? "?"}${r.source ? `, fonte: ${r.source}` : ""})`,
`${i + 1}. [${r.kind}/${r.scope} score=${r.score}] ${r.text}\n (id: ${r.memory_id}, agente: ${r.agent_id ?? "?"}, creato: ${r.created_at ?? "?"}${r.source ? `, fonte: ${r.source}` : ""}${r.supersedes_id ? `, supersede ${r.supersedes_id}` : ""}${r.superseded_by ? `, ⚠️ superseduto da ${r.superseded_by}` : ""})`,
);
return {
content: [{ type: "text", text: lines.join("\n") }],
@@ -256,6 +262,111 @@ export default function qmemExtension(pi: ExtensionAPI) {
},
});
// =========================================================================
// TOOL: qmem_correct — supersede di una memoria falsa o superata
// =========================================================================
pi.registerTool({
name: "qmem_correct",
label: "Qmem memory correct",
description:
"Corregge una memoria falsa o superata: crea un NUOVO record che supersede il vecchio (che resta " +
"in archivio marcato superseded, mai eliminato). Passa memory_id se lo conosci (dalla risposta di " +
"qmem_search), oppure query per individuare automaticamente il record attivo più rilevante. Il testo " +
"corretto sostituisce quello vecchio nella ricerca semantica. Usalo quando hai evidenza verificata che " +
"una memoria è falsa: contraddizione con fonte autorevole, conferma dell'utente o esito di un'azione.",
parameters: Type.Object({
memory_id: Type.Optional(Type.String({ description: "UUID del record attivo da supersedere (dalla risposta di qmem_search)." })),
query: Type.Optional(Type.String({ description: "Query per trovare il record da correggere (usata solo se memory_id non è fornito)." })),
corrected_text: Type.String({ description: "Il testo corretto e verificato che sostituisce quello falso." }),
reason: Type.Optional(Type.String({ description: "Motivo della correzione (visibile in audit e sul vecchio record)." })),
kind: Type.Optional(
Type.Union(
[Type.Literal("decision"), Type.Literal("fact"), Type.Literal("episode"), Type.Literal("preference")],
{ description: "Tipo del nuovo record (default: eredita dal record superseduto)." },
),
),
project_id: Type.Optional(Type.String({ description: "Progetto del nuovo record (default: eredita dal record superseduto)." })),
agent_id: Type.Optional(Type.String({ description: "Nome dell'agente che corregge (solo provenienza)." })),
}),
async execute(toolCallId, params, signal, onUpdate, ctx) {
const cfg = loadConfig();
if (!cfg.apiKey) {
return {
content: [{ type: "text", text: "Config mancante: esegui /qmem:config per impostare url e apiKey." }],
details: { error: "missing_config" },
};
}
const p = params as any;
if (!p.memory_id && !p.query) {
return {
content: [{ type: "text", text: "Serve memory_id (da qmem_search) oppure query per trovare il record da correggere." }],
details: { error: "missing_target" },
};
}
let memoryId = p.memory_id;
let orig: any = {};
if (!memoryId) {
onUpdate?.({ content: [{ type: "text", text: `qmem: ricerca del record da correggere ("${p.query}")...` }] });
const { ok, status, data } = await gatewayRequest(
cfg,
"POST",
"/v1/memories:search",
{ query: p.query, top_k: 1, include_superseded: false },
signal,
);
if (!ok) {
return {
content: [{ type: "text", text: `Errore ${status}: ${JSON.stringify(data)}` }],
details: { error: "gateway_error", status },
};
}
const results = data.results ?? [];
if (results.length === 0) {
return {
content: [{ type: "text", text: "Nessun record attivo trovato per la query. Nessuna correzione applicata." }],
details: { error: "not_found" },
};
}
memoryId = results[0].memory_id;
orig = results[0];
}
onUpdate?.({ content: [{ type: "text", text: `qmem: supersede di ${memoryId}...` }] });
const { ok, status, data } = await gatewayRequest(
cfg,
"POST",
"/v1/memories",
{
text: p.corrected_text,
kind: p.kind ?? orig.kind ?? "fact",
agent_id: p.agent_id,
project_id: p.project_id ?? orig.project_id,
scope: orig.scope ?? "agent",
source: "qmem_correct",
supersedes_id: memoryId,
supersede_reason: p.reason,
},
signal,
);
if (!ok) {
return {
content: [{ type: "text", text: `Errore ${status}: ${JSON.stringify(data)}` }],
details: { error: "gateway_error", status },
};
}
return {
content: [
{
type: "text",
text: `Correzione applicata: nuovo record ${data.memory_id} supersede ${memoryId}${p.reason ? ` (motivo: ${p.reason})` : ""}. Il vecchio record resta in archivio marcato superseded.`,
},
],
details: { new_id: data.memory_id, superseded_id: memoryId },
};
},
});
// =========================================================================
// COMANDO: /qmem:config — menu interattivo + modalità CLI rapida
// =========================================================================
+40 -7
View File
@@ -8,8 +8,8 @@ agente (attuale o futuro) con la chiave può consultare e aggiungere
informazioni liberamente. L'agent_id è solo metadata di provenienza.
Endpoints:
POST /v1/memories → crea un record di memoria
POST /v1/memories:search → ricerca semantica con filtri
POST /v1/memories → crea un record (con supersedes_id corregge un record esistente)
POST /v1/memories:search → ricerca semantica con filtri (include_superseded per la lineage)
GET /v1/memories/{id} → recupera per UUID
DELETE /v1/memories/{id} → elimina per UUID
GET /v1/status → health + statistiche
@@ -48,7 +48,7 @@ MAX_TEXT_LEN = int(os.environ.get("MAX_TEXT_LEN", "8000"))
logging.basicConfig(level=logging.INFO, format="%(asctime)s %(levelname)s %(message)s")
log = logging.getLogger("memory-gateway")
app = FastAPI(title="Memory Gateway", version="2.0.0")
app = FastAPI(title="Memory Gateway", version="2.1.0")
qdrant = QdrantClient(url=QDRANT_URL, api_key=QDRANT_API_KEY)
# Rate limit in-memory: {key: [timestamps]}
@@ -67,6 +67,7 @@ class MemoryIn(BaseModel):
source: Optional[str] = Field(default=None, max_length=256)
expires_at: Optional[str] = None # ISO 8601
supersedes_id: Optional[str] = None
supersede_reason: Optional[str] = Field(default=None, max_length=512)
class SearchIn(BaseModel):
@@ -74,6 +75,7 @@ class SearchIn(BaseModel):
kind: Optional[Literal["decision", "fact", "episode", "preference"]] = None
project_id: Optional[str] = None
scope: Optional[Literal["agent", "project", "org"]] = None
include_superseded: bool = False
top_k: int = Field(default=5, ge=1, le=20)
@@ -142,7 +144,7 @@ def startup() -> None:
collection_name=COLLECTION,
vectors_config=qm.VectorParams(size=EMBED_DIM, distance=qm.Distance.COSINE),
)
for field in ("agent_id", "project_id", "scope", "kind"):
for field in ("agent_id", "project_id", "scope", "kind", "supersedes_id", "superseded_by"):
qdrant.create_payload_index(
collection_name=COLLECTION,
field_name=field,
@@ -159,6 +161,17 @@ def startup() -> None:
@app.post("/v1/memories")
async def add_memory(body: MemoryIn, key: str = Depends(require_auth)) -> dict:
memory_id = str(uuid.uuid4())
# Supersede: il nuovo record corregge uno esistente, che resta in archivio marcato
superseded_id: Optional[str] = None
if body.supersedes_id:
old = qdrant.retrieve(collection_name=COLLECTION, ids=[body.supersedes_id], with_payload=True)
if not old:
raise HTTPException(status_code=404, detail="Memoria da supersedere non trovata")
if old[0].payload.get("superseded_by"):
raise HTTPException(status_code=409, detail="La memoria è già stata superseduta: correggi la versione attiva")
superseded_id = body.supersedes_id
vector = await embed(body.text)
payload: dict[str, Any] = {
@@ -170,15 +183,29 @@ async def add_memory(body: MemoryIn, key: str = Depends(require_auth)) -> dict:
"source": body.source,
"created_at": _now_iso(),
"expires_at": _parse_ts(body.expires_at),
"supersedes_id": body.supersedes_id,
"supersedes_id": superseded_id,
"supersede_reason": body.supersede_reason,
"embedding_model": EMBED_MODEL,
}
qdrant.upsert(
collection_name=COLLECTION,
points=[qm.PointStruct(id=memory_id, vector=vector, payload=payload)],
)
_audit(key, "create", memory_id=memory_id, kind=body.kind, agent_id=payload["agent_id"])
return {"memory_id": memory_id, "created_at": payload["created_at"]}
if superseded_id:
qdrant.set_payload(
collection_name=COLLECTION,
payload={
"superseded_by": memory_id,
"superseded_at": _now_iso(),
"supersede_reason": body.supersede_reason,
},
points=[superseded_id],
)
_audit(key, "supersede", old_id=superseded_id, new_id=memory_id, kind=body.kind, agent_id=payload["agent_id"])
else:
_audit(key, "create", memory_id=memory_id, kind=body.kind, agent_id=payload["agent_id"])
return {"memory_id": memory_id, "created_at": payload["created_at"], "supersedes_id": superseded_id}
@app.post("/v1/memories:search")
@@ -192,6 +219,9 @@ async def search_memories(body: SearchIn, key: str = Depends(require_auth)) -> d
must.append(qm.FieldCondition(key="project_id", match=qm.MatchValue(value=body.project_id)))
if body.scope:
must.append(qm.FieldCondition(key="scope", match=qm.MatchValue(value=body.scope)))
if not body.include_superseded:
# default: esclude i record già corretti (superseded_by presente)
must.append(qm.IsEmptyCondition(is_empty=qm.PayloadField(key="superseded_by")))
hits = qdrant.search(
collection_name=COLLECTION,
@@ -211,6 +241,9 @@ async def search_memories(body: SearchIn, key: str = Depends(require_auth)) -> d
"project_id": h.payload.get("project_id"),
"created_at": h.payload.get("created_at"),
"source": h.payload.get("source"),
"supersedes_id": h.payload.get("supersedes_id"),
"superseded_by": h.payload.get("superseded_by"),
"supersede_reason": h.payload.get("supersede_reason"),
}
for h in hits
]
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "pi-qmem",
"version": "1.0.0",
"version": "1.1.0",
"description": "Memoria centralizzata e condivisa per agenti AI: salva e cerca record semantici (Qdrant + BGE-M3) via Memory Gateway.",
"keywords": ["pi-package", "memory", "agent", "qdrant", "rag"],
"license": "MIT",