Files
flowdeck/app/services/semantic_search.py
T
bruno 1706ad1ee9
FlowDeck CI / lint (push) Successful in 1m48s
FlowDeck CI / test (push) Failing after 21m19s
FlowDeck CI / docker (push) Skipped
feat: v7.3.0 — cycle v6.8.0→v7.3.0 (Sites, Search, Automations, Calendar, SCIM, Wiki) + audit A9
- v6.8.0 Sites & Forms publics (migrations 24)
- v6.9.0 Recherche sémantique hybride + Ask AI (migration 25)
- v7.0.0 Automations v2 multi-étapes + Workers sandboxés (migration 26)
- v7.1.0 Calendar sync Google/CalDAV + Meeting Notes (migration 27)
- v7.2.0 Enterprise : SCIM 2.0, 2FA TOTP/passkeys, audit UI, agent approvals (migration 28)
- v7.3.0 Wiki/Teamspaces, verified pages, collab polish, charts, unfurl (migration 29)
- docs V68→V73, ROADMAP/CHANGELOG/WORKLOAD à jour, VERSION 7.3.0
- A9 : flowdeck.db, flowdeck_dev.db, test-commit.md, upload_test.txt et e2e/{node_modules,shots,test-results} désindexés + ignorés (.gitignore/.dockerignore)
2026-09-30 20:02:57 -04:00

536 lines
21 KiB
Python

"""FlowDeck — semantic (vector) search + Ask AI (v6.9.0).
Hybrid retrieval = lexical (FTS5/LIKE via :mod:`app.services.search`) fused
with vector cosine similarity via Reciprocal Rank Fusion, then filtered
through :class:`PermissionManager` so unauthorized chunks never surface
(and never enter an LLM prompt).
Vectors use a dependency-free **hashed TF** encoder (``hash-256``): token →
``md5 % 256`` with L2 normalization. Deterministic, offline-first, good
enough for recall on small workspaces; the ``embed_texts`` entry point is
pluggable should an LLM ``/embeddings`` provider be wired later.
See ``docs/V69_Search_Ask_AI.md``.
"""
from __future__ import annotations
import asyncio
import hashlib
import json
import logging
import math
import re
import struct
import time
from app.db import get_conn
logger = logging.getLogger(__name__)
DIM = 256
MODEL = "hash-256"
CHUNK_SIZE = 1200
CHUNK_OVERLAP = 150
MAX_CHUNKS_PER_RESOURCE = 50
RRF_K = 60
_TOKEN_RE = re.compile(r"[\wÀ-ÿ]+", flags=re.UNICODE)
# Ask cache: (question_hash, workspace_id, user_id) -> (expires_at, payload)
_ask_cache: dict[tuple[str, int | None, int], tuple[float, dict]] = {}
_ASK_CACHE_TTL = 600.0
# Ask rate limit: user_id -> (window_start, count)
_ask_rate: dict[int, tuple[float, int]] = {}
_ASK_RATE_MAX = 30
_ASK_RATE_WINDOW = 60.0
# ── text extraction & chunking ─────────────────────────────────────────────
def _blocks_to_text(blocks) -> list[str]:
parts: list[str] = []
def _walk(items) -> None:
for b in items or []:
if not isinstance(b, dict):
continue
for key in ("content", "text", "title"):
val = b.get(key)
if isinstance(val, str) and val.strip():
parts.append(val.strip())
break
children = b.get("children")
if isinstance(children, list):
_walk(children)
_walk(blocks if isinstance(blocks, list) else [])
return parts
def extract_page_text(content: str | None, content_format: str | None) -> str:
"""Full searchable text of a ``pages`` row (all blocks, recursive)."""
if not content:
return ""
if (content_format or "blocks") == "blocks":
try:
blocks = json.loads(content)
return "\n".join(_blocks_to_text(blocks))
except Exception:
return content
return content
def chunk_text(text: str, size: int = CHUNK_SIZE, overlap: int = CHUNK_OVERLAP) -> list[str]:
"""Split text into overlapping chunks (char-based, word-boundary aware)."""
text = (text or "").strip()
if not text:
return []
if len(text) <= size:
return [text]
chunks: list[str] = []
start = 0
while start < len(text):
end = min(start + size, len(text))
if end < len(text):
space = text.rfind(" ", start, end)
if space > start + size // 2:
end = space
chunks.append(text[start:end].strip())
if end >= len(text):
break
start = max(end - overlap, start + 1)
if len(chunks) >= MAX_CHUNKS_PER_RESOURCE:
break
return [c for c in chunks if c]
# ── hashed-TF embeddings ───────────────────────────────────────────────────
def _tokens(text: str) -> list[str]:
return [t.lower() for t in _TOKEN_RE.findall(text or "") if t]
def embed_text(text: str, dim: int = DIM) -> bytes:
"""Deterministic L2-normalized hashed-TF vector, struct-packed float32."""
vec = [0.0] * dim
for tok in _tokens(text):
idx = int(hashlib.md5(tok.encode()).hexdigest(), 16) % dim
vec[idx] += 1.0
norm = math.sqrt(sum(v * v for v in vec))
if norm > 0:
vec = [v / norm for v in vec]
return struct.pack(f"<{dim}f", *vec)
def embed_texts(texts: list[str], dim: int = DIM) -> list[bytes]:
"""Batch entry point (pluggable: LLM /embeddings can replace hashing)."""
return [embed_text(t, dim) for t in texts]
def cosine(a: bytes, b: bytes, dim: int = DIM) -> float:
"""Cosine similarity of two packed normalized vectors (== dot product)."""
try:
va = struct.unpack(f"<{dim}f", a)
vb = struct.unpack(f"<{dim}f", b)
except struct.error:
return 0.0
return sum(x * y for x, y in zip(va, vb, strict=True))
# ── indexing ───────────────────────────────────────────────────────────────
def _resource_text(conn, resource_type: str, resource_id: int) -> str | None:
"""Return indexable text, or None when the resource must not be indexed."""
if resource_type == "page":
row = conn.execute(
"SELECT title, content, content_format FROM pages "
"WHERE id=? AND (deleted_at IS NULL OR deleted_at='') "
"AND COALESCE(search_excluded, 0)=0",
(resource_id,),
).fetchone()
if not row:
return None
body = extract_page_text(row["content"], row["content_format"])
return f"{row['title'] or ''}\n{body}".strip()
if resource_type == "collection":
row = conn.execute(
"SELECT name, description FROM collections WHERE id=?", (resource_id,)
).fetchone()
if not row:
return None
return f"{row['name'] or ''}\n{row['description'] or ''}".strip()
return None
def index_resource(resource_type: str, resource_id: int) -> int:
"""(Re)index one resource. Returns the number of chunks stored."""
with get_conn() as conn:
text = _resource_text(conn, resource_type, resource_id)
conn.execute(
"DELETE FROM semantic_embeddings WHERE resource_type=? AND resource_id=?",
(resource_type, resource_id),
)
n = 0
if text:
for i, chunk in enumerate(chunk_text(text)[:MAX_CHUNKS_PER_RESOURCE]):
conn.execute(
"""INSERT INTO semantic_embeddings
(resource_type, resource_id, chunk_id, chunk_text, embedding, model)
VALUES (?, ?, ?, ?, ?, ?)""",
(resource_type, resource_id, i, chunk, embed_text(chunk), MODEL),
)
n += 1
conn.execute(
"""INSERT INTO semantic_index_state (resource_type, resource_id, indexed_at)
VALUES (?, ?, CURRENT_TIMESTAMP)
ON CONFLICT(resource_type, resource_id)
DO UPDATE SET indexed_at=CURRENT_TIMESTAMP""",
(resource_type, resource_id),
)
conn.commit()
return n
def _stale_resources(conn, limit: int) -> list[tuple[str, int]]:
out: list[tuple[str, int]] = []
rows = conn.execute(
"""SELECT p.id, p.updated_at FROM pages p
LEFT JOIN semantic_index_state s
ON s.resource_type='page' AND s.resource_id=p.id
WHERE (p.deleted_at IS NULL OR p.deleted_at='')
AND COALESCE(p.search_excluded, 0)=0
AND (s.indexed_at IS NULL OR p.updated_at > s.indexed_at)
ORDER BY p.updated_at DESC LIMIT ?""",
(limit,),
).fetchall()
out += [("page", r["id"]) for r in rows]
if len(out) < limit:
rows = conn.execute(
"""SELECT c.id, c.updated_at FROM collections c
LEFT JOIN semantic_index_state s
ON s.resource_type='collection' AND s.resource_id=c.id
WHERE s.indexed_at IS NULL OR c.updated_at > s.indexed_at
ORDER BY c.updated_at DESC LIMIT ?""",
(limit - len(out),),
).fetchall()
out += [("collection", r["id"]) for r in rows]
return out
def purge_orphans() -> int:
"""Drop vectors for deleted/excluded resources. Returns rows removed."""
with get_conn() as conn:
cur = conn.execute(
"""DELETE FROM semantic_embeddings
WHERE (resource_type='page' AND resource_id NOT IN (
SELECT id FROM pages WHERE (deleted_at IS NULL OR deleted_at='')
AND COALESCE(search_excluded, 0)=0))
OR (resource_type='collection' AND resource_id NOT IN (
SELECT id FROM collections))"""
)
conn.execute(
"""DELETE FROM semantic_index_state
WHERE (resource_type='page' AND resource_id NOT IN (
SELECT id FROM pages WHERE (deleted_at IS NULL OR deleted_at='')
AND COALESCE(search_excluded, 0)=0))
OR (resource_type='collection' AND resource_id NOT IN (
SELECT id FROM collections))"""
)
conn.commit()
return cur.rowcount or 0
def index_pending(limit: int = 50) -> dict:
"""Index up to ``limit`` stale resources + purge orphans (scheduler job)."""
with get_conn() as conn:
stale = _stale_resources(conn, limit)
indexed = 0
for rtype, rid in stale:
try:
index_resource(rtype, rid)
indexed += 1
except Exception as exc: # never break the scheduler loop
logger.debug("semantic index failed for %s %s: %s", rtype, rid, exc)
purged = purge_orphans()
return {"checked": len(stale), "indexed": indexed, "purged": purged}
async def semantic_index_scheduler(interval_seconds: int = 300) -> None:
"""Background task: incremental indexing (wired in app lifespan)."""
while True:
try:
await asyncio.to_thread(index_pending)
except Exception as exc: # noqa: BLE001 — scheduler must survive
logger.debug("semantic index scheduler: %s", exc)
await asyncio.sleep(interval_seconds)
# ── vector search ──────────────────────────────────────────────────────────
def vector_search(query: str, *, limit: int = 20,
resource_types: tuple[str, ...] = ("page", "collection")) -> list[dict]:
"""Brute-force cosine scan (fine at this scale). Returns ranked chunks."""
q = (query or "").strip()
if not q:
return []
qvec = embed_text(q)
with get_conn() as conn:
placeholders = ",".join("?" for _ in resource_types)
rows = conn.execute(
f"""SELECT resource_type, resource_id, chunk_id, chunk_text
FROM semantic_embeddings WHERE resource_type IN ({placeholders})""",
list(resource_types),
).fetchall()
scored = []
for r in rows:
row = conn.execute(
"SELECT embedding FROM semantic_embeddings "
"WHERE resource_type=? AND resource_id=? AND chunk_id=?",
(r["resource_type"], r["resource_id"], r["chunk_id"]),
).fetchone()
s = cosine(qvec, row["embedding"]) if row else 0.0
if s > 0:
scored.append({
"resource_type": r["resource_type"],
"resource_id": r["resource_id"],
"chunk_id": r["chunk_id"],
"chunk_text": r["chunk_text"],
"score": s,
})
scored.sort(key=lambda d: d["score"], reverse=True)
return scored[:limit]
# ── hybrid (lexical + vector, RRF) + ACL ───────────────────────────────────
def _rrf_fuse(ranked_lists: list[list[tuple[str, int]]], k: int = RRF_K) -> list[tuple[str, int, float]]:
scores: dict[tuple[str, int], float] = {}
for ranked in ranked_lists:
for rank, key in enumerate(ranked):
scores[key] = scores.get(key, 0.0) + 1.0 / (k + rank + 1)
fused = [(t, i, s) for (t, i), s in scores.items()]
fused.sort(key=lambda x: x[2], reverse=True)
return fused
def hybrid_search(query: str, user: dict, *, limit: int = 20,
workspace_id: int | None = None,
resource_types: tuple[str, ...] = ("page", "collection")) -> tuple[list[dict], int]:
"""Lexical + vector fusion, workspace-scoped, ACL-filtered.
Returns (results, total). Each result: {type, id, title, excerpt, url, score}.
"""
from app.services.permission_manager import PermissionManager
q = (query or "").strip()
if not q:
return [], 0
limit = max(1, min(int(limit or 20), 100))
user_id = user.get("id")
pm = PermissionManager(user_id, bool(user.get("is_admin")))
# 1) lexical candidates (already workspace-membership scoped)
from app.services import search as search_service
lex = search_service.search(q, user_id, limit * 3)
lex_ranked: list[tuple[str, int]] = []
lex_by_key: dict[tuple[str, int], dict] = {}
for item in (lex.get("pages") or []) + (lex.get("collections") or []):
key = (item["type"], item["id"])
if key not in lex_by_key:
lex_by_key[key] = item
lex_ranked.append(key)
# 2) vector candidates
vec = vector_search(q, limit=limit * 3, resource_types=resource_types)
vec_ranked = [(d["resource_type"], d["resource_id"]) for d in vec]
# 3) fuse
fused = _rrf_fuse([lex_ranked, vec_ranked])
# 4) ACL + workspace filter, enrich
results: list[dict] = []
with get_conn() as conn:
for rtype, rid, score in fused:
if rtype == "page":
if not pm.can_view_page(rid):
continue
row = conn.execute(
"SELECT id, title, content, content_format, workspace_id, "
"COALESCE(search_excluded, 0) AS excluded "
"FROM pages WHERE id=?", (rid,)).fetchone()
if not row or row["excluded"]:
continue
if workspace_id and row["workspace_id"] != workspace_id:
continue
excerpt = (lex_by_key.get((rtype, rid), {}).get("excerpt")
or extract_page_text(row["content"], row["content_format"])[:160])
results.append({"type": "page", "id": rid,
"title": (row["title"] or "Untitled"),
"excerpt": excerpt, "url": f"/pages/{rid}",
"score": round(score, 5)})
else:
if not pm.can_view_collection(rid):
continue
row = conn.execute(
"SELECT id, name, description, workspace_id FROM collections WHERE id=?",
(rid,)).fetchone()
if not row:
continue
if workspace_id and row["workspace_id"] != workspace_id:
continue
excerpt = (lex_by_key.get((rtype, rid), {}).get("subtitle")
or (row["description"] or "")[:160])
results.append({"type": "collection", "id": rid,
"title": (row["name"] or "Untitled"),
"excerpt": excerpt, "url": f"/db/{rid}",
"score": round(score, 5)})
if len(results) >= limit:
break
return results, len(results)
# ── Ask AI ─────────────────────────────────────────────────────────────────
def _check_ask_rate(user_id: int) -> None:
from fastapi import HTTPException
now = time.time()
start, count = _ask_rate.get(user_id, (now, 0))
if now - start > _ASK_RATE_WINDOW:
_ask_rate[user_id] = (now, 1)
return
if count >= _ASK_RATE_MAX:
raise HTTPException(429, "Too many questions. Slow down.")
_ask_rate[user_id] = (start, count + 1)
def _offline_answer(question: str, chunks: list[dict]) -> str:
"""Extractive fallback: top sentences sharing query terms + citations."""
qterms = {t.lower() for t in _tokens(question)}
picked: list[str] = []
for ch in chunks[:8]:
for sent in re.split(r"(?<=[.!?])\s+", ch["chunk_text"] or ""):
words = {t.lower() for t in _tokens(sent)}
if qterms & words and len(sent.strip()) > 20:
picked.append((sent.strip(), ch))
if len(picked) >= 4:
break
if len(picked) >= 4:
break
if not picked:
# No lexical overlap: still cite the top vector matches.
lines = []
for ch in chunks[:3]:
snippet = (ch["chunk_text"] or "")[:200].replace("\n", " ")
lines.append(f"- {snippet} [[fdpage:{ch['resource_id']}]]"
if ch["resource_type"] == "page" else f"- {snippet}")
return ("Je n'ai pas trouvé de passage répondant directement, "
"mais voici les passages les plus proches :\n" + "\n".join(lines))
lines = []
for sent, ch in picked:
if ch["resource_type"] == "page":
lines.append(f"- {sent} [[fdpage:{ch['resource_id']}]]")
else:
lines.append(f"- {sent}")
return "Voici ce que j'ai trouvé dans votre workspace :\n" + "\n".join(lines)
def _resolve_citations(conn, chunks: list[dict]) -> list[dict]:
seen: list[dict] = []
done: set[tuple[str, int]] = set()
for ch in chunks:
key = (ch["resource_type"], ch["resource_id"])
if key in done:
continue
done.add(key)
if ch["resource_type"] == "page":
row = conn.execute("SELECT title FROM pages WHERE id=?", (ch["resource_id"],)).fetchone()
seen.append({"type": "page", "id": ch["resource_id"],
"title": (row["title"] if row else "Deleted page") or "Untitled"})
else:
row = conn.execute("SELECT name FROM collections WHERE id=?",
(ch["resource_id"],)).fetchone()
seen.append({"type": "collection", "id": ch["resource_id"],
"title": (row["name"] if row else "Deleted") or "Untitled"})
return seen
async def ask(question: str, user: dict, workspace_id: int | None = None) -> dict:
"""RAG answer over the user's authorized chunks (LLM or offline fallback)."""
from app.services.permission_manager import PermissionManager
q = (question or "").strip()
if not q:
from fastapi import HTTPException
raise HTTPException(400, "question is required")
_check_ask_rate(user.get("id") or 0)
cache_key = (hashlib.sha256(q.encode()).hexdigest(), workspace_id, user.get("id"))
now = time.time()
hit = _ask_cache.get(cache_key)
if hit and hit[0] > now:
out = dict(hit[1])
out["cached"] = True
return out
pm = PermissionManager(user.get("id"), bool(user.get("is_admin")))
vec = vector_search(q, limit=24)
allowed = []
with get_conn() as conn:
for ch in vec:
if ch["resource_type"] == "page":
row = conn.execute(
"SELECT workspace_id, COALESCE(search_excluded, 0) AS excluded "
"FROM pages WHERE id=?", (ch["resource_id"],)).fetchone()
if not row or row["excluded"] or not pm.can_view_page(ch["resource_id"]):
continue
if workspace_id and row["workspace_id"] != workspace_id:
continue
else:
row = conn.execute(
"SELECT workspace_id FROM collections WHERE id=?",
(ch["resource_id"],)).fetchone()
if not row or not pm.can_view_collection(ch["resource_id"]):
continue
if workspace_id and row["workspace_id"] != workspace_id:
continue
allowed.append(ch)
if len(allowed) >= 8:
break
answer = ""
offline = True
if allowed:
try:
from app.services.llm_client import LLMClient
llm = LLMClient()
if await llm.is_available():
ctx = "\n\n".join(
f"[doc {i+1} page_id={c['resource_id']}]\n{c['chunk_text'][:1500]}"
for i, c in enumerate(allowed)
)
resp = await llm.complete([
{"role": "system",
"content": "Réponds en français en citant les sources avec "
"[[fdpage:ID]] (ID = page_id indiqué). Concis."},
{"role": "user", "content": f"Question : {q}\n\nContexte :\n{ctx}"},
])
answer = (resp.text or "").strip()
offline = False
except Exception as exc: # noqa: BLE001 — fall back to extractive
logger.debug("ask LLM failed, offline fallback: %s", exc)
if not answer:
answer = ("Aucun contenu accessible ne correspond à votre question."
if not allowed else _offline_answer(q, allowed))
with get_conn() as conn:
citations = _resolve_citations(conn, allowed[:8])
out = {"answer_markdown": answer, "citations": citations,
"offline": offline, "cached": False}
_ask_cache[cache_key] = (now + _ASK_CACHE_TTL, out)
return out
def reset_state() -> None:
"""Test helper: clear ask cache + rate limiter."""
_ask_cache.clear()
_ask_rate.clear()