← Files Portable ResumeARCHIVED FILE
skills/.portable-resume/runtime/portable_resume/adapters/kimi.py
57.4 KB · Oct 2, 2026 · 00:33 UTC
"""Read-only adapters for current Kimi Code and legacy Kimi CLI stores.
The current store is discovered from ``session_index.jsonl`` and reads each
session's ``agents/main/wire.jsonl``. The legacy store uses ``kimi.json`` to
map work-directory hashes and reads ``context.jsonl`` (or ``wire.jsonl`` as a
fallback). Source files are never opened directly: all content reads go
through the shared stable, no-follow snapshot primitive.
"""
from __future__ import annotations
import hashlib
import json
import math
import os
import re
import stat
import uuid
from dataclasses import replace
from datetime import datetime, timedelta, timezone
from typing import Any, Iterable, Mapping
from .base import CapabilityReport, ResolvedRef
from ..bounds import DEFAULT_BOUNDS, ReadBudget
from ..diagnostics import DiagnosticError
from ..model import Query, Session, SessionSummary, Turn
from ..paths import (
canonical_root,
canonicalize_cwd,
is_within,
require_regular_no_symlinks,
same_cwd,
)
from ..sanitize import sanitize_turn_record
from ..snapshot import stable_read_bytes, stable_scan_lines
FORMAT_ID = "kimi-code-wire-jsonl-v1"
LEGACY_FORMAT_ID = "kimi-legacy-context-jsonl-v1"
_ID = re.compile(r"^[A-Za-z0-9][A-Za-z0-9._-]{0,1023}$")
_HASH = re.compile(r"^[0-9a-fA-F]{32}$")
_FILTERED = frozenset(
{
"system",
"developer",
"reasoning",
"thinking",
"think",
"thought",
"analysis",
"internal",
"control",
"policy",
}
)
_ROLE_ALIASES = {
"human": "user",
"user": "user",
"assistant": "assistant",
"agent": "assistant",
"model": "assistant",
"ai": "assistant",
"tool": "tool",
"function": "tool",
}
_CURRENT_CONTROL_TYPES = frozenset(
{
# Protocol 1.5 persisted records generated by Kimi Code's
# agent-core-v2 wire manifest.
"metadata",
"config.update",
"context.append_message",
"context.apply_compaction",
"context.clear",
"context.undo",
"forked",
"full_compaction.begin",
"full_compaction.cancel",
"full_compaction.complete",
"goal.clear",
"goal.create",
"goal.update",
"interaction.request",
"interaction.resolved",
"llm.request",
"llm.tools_snapshot",
"mcp.tools_discovered",
"permission.record_approval_result",
"permission.set_mode",
"plan_mode.cancel",
"plan_mode.enter",
"plan_mode.exit",
"plan.revision",
"profile.bind",
"swarm_mode.enter",
"swarm_mode.exit",
"task.started",
"task.terminated",
"tools.register_user_tool",
"tools.reset_active_tools",
"tools.set_active_tools",
"tools.unregister_user_tool",
"tools.update_store",
"turn.cancel",
"turn.prompt",
"turn.steer",
"usage.record",
# Persisted names from the immediately preceding agent-core journal
# that can still exist in otherwise current ~/.kimi-code sessions.
"context.update_token_count",
"micro_compaction.apply",
}
)
_CURRENT_LOOP_EVENT_TYPES = frozenset(
{"step.begin", "step.end", "content.part", "tool.call", "tool.result"}
)
class _DuplicateKey(ValueError):
pass
def _object(pairs: list[tuple[str, Any]]) -> dict[str, Any]:
result: dict[str, Any] = {}
for key, value in pairs:
if key in result:
raise _DuplicateKey(key)
result[key] = value
return result
def _shape(value: Any, depth: int = 0) -> None:
if depth > 32:
raise DiagnosticError.limit_exceeded()
if isinstance(value, Mapping):
if len(value) > 512:
raise DiagnosticError.limit_exceeded()
for item in value.values():
_shape(item, depth + 1)
elif isinstance(value, list):
if len(value) > DEFAULT_BOUNDS.scanned_records:
raise DiagnosticError.limit_exceeded()
for item in value:
_shape(item, depth + 1)
def _loads(data: bytes, *, optional: bool = False) -> Any:
def reject_constant(value: str) -> None:
raise ValueError(value)
def finite_float(value: str) -> float:
parsed = float(value)
if not math.isfinite(parsed):
raise ValueError(value)
return parsed
try:
value = json.loads(
data.decode("utf-8"),
object_pairs_hook=_object,
parse_constant=reject_constant,
parse_float=finite_float,
)
except (
UnicodeDecodeError,
json.JSONDecodeError,
_DuplicateKey,
RecursionError,
ValueError,
) as error:
code = "E_UNSUPPORTED_FORMAT" if optional else "E_CORRUPT_RECORD"
raise DiagnosticError(code, source="kimi") from error
_shape(value)
return value
def _identifier(value: object, *, provider: str) -> str:
if not isinstance(value, str) or value in {".", ".."} or _ID.fullmatch(value) is None:
raise DiagnosticError("E_CORRUPT_RECORD", source="kimi", provider=provider)
return value
def _timestamp(value: object) -> str | None:
if isinstance(value, bool):
return None
if isinstance(value, (int, float)):
seconds = float(value)
if abs(seconds) >= 100_000_000_000:
seconds /= 1000.0
try:
parsed = datetime.fromtimestamp(seconds, timezone.utc)
except (OSError, OverflowError, ValueError):
return None
return parsed.isoformat(timespec="seconds").replace("+00:00", "Z")
if isinstance(value, str):
try:
parsed = datetime.fromisoformat(value.replace("Z", "+00:00"))
except ValueError:
return None
if parsed.tzinfo is None:
return None
return parsed.isoformat(timespec="seconds").replace("+00:00", "Z")
return None
def _first(mapping: Mapping[str, Any], names: tuple[str, ...]) -> Any:
for name in names:
if name in mapping:
return mapping[name]
for container_name in ("session", "metadata", "info", "state"):
nested = mapping.get(container_name)
if isinstance(nested, Mapping):
for name in names:
if name in nested:
return nested[name]
return None
def _metadata(value: Any, *, fallback_cwd: str | None = None) -> dict[str, str | None]:
if not isinstance(value, Mapping):
return {
"title": None,
"cwd": fallback_cwd,
"branch": None,
"created_at": None,
"updated_at": None,
"source_repo_root": None,
}
cwd = _first(value, ("workDir", "work_dir", "cwd", "directory", "workingDirectory"))
return {
"title": _string(_first(value, ("customTitle", "custom_title", "title", "name", "summary"))),
"cwd": canonicalize_cwd(cwd) if isinstance(cwd, str) and os.path.isabs(cwd) else fallback_cwd,
"branch": _string(_first(value, ("branch", "gitBranch", "git_branch"))),
"created_at": _timestamp(_first(value, ("createdAt", "created_at", "created", "startTime"))),
"updated_at": _timestamp(
_first(
value,
("updatedAt", "updated_at", "lastActiveAt", "last_active_at", "modifiedAt", "wire_mtime", "mtime"),
)
),
"source_repo_root": _absolute(_first(value, ("repoRoot", "repo_root", "gitRoot", "git_root"))),
}
def _string(value: object) -> str | None:
return value if isinstance(value, str) else None
def _absolute(value: object) -> str | None:
return canonicalize_cwd(value) if isinstance(value, str) and os.path.isabs(value) else None
def _eligible(summary: SessionSummary, query: Query) -> bool:
ref = query.ref.strip() if query.ref else None
if ref == summary.session_id:
return True
if ref:
try:
if str(uuid.UUID(ref)) == str(uuid.UUID(summary.session_id)):
return True
except ValueError:
pass
if ref and os.path.isabs(ref) and summary.source_path is not None:
if canonicalize_cwd(ref) == canonicalize_cwd(summary.source_path):
return True
if query.cwd is not None and (summary.cwd is None or not same_cwd(summary.cwd, query.cwd)):
return False
minutes = query.within_min if query.within_min is not None else DEFAULT_BOUNDS.listing_age_minutes
if minutes is not None and minutes <= 0:
return True
if summary.updated_at is None:
return False
try:
updated = datetime.fromisoformat(summary.updated_at.replace("Z", "+00:00"))
except ValueError:
return False
return updated >= datetime.now(timezone.utc) - timedelta(minutes=minutes)
def _session_from(summary: SessionSummary, turns: Iterable[Turn], warnings: Iterable[str]) -> Session:
values = tuple(turns)
return Session(
source=summary.source,
session_id=summary.session_id,
source_path=summary.source_path,
title=summary.title,
cwd=summary.cwd,
branch=summary.branch,
created_at=summary.created_at,
updated_at=summary.updated_at,
source_repo_root=summary.source_repo_root,
last_user_request=next((turn.content for turn in reversed(values) if turn.role == "user"), None),
last_assistant_action=next((turn.content for turn in reversed(values) if turn.role == "assistant"), None),
turns=values,
warnings=tuple(dict.fromkeys((*summary.warnings, *warnings))),
)
class KimiAdapter:
key = "kimi"
def __init__(self, *, root: str | None = None, read_hook: Any = None):
self._configured_root = root
self._read_hook = read_hook
def _explicit_root(self, query: Query) -> str | None:
candidate = query.source_root or self._configured_root
if candidate is None:
return None
if not os.path.isdir(candidate):
return None
return canonical_root(candidate)
def _roots(self, query: Query) -> tuple[tuple[str, str], ...]:
explicit = self._explicit_root(query)
if explicit is not None:
has_index = os.path.isfile(os.path.join(explicit, "session_index.jsonl"))
has_legacy = os.path.isfile(os.path.join(explicit, "kimi.json"))
has_sessions = os.path.isdir(os.path.join(explicit, "sessions"))
# Current store may lack session_index.jsonl; FS fallback still applies.
if has_index or (has_sessions and not has_legacy):
return ((explicit, FORMAT_ID),)
if has_legacy:
return ((explicit, LEGACY_FORMAT_ID),)
return ()
current_path = os.environ.get("KIMI_CODE_HOME") or os.path.expanduser("~/.kimi-code")
legacy_path = os.environ.get("KIMI_SHARE_DIR") or os.path.expanduser("~/.kimi")
output: list[tuple[str, str]] = []
if os.path.isdir(current_path):
output.append((canonical_root(current_path), FORMAT_ID))
if os.path.isdir(legacy_path):
legacy_root = canonical_root(legacy_path)
if not output or legacy_root != output[0][0]:
output.append((legacy_root, LEGACY_FORMAT_ID))
return tuple(output)
def approved_roots(self, query: Query) -> tuple[str, ...]:
return tuple(dict.fromkeys(root for root, _ in self._roots(query)))
def _selected(self, query: Query, budget: ReadBudget | None) -> tuple[str, str, list[dict[str, Any]], list[str]]:
roots = self._roots(query)
ref = query.ref.strip() if query.ref else None
if ref and os.path.isabs(ref):
for root, provider in roots:
try:
safe, _ = require_regular_no_symlinks(ref, root)
path = canonicalize_cwd(safe)
session_id = _identifier(
_session_id_from_transcript(path, provider),
provider=provider,
)
self._validate_transcript_shape(path, root, provider, session_id)
except DiagnosticError:
continue
cwd: str | None = None
title: str | None = None
warnings: list[str] = []
if provider == FORMAT_ID:
cwd, title, warnings = self._current_exact_path_hints(
root,
session_id,
path,
budget,
)
return (
root,
provider,
[
{
"session_id": session_id,
"path": path,
"cwd": cwd,
"title": title,
}
],
warnings,
)
raise DiagnosticError.unsafe_path()
if ref:
try:
exact_uuid = str(uuid.UUID(ref))
except ValueError:
exact_uuid = None
if exact_uuid is not None:
exact_empty_selection: tuple[
str, str, list[dict[str, Any]], list[str]
] | None = None
for root, provider in roots:
try:
if provider == FORMAT_ID:
records, warnings = self._current_exact_id_records(
root,
exact_uuid,
budget,
)
else:
all_records, warnings = self._legacy_records(root, budget)
records = [
record
for record in all_records
if record["session_id"] == exact_uuid
]
except DiagnosticError as error:
if (
error.code == "E_UNSUPPORTED_FORMAT"
and query.source_root is None
and self._configured_root is None
):
continue
raise
# Miss on this root: try the next provider/home (e.g. empty
# current store, session only in legacy).
if records:
return root, provider, records, warnings
if exact_empty_selection is None:
exact_empty_selection = (root, provider, records, warnings)
# Exact UUID path finished without a hit. Do not fall through into
# a second full index reduce on the same ReadBudget (large
# append-only indexes would double-charge transcript_records and
# can turn a clean miss into E_LIMIT_EXCEEDED).
if exact_empty_selection is not None:
return exact_empty_selection
raise DiagnosticError(
"E_CAPABILITY_UNAVAILABLE" if not roots else "E_UNSUPPORTED_FORMAT",
source=self.key,
)
empty_selection: tuple[str, str, list[dict[str, Any]], list[str]] | None = None
for root, provider in roots:
try:
if provider == FORMAT_ID:
records, warnings = self._current_records(root, budget)
else:
records, warnings = self._legacy_records(root, budget)
except DiagnosticError as error:
if error.code == "E_UNSUPPORTED_FORMAT" and query.source_root is None and self._configured_root is None:
continue
raise
selection = (root, provider, records, warnings)
if records or query.source_root is not None or self._configured_root is not None:
return selection
# For implicit homes, an empty current provider does not hide a
# populated legacy provider. If every provider is empty, retain the
# first supported empty result so probe/list semantics stay stable.
if empty_selection is None:
empty_selection = selection
if empty_selection is not None:
return empty_selection
raise DiagnosticError("E_CAPABILITY_UNAVAILABLE" if not roots else "E_UNSUPPORTED_FORMAT", source=self.key)
def probe(self, query: Query) -> CapabilityReport:
roots = self._roots(query)
if not roots:
return CapabilityReport(self.key, None, "unavailable")
try:
root, provider, records, warnings = self._selected(query, None)
except DiagnosticError as error:
if error.code in {"E_UNSUPPORTED_FORMAT", "E_CAPABILITY_UNAVAILABLE"}:
return CapabilityReport(self.key, roots[0][1], "unsupported", root=roots[0][0])
raise
return CapabilityReport(
self.key,
provider,
"supported",
root=root,
evidence=(
"session_index.jsonl:sessions/*/*/agents/main/wire.jsonl"
if provider == FORMAT_ID
else "kimi.json:sessions/md5-workdir/*/context.jsonl"
,),
warnings=tuple(warnings),
)
def list(self, query: Query, budget: ReadBudget) -> list[SessionSummary]:
root, provider, records, discovery_warnings = self._selected(query, budget)
output: list[SessionSummary] = []
for record in records:
path = record["path"]
session_id = record["session_id"]
state, state_warnings = self._read_state(os.path.dirname(path), root, budget)
meta = _metadata(state, fallback_cwd=record.get("cwd"))
file_time = self._transcript_mtime(path)
summary = SessionSummary(
source=self.key,
session_id=session_id,
source_path=path,
title=meta["title"] or record.get("title"),
cwd=meta["cwd"],
branch=meta["branch"],
created_at=meta["created_at"] or file_time,
updated_at=meta["updated_at"] or file_time,
source_repo_root=meta["source_repo_root"],
provider=provider,
warnings=tuple(dict.fromkeys((*discovery_warnings, *state_warnings))),
)
if _eligible(summary, query):
output.append(summary)
output.sort(key=lambda item: (item.updated_at or "", item.session_id), reverse=True)
return output[: DEFAULT_BOUNDS.listed_sessions]
def show(self, ref: ResolvedRef, query: Query, budget: ReadBudget) -> Session:
if ref.provider not in {FORMAT_ID, LEGACY_FORMAT_ID}:
raise DiagnosticError("E_UNSUPPORTED_FORMAT", source=self.key, provider=ref.provider)
roots = [root for root, provider in self._roots(query) if provider == ref.provider]
if not roots or ref.source_path is None:
raise DiagnosticError.unsafe_path()
root = next((item for item in roots if is_within(ref.source_path, item)), None)
if root is None:
raise DiagnosticError.unsafe_path()
path, _ = require_regular_no_symlinks(ref.source_path, root)
if _session_id_from_transcript(path, ref.provider) != ref.session_id:
raise DiagnosticError("E_CORRUPT_RECORD", source=self.key, provider=ref.provider)
self._validate_transcript_shape(path, root, ref.provider, ref.session_id)
# Exact path is authoritative: do not scan index/FS here. That keeps
# the caller's single source_read_bytes / transcript_records budget for
# the selected wire and prevents optional discovery from blocking show
# or doubling aggregate admitted source bytes.
discovery_warnings: list[str] = list(ref.warnings)
if ref.provider == FORMAT_ID:
discovery_warnings.extend(self._current_index_warnings(root))
self._require_regular_transcript(path, root)
state, state_warnings = self._read_state(os.path.dirname(path), root, budget)
# Legacy cwd/title often live only in kimi.json (not state.json). Load
# that small metadata file only — never re-scan the full session tree.
fallback_cwd = ref.cwd
fallback_title = ref.title
if ref.provider == LEGACY_FORMAT_ID:
legacy_cwd, legacy_title = self._legacy_hints_for_show(
root, path, ref.session_id, budget
)
fallback_cwd = legacy_cwd or fallback_cwd
fallback_title = legacy_title or fallback_title
meta = _metadata(state, fallback_cwd=fallback_cwd)
if meta["title"] is None and fallback_title is not None:
meta = {**meta, "title": fallback_title}
turns, transcript_warnings, transcript_times = self._parse_transcript(
path, root, query, budget, include_turns=True
)
summary = SessionSummary(
source=self.key,
session_id=ref.session_id,
source_path=path,
title=meta["title"],
cwd=meta["cwd"],
branch=meta["branch"],
created_at=meta["created_at"] or transcript_times[0],
updated_at=meta["updated_at"] or transcript_times[1],
source_repo_root=meta["source_repo_root"],
provider=ref.provider,
warnings=tuple(dict.fromkeys((*discovery_warnings, *state_warnings, *transcript_warnings))),
)
return _session_from(summary, turns, (*discovery_warnings, *state_warnings, *transcript_warnings))
def _current_exact_id_records(
self,
root: str,
session_id: str,
budget: ReadBudget | None,
) -> tuple[list[dict[str, Any]], list[str]]:
effective_budget = budget if budget is not None else ReadBudget()
index = os.path.join(root, "session_index.jsonl")
warnings: list[str] = []
index_state = self._index_presence(index)
if index_state == "unreadable":
return [], ["W_STALE_INDEX"]
if index_state == "regular":
try:
reduced, index_warnings, tombstones = self._reduce_current_index(
index,
root,
effective_budget,
)
except DiagnosticError as error:
if error.code == "E_LIMIT_EXCEEDED":
raise
return [], ["W_STALE_INDEX"]
warnings.extend(index_warnings)
if session_id in tombstones:
return [], list(dict.fromkeys(warnings))
match = next(
(record for record in reduced if record["session_id"] == session_id),
None,
)
if match is not None:
return [match], list(dict.fromkeys(warnings))
else:
warnings.append("W_STALE_INDEX")
sessions_root = os.path.join(root, "sessions")
buckets = self._safe_entries(sessions_root, allow_missing=True)
if index_state == "absent" and not buckets:
raise DiagnosticError("E_UNSUPPORTED_FORMAT", source=self.key, provider=FORMAT_ID)
for bucket in buckets:
effective_budget.consume_records()
if not bucket.is_dir(follow_symlinks=False):
continue
session_dir = os.path.join(bucket.path, session_id)
try:
transcript = self._choose_transcript(
[session_dir],
root,
current=True,
)
except DiagnosticError as error:
if error.code == "E_UNSAFE_PATH":
continue
raise
if transcript is None:
continue
warnings.append("W_STALE_INDEX")
return (
[
{
"session_id": session_id,
"path": transcript,
"cwd": None,
"title": None,
}
],
list(dict.fromkeys(warnings)),
)
return [], list(dict.fromkeys(warnings))
def _current_exact_path_hints(
self,
root: str,
session_id: str,
path: str,
budget: ReadBudget | None,
) -> tuple[str | None, str | None, list[str]]:
"""Optionally reconcile an exact path with the append-only index.
Exact path recovery remains authoritative. Index absence, staleness, or
limits only remove optional metadata and add W_STALE_INDEX.
"""
index = os.path.join(root, "session_index.jsonl")
if self._index_presence(index) != "regular":
return None, None, ["W_STALE_INDEX"]
try:
reduced, warnings, tombstones = self._reduce_current_index(
index,
root,
budget if budget is not None else ReadBudget(),
)
except DiagnosticError:
return None, None, ["W_STALE_INDEX"]
if session_id in tombstones:
return None, None, list(dict.fromkeys((*warnings, "W_STALE_INDEX")))
match = next(
(record for record in reduced if record["session_id"] == session_id),
None,
)
if match is None or canonicalize_cwd(match["path"]) != path:
return None, None, list(dict.fromkeys((*warnings, "W_STALE_INDEX")))
return match.get("cwd"), match.get("title"), list(dict.fromkeys(warnings))
def _index_presence(self, index: str) -> str:
"""Return absent | regular | unreadable for session_index.jsonl."""
try:
mode = os.lstat(index).st_mode
except FileNotFoundError:
return "absent"
except OSError as error:
raise DiagnosticError.source_busy() from error
if stat.S_ISLNK(mode) or not stat.S_ISREG(mode):
return "unreadable"
return "regular"
def _current_records(self, root: str, budget: ReadBudget | None) -> tuple[list[dict[str, Any]], list[str]]:
effective_budget = budget if budget is not None else ReadBudget()
index = os.path.join(root, "session_index.jsonl")
warnings: list[str] = []
output_by_id: dict[str, dict[str, Any]] = {}
tombstones: set[str] = set()
index_state = self._index_presence(index)
if index_state == "unreadable":
# Present but not a regular no-follow file: fail closed for FS union.
warnings.append("W_STALE_INDEX")
return [], list(dict.fromkeys(warnings))
if index_state == "regular":
try:
reduced, index_warnings, tombstones = self._reduce_current_index(
index, root, effective_budget
)
warnings.extend(index_warnings)
for record in reduced:
output_by_id[record["session_id"]] = record
except DiagnosticError as error:
if error.code == "E_LIMIT_EXCEEDED":
raise
# Index present but not fully reduced (busy/unreadable). Do not
# FS-resurrect sessions whose deleted tombstones may be lost.
warnings.append("W_STALE_INDEX")
return list(output_by_id.values()), list(dict.fromkeys(warnings))
# Union read-only FS discovery so unindexed / partial-index sessions surface.
# Index tombstones (deleted: true) must not be revived from leftover dirs.
try:
fallback, fallback_warnings = self._fs_fallback_current(
root, effective_budget
)
warnings.extend(fallback_warnings)
for record in fallback:
session_id = record["session_id"]
if session_id in tombstones:
continue
if session_id not in output_by_id:
output_by_id[session_id] = record
warnings.append("W_STALE_INDEX")
if index_state == "absent" and output_by_id:
warnings.append("W_STALE_INDEX")
except DiagnosticError as error:
if not output_by_id:
if error.code == "E_UNSAFE_PATH":
raise
if index_state == "regular":
return [], list(dict.fromkeys(warnings + ["W_STALE_INDEX"]))
raise
warnings.append("W_STALE_INDEX")
if output_by_id:
return list(output_by_id.values()), list(dict.fromkeys(warnings))
if index_state == "regular":
return [], list(dict.fromkeys(warnings + ["W_STALE_INDEX"]))
raise DiagnosticError("E_UNSUPPORTED_FORMAT", source=self.key, provider=FORMAT_ID)
def _reduce_current_index(
self, index: str, root: str, budget: ReadBudget | None
) -> tuple[list[dict[str, Any]], list[str], set[str]]:
sessions_root = os.path.join(root, "sessions")
output_by_id: dict[str, dict[str, Any]] = {}
tombstones: set[str] = set()
warnings: list[str] = []
# Index events are physical history lines (≤ transcript_records) on the
# list/probe discovery budget only — show does not call this path.
for value, line_warnings, _terminated in self._iter_jsonl(
index,
root,
budget,
charge_transcript=True,
soft_corrupt=True,
):
warnings.extend(line_warnings)
if value is None:
continue
if not isinstance(value, Mapping):
raise DiagnosticError("E_CORRUPT_RECORD", source=self.key, provider=FORMAT_ID)
session_id = _identifier(value.get("sessionId"), provider=FORMAT_ID)
if value.get("deleted") is True:
output_by_id.pop(session_id, None)
tombstones.add(session_id)
continue
cwd = value.get("workDir")
session_dir = value.get("sessionDir")
if not isinstance(cwd, str) or not isinstance(session_dir, str):
raise DiagnosticError("E_CORRUPT_RECORD", source=self.key, provider=FORMAT_ID)
# Current Kimi Code treats the append-only index as an untrusted
# hint: only absolute session directories immediately keyed by the
# indexed session ID and contained under <home>/sessions survive.
if not os.path.isabs(session_dir):
continue
canonical_dir = canonicalize_cwd(session_dir)
if not is_within(canonical_dir, sessions_root) or os.path.basename(canonical_dir) != session_id:
continue
try:
transcript = self._choose_transcript([session_dir], root, current=True)
except DiagnosticError as error:
if error.code == "E_UNSAFE_PATH":
warnings.append("W_MISSING_BLOB")
output_by_id.pop(session_id, None)
continue
raise
if transcript is None:
warnings.append("W_MISSING_BLOB")
output_by_id.pop(session_id, None)
continue
tombstones.discard(session_id)
output_by_id[session_id] = {
"session_id": session_id,
"path": transcript,
"cwd": canonicalize_cwd(cwd) if os.path.isabs(cwd) else None,
}
return list(output_by_id.values()), list(dict.fromkeys(warnings)), tombstones
def _fs_fallback_current(
self, root: str, budget: ReadBudget | None
) -> tuple[list[dict[str, Any]], list[str]]:
effective_budget = budget if budget is not None else ReadBudget()
sessions_root = os.path.join(root, "sessions")
output: list[dict[str, Any]] = []
warnings: list[str] = []
for bucket in self._safe_entries(sessions_root, allow_missing=True):
if not bucket.is_dir(follow_symlinks=False):
continue
try:
sessions = self._safe_entries(bucket.path)
except DiagnosticError as error:
if error.code == "E_UNSAFE_PATH":
warnings.append("W_STALE_INDEX")
continue
raise
for session in sessions:
effective_budget.consume_records()
if not session.is_dir(follow_symlinks=False):
continue
try:
session_id = _identifier(session.name, provider=FORMAT_ID)
except DiagnosticError:
continue
if os.path.basename(session.path) != session_id:
continue
try:
transcript = self._choose_transcript([session.path], root, current=True)
except DiagnosticError as error:
if error.code == "E_UNSAFE_PATH":
warnings.append("W_MISSING_BLOB")
continue
raise
if transcript is None:
continue
output.append(
{
"session_id": session_id,
"path": transcript,
"cwd": None,
"title": None,
}
)
return output, list(dict.fromkeys(warnings))
def _load_legacy_metadata(
self, root: str, budget: ReadBudget | None
) -> tuple[dict[str, str], dict[str, Mapping[str, Any]]]:
metadata_path = os.path.join(root, "kimi.json")
if not os.path.isfile(metadata_path):
raise DiagnosticError("E_UNSUPPORTED_FORMAT", source=self.key, provider=LEGACY_FORMAT_ID)
read = stable_read_bytes(
metadata_path,
root=root,
max_bytes=DEFAULT_BOUNDS.record_bytes,
budget=budget,
hook=self._read_hook,
)
value = _loads(read.data)
return self._legacy_metadata(value)
def _legacy_hints_for_show(
self, root: str, path: str, session_id: str, budget: ReadBudget | None
) -> tuple[str | None, str | None]:
try:
cwd_by_hash, session_meta = self._load_legacy_metadata(root, budget)
except DiagnosticError:
return None, None
cwd: str | None = None
sessions_root = os.path.join(root, "sessions")
if is_within(path, sessions_root):
relative = os.path.relpath(path, sessions_root)
bucket = relative.split(os.sep, 1)[0]
if _HASH.fullmatch(bucket) is not None:
cwd = cwd_by_hash.get(bucket.casefold())
title: str | None = None
meta = session_meta.get(session_id)
if isinstance(meta, Mapping):
title = _string(_first(meta, ("title", "name", "summary", "customTitle")))
return cwd, title
def _legacy_records(self, root: str, budget: ReadBudget | None) -> tuple[list[dict[str, Any]], list[str]]:
cwd_by_hash, session_meta = self._load_legacy_metadata(root, budget)
sessions_root = os.path.join(root, "sessions")
entries = self._safe_entries(sessions_root, allow_missing=True)
output: list[dict[str, Any]] = []
for bucket in entries:
if _HASH.fullmatch(bucket.name) is None or not bucket.is_dir(follow_symlinks=False):
continue
cwd = cwd_by_hash.get(bucket.name.casefold())
for session in self._safe_entries(bucket.path):
if len(output) >= DEFAULT_BOUNDS.scanned_records:
raise DiagnosticError.limit_exceeded()
if session.is_file(follow_symlinks=False) and session.name.endswith(".jsonl"):
session_id = session.name[: -len(".jsonl")]
try:
session_id = _identifier(session_id, provider=LEGACY_FORMAT_ID)
except DiagnosticError:
continue
meta = session_meta.get(session_id, {})
output.append(
{
"session_id": session_id,
"path": canonicalize_cwd(session.path),
"cwd": cwd,
"title": meta.get("title") if isinstance(meta, Mapping) else None,
}
)
continue
if not session.is_dir(follow_symlinks=False):
continue
try:
session_id = _identifier(session.name, provider=LEGACY_FORMAT_ID)
except DiagnosticError:
continue
transcript = self._choose_transcript([session.path], root, current=False)
if transcript is None:
continue
meta = session_meta.get(session_id, {})
output.append(
{
"session_id": session_id,
"path": transcript,
"cwd": cwd,
"title": meta.get("title") if isinstance(meta, Mapping) else None,
}
)
return output, []
def _legacy_metadata(self, value: Any) -> tuple[dict[str, str], dict[str, Mapping[str, Any]]]:
if not isinstance(value, Mapping):
raise DiagnosticError("E_CORRUPT_RECORD", source=self.key, provider=LEGACY_FORMAT_ID)
work_dirs = value.get("work_dirs")
if not isinstance(work_dirs, (Mapping, list)):
raise DiagnosticError("E_UNSUPPORTED_FORMAT", source=self.key, provider=LEGACY_FORMAT_ID)
cwd_by_hash: dict[str, str] = {}
session_meta: dict[str, Mapping[str, Any]] = {}
items: Iterable[tuple[object, object]]
items = work_dirs.items() if isinstance(work_dirs, Mapping) else enumerate(work_dirs)
for raw_key, raw_value in items:
candidate = raw_key if isinstance(raw_key, str) and os.path.isabs(raw_key) else None
explicit_hash = raw_key.casefold() if isinstance(raw_key, str) and _HASH.fullmatch(raw_key) else None
if isinstance(raw_value, str):
if os.path.isabs(raw_value):
candidate = raw_value
elif candidate is not None and _HASH.fullmatch(raw_value):
cwd_by_hash[raw_value.casefold()] = canonicalize_cwd(candidate)
if isinstance(raw_value, Mapping):
path_value = _first(raw_value, ("path", "work_dir", "workDir", "cwd", "directory"))
if isinstance(path_value, str) and os.path.isabs(path_value):
candidate = path_value
hash_value = _first(raw_value, ("hash", "work_dir_hash", "workDirHash", "key"))
if candidate is not None and isinstance(hash_value, str) and _HASH.fullmatch(hash_value):
cwd_by_hash[hash_value.casefold()] = canonicalize_cwd(candidate)
sessions = raw_value.get("sessions")
if isinstance(sessions, Mapping):
for session_id, meta in sessions.items():
if isinstance(session_id, str) and isinstance(meta, Mapping):
session_meta[session_id] = meta
elif isinstance(sessions, list):
for meta in sessions:
if isinstance(meta, Mapping):
sid = _first(meta, ("id", "session_id", "sessionId"))
if isinstance(sid, str):
session_meta[sid] = meta
if candidate is not None:
cwd = canonicalize_cwd(candidate)
if explicit_hash is not None:
cwd_by_hash[explicit_hash] = cwd
cwd_by_hash[hashlib.md5(candidate.encode("utf-8")).hexdigest()] = cwd
cwd_by_hash[hashlib.md5(cwd.encode("utf-8")).hexdigest()] = cwd
return cwd_by_hash, session_meta
def _safe_dirs(self, path: str, *, allow_missing: bool = False) -> list[os.DirEntry[str]]:
return [
entry
for entry in self._safe_entries(path, allow_missing=allow_missing)
if entry.is_dir(follow_symlinks=False)
]
def _safe_entries(self, path: str, *, allow_missing: bool = False) -> list[os.DirEntry[str]]:
try:
mode = os.lstat(path).st_mode
except FileNotFoundError:
if allow_missing:
return []
raise DiagnosticError("E_UNSUPPORTED_FORMAT", source=self.key)
except OSError as error:
raise DiagnosticError.source_busy() from error
if stat.S_ISLNK(mode) or not stat.S_ISDIR(mode):
raise DiagnosticError.unsafe_path()
try:
entries = sorted(os.scandir(path), key=lambda item: item.name)
except OSError as error:
raise DiagnosticError.source_busy() from error
output: list[os.DirEntry[str]] = []
for entry in entries:
if len(output) >= DEFAULT_BOUNDS.scanned_records:
raise DiagnosticError.limit_exceeded()
mode = entry.stat(follow_symlinks=False).st_mode
if stat.S_ISLNK(mode):
raise DiagnosticError.unsafe_path()
if stat.S_ISDIR(mode) or stat.S_ISREG(mode):
output.append(entry)
else:
raise DiagnosticError.unsafe_path()
return output
def _choose_transcript(self, directories: Iterable[str], root: str, *, current: bool) -> str | None:
relatives = ("agents/main/wire.jsonl",) if current else ("context.jsonl", "wire.jsonl")
for directory in directories:
if not is_within(directory, root):
raise DiagnosticError.unsafe_path()
for relative in relatives:
candidate = os.path.join(directory, *relative.split("/"))
try:
mode = os.lstat(candidate).st_mode
except FileNotFoundError:
continue
except OSError as error:
raise DiagnosticError.source_busy() from error
if stat.S_ISLNK(mode) or not stat.S_ISREG(mode):
raise DiagnosticError.unsafe_path()
safe, _ = require_regular_no_symlinks(candidate, root)
return safe
return None
def _read_state(
self, transcript_dir: str, root: str, budget: ReadBudget
) -> tuple[Mapping[str, Any], list[str]]:
if os.path.basename(transcript_dir).endswith(".jsonl"):
return {}, ["W_MISSING_BLOB"]
session_dir = transcript_dir
if os.path.basename(transcript_dir) == "main":
session_dir = os.path.dirname(os.path.dirname(transcript_dir))
path = os.path.join(session_dir, "state.json")
if not os.path.exists(path):
return {}, ["W_MISSING_BLOB"]
read = stable_read_bytes(
path,
root=root,
max_bytes=DEFAULT_BOUNDS.record_bytes,
budget=budget,
hook=self._read_hook,
)
value = _loads(read.data)
if not isinstance(value, Mapping):
raise DiagnosticError("E_CORRUPT_RECORD", source=self.key)
return value, []
def _transcript_mtime(self, path: str) -> str | None:
try:
mtime_ns = os.lstat(path).st_mtime_ns
except OSError:
return None
return (
datetime.fromtimestamp(mtime_ns / 1_000_000_000, timezone.utc)
.isoformat(timespec="seconds")
.replace("+00:00", "Z")
)
def _require_regular_transcript(self, path: str, root: str) -> None:
if not is_within(path, root):
raise DiagnosticError.unsafe_path()
try:
mode = os.lstat(path).st_mode
except OSError as error:
raise DiagnosticError.source_busy() from error
if stat.S_ISLNK(mode) or not stat.S_ISREG(mode):
raise DiagnosticError.unsafe_path()
def _current_index_warnings(self, root: str) -> list[str]:
index = os.path.join(root, "session_index.jsonl")
try:
mode = os.lstat(index).st_mode
except OSError:
return ["W_STALE_INDEX"]
if stat.S_ISLNK(mode) or not stat.S_ISREG(mode):
return ["W_STALE_INDEX"]
return []
def _validate_transcript_shape(
self,
path: str,
root: str,
provider: str,
session_id: str,
) -> None:
relative = os.path.relpath(path, root)
parts = tuple(part for part in relative.split(os.sep) if part not in {"", "."})
if any(part == os.pardir for part in parts):
raise DiagnosticError.unsafe_path()
if provider == FORMAT_ID:
valid = (
len(parts) == 6
and parts[0] == "sessions"
and parts[2] == session_id
and parts[3:] == ("agents", "main", "wire.jsonl")
)
else:
valid = (
len(parts) == 4
and parts[0] == "sessions"
and _HASH.fullmatch(parts[1]) is not None
and parts[2] == session_id
and parts[3] in {"context.jsonl", "wire.jsonl"}
) or (
len(parts) == 3
and parts[0] == "sessions"
and _HASH.fullmatch(parts[1]) is not None
and parts[2] == f"{session_id}.jsonl"
)
if not valid:
raise DiagnosticError.unsafe_path()
def _iter_jsonl(
self,
path: str,
root: str,
budget: ReadBudget | None,
*,
charge_transcript: bool,
soft_corrupt: bool = False,
) -> Iterable[tuple[Any, list[str], bool]]:
"""Yield decoded JSONL records under aggregate source_read_bytes.
Per-line size uses record_bytes via stable_scan_lines. Physical line
counts charge transcript_records when charge_transcript is True, else
scanned_records. Does not retain the raw file body.
When ``soft_corrupt`` is True (index reduce), malformed terminated rows
are skipped with warnings so earlier tombstones stay authoritative.
"""
for line in stable_scan_lines(
path,
root=root,
budget=budget,
charge_transcript=charge_transcript,
hook=self._read_hook,
):
if not line.utf8_valid:
if not line.terminated:
yield None, ["W_PARTIAL_TAIL"], False
continue
if soft_corrupt:
yield None, ["W_STALE_INDEX"], True
continue
raise DiagnosticError("E_CORRUPT_RECORD", source=self.key)
text = line.text.strip()
if not text:
continue
try:
value = _loads(text.encode("utf-8"))
except DiagnosticError:
if not line.terminated:
yield None, ["W_PARTIAL_TAIL"], False
continue
if soft_corrupt:
yield None, ["W_STALE_INDEX"], True
continue
raise
yield value, [], line.terminated
def _parse_transcript(
self,
path: str,
root: str,
query: Query,
budget: ReadBudget,
*,
include_turns: bool,
) -> tuple[list[Turn], list[str], tuple[str | None, str | None]]:
turns: list[Turn] = []
pending_turns: list[tuple[dict[str, Any], str | None]] = []
timestamps: list[str] = []
warnings: list[str] = []
recognized = 0
file_time = self._transcript_mtime(path)
for value, line_warnings, _terminated in self._iter_jsonl(
path, root, budget, charge_transcript=True
):
warnings.extend(line_warnings)
if value is None:
continue
if not isinstance(value, Mapping):
raise DiagnosticError("E_CORRUPT_RECORD", source=self.key)
timestamp = _timestamp(_first(value, ("timestamp", "createdAt", "created_at", "time")))
if timestamp is not None:
timestamps.append(timestamp)
current = _extract_current_record(value)
if current is not None:
role, content, tool_name, coalesce_key = current
extracted = (role, content, tool_name)
else:
extracted = _extract_wire_event(value) or _extract_message(value)
coalesce_key = None
if extracted is None:
if _known_current_record(value) or _looks_filtered(value):
recognized += 1
else:
warnings.append("W_UNKNOWN_RECORD_SKIPPED")
continue
recognized += 1
role, content, tool_name = extracted
if not include_turns:
continue
record = {"role": role, "content": content, "timestamp": timestamp, "tool_name": tool_name}
if (
coalesce_key is not None
and pending_turns
and pending_turns[-1][1] == coalesce_key
):
previous, _ = pending_turns[-1]
previous["content"] = str(previous["content"]) + content
if timestamp is not None:
previous["timestamp"] = timestamp
else:
pending_turns.append((record, coalesce_key))
bounds = replace(DEFAULT_BOUNDS, tool_output_chars=query.max_tool_chars)
for record, _ in pending_turns:
turn, found = sanitize_turn_record(record, ordinal=len(turns), bounds=bounds)
warnings.extend(found)
if turn is not None:
budget.consume_turns()
turns.append(turn)
if recognized == 0:
raise DiagnosticError("E_UNSUPPORTED_FORMAT", source=self.key)
return turns, list(dict.fromkeys(warnings)), (
min(timestamps) if timestamps else file_time,
max(timestamps) if timestamps else file_time,
)
def _looks_filtered(value: Mapping[str, Any]) -> bool:
pending: list[tuple[Any, int]] = [(value, 0)]
seen: set[int] = set()
while pending:
candidate, depth = pending.pop(0)
if isinstance(candidate, Mapping):
identity = id(candidate)
if identity in seen:
continue
seen.add(identity)
for name in ("role", "type", "kind", "event"):
kind = candidate.get(name)
if isinstance(kind, str):
lowered = kind.casefold()
if lowered == "_usage" or any(token in lowered for token in _FILTERED):
return True
if depth < 6:
pending.extend((nested, depth + 1) for nested in candidate.values())
elif isinstance(candidate, list) and depth < 6:
pending.extend((nested, depth + 1) for nested in candidate)
return False
def _known_current_record(value: Mapping[str, Any]) -> bool:
kind = value.get("type")
if not isinstance(kind, str):
return False
if kind != "context.append_loop_event":
return kind in _CURRENT_CONTROL_TYPES
event = value.get("event")
return (
isinstance(event, Mapping)
and isinstance(event.get("type"), str)
and event["type"] in _CURRENT_LOOP_EVENT_TYPES
)
def _extract_current_record(
value: Mapping[str, Any],
) -> tuple[str, str, str | None, str | None] | None:
kind = value.get("type")
if kind == "context.append_message":
extracted = _extract_message(value)
return (*extracted, None) if extracted is not None else None
if kind != "context.append_loop_event":
return None
event = value.get("event")
if not isinstance(event, Mapping):
return None
event_type = event.get("type")
if event_type == "content.part":
text = _content(event.get("part"))
if text is None:
return None
step_uuid = event.get("stepUuid")
group = (
f"assistant:{step_uuid}"
if isinstance(step_uuid, str) and step_uuid
else None
)
return "assistant", text, None, group
if event_type == "tool.result":
result = event.get("result")
if not isinstance(result, Mapping):
return None
text = _content(
_first(result, ("output", "content", "message", "value"))
)
if text is None:
return None
return "tool", text, "tool_result", None
return None
def _session_id_from_transcript(path: str, provider: str) -> str:
if provider == FORMAT_ID:
# <session>/agents/main/wire.jsonl
return os.path.basename(os.path.dirname(os.path.dirname(os.path.dirname(path))))
basename = os.path.basename(path)
if basename.endswith(".jsonl") and basename not in {"context.jsonl", "wire.jsonl"}:
return basename[: -len(".jsonl")]
return os.path.basename(os.path.dirname(path))
def _extract_wire_event(value: Mapping[str, Any]) -> tuple[str, str, str | None] | None:
"""Normalize the replay event shape used before current AgentRecord logs."""
event: Mapping[str, Any] = value
if value.get("type") == "agent":
if value.get("agentId") != "main" or not isinstance(value.get("event"), Mapping):
return None
event = value["event"]
message = event.get("message")
if isinstance(message, Mapping):
event = message
event_type = event.get("type")
payload = event.get("payload")
if not isinstance(event_type, str) or not isinstance(payload, Mapping):
return None
if event_type == "TurnBegin":
text = _content(payload.get("user_input"))
return ("user", text, None) if text is not None else None
if event_type == "ContentPart":
part_type = payload.get("type")
if isinstance(part_type, str) and part_type.casefold() in _FILTERED:
return None
text = _content(payload)
return ("assistant", text, None) if text is not None else None
if event_type == "ToolResult":
result = payload.get("return_value")
text = _content(result)
if text is None and isinstance(result, Mapping):
text = _content(_first(result, ("output", "message")))
return ("tool", text, None) if text is not None else None
return None
def _extract_message(value: Mapping[str, Any]) -> tuple[str, str, str | None] | None:
candidates: list[Mapping[str, Any]] = []
pending: list[tuple[Mapping[str, Any], int]] = [(value, 0)]
seen: set[int] = set()
while pending:
candidate, depth = pending.pop(0)
identity = id(candidate)
if identity in seen:
continue
seen.add(identity)
candidates.append(candidate)
if depth >= 6:
continue
for key in ("message", "data", "payload", "record", "event", "params", "update"):
nested = candidate.get(key)
if isinstance(nested, Mapping):
pending.append((nested, depth + 1))
for candidate in candidates:
raw_role = _first(candidate, ("role", "sender", "author"))
if isinstance(raw_role, Mapping):
raw_role = _first(raw_role, ("role", "type", "name"))
if not isinstance(raw_role, str):
kind = _first(candidate, ("type", "kind", "event", "sessionUpdate"))
if isinstance(kind, str):
lowered = kind.casefold()
if lowered in {"user", "human"} or "user_message" in lowered:
raw_role = "user"
elif lowered in {"assistant", "agent", "model", "ai"} or "assistant_message" in lowered or "agent_message" in lowered:
raw_role = "assistant"
elif "tool_result" in lowered:
raw_role = "tool"
if not isinstance(raw_role, str):
continue
lowered_role = raw_role.casefold()
if any(token in lowered_role for token in _FILTERED):
return None
role = _ROLE_ALIASES.get(lowered_role)
if role is None:
continue
content = _content(_first(candidate, ("content", "text", "message", "output", "result")))
if content is None:
continue
tool_name = _string(_first(candidate, ("tool_name", "toolName", "name"))) if role == "tool" else None
return role, content, tool_name
return None
def _content(value: Any) -> str | None:
if isinstance(value, str):
return value
if isinstance(value, Mapping):
kind = value.get("type")
if isinstance(kind, str) and any(token in kind.casefold() for token in _FILTERED):
return None
for key in ("text", "content", "value"):
text = value.get(key)
if isinstance(text, str):
return text
return None
if isinstance(value, list):
pieces: list[str] = []
for item in value:
if isinstance(item, str):
pieces.append(item)
elif isinstance(item, Mapping):
kind = item.get("type")
if isinstance(kind, str) and any(token in kind.casefold() for token in _FILTERED):
continue
text = _content(item)
if text is not None:
pieces.append(text)
return "".join(pieces) if pieces else None
return None
ADAPTER = KimiAdapter()
def get_adapter() -> KimiAdapter:
return ADAPTER
SHA-256: 38ee30f468e1d8d0d92bc77f38ce2f65f936c14b17640f39a911678043f2c23d