mirror of
https://github.com/EverMind-AI/EverOS.git
synced 2026-09-28 21:35:19 +08:00
fix(review): close 3 blockers surfaced by round-4 review
N1. cluster_repo.find_cluster_id_for_member was cross-owner-unsafe.
Its reverse index (member_type, member_id) alone cannot disambiguate
two owners whose entry_id happens to collide — entry_id is
deliberately only per-owner unique (see entries.py:47:
'Cross-user uniqueness is handled at the database layer via a
composite <user_id>_<entry_id> field; it is not encoded into the
EntryId string itself'). Phase 2's _scan_all_rows crosses all
owners, so on any multi-owner root, same-day seq=1 episodes under
different owners would either false-hit each other's cluster or
be silently skipped from clustering. Add required (app_id,
project_id, owner_id) keyword args + JOIN Cluster (which already
carries scope) to filter by parent scope. Prior signature had zero
production callers except two the same PR just added, so the
API break is contained. Regression test: two owners persist a
cluster each around the same entry_id, each lookup resolves to
its own owner's cluster, a third owner's lookup returns None.
N2. Ctrl-C / EOF at the y/N prompt was landing on the generic
except Exception branch (exit 2 with rich traceback) instead of
the exit-130 interrupt path. Root cause: typer 0.15+ vendored
click under typer._click, so typer.Abort and the standalone
click.exceptions.Abort are distinct classes. The interrupt-branch
catch only listed the standalone one; every existing 'abort'
test was manually raising click.exceptions.Abort so the miss
was a false-positive guard rail. Widen the catch to
(typer.Abort, click.exceptions.Abort) and declare click as a
first-class dependency (it was only pulled in via uvicorn).
Regression test: raise real typer.Abort() at the confirm step
and assert exit 130 + INTERRUPTED banner.
N3. _looks_like_utf8_text used mime.startswith('text/'), which
caught text/html as well. HTML uploads then bypassed everalgo's
_aparse_html — losing clean_html_for_llm (strips <script>/<style>
/<nav>/<iframe> + HTML comments) and the 1 MiB output cap. A
40 MiB .html with <script> bodies and <!-- prompt injection -->
comments would flow straight into the extraction LLM. Replace
with explicit allowlist {text/plain, text/markdown, text/x-rst,
text/x-markdown}; text/html and any future text/* mime now
default to the parser path. Test matrix asserts text/html →
False (was regressed as True by the earlier commit).
Hermetic env full pytest: 2033 passed / 7 deselected (+6 tests
from these regressions).
Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 4.7
parent
07a7cd51d5
commit
561b5fec77
@@ -50,6 +50,14 @@ dependencies = [
|
||||
|
||||
# CLI + TUI
|
||||
"typer>=0.12.0",
|
||||
# click is a transitive dep of uvicorn today, but our CLI raises
|
||||
# ``click.exceptions.Abort`` directly (e.g. tests that inject aborts)
|
||||
# and typer 0.15+ vendored a copy of click under ``typer._click`` — so
|
||||
# the two classes are DIFFERENT even though users only ever see
|
||||
# ``typer.Abort``. We depend on the standalone click package so
|
||||
# ``click.exceptions.Abort`` stays importable and both classes are
|
||||
# covered in the Ctrl-C catch (see backfill exit-130 path).
|
||||
"click>=8.1",
|
||||
"textual>=8.2.7",
|
||||
|
||||
# Tokenization (BM25 Chinese support)
|
||||
|
||||
@@ -325,21 +325,35 @@ def _reject_oversized_upload(file: UploadFile) -> None:
|
||||
_PLAIN_TEXT_EXTENSIONS: frozenset[str] = frozenset(
|
||||
{"md", "txt", "rst", "markdown", "text"}
|
||||
)
|
||||
# Explicit mime allowlist. ``text/*`` prefix matching would let
|
||||
# ``text/html`` (and any future ``text/xml`` etc.) bypass the parser and
|
||||
# feed raw markup — script tags, HTML comments, style blocks — straight
|
||||
# into the knowledge-extraction LLM. everalgo's HTML parser runs
|
||||
# ``clean_html_for_llm`` (strips ``<script>/<style>/<nav>/<iframe>`` +
|
||||
# comments) and caps output at 1 MiB — semantics we must NOT skip.
|
||||
_PLAIN_TEXT_MIMES: frozenset[str] = frozenset(
|
||||
{"text/plain", "text/markdown", "text/x-rst", "text/x-markdown"}
|
||||
)
|
||||
|
||||
|
||||
def _looks_like_utf8_text(file: UploadFile) -> bool:
|
||||
"""Decide whether to skip the parser and go straight to UTF-8 decode.
|
||||
|
||||
A file is treated as plain UTF-8 text when its ``content_type`` starts
|
||||
with ``text/`` (browsers set this for ``.md`` / ``.txt`` uploads), or
|
||||
when the mime is missing / ``application/octet-stream`` (typical for
|
||||
``curl -F file=@x.md`` without an explicit ``type=`` hint) AND the
|
||||
filename extension is on a small allowlist. This short-circuits the
|
||||
parser path — which depends on ``[multimodal]`` being configured —
|
||||
for the common case of uploading a markdown/plaintext knowledge doc.
|
||||
A file is treated as plain UTF-8 text when its ``content_type`` is on
|
||||
``_PLAIN_TEXT_MIMES`` (browsers set ``text/markdown`` / ``text/plain``
|
||||
for ``.md`` / ``.txt`` uploads), or when the mime is missing /
|
||||
``application/octet-stream`` (typical for ``curl -F file=@x.md``
|
||||
without an explicit ``type=`` hint) AND the filename extension is on
|
||||
the extension allowlist.
|
||||
|
||||
``text/html`` is deliberately EXCLUDED — HTML uploads must go through
|
||||
the parser so ``clean_html_for_llm`` can strip active markup and cap
|
||||
the payload; the ``text/*`` prefix match this file used to carry
|
||||
would have let raw HTML (including ``<script>`` bodies and
|
||||
``<!-- prompt injection -->`` comments) reach the LLM as-is.
|
||||
"""
|
||||
mime = (file.content_type or "").lower()
|
||||
if mime.startswith("text/"):
|
||||
if mime in _PLAIN_TEXT_MIMES:
|
||||
return True
|
||||
if mime and mime != "application/octet-stream":
|
||||
return False
|
||||
|
||||
@@ -431,13 +431,21 @@ async def run_backfill(*, phase: str, auto_yes: bool) -> int:
|
||||
_print_summary(summary, exit_code=1)
|
||||
return 1
|
||||
continue
|
||||
# ``click.exceptions.Abort`` fires when ``typer.confirm`` (== click.confirm)
|
||||
# catches KeyboardInterrupt at the y/N prompt and re-raises it as
|
||||
# ``Abort`` (a RuntimeError subclass, NOT a KeyboardInterrupt). Catching
|
||||
# it alongside KeyboardInterrupt / CancelledError keeps Ctrl-C at the
|
||||
# prompt on the interrupt path (exit 130 with resume hint) rather than
|
||||
# falling through to the generic ``except Exception`` (exit 2).
|
||||
except (KeyboardInterrupt, asyncio.CancelledError, click.exceptions.Abort):
|
||||
# ``typer.Abort`` fires when ``typer.confirm`` catches KeyboardInterrupt
|
||||
# or EOF at the y/N prompt and re-raises it as an ``Abort`` (a
|
||||
# ``RuntimeError`` subclass, NOT a ``KeyboardInterrupt``). Typer 0.15+
|
||||
# vendored click under ``typer._click`` — ``typer.Abort`` and the
|
||||
# standalone ``click.exceptions.Abort`` are DISTINCT classes. Catching
|
||||
# only ``click.exceptions.Abort`` would let the real typer-raised
|
||||
# abort fall through to the generic ``except Exception`` branch (exit
|
||||
# 2, rich traceback). Both are listed so the Ctrl-C-at-prompt path
|
||||
# (exit 130 with resume hint) fires regardless of which one arrives.
|
||||
except (
|
||||
KeyboardInterrupt,
|
||||
asyncio.CancelledError,
|
||||
typer.Abort,
|
||||
click.exceptions.Abort,
|
||||
):
|
||||
logger.warning(
|
||||
"cascade_backfill_interrupted",
|
||||
phase=current_phase.slug if current_phase else phase,
|
||||
|
||||
@@ -113,17 +113,43 @@ class _ClusterRepo(RepoBase[Cluster]):
|
||||
self,
|
||||
member_type: str,
|
||||
member_id: str,
|
||||
*,
|
||||
app_id: str,
|
||||
project_id: str,
|
||||
owner_id: str,
|
||||
) -> str | None:
|
||||
"""Reverse lookup: ``(member_type, member_id) → cluster_id``.
|
||||
"""Reverse lookup: ``(member_type, member_id) → cluster_id`` scoped
|
||||
to ``(app_id, project_id, owner_id)``.
|
||||
|
||||
Returns ``None`` when the entity is not yet attached to any cluster.
|
||||
Backed by ``ix_cluster_member_reverse`` so it is O(log N).
|
||||
``member_id`` (e.g. episode ``entry_id`` like
|
||||
``ep_20260517_00000001``) is only per-owner unique — see
|
||||
``core/persistence/markdown/entries.py``: "Cross-user uniqueness
|
||||
is handled at the database layer via a composite
|
||||
``<user_id>_<entry_id>`` field; it is not encoded into the
|
||||
EntryId string itself." Without the scope filter, two owners
|
||||
writing on the same day would share ``entry_id`` and either
|
||||
collide on the reverse index (false hit → the second owner's
|
||||
row silently drops from a cluster it was never part of) or find
|
||||
a foreign cluster.
|
||||
|
||||
Filter cascades: reverse index on ``(member_type, member_id)``
|
||||
narrows to O(k) rows where k = number of owners sharing that
|
||||
entry_id (usually 1), then the join to :class:`Cluster` filters
|
||||
by scope. ``Cluster`` already carries ``app_id / project_id /
|
||||
owner_id`` as of the initial schema.
|
||||
|
||||
Returns ``None`` when the entity is not attached to a cluster
|
||||
under the specified scope.
|
||||
"""
|
||||
async with session_scope(self._factory) as s:
|
||||
stmt = (
|
||||
select(ClusterMember.cluster_id)
|
||||
.join(Cluster, Cluster.cluster_id == ClusterMember.cluster_id)
|
||||
.where(ClusterMember.member_type == member_type)
|
||||
.where(ClusterMember.member_id == member_id)
|
||||
.where(Cluster.app_id == app_id)
|
||||
.where(Cluster.project_id == project_id)
|
||||
.where(Cluster.owner_id == owner_id)
|
||||
.limit(1)
|
||||
)
|
||||
return (await s.execute(stmt)).scalar_one_or_none()
|
||||
|
||||
@@ -928,9 +928,15 @@ async def _emit_synthetic_events(
|
||||
total = len(episodes) + len(cases)
|
||||
emitted = 0
|
||||
for raw in episodes:
|
||||
# entry_id is only per-owner unique — scope the reverse lookup
|
||||
# so a same-day, same-seq episode under a different owner does
|
||||
# not falsely match this owner's cluster (or vice versa).
|
||||
existing = await cluster_repo.find_cluster_id_for_member(
|
||||
member_type="episode",
|
||||
member_id=raw["entry_id"],
|
||||
app_id=raw["app_id"],
|
||||
project_id=raw["project_id"],
|
||||
owner_id=raw["owner_id"],
|
||||
)
|
||||
if existing is None:
|
||||
await engine.emit(_episode_row_to_event(raw))
|
||||
@@ -940,6 +946,9 @@ async def _emit_synthetic_events(
|
||||
existing = await cluster_repo.find_cluster_id_for_member(
|
||||
member_type="case",
|
||||
member_id=raw["entry_id"],
|
||||
app_id=raw["app_id"],
|
||||
project_id=raw["project_id"],
|
||||
owner_id=raw["owner_id"],
|
||||
)
|
||||
if existing is None:
|
||||
await engine.emit(_agent_case_row_to_event(raw))
|
||||
|
||||
@@ -573,17 +573,28 @@ class _UploadStub:
|
||||
@pytest.mark.parametrize(
|
||||
"mime, filename, expected",
|
||||
[
|
||||
# Explicit plain-text mimes: short-circuit
|
||||
("text/markdown", "note.md", True),
|
||||
("text/plain", "note.txt", True),
|
||||
("TEXT/HTML", "page.html", True), # case-insensitive
|
||||
("application/octet-stream", "note.md", True), # curl default + .md
|
||||
(None, "note.md", True), # missing mime + known extension
|
||||
("TEXT/MARKDOWN", "note.md", True), # case-insensitive
|
||||
("text/x-rst", "readme.rst", True),
|
||||
# HTML must NOT short-circuit — needs everalgo's clean_html_for_llm
|
||||
# (strips <script>/<style>/<nav>/comments) + 1 MiB cap. Was a
|
||||
# regression when the whitelist used ``startswith("text/")``.
|
||||
("text/html", "page.html", False),
|
||||
("TEXT/HTML", "page.html", False),
|
||||
# Missing / octet-stream mime + known plaintext extension: OK
|
||||
("application/octet-stream", "note.md", True),
|
||||
(None, "note.md", True),
|
||||
("", "notes.rst", True),
|
||||
("application/octet-stream", "note.bin", False), # unknown extension
|
||||
# Unknown extension: fall through
|
||||
("application/octet-stream", "note.bin", False),
|
||||
(None, "note.bin", False),
|
||||
# Non-text mimes always route to parser
|
||||
("application/pdf", "doc.pdf", False),
|
||||
("image/png", "img.png", False),
|
||||
("application/json", "data.json", False), # explicit non-text mime
|
||||
("application/json", "data.json", False),
|
||||
("text/xml", "data.xml", False), # future text/* additions default to safe
|
||||
],
|
||||
)
|
||||
def test_looks_like_utf8_text_matrix(
|
||||
|
||||
@@ -286,16 +286,130 @@ async def test_find_cluster_id_for_member_reverse_lookup(
|
||||
)
|
||||
|
||||
assert (
|
||||
await repo.find_cluster_id_for_member("memcell", "mc_one") == "cl_user0000001"
|
||||
await repo.find_cluster_id_for_member(
|
||||
"memcell",
|
||||
"mc_one",
|
||||
app_id="default",
|
||||
project_id="default",
|
||||
owner_id="u_alice",
|
||||
)
|
||||
== "cl_user0000001"
|
||||
)
|
||||
assert (
|
||||
await repo.find_cluster_id_for_member("case", "ac_20260517_0001")
|
||||
await repo.find_cluster_id_for_member(
|
||||
"case",
|
||||
"ac_20260517_0001",
|
||||
app_id="default",
|
||||
project_id="default",
|
||||
owner_id="agent_42",
|
||||
)
|
||||
== "cl_case0000001"
|
||||
)
|
||||
# Type-discriminated: same id under wrong type misses.
|
||||
assert await repo.find_cluster_id_for_member("case", "mc_one") is None
|
||||
assert await repo.find_cluster_id_for_member("memcell", "ac_20260517_0001") is None
|
||||
assert await repo.find_cluster_id_for_member("memcell", "mc_missing") is None
|
||||
assert (
|
||||
await repo.find_cluster_id_for_member(
|
||||
"case",
|
||||
"mc_one",
|
||||
app_id="default",
|
||||
project_id="default",
|
||||
owner_id="u_alice",
|
||||
)
|
||||
is None
|
||||
)
|
||||
assert (
|
||||
await repo.find_cluster_id_for_member(
|
||||
"memcell",
|
||||
"ac_20260517_0001",
|
||||
app_id="default",
|
||||
project_id="default",
|
||||
owner_id="agent_42",
|
||||
)
|
||||
is None
|
||||
)
|
||||
assert (
|
||||
await repo.find_cluster_id_for_member(
|
||||
"memcell",
|
||||
"mc_missing",
|
||||
app_id="default",
|
||||
project_id="default",
|
||||
owner_id="u_alice",
|
||||
)
|
||||
is None
|
||||
)
|
||||
|
||||
|
||||
async def test_find_cluster_id_for_member_scoped_to_owner(
|
||||
repo: _ClusterRepo,
|
||||
) -> None:
|
||||
"""Same ``member_id`` under two owners resolves to each owner's own cluster.
|
||||
|
||||
Regression: prior signature took only ``(member_type, member_id)`` and
|
||||
would falsely resolve owner B's episode ``ep_20260517_00000001`` to
|
||||
owner A's cluster (or vice versa) whenever ``entry_id`` collided —
|
||||
which is by design for entry_id (see
|
||||
``core/persistence/markdown/entries.py``: entry_id is per-owner
|
||||
unique, cross-owner uniqueness comes from the composite
|
||||
``<owner>_<entry_id>`` at the storage layer).
|
||||
"""
|
||||
entry_id = "ep_20260517_00000001"
|
||||
alice_cluster = _make_cluster(
|
||||
cluster_id="cl_alice000001",
|
||||
centroid_vals=[1.0, 0.0],
|
||||
members=[entry_id],
|
||||
)
|
||||
bob_cluster = _make_cluster(
|
||||
cluster_id="cl_bob00000001",
|
||||
centroid_vals=[0.0, 1.0],
|
||||
members=[entry_id], # SAME entry_id, different owner
|
||||
)
|
||||
await repo.upsert_with_members(
|
||||
alice_cluster,
|
||||
owner_id="u_alice",
|
||||
owner_type="user",
|
||||
kind="user_memory",
|
||||
member_type="episode",
|
||||
)
|
||||
await repo.upsert_with_members(
|
||||
bob_cluster,
|
||||
owner_id="u_bob",
|
||||
owner_type="user",
|
||||
kind="user_memory",
|
||||
member_type="episode",
|
||||
)
|
||||
|
||||
# Each owner resolves to their own cluster — no cross-owner bleed.
|
||||
assert (
|
||||
await repo.find_cluster_id_for_member(
|
||||
"episode",
|
||||
entry_id,
|
||||
app_id="default",
|
||||
project_id="default",
|
||||
owner_id="u_alice",
|
||||
)
|
||||
== "cl_alice000001"
|
||||
)
|
||||
assert (
|
||||
await repo.find_cluster_id_for_member(
|
||||
"episode",
|
||||
entry_id,
|
||||
app_id="default",
|
||||
project_id="default",
|
||||
owner_id="u_bob",
|
||||
)
|
||||
== "cl_bob00000001"
|
||||
)
|
||||
# A third owner with no cluster attached returns None even though the
|
||||
# entry_id exists under two other owners.
|
||||
assert (
|
||||
await repo.find_cluster_id_for_member(
|
||||
"episode",
|
||||
entry_id,
|
||||
app_id="default",
|
||||
project_id="default",
|
||||
owner_id="u_carol",
|
||||
)
|
||||
is None
|
||||
)
|
||||
|
||||
|
||||
# ── remove_members ─────────────────────────────────────────────────────
|
||||
|
||||
@@ -80,7 +80,13 @@ async def test_emit_synthetic_events_skips_already_clustered_rows(
|
||||
|
||||
class _StubClusterRepo:
|
||||
async def find_cluster_id_for_member(
|
||||
self, *, member_type: str, member_id: str
|
||||
self,
|
||||
member_type: str,
|
||||
member_id: str,
|
||||
*,
|
||||
app_id: str,
|
||||
project_id: str,
|
||||
owner_id: str,
|
||||
) -> str | None:
|
||||
lookups.append((member_type, member_id))
|
||||
return clustered.get((member_type, member_id))
|
||||
@@ -131,7 +137,13 @@ async def test_emit_synthetic_events_uses_write_path_member_type_strings(
|
||||
|
||||
class _StubClusterRepo:
|
||||
async def find_cluster_id_for_member(
|
||||
self, *, member_type: str, member_id: str
|
||||
self,
|
||||
member_type: str,
|
||||
member_id: str,
|
||||
*,
|
||||
app_id: str,
|
||||
project_id: str,
|
||||
owner_id: str,
|
||||
) -> str | None:
|
||||
seen_member_types.add(member_type)
|
||||
return None
|
||||
@@ -160,7 +172,13 @@ async def test_emit_synthetic_events_no_skip_when_no_prior_clusters(
|
||||
|
||||
class _StubClusterRepo:
|
||||
async def find_cluster_id_for_member(
|
||||
self, *, member_type: str, member_id: str
|
||||
self,
|
||||
member_type: str,
|
||||
member_id: str,
|
||||
*,
|
||||
app_id: str,
|
||||
project_id: str,
|
||||
owner_id: str,
|
||||
) -> str | None:
|
||||
return None
|
||||
|
||||
|
||||
@@ -672,3 +672,65 @@ async def test_run_backfill_ctrl_c_at_confirm_returns_130(
|
||||
assert "INTERRUPTED" in combined
|
||||
# Regression guard: must not have taken the generic-error path.
|
||||
assert "Backfill failed — see logs for details." not in combined
|
||||
|
||||
|
||||
async def test_run_backfill_typer_abort_at_confirm_returns_130(
|
||||
_isolated_root: Path, monkeypatch, capsys
|
||||
) -> None:
|
||||
"""Same interrupt semantics as the click-Abort test above, but with
|
||||
the real ``typer.Abort`` class typer 0.15+ actually raises.
|
||||
|
||||
Regression: typer vendored click under ``typer._click`` in 0.15+, so
|
||||
``typer.Abort`` and the standalone ``click.exceptions.Abort`` are
|
||||
DISTINCT classes — a catch of only ``click.exceptions.Abort`` would
|
||||
silently miss the real typer-raised abort (letting it fall through
|
||||
to ``except Exception`` → exit 2 with rich traceback). This test
|
||||
proves both classes are covered by the interrupt-branch tuple.
|
||||
"""
|
||||
import typer as _typer
|
||||
|
||||
assert _typer.Abort is not click.exceptions.Abort, (
|
||||
"typer.Abort must be distinct from click.exceptions.Abort "
|
||||
"for this test to exercise the regression scenario"
|
||||
)
|
||||
|
||||
monkeypatch.setattr(
|
||||
_backfill, "get_embedding_capability", lambda: _FakeCapabilityAvailable()
|
||||
)
|
||||
|
||||
async def _fake_scan():
|
||||
from everos.memory.cascade._backfill import (
|
||||
_TABLE_SPECS,
|
||||
_NullVectorRow,
|
||||
_TableBacklog,
|
||||
)
|
||||
|
||||
return [
|
||||
_TableBacklog(
|
||||
spec=_TABLE_SPECS[0],
|
||||
rows=[
|
||||
_NullVectorRow(
|
||||
id="row_1",
|
||||
text="hello",
|
||||
subject_text=None,
|
||||
tokens=5,
|
||||
)
|
||||
],
|
||||
)
|
||||
], 0
|
||||
|
||||
monkeypatch.setattr(_backfill, "_scan_null_vector_backlog", _fake_scan)
|
||||
|
||||
def _raise_typer_abort(*args, **kwargs):
|
||||
raise _typer.Abort()
|
||||
|
||||
monkeypatch.setattr(_backfill_cmd, "_confirm", _raise_typer_abort)
|
||||
|
||||
exit_code = await _backfill_cmd.run_backfill(phase="vectors", auto_yes=False)
|
||||
assert exit_code == 130
|
||||
|
||||
captured = capsys.readouterr()
|
||||
combined = captured.out + captured.err
|
||||
assert "Interrupted — partial progress was written." in combined
|
||||
assert "INTERRUPTED" in combined
|
||||
assert "Backfill failed — see logs for details." not in combined
|
||||
|
||||
@@ -567,6 +567,7 @@ dependencies = [
|
||||
{ name = "alembic" },
|
||||
{ name = "anyio" },
|
||||
{ name = "apscheduler" },
|
||||
{ name = "click" },
|
||||
{ name = "everalgo-agent-memory" },
|
||||
{ name = "everalgo-knowledge" },
|
||||
{ name = "everalgo-rank" },
|
||||
@@ -621,6 +622,7 @@ requires-dist = [
|
||||
{ name = "alembic", specifier = ">=1.13.0" },
|
||||
{ name = "anyio", specifier = ">=4.0" },
|
||||
{ name = "apscheduler", specifier = ">=3.10.4,<4.0" },
|
||||
{ name = "click", specifier = ">=8.1" },
|
||||
{ name = "everalgo-agent-memory", specifier = "==0.3.1" },
|
||||
{ name = "everalgo-knowledge", specifier = "==0.1.1" },
|
||||
{ name = "everalgo-parser", extras = ["svg"], marker = "extra == 'multimodal'", specifier = ">=0.2.1" },
|
||||
|
||||
Reference in New Issue
Block a user