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:
Jiayao Song
2026-07-28 18:31:17 +08:00
co-authored by Claude Opus 4.7
parent 07a7cd51d5
commit 561b5fec77
10 changed files with 303 additions and 31 deletions
+8
View File
@@ -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)
+22 -8
View File
@@ -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()
+9
View File
@@ -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
Generated
+2
View File
@@ -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" },