mirror of
https://github.com/langchain-ai/deepagents.git
synced 2026-08-28 05:00:04 -04:00
dfcd9417a8
A transient provider failure after partial output no longer ends the turn: the model call is retried, and already-streamed output is marked as incomplete instead of being joined to the replay. --- Previously, retries stopped once any message chunk had been emitted. This change separates retry policy from presentation: - The middleware emits attempt lifecycle and retry correlation events. - Streaming clients preserve and annotate superseded partial output. - Buffered `--no-stream` mode drops failed-attempt text before writing to stdout. - Transcript, usage, and tool-hook state are scoped by attempt so replayed output is not merged or suppressed. Failed attempts still count toward usage, and clients tolerate missing or malformed lifecycle events for compatibility. Design details and alternatives are documented in `libs/code/STREAMING_RETRY_DESIGN.md`. <details> <summary>Test plan</summary> - Retry behavior before and after output begins - Streaming, buffered, and TUI presentation - Transcript, usage, and tool-hook settlement across retries - Missing, malformed, duplicated, and out-of-order lifecycle events </details> ### Screenshot  Made by [Open SWE](https://openswe.vercel.app/agents/722f8e40-29e4-5eb0-9fd5-368416105bee) --------- Co-authored-by: open-swe[bot] <open-swe@users.noreply.github.com>
817 lines
28 KiB
Python
817 lines
28 KiB
Python
"""Client-owned conversation transcript projections for Hooks v2.
|
|
|
|
Materializes versioned per-thread and per-subagent JSONL files that hook
|
|
commands can read via `transcript_path` / `agent_transcript_path`.
|
|
|
|
Lag semantics:
|
|
The on-disk JSONL may lag behind live checkpoint/UI state. Callers that need
|
|
the just-finished assistant turn must prefer `last_assistant_message` on
|
|
Stop/SubagentStop. `materialize()` flushes pending records immediately before
|
|
returning a path so hooks see a consistent snapshot of what the store has
|
|
accepted so far, not a live tail of the server checkpoint.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import hashlib
|
|
import logging
|
|
import os
|
|
import re
|
|
import tempfile
|
|
import threading
|
|
import unicodedata
|
|
from contextlib import contextmanager, suppress
|
|
from dataclasses import dataclass, field
|
|
from pathlib import Path
|
|
from typing import TYPE_CHECKING, Literal
|
|
from urllib.parse import parse_qsl, urlencode, urlsplit, urlunsplit
|
|
|
|
if os.name != "nt":
|
|
import fcntl
|
|
else:
|
|
import msvcrt
|
|
|
|
from langchain_core.messages import (
|
|
AIMessage,
|
|
BaseMessage,
|
|
BaseMessageChunk,
|
|
HumanMessage,
|
|
SystemMessage,
|
|
ToolMessage,
|
|
message_chunk_to_message,
|
|
)
|
|
from pydantic import BaseModel, ConfigDict
|
|
|
|
from deepagents_code._constants import LOCAL_CONTEXT_MESSAGE_SOURCE
|
|
from deepagents_code.config_manifest import _is_secret_env
|
|
from deepagents_code.json_types import JSON_VALUE_ADAPTER, JsonValue
|
|
|
|
if TYPE_CHECKING:
|
|
from collections.abc import Iterator, Mapping, Sequence
|
|
from typing import Protocol
|
|
|
|
class _TranscriptRuntime(Protocol):
|
|
def append_messages(
|
|
self,
|
|
thread_id: str,
|
|
messages: Sequence[BaseMessage],
|
|
*,
|
|
agent_id: str | None = None,
|
|
) -> None: ...
|
|
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
SUBAGENT_TRANSCRIPT_ID_METADATA_KEY = "dcode_subagent_id"
|
|
_INTERNAL_STREAM_SOURCES = frozenset(
|
|
{LOCAL_CONTEXT_MESSAGE_SOURCE, "summarization", "auto_mode_classifier"}
|
|
)
|
|
|
|
TRANSCRIPT_SCHEMA_VERSION = 1
|
|
DEFAULT_RETENTION_REVISIONS = 20
|
|
_FILE_MODE = 0o600
|
|
_DIR_MODE = 0o700
|
|
# Credential-style assignments. Bare names like PASSWORD= are matched via the
|
|
# trailing keyword alternatives when preceded by an underscore-separated prefix
|
|
# (for example OPENAI_API_KEY=), matching the repository secret-name policy.
|
|
_SECRET_ASSIGNMENT_RE = re.compile(
|
|
r"(?i)\b([A-Z][A-Z0-9_]*(?:API[_-]?KEY|KEY|TOKEN|SECRET|PASSWORD|CREDENTIAL)"
|
|
r"[A-Z0-9_]*)\s*=\s*([^\s,;]+)"
|
|
)
|
|
_BEARER_RE = re.compile(r"(?i)\b(Bearer)\s+[A-Za-z0-9._~+/=-]{8,}")
|
|
_PREFIXED_TOKEN_RE = re.compile(
|
|
r"(?<![A-Za-z0-9])(?:"
|
|
r"sk-(?:ant-)?|sk_(?:live|test)_|pk_(?:live|test)_|gh[pousr]_|"
|
|
r"github_pat_|glpat-|xox[baprs]-|hf_|npm_|AIza|AKIA"
|
|
r")[A-Za-z0-9._-]{8,}"
|
|
)
|
|
_JWT_RE = re.compile(
|
|
r"(?<![A-Za-z0-9_-])eyJ[A-Za-z0-9_-]{6,}\."
|
|
r"[A-Za-z0-9_-]{6,}\.[A-Za-z0-9_-]{6,}"
|
|
)
|
|
_URL_RE = re.compile(r"https?://[^\s<>\"']+", re.IGNORECASE)
|
|
_SAFE_PREFIX_RE = re.compile(r"[^a-z0-9]+")
|
|
_SAFE_PREFIX_LENGTH = 32
|
|
_EMPTY_REVISION = hashlib.sha256(b"").hexdigest()
|
|
|
|
|
|
class TranscriptRecord(BaseModel):
|
|
"""One JSONL record in a materialized transcript projection."""
|
|
|
|
model_config = ConfigDict(extra="forbid")
|
|
|
|
schema_version: Literal[1] = TRANSCRIPT_SCHEMA_VERSION
|
|
"""Transcript schema version used to interpret this record."""
|
|
|
|
sequence: int
|
|
"""Zero-based position of this record within its transcript."""
|
|
|
|
record_id: str
|
|
"""Message identifier, or a deterministic role-and-sequence fallback."""
|
|
|
|
timestamp: str | None = None
|
|
"""Source message timestamp when one is available."""
|
|
|
|
thread_id: str
|
|
"""Conversation thread that owns this record."""
|
|
|
|
agent_id: str | None = None
|
|
"""Subagent scope for an agent transcript, otherwise `None`."""
|
|
|
|
role: Literal["user", "assistant", "tool", "system"]
|
|
"""Normalized conversation role for the projected message."""
|
|
|
|
message_id: str | None = None
|
|
"""Original LangChain message identifier when one is available."""
|
|
|
|
content: JsonValue
|
|
"""Redacted JSON-compatible message content."""
|
|
|
|
name: str | None = None
|
|
"""Tool or message name when one is available."""
|
|
|
|
|
|
@dataclass(frozen=True, slots=True)
|
|
class TranscriptHandle:
|
|
"""Identity of a materialized transcript file."""
|
|
|
|
path: Path
|
|
revision: str
|
|
thread_id: str
|
|
agent_id: str | None = None
|
|
|
|
|
|
@dataclass
|
|
class _TranscriptBuffer:
|
|
records: list[TranscriptRecord] = field(default_factory=list)
|
|
record_ids: set[str] = field(default_factory=set)
|
|
dirty: bool = False
|
|
revision: str = _EMPTY_REVISION
|
|
|
|
|
|
class TranscriptStore:
|
|
"""Append-only JSONL transcript projections owned by the client process."""
|
|
|
|
def __init__(
|
|
self,
|
|
root: Path,
|
|
*,
|
|
retention_revisions: int = DEFAULT_RETENTION_REVISIONS,
|
|
) -> None:
|
|
"""Create a store rooted at `root`.
|
|
|
|
Args:
|
|
root: Directory that will contain per-thread transcript files.
|
|
retention_revisions: Maximum prior `.bak-*` revisions retained per
|
|
transcript after each rewrite.
|
|
|
|
Raises:
|
|
ValueError: If `retention_revisions` is negative.
|
|
"""
|
|
if retention_revisions < 0:
|
|
msg = "retention_revisions must be nonnegative"
|
|
raise ValueError(msg)
|
|
self.root = root.expanduser().resolve()
|
|
self.retention_revisions = retention_revisions
|
|
self._buffers: dict[tuple[str, str | None], _TranscriptBuffer] = {}
|
|
self._lock = threading.RLock()
|
|
_ensure_private_directories(self.root, self.root)
|
|
|
|
def thread_path(self, thread_id: str) -> Path:
|
|
"""Return the materialized path for a thread transcript.
|
|
|
|
Args:
|
|
thread_id: Conversation thread identifier.
|
|
|
|
Returns:
|
|
Absolute JSONL path for the thread.
|
|
"""
|
|
return self.root / f"{_safe_component(thread_id)}.jsonl"
|
|
|
|
def agent_path(self, thread_id: str, agent_id: str) -> Path:
|
|
"""Return the materialized path for a subagent transcript.
|
|
|
|
Args:
|
|
thread_id: Parent conversation thread identifier.
|
|
agent_id: Subagent identifier.
|
|
|
|
Returns:
|
|
Absolute JSONL path nested under the thread.
|
|
"""
|
|
return (
|
|
self.root
|
|
/ _safe_component(thread_id)
|
|
/ "agents"
|
|
/ f"{_safe_component(agent_id)}.jsonl"
|
|
)
|
|
|
|
def append_messages(
|
|
self,
|
|
thread_id: str,
|
|
messages: Sequence[BaseMessage],
|
|
*,
|
|
agent_id: str | None = None,
|
|
) -> None:
|
|
"""Append redacted message projections to the in-memory buffer.
|
|
|
|
Args:
|
|
thread_id: Conversation thread identifier.
|
|
messages: LangChain messages to project.
|
|
agent_id: Optional subagent scope.
|
|
"""
|
|
with self._lock:
|
|
buffer = self._buffer(thread_id, agent_id)
|
|
for message in messages:
|
|
if _is_local_context_message(message):
|
|
continue
|
|
record = _record_from_message(
|
|
message,
|
|
thread_id=thread_id,
|
|
agent_id=agent_id,
|
|
sequence=len(buffer.records),
|
|
)
|
|
if record is None:
|
|
continue
|
|
if (
|
|
record.message_id is not None
|
|
and record.record_id in buffer.record_ids
|
|
):
|
|
continue
|
|
buffer.records.append(record)
|
|
buffer.record_ids.add(record.record_id)
|
|
buffer.dirty = True
|
|
|
|
def materialize(
|
|
self,
|
|
thread_id: str,
|
|
*,
|
|
agent_id: str | None = None,
|
|
) -> TranscriptHandle:
|
|
"""Flush pending records and return the client-readable path.
|
|
|
|
Materialization rewrites the whole per-thread JSONL file, so it runs
|
|
under a cross-process advisory lock and re-reads the on-disk records
|
|
first, merging in anything another dcode process appended since this
|
|
process last loaded the file.
|
|
|
|
Args:
|
|
thread_id: Conversation thread identifier.
|
|
agent_id: Optional subagent scope.
|
|
|
|
Returns:
|
|
Handle with path and content revision identity.
|
|
"""
|
|
with self._lock:
|
|
buffer = self._buffer(thread_id, agent_id)
|
|
path = (
|
|
self.agent_path(thread_id, agent_id)
|
|
if agent_id is not None
|
|
else self.thread_path(thread_id)
|
|
)
|
|
_ensure_private_directories(self.root, path.parent)
|
|
if path.is_file() and os.name != "nt":
|
|
path.chmod(_FILE_MODE)
|
|
with _file_lock(path.with_suffix(path.suffix + ".lock")):
|
|
_merge_disk_records(path, buffer)
|
|
if buffer.dirty or not path.is_file():
|
|
revision = _write_transcript(
|
|
self.root,
|
|
path,
|
|
buffer.records,
|
|
self.retention_revisions,
|
|
)
|
|
buffer.revision = revision
|
|
buffer.dirty = False
|
|
return TranscriptHandle(
|
|
path=path,
|
|
revision=buffer.revision,
|
|
thread_id=thread_id,
|
|
agent_id=agent_id,
|
|
)
|
|
|
|
def revision(self, thread_id: str, *, agent_id: str | None = None) -> str:
|
|
"""Return the current revision id without forcing a flush.
|
|
|
|
Args:
|
|
thread_id: Conversation thread identifier.
|
|
agent_id: Optional subagent scope.
|
|
|
|
Returns:
|
|
Content revision string for the buffered projection.
|
|
"""
|
|
with self._lock:
|
|
buffer = self._buffer(thread_id, agent_id)
|
|
if buffer.dirty:
|
|
return _revision_for_records(buffer.records)
|
|
return buffer.revision
|
|
|
|
def _buffer(self, thread_id: str, agent_id: str | None) -> _TranscriptBuffer:
|
|
key = (thread_id, agent_id)
|
|
buffer = self._buffers.get(key)
|
|
if buffer is None:
|
|
buffer = _TranscriptBuffer()
|
|
path = (
|
|
self.agent_path(thread_id, agent_id)
|
|
if agent_id is not None
|
|
else self.thread_path(thread_id)
|
|
)
|
|
if path.is_file():
|
|
buffer.records, valid = _read_transcript(path)
|
|
buffer.record_ids = {record.record_id for record in buffer.records}
|
|
buffer.revision = _revision_for_records(buffer.records)
|
|
buffer.dirty = not valid
|
|
self._buffers[key] = buffer
|
|
return buffer
|
|
|
|
|
|
@dataclass(slots=True)
|
|
class _AttemptScope:
|
|
"""Staged records for one model attempt of one agent.
|
|
|
|
The owning `agent_id` is the `TranscriptRecorder._attempts` key rather than
|
|
a field here, so a scope cannot disagree with where it is stored. "Open" is
|
|
likewise dict membership, not a flag: `complete_attempt` and
|
|
`discard_attempt` both pop, which makes a committed-and-discarded scope
|
|
unrepresentable instead of merely invalid.
|
|
"""
|
|
|
|
call_id: str
|
|
attempt: int
|
|
staged: list[BaseMessage] = field(default_factory=list)
|
|
|
|
|
|
@dataclass(slots=True)
|
|
class TranscriptRecorder:
|
|
"""Collect completed stream messages into a Hooks transcript runtime."""
|
|
|
|
runtime: _TranscriptRuntime
|
|
thread_id: str
|
|
_chunks: dict[tuple[str | None, str], BaseMessageChunk] = field(
|
|
default_factory=dict
|
|
)
|
|
_attempts: dict[str | None, _AttemptScope] = field(default_factory=dict)
|
|
|
|
def record(
|
|
self,
|
|
message: object,
|
|
metadata: Mapping[str, object] | None,
|
|
*,
|
|
main_agent: bool,
|
|
) -> None:
|
|
"""Record one streamed message when its transcript identity is stable.
|
|
|
|
While an attempt scope is open for the resolved `agent_id`, the record
|
|
is staged rather than appended, and only reaches the transcript when
|
|
`complete_attempt` commits it.
|
|
|
|
Args:
|
|
message: Streamed LangChain message or chunk.
|
|
metadata: Stream metadata carrying optional subagent identity.
|
|
main_agent: Whether the message belongs to the root graph.
|
|
"""
|
|
if (
|
|
metadata is not None
|
|
and metadata.get("lc_source") in _INTERNAL_STREAM_SOURCES
|
|
) or _is_local_context_message(message):
|
|
return
|
|
agent_id = None if main_agent else _stream_agent_id(metadata)
|
|
if not main_agent and agent_id is None:
|
|
return
|
|
if isinstance(message, BaseMessageChunk):
|
|
key = (agent_id, message.id or type(message).__name__)
|
|
previous = self._chunks.get(key)
|
|
combined = message if previous is None else previous + message
|
|
self._chunks[key] = combined
|
|
if getattr(message, "chunk_position", None) != "last":
|
|
return
|
|
self._chunks.pop(key, None)
|
|
self._append(message_chunk_to_message(combined), agent_id=agent_id)
|
|
return
|
|
if isinstance(message, BaseMessage):
|
|
self._append(message, agent_id=agent_id)
|
|
|
|
def start_attempt(
|
|
self,
|
|
*,
|
|
agent_id: str | None,
|
|
call_id: str,
|
|
attempt: int,
|
|
) -> None:
|
|
"""Open (or replace) the active staging scope for an agent.
|
|
|
|
Staged messages buffer in memory until `complete_attempt` commits
|
|
them, so retried attempts do not duplicate transcript records.
|
|
|
|
Args:
|
|
agent_id: Subagent scope, or `None` for the main agent.
|
|
call_id: LLM call identifier for this attempt.
|
|
attempt: Zero-based attempt number within the call.
|
|
"""
|
|
if self._active_scope(agent_id, call_id=call_id, attempt=attempt) is not None:
|
|
return
|
|
existing = self._attempts.get(agent_id)
|
|
if existing is not None and existing.staged:
|
|
# A start for a different attempt replaces the open scope outright,
|
|
# which drops whatever it staged. Callers are expected to discard
|
|
# the superseded scope first; reaching here with staged records
|
|
# means a lifecycle event went missing, so leave a trace rather
|
|
# than losing transcript records silently.
|
|
logger.warning(
|
|
"Replacing open transcript attempt %s/%d for agent %s while it "
|
|
"held %d staged record(s); they are dropped",
|
|
existing.call_id,
|
|
existing.attempt,
|
|
agent_id,
|
|
len(existing.staged),
|
|
)
|
|
self._drop_chunks(agent_id)
|
|
self._attempts[agent_id] = _AttemptScope(call_id=call_id, attempt=attempt)
|
|
|
|
def complete_attempt(
|
|
self,
|
|
*,
|
|
agent_id: str | None,
|
|
call_id: str,
|
|
attempt: int,
|
|
) -> None:
|
|
"""Commit staged messages for a matching attempt, then clear it.
|
|
|
|
Duplicate or mismatched calls are ignored idempotently.
|
|
|
|
Args:
|
|
agent_id: Subagent scope, or `None` for the main agent.
|
|
call_id: LLM call identifier for this attempt.
|
|
attempt: Zero-based attempt number within the call.
|
|
"""
|
|
scope = self._active_scope(agent_id, call_id=call_id, attempt=attempt)
|
|
if scope is None:
|
|
return
|
|
staged = scope.staged
|
|
self._attempts.pop(agent_id, None)
|
|
self._drop_chunks(agent_id)
|
|
for message in staged:
|
|
self._append(message, agent_id=agent_id)
|
|
|
|
def discard_attempt(
|
|
self,
|
|
*,
|
|
agent_id: str | None,
|
|
call_id: str,
|
|
attempt: int,
|
|
) -> None:
|
|
"""Drop staged messages and pending chunks for a matching attempt.
|
|
|
|
Duplicate or mismatched calls are ignored idempotently.
|
|
|
|
Args:
|
|
agent_id: Subagent scope, or `None` for the main agent.
|
|
call_id: LLM call identifier for this attempt.
|
|
attempt: Zero-based attempt number within the call.
|
|
"""
|
|
scope = self._active_scope(agent_id, call_id=call_id, attempt=attempt)
|
|
if scope is None:
|
|
return
|
|
self._attempts.pop(agent_id, None)
|
|
self._drop_chunks(agent_id)
|
|
|
|
def drop_uncommitted(self) -> None:
|
|
"""Clear all attempt scopes and chunk accumulators without appending.
|
|
|
|
Terminal teardown: these attempts never completed (abort, retry-budget
|
|
exhaustion, cancellation), so their records are not transcript history.
|
|
Dropping non-empty staging makes the on-screen conversation and the
|
|
persisted transcript diverge, so count it rather than clearing silently.
|
|
"""
|
|
dropped = sum(len(scope.staged) for scope in self._attempts.values())
|
|
if dropped:
|
|
logger.warning(
|
|
"Dropping %d staged transcript record(s) from %d uncommitted "
|
|
"model attempt(s)",
|
|
dropped,
|
|
len(self._attempts),
|
|
)
|
|
self._attempts.clear()
|
|
self._chunks.clear()
|
|
|
|
def _drop_chunks(self, agent_id: str | None) -> None:
|
|
"""Drop every partial chunk accumulator belonging to one agent."""
|
|
for key in [key for key in self._chunks if key[0] == agent_id]:
|
|
del self._chunks[key]
|
|
|
|
def _active_scope(
|
|
self,
|
|
agent_id: str | None,
|
|
*,
|
|
call_id: str,
|
|
attempt: int,
|
|
) -> _AttemptScope | None:
|
|
scope = self._attempts.get(agent_id)
|
|
if scope is None or scope.call_id != call_id or scope.attempt != attempt:
|
|
return None
|
|
return scope
|
|
|
|
def _append(self, message: BaseMessage, *, agent_id: str | None) -> None:
|
|
scope = self._attempts.get(agent_id)
|
|
if scope is not None:
|
|
scope.staged.append(message)
|
|
return
|
|
append_messages = getattr(self.runtime, "append_messages", None)
|
|
if not callable(append_messages):
|
|
return
|
|
try:
|
|
append_messages(self.thread_id, [message], agent_id=agent_id)
|
|
except (TypeError, ValueError):
|
|
logger.warning(
|
|
"Skipping invalid streamed transcript message",
|
|
exc_info=True,
|
|
)
|
|
|
|
def append(self, messages: Sequence[BaseMessage]) -> None:
|
|
"""Append checkpoint or input messages to the root transcript.
|
|
|
|
Deliberately bypasses attempt staging: these messages are turn inputs
|
|
and checkpoint history, not model output, so no attempt owns them and
|
|
none should be able to discard them. Callers append them between turns,
|
|
never while a scope is open -- doing so would order them ahead of model
|
|
output that was staged earlier.
|
|
"""
|
|
append_messages = getattr(self.runtime, "append_messages", None)
|
|
if callable(append_messages):
|
|
append_messages(self.thread_id, messages)
|
|
|
|
|
|
def _stream_agent_id(metadata: Mapping[str, object] | None) -> str | None:
|
|
if metadata is None:
|
|
return None
|
|
value = metadata.get(SUBAGENT_TRANSCRIPT_ID_METADATA_KEY)
|
|
return value if isinstance(value, str) and value else None
|
|
|
|
|
|
def _is_local_context_message(message: object) -> bool:
|
|
"""Return whether a local or serialized message is internal context."""
|
|
if isinstance(message, BaseMessage):
|
|
metadata = getattr(message, "additional_kwargs", None)
|
|
elif isinstance(message, dict):
|
|
metadata = message.get("additional_kwargs")
|
|
else:
|
|
return False
|
|
return (
|
|
isinstance(metadata, dict)
|
|
and metadata.get("lc_source") == LOCAL_CONTEXT_MESSAGE_SOURCE
|
|
)
|
|
|
|
|
|
def _record_from_message(
|
|
message: BaseMessage,
|
|
*,
|
|
thread_id: str,
|
|
agent_id: str | None,
|
|
sequence: int,
|
|
) -> TranscriptRecord | None:
|
|
if isinstance(message, HumanMessage):
|
|
role: Literal["user", "assistant", "tool", "system"] = "user"
|
|
elif isinstance(message, AIMessage):
|
|
role = "assistant"
|
|
elif isinstance(message, ToolMessage):
|
|
role = "tool"
|
|
elif isinstance(message, SystemMessage):
|
|
role = "system"
|
|
else:
|
|
return None
|
|
|
|
raw = message.model_dump(mode="json")
|
|
content = JSON_VALUE_ADAPTER.validate_python(
|
|
redact_transcript_value(raw.get("content"))
|
|
)
|
|
message_id = message.id if isinstance(message.id, str) else None
|
|
name = getattr(message, "name", None)
|
|
tool_name = name if isinstance(name, str) else None
|
|
record_id = message_id or f"{role}:{sequence}"
|
|
return TranscriptRecord(
|
|
sequence=sequence,
|
|
record_id=record_id,
|
|
thread_id=thread_id,
|
|
agent_id=agent_id,
|
|
role=role,
|
|
message_id=message_id,
|
|
content=content,
|
|
name=tool_name,
|
|
)
|
|
|
|
|
|
def redact_transcript_value(value: object) -> JsonValue:
|
|
"""Redact secret-like strings inside transcript content.
|
|
|
|
Args:
|
|
value: Arbitrary message content.
|
|
|
|
Returns:
|
|
JSON-compatible content with URLs/credentials scrubbed.
|
|
"""
|
|
if isinstance(value, str):
|
|
return _redact_text(value)
|
|
if isinstance(value, list):
|
|
return [redact_transcript_value(item) for item in value]
|
|
if isinstance(value, dict):
|
|
return {
|
|
str(key): (
|
|
"[redacted]"
|
|
if _is_secret_env(str(key))
|
|
else redact_transcript_value(item)
|
|
)
|
|
for key, item in value.items()
|
|
}
|
|
return JSON_VALUE_ADAPTER.validate_python(value)
|
|
|
|
|
|
def _redact_text(text: str) -> str:
|
|
redacted = _SECRET_ASSIGNMENT_RE.sub(
|
|
lambda match: f"{match.group(1)}=[redacted]",
|
|
text,
|
|
)
|
|
redacted = _BEARER_RE.sub(lambda match: f"{match.group(1)} [redacted]", redacted)
|
|
redacted = _PREFIXED_TOKEN_RE.sub("[redacted]", redacted)
|
|
redacted = _JWT_RE.sub("[redacted]", redacted)
|
|
return _URL_RE.sub(lambda match: _redact_url(match.group(0)), redacted)
|
|
|
|
|
|
def _redact_url(value: str) -> str:
|
|
try:
|
|
parsed = urlsplit(value)
|
|
hostname = parsed.hostname or ""
|
|
port = parsed.port
|
|
except ValueError:
|
|
return "[redacted URL]"
|
|
if ":" in hostname and not hostname.startswith("["):
|
|
hostname = f"[{hostname}]"
|
|
netloc = f"{hostname}:{port}" if port is not None else hostname
|
|
path = "/[redacted]" if parsed.path else ""
|
|
query_items = parse_qsl(parsed.query, keep_blank_values=True)
|
|
query = urlencode([(key, "[redacted]") for key, _value in query_items])
|
|
fragment = "[redacted]" if parsed.fragment else ""
|
|
return urlunsplit((parsed.scheme, netloc, path, query, fragment))
|
|
|
|
|
|
@contextmanager
|
|
def _file_lock(lock_path: Path) -> Iterator[None]:
|
|
"""Hold an advisory cross-process lock on `lock_path`.
|
|
|
|
Materialization rewrites the whole transcript file, so concurrent dcode
|
|
processes resuming the same thread must serialize their read-merge-write
|
|
cycle on something stronger than the in-process `threading.RLock`.
|
|
"""
|
|
fd = os.open(lock_path, os.O_RDWR | os.O_CREAT, _FILE_MODE)
|
|
try:
|
|
if os.name != "nt":
|
|
fcntl.flock(fd, fcntl.LOCK_EX)
|
|
else:
|
|
msvcrt.locking(fd, msvcrt.LK_LOCK, 1)
|
|
try:
|
|
yield
|
|
finally:
|
|
if os.name != "nt":
|
|
fcntl.flock(fd, fcntl.LOCK_UN)
|
|
else:
|
|
with suppress(OSError):
|
|
os.lseek(fd, 0, os.SEEK_SET)
|
|
msvcrt.locking(fd, msvcrt.LK_UNLCK, 1)
|
|
finally:
|
|
os.close(fd)
|
|
|
|
|
|
def _merge_disk_records(path: Path, buffer: _TranscriptBuffer) -> None:
|
|
"""Fold on-disk records missing from `buffer` back into it.
|
|
|
|
Another process may have materialized the shared transcript since this
|
|
process last loaded it; re-reading under the file lock prevents the next
|
|
rewrite from silently dropping that process's records. Buffer records win
|
|
ordering ties; disk-only records keep their relative order appended after.
|
|
"""
|
|
if not path.is_file():
|
|
return
|
|
disk_records, valid = _read_transcript(path)
|
|
if not valid:
|
|
return
|
|
buffer_ids = {record.record_id for record in buffer.records}
|
|
if all(record.record_id in buffer_ids for record in disk_records):
|
|
return
|
|
merged = list(buffer.records)
|
|
merged_ids = set(buffer_ids)
|
|
for record in disk_records:
|
|
if record.record_id in merged_ids:
|
|
continue
|
|
merged.append(record.model_copy(update={"sequence": len(merged)}))
|
|
merged_ids.add(record.record_id)
|
|
buffer.records = merged
|
|
buffer.record_ids = merged_ids
|
|
buffer.dirty = True
|
|
|
|
|
|
def _write_transcript(
|
|
root: Path,
|
|
path: Path,
|
|
records: Sequence[TranscriptRecord],
|
|
retention_revisions: int,
|
|
) -> str:
|
|
_ensure_private_directories(root, path.parent)
|
|
payload = "".join(
|
|
record.model_dump_json(exclude_none=True) + "\n" for record in records
|
|
)
|
|
revision = hashlib.sha256(payload.encode("utf-8")).hexdigest()
|
|
fd, raw_tmp = tempfile.mkstemp(dir=path.parent, suffix=".tmp")
|
|
tmp_path = Path(raw_tmp)
|
|
try:
|
|
with os.fdopen(fd, "w", encoding="utf-8") as handle:
|
|
handle.write(payload)
|
|
handle.flush()
|
|
os.fsync(handle.fileno())
|
|
if os.name != "nt":
|
|
tmp_path.chmod(_FILE_MODE)
|
|
# Copy the previous revision aside first, then atomically replace the
|
|
# live path so concurrent readers never observe a missing file.
|
|
if path.exists():
|
|
prior_payload = path.read_bytes()
|
|
prior_revision = hashlib.sha256(prior_payload).hexdigest()
|
|
backup = path.with_suffix(path.suffix + f".bak-{prior_revision}")
|
|
_write_backup(backup, prior_payload)
|
|
_prune_backups(path, retention_revisions)
|
|
tmp_path.replace(path)
|
|
if os.name != "nt":
|
|
path.chmod(_FILE_MODE)
|
|
except OSError:
|
|
logger.warning("Failed to materialize transcript at %s", path, exc_info=True)
|
|
with suppress(OSError):
|
|
tmp_path.unlink(missing_ok=True)
|
|
raise
|
|
return revision
|
|
|
|
|
|
def _write_backup(path: Path, payload: bytes) -> None:
|
|
fd, raw_tmp = tempfile.mkstemp(dir=path.parent, suffix=".bak.tmp")
|
|
tmp_path = Path(raw_tmp)
|
|
try:
|
|
with os.fdopen(fd, "wb") as handle:
|
|
handle.write(payload)
|
|
handle.flush()
|
|
os.fsync(handle.fileno())
|
|
if os.name != "nt":
|
|
tmp_path.chmod(_FILE_MODE)
|
|
tmp_path.replace(path)
|
|
except OSError:
|
|
with suppress(OSError):
|
|
tmp_path.unlink(missing_ok=True)
|
|
raise
|
|
|
|
|
|
def _safe_component(identifier: str) -> str:
|
|
normalized = unicodedata.normalize("NFKD", identifier)
|
|
readable = normalized.encode("ascii", errors="ignore").decode("ascii").lower()
|
|
prefix = _SAFE_PREFIX_RE.sub("-", readable).strip("-")[:_SAFE_PREFIX_LENGTH]
|
|
digest = hashlib.sha256(identifier.encode("utf-8")).hexdigest()
|
|
return f"{prefix or 'id'}--{digest}"
|
|
|
|
|
|
def _ensure_private_directories(root: Path, target: Path) -> None:
|
|
root.mkdir(parents=True, exist_ok=True, mode=_DIR_MODE)
|
|
target.mkdir(parents=True, exist_ok=True, mode=_DIR_MODE)
|
|
if os.name == "nt":
|
|
return
|
|
root.chmod(_DIR_MODE)
|
|
relative = target.relative_to(root)
|
|
current = root
|
|
for part in relative.parts:
|
|
current /= part
|
|
current.chmod(_DIR_MODE)
|
|
|
|
|
|
def _prune_backups(path: Path, retention_revisions: int) -> None:
|
|
pattern = f"{path.name}.bak-*"
|
|
backups = sorted(
|
|
path.parent.glob(pattern),
|
|
key=lambda item: (item.stat().st_mtime_ns, item.name),
|
|
)
|
|
excess = len(backups) - retention_revisions
|
|
for stale in backups[: max(0, excess)]:
|
|
with suppress(OSError):
|
|
stale.unlink()
|
|
|
|
|
|
def _read_transcript(path: Path) -> tuple[list[TranscriptRecord], bool]:
|
|
records: list[TranscriptRecord] = []
|
|
try:
|
|
for line in path.read_text(encoding="utf-8").splitlines():
|
|
if not line.strip():
|
|
continue
|
|
records.append(TranscriptRecord.model_validate_json(line))
|
|
except (OSError, ValueError):
|
|
logger.warning("Could not read transcript at %s", path, exc_info=True)
|
|
return [], False
|
|
return records, True
|
|
|
|
|
|
def _revision_for_records(records: Sequence[TranscriptRecord]) -> str:
|
|
payload = "".join(
|
|
record.model_dump_json(exclude_none=True) + "\n" for record in records
|
|
)
|
|
return hashlib.sha256(payload.encode("utf-8")).hexdigest()
|