Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 2 additions & 1 deletion loopx/capabilities/periodic_report/pending_intent.py
Original file line number Diff line number Diff line change
Expand Up @@ -760,7 +760,8 @@ def _actual_work_window(
run_index = runtime_root / "goals" / goal_id / "runs" / "index.jsonl"
if run_index.is_file():
try:
rows = run_index.read_text(encoding="utf-8").splitlines()
# LF framing, not `splitlines()`: see `loopx/history.py` for the same reason.
rows = run_index.read_text(encoding="utf-8").split("\n")
except OSError:
rows = []
for raw_row in rows:
Expand Down
15 changes: 2 additions & 13 deletions loopx/control_plane/quota/monitor_poll.py
Original file line number Diff line number Diff line change
Expand Up @@ -387,7 +387,8 @@ def _find_monitor_poll_turn(
normalized_todo_id = normalize_todo_id(todo_id) if todo_id else None
normalized_target_key = str(target_key or "").strip() or None
try:
lines = index_path.read_text(encoding="utf-8").splitlines()
# LF framing keeps one record one record when a value carries U+0085.
lines = index_path.read_text(encoding="utf-8").split("\n")
except OSError:
return None
for line in reversed(lines):
Expand Down Expand Up @@ -698,7 +699,6 @@ def record_quota_monitor_poll_for_decision(
task_lease_idempotency_key: str | None = None,
task_lease_expected_version: int | None = None,
use_current_task_lease: bool = False,
auxiliary_settlement_todo: Mapping[str, Any] | None = None,
turn_instance_id: str | None = None,
_index_lock_held: bool = False,
status_reloader: Callable[[], dict[str, Any]] | None = None,
Expand Down Expand Up @@ -751,17 +751,6 @@ def record_quota_monitor_poll_for_decision(
registry_path=registry_path,
runtime_root=runtime_root,
)
if auxiliary_settlement_todo is not None:
decision["auxiliary_settlement_todo"] = {
key: auxiliary_settlement_todo.get(key)
for key in (
"todo_id",
"task_class",
"status",
"claimed_by",
"excluded_agents",
)
}
observation = _observation_packet(
before=before,
agent_id=agent_id,
Expand Down
4 changes: 3 additions & 1 deletion loopx/control_plane/runtime/run_index_rebuild.py
Original file line number Diff line number Diff line change
Expand Up @@ -38,7 +38,9 @@ def _event_identity(record: dict[str, Any]) -> dict[str, Any]:


def read_index_rows(index_path: Path) -> tuple[list[str], list[tuple[int, dict[str, Any]]]]:
raw_lines = index_path.read_text(encoding="utf-8").splitlines()
# One JSON document per LF: the writer keeps non-ASCII verbatim, so a value
# carrying U+0085 must not be treated as a line break here.
raw_lines = index_path.read_text(encoding="utf-8").split("\n")
rows: list[tuple[int, dict[str, Any]]] = []
for line_number, line in enumerate(raw_lines, start=1):
if not line.strip():
Expand Down
6 changes: 5 additions & 1 deletion loopx/history.py
Original file line number Diff line number Diff line change
Expand Up @@ -699,7 +699,11 @@ def repair_index_duplicates(
else nullcontext()
)
with lock:
raw_lines = index_path.read_text(encoding="utf-8").splitlines()
# The index is one JSON document per LF. `json.dumps(..., ensure_ascii=False)`
# leaves U+0085/U+2028/U+2029 in a value verbatim, and `str.splitlines()`
# treats those as line breaks, which tore one record into two unparsable
# fragments. Frame on LF instead.
raw_lines = index_path.read_text(encoding="utf-8").split("\n")
grouped: dict[tuple[str, str, str], list[tuple[int, dict[str, Any]]]] = {}
for line_number, line in enumerate(raw_lines, start=1):
if not line.strip():
Expand Down
40 changes: 40 additions & 0 deletions tests/control_plane/test_run_index_jsonl_framing.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,40 @@
"""Record framing for the JSONL run index."""

from __future__ import annotations

import json
from pathlib import Path

from loopx.control_plane.runtime.run_index_rebuild import read_index_rows

NEL = "\x85"


def test_run_index_record_with_a_raw_separator_stays_one_record(tmp_path: Path) -> None:
"""`json.dumps(..., ensure_ascii=False)` keeps U+0085 verbatim in a value.

`str.splitlines()` treats U+0085, U+2028 and U+2029 as line breaks, so the
record arrived as two fragments that both failed to parse and the row was
dropped instead of being read.
"""

index = tmp_path / "run_index.jsonl"
record = {"agent_id": "agent-a", "text": f"note{NEL}with-nel"}
index.write_text(json.dumps(record, ensure_ascii=False) + "\n", encoding="utf-8")

_, rows = read_index_rows(index)

assert [row for _, row in rows] == [record]


def test_run_index_ordinary_records_are_unaffected(tmp_path: Path) -> None:
index = tmp_path / "run_index.jsonl"
records = [{"agent_id": "agent-a", "index": value} for value in range(3)]
index.write_text(
"".join(json.dumps(row, ensure_ascii=False) + "\n" for row in records),
encoding="utf-8",
)

_, rows = read_index_rows(index)

assert [row for _, row in rows] == records
Loading