← Files Portable ResumeARCHIVED FILE

skills/.portable-resume/runtime/portable_resume/adapters/kimi.py

57.4 KB · Oct 2, 2026 · 00:33 UTC

↓ Download file

"""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