refactor(state): collapse multi-line SQL assembly and log calls in hermes_state_common to one-liners
This commit is contained in:
@@ -14,9 +14,8 @@ import time
|
||||
from typing import Any
|
||||
|
||||
from agent.skill_commands import SKILL_EXCERPT_JOINT, SKILL_SCAFFOLD_SQL_LIKE, describe_skill_invocation
|
||||
from agent.context_compressor import (
|
||||
LEGACY_SUMMARY_PREFIX, SUMMARY_PREFIX, _MERGED_PRIOR_CONTEXT_HEADER, _MERGED_SUMMARY_DELIMITER,
|
||||
_SUMMARY_END_MARKER)
|
||||
from agent.context_compressor import (LEGACY_SUMMARY_PREFIX, SUMMARY_PREFIX, _MERGED_PRIOR_CONTEXT_HEADER,
|
||||
_MERGED_SUMMARY_DELIMITER, _SUMMARY_END_MARKER)
|
||||
|
||||
|
||||
# Session preview = head of the first user message, shown wherever a session has no title. A /skill
|
||||
@@ -68,49 +67,35 @@ _PREVIEW_LONG_FORM_PREFIX = SUMMARY_PREFIX.split("Do NOT answer", 1)[0]
|
||||
_PREVIEW_SUMMARY_PREFIXES = (_PREVIEW_LONG_FORM_PREFIX, LEGACY_SUMMARY_PREFIX)
|
||||
_PREVIEW_STANDALONE_SUMMARY_SQL = _sql_starts_with("m.content", _PREVIEW_SUMMARY_PREFIXES)
|
||||
_PREVIEW_MERGED_AFTER_SQL = _sql_after_marker(_MERGED_SUMMARY_DELIMITER)
|
||||
_PREVIEW_MERGED_SUMMARY_SQL = (
|
||||
f"(INSTR(m.content, {_sql_literal(_MERGED_SUMMARY_DELIMITER)}) > 0"
|
||||
f" AND {_sql_starts_with(_PREVIEW_MERGED_AFTER_SQL, _PREVIEW_SUMMARY_PREFIXES)})"
|
||||
)
|
||||
_PREVIEW_MERGED_SUMMARY_SQL = (f"(INSTR(m.content, {_sql_literal(_MERGED_SUMMARY_DELIMITER)}) > 0"
|
||||
f" AND {_sql_starts_with(_PREVIEW_MERGED_AFTER_SQL, _PREVIEW_SUMMARY_PREFIXES)})")
|
||||
_PREVIEW_MERGED_PRIOR_SQL = _sql_trim_whitespace(
|
||||
f"SUBSTR(m.content, 1, INSTR(m.content, {_sql_literal(_MERGED_SUMMARY_DELIMITER)}) - 1)"
|
||||
)
|
||||
f"SUBSTR(m.content, 1, INSTR(m.content, {_sql_literal(_MERGED_SUMMARY_DELIMITER)}) - 1)")
|
||||
_PREVIEW_MERGED_PRIOR_LTRIMMED_SQL = _sql_ltrim_whitespace(_PREVIEW_MERGED_PRIOR_SQL)
|
||||
_PREVIEW_MERGED_PRIOR_UNWRAPPED_SQL = (
|
||||
f"CASE WHEN SUBSTR({_PREVIEW_MERGED_PRIOR_LTRIMMED_SQL}, 1,"
|
||||
_PREVIEW_MERGED_PRIOR_UNWRAPPED_SQL = (f"CASE WHEN SUBSTR({_PREVIEW_MERGED_PRIOR_LTRIMMED_SQL}, 1,"
|
||||
f" {len(_MERGED_PRIOR_CONTEXT_HEADER)}) = {_sql_literal(_MERGED_PRIOR_CONTEXT_HEADER)}"
|
||||
f" THEN {_sql_ltrim_whitespace(f'SUBSTR({_PREVIEW_MERGED_PRIOR_LTRIMMED_SQL}, {len(_MERGED_PRIOR_CONTEXT_HEADER) + 1})')}"
|
||||
f" ELSE {_PREVIEW_MERGED_PRIOR_SQL} END"
|
||||
)
|
||||
f" ELSE {_PREVIEW_MERGED_PRIOR_SQL} END")
|
||||
_PREVIEW_FORCE_USER_REMAINDER_SQL = _sql_after_marker(_SUMMARY_END_MARKER)
|
||||
|
||||
# Pure compaction rows are ineligible for previews; force-user-leading and merged
|
||||
# carriers are eligible only when authentic content survives.
|
||||
_PREVIEW_ELIGIBLE_SQL = (
|
||||
f"((NOT {_PREVIEW_STANDALONE_SUMMARY_SQL} AND NOT {_PREVIEW_MERGED_SUMMARY_SQL})"
|
||||
f" OR ({_PREVIEW_STANDALONE_SUMMARY_SQL}"
|
||||
f" AND INSTR(m.content, {_sql_literal(_SUMMARY_END_MARKER)}) > 0"
|
||||
_PREVIEW_ELIGIBLE_SQL = (f"((NOT {_PREVIEW_STANDALONE_SUMMARY_SQL} AND NOT {_PREVIEW_MERGED_SUMMARY_SQL})"
|
||||
f" OR ({_PREVIEW_STANDALONE_SUMMARY_SQL} AND INSTR(m.content, {_sql_literal(_SUMMARY_END_MARKER)}) > 0"
|
||||
f" AND LENGTH({_sql_trim_whitespace(_PREVIEW_FORCE_USER_REMAINDER_SQL)}) > 0)"
|
||||
f" OR ({_PREVIEW_MERGED_SUMMARY_SQL}"
|
||||
f" AND LENGTH({_sql_trim_whitespace(_PREVIEW_MERGED_PRIOR_UNWRAPPED_SQL)}) > 0))"
|
||||
)
|
||||
f" AND LENGTH({_sql_trim_whitespace(_PREVIEW_MERGED_PRIOR_UNWRAPPED_SQL)}) > 0))")
|
||||
|
||||
# Shared ``_preview_raw`` SELECT expression for every listing query (scaffolded rows:
|
||||
# head + tail spliced around SKILL_EXCERPT_JOINT when over budget).
|
||||
_PREVIEW_RAW_SELECT = (
|
||||
f"CASE WHEN {_PREVIEW_STANDALONE_SUMMARY_SQL}"
|
||||
f" THEN {_PREVIEW_FORCE_USER_REMAINDER_SQL}"
|
||||
f" WHEN {_PREVIEW_MERGED_SUMMARY_SQL}"
|
||||
f" THEN {_PREVIEW_MERGED_PRIOR_UNWRAPPED_SQL}"
|
||||
f" WHEN {_PREVIEW_SCAFFOLDED_SQL}"
|
||||
f" AND LENGTH(m.content) > {_PREVIEW_SCAFFOLD_WINDOW * 2}"
|
||||
f" THEN SUBSTR({_PREVIEW_CONTENT_SQL}, 1, {_PREVIEW_SCAFFOLD_WINDOW})"
|
||||
f" || '{SKILL_EXCERPT_JOINT}'"
|
||||
f"CASE WHEN {_PREVIEW_STANDALONE_SUMMARY_SQL} THEN {_PREVIEW_FORCE_USER_REMAINDER_SQL}"
|
||||
f" WHEN {_PREVIEW_MERGED_SUMMARY_SQL} THEN {_PREVIEW_MERGED_PRIOR_UNWRAPPED_SQL}"
|
||||
f" WHEN {_PREVIEW_SCAFFOLDED_SQL} AND LENGTH(m.content) > {_PREVIEW_SCAFFOLD_WINDOW * 2}"
|
||||
f" THEN SUBSTR({_PREVIEW_CONTENT_SQL}, 1, {_PREVIEW_SCAFFOLD_WINDOW}) || '{SKILL_EXCERPT_JOINT}'"
|
||||
f" || SUBSTR({_PREVIEW_CONTENT_SQL}, -{_PREVIEW_SCAFFOLD_WINDOW})"
|
||||
f" WHEN {_PREVIEW_SCAFFOLDED_SQL}"
|
||||
f" THEN SUBSTR({_PREVIEW_CONTENT_SQL}, 1, {_PREVIEW_SCAFFOLD_WINDOW * 2})"
|
||||
f" ELSE SUBSTR({_PREVIEW_CONTENT_SQL}, 1, {_PREVIEW_HEAD_CHARS}) END"
|
||||
)
|
||||
f" WHEN {_PREVIEW_SCAFFOLDED_SQL} THEN SUBSTR({_PREVIEW_CONTENT_SQL}, 1, {_PREVIEW_SCAFFOLD_WINDOW * 2})"
|
||||
f" ELSE SUBSTR({_PREVIEW_CONTENT_SQL}, 1, {_PREVIEW_HEAD_CHARS}) END")
|
||||
|
||||
|
||||
def _shape_preview(raw: Any) -> str:
|
||||
@@ -125,50 +110,31 @@ def _shape_preview(raw: Any) -> str:
|
||||
|
||||
|
||||
# Correlated ``_preview_raw`` column for a ``sessions s`` row.
|
||||
_PREVIEW_RAW_SUBQUERY_SQL = (
|
||||
f"COALESCE((SELECT {_PREVIEW_RAW_SELECT} FROM messages m"
|
||||
f" WHERE m.session_id = s.id AND m.role = 'user' AND m.content IS NOT NULL"
|
||||
f" AND {_PREVIEW_ELIGIBLE_SQL}"
|
||||
f" ORDER BY m.timestamp, m.id LIMIT 1), '') AS _preview_raw"
|
||||
)
|
||||
_PREVIEW_RAW_SUBQUERY_SQL = (f"COALESCE((SELECT {_PREVIEW_RAW_SELECT} FROM messages m"
|
||||
f" WHERE m.session_id = s.id AND m.role = 'user' AND m.content IS NOT NULL AND {_PREVIEW_ELIGIBLE_SQL}"
|
||||
f" ORDER BY m.timestamp, m.id LIMIT 1), '') AS _preview_raw")
|
||||
|
||||
# ── Session lineage predicates ({a} = sessions alias) ───────────────────────
|
||||
|
||||
# A /branch child (kept visible, never cascade-deleted): stable marker OR the legacy end_reason heuristic.
|
||||
_BRANCH_CHILD_SQL = (
|
||||
"json_extract(COALESCE({a}.model_config, '{{}}'), '$._branched_from') IS NOT NULL"
|
||||
_BRANCH_CHILD_SQL = ("json_extract(COALESCE({a}.model_config, '{{}}'), '$._branched_from') IS NOT NULL"
|
||||
" OR EXISTS (SELECT 1 FROM sessions p WHERE p.id = {a}.parent_session_id"
|
||||
" AND p.end_reason = 'branched' AND {a}.started_at >= p.ended_at)"
|
||||
)
|
||||
_COMPRESSION_CHILD_SQL = (
|
||||
"EXISTS (SELECT 1 FROM sessions p WHERE p.id = {a}.parent_session_id"
|
||||
" AND p.end_reason = 'compression')"
|
||||
)
|
||||
" AND p.end_reason = 'branched' AND {a}.started_at >= p.ended_at)")
|
||||
_COMPRESSION_CHILD_SQL = ("EXISTS (SELECT 1 FROM sessions p WHERE p.id = {a}.parent_session_id"
|
||||
" AND p.end_reason = 'compression')")
|
||||
|
||||
_RESET_END_REASONS = (
|
||||
"session_reset",
|
||||
# switch_session() creates no child row, but pre-marker DBs hold legacy reset
|
||||
# children whose parent later ended 'session_switch'. Must stay identical to the
|
||||
# recovery fence in find_latest_gateway_session_for_peer (interpolates the SQL form).
|
||||
"session_switch",
|
||||
"idle",
|
||||
"daily",
|
||||
"suspended",
|
||||
"resume_pending_expired",
|
||||
)
|
||||
# 'session_switch': switch_session() creates no child row, but pre-marker DBs hold legacy reset children
|
||||
# whose parent later ended that way. Must stay identical to the recovery fence in
|
||||
# find_latest_gateway_session_for_peer (interpolates the SQL form).
|
||||
_RESET_END_REASONS = ("session_reset", "session_switch", "idle", "daily", "suspended", "resume_pending_expired")
|
||||
_RESET_END_REASONS_SQL = ", ".join(f"'{reason}'" for reason in _RESET_END_REASONS)
|
||||
|
||||
# Accidental end reasons recovery treats as resumable (docs/session-lifecycle.md). Single source of truth:
|
||||
# interpolated into recovery SQL AND exposed as SessionDB.RECOVERABLE_END_REASONS.
|
||||
_RECOVERABLE_END_REASONS = (
|
||||
"agent_close",
|
||||
"ws_orphan_reap",
|
||||
# Stale sentinel-parked runtime superseded by a fresh session.resume.
|
||||
"superseded_by_resume",
|
||||
# Startup sweep of rows orphaned by a dead gateway process: same accident class as
|
||||
# ws_orphan_reap, kept distinct for forensics.
|
||||
"startup_orphan_reap",
|
||||
)
|
||||
# superseded_by_resume: stale sentinel-parked runtime superseded by a fresh session.resume.
|
||||
# startup_orphan_reap: startup sweep of rows orphaned by a dead gateway process (same accident class as
|
||||
# ws_orphan_reap, kept distinct for forensics).
|
||||
_RECOVERABLE_END_REASONS = ("agent_close", "ws_orphan_reap", "superseded_by_resume", "startup_orphan_reap")
|
||||
_RECOVERABLE_END_REASONS_SQL = ", ".join(f"'{reason}'" for reason in _RECOVERABLE_END_REASONS)
|
||||
|
||||
# End reasons written by AUTOMATIC cleanup (shutdown, orphan reapers, idle/LRU eviction), not by a
|
||||
@@ -189,39 +155,26 @@ def _legacy_reset_child_sql(alias: str, reasons_sql: str) -> str:
|
||||
"""Pre-marker reset-continuation heuristic: child rides its parent's exact non-empty routing key and the
|
||||
parent ended at a reset boundary. Shared by ``_RESET_CHILD_SQL`` and ``reopen_session()``'s
|
||||
marker-stamping UPDATE so the two cannot drift; ``reasons_sql`` is a literal or placeholder list."""
|
||||
return (
|
||||
f"EXISTS (SELECT 1 FROM sessions p"
|
||||
f" WHERE p.id = {alias}.parent_session_id"
|
||||
f" AND p.end_reason IN ({reasons_sql})"
|
||||
f" AND {alias}.session_key IS NOT NULL"
|
||||
f" AND {alias}.session_key != ''"
|
||||
f" AND {alias}.session_key = p.session_key)"
|
||||
)
|
||||
return (f"EXISTS (SELECT 1 FROM sessions p WHERE p.id = {alias}.parent_session_id"
|
||||
f" AND p.end_reason IN ({reasons_sql}) AND {alias}.session_key IS NOT NULL"
|
||||
f" AND {alias}.session_key != '' AND {alias}.session_key = p.session_key)")
|
||||
|
||||
|
||||
# A reset starts a separate user-visible conversation though rows keep parent_session_id
|
||||
# for lineage. Stable marker, or the same-key fallback for pre-marker rows (the
|
||||
# exact-key requirement keeps subagent children out).
|
||||
_RESET_CHILD_SQL = (
|
||||
"json_extract(COALESCE({a}.model_config, '{{}}'), '$._reset_from') IS NOT NULL"
|
||||
" OR " + _legacy_reset_child_sql("{a}", _RESET_END_REASONS_SQL)
|
||||
)
|
||||
_RESET_CHILD_SQL = ("json_extract(COALESCE({a}.model_config, '{{}}'), '$._reset_from') IS NOT NULL"
|
||||
" OR " + _legacy_reset_child_sql("{a}", _RESET_END_REASONS_SQL))
|
||||
|
||||
# Picker-visible rows: roots + branch/reset children (not subagent runs or compression continuations).
|
||||
_LISTABLE_CHILD_SQL = (
|
||||
f"(s.parent_session_id IS NULL OR {_BRANCH_CHILD_SQL.format(a='s')}"
|
||||
f" OR {_RESET_CHILD_SQL.format(a='s')})"
|
||||
)
|
||||
_LISTABLE_CHILD_SQL = (f"(s.parent_session_id IS NULL OR {_BRANCH_CHILD_SQL.format(a='s')}"
|
||||
f" OR {_RESET_CHILD_SQL.format(a='s')})")
|
||||
|
||||
|
||||
def _ephemeral_child_sql(alias: str = "s") -> str:
|
||||
"""Subagent runs, not branch, reset, or compression children."""
|
||||
return (
|
||||
f"({alias}.parent_session_id IS NOT NULL"
|
||||
f" AND NOT ({_BRANCH_CHILD_SQL.format(a=alias)})"
|
||||
f" AND NOT ({_COMPRESSION_CHILD_SQL.format(a=alias)})"
|
||||
f" AND NOT ({_RESET_CHILD_SQL.format(a=alias)}))"
|
||||
)
|
||||
return (f"({alias}.parent_session_id IS NOT NULL AND NOT ({_BRANCH_CHILD_SQL.format(a=alias)})"
|
||||
f" AND NOT ({_COMPRESSION_CHILD_SQL.format(a=alias)}) AND NOT ({_RESET_CHILD_SQL.format(a=alias)}))")
|
||||
|
||||
|
||||
def _sql_freshest_of(activity: str, session_id_expr: str, started: str) -> str:
|
||||
@@ -229,15 +182,8 @@ def _sql_freshest_of(activity: str, session_id_expr: str, started: str) -> str:
|
||||
else *started*. Heartbeats are rate-limited (~60s) so ``last_activity_at`` can lag
|
||||
a newer message; never prefer it alone."""
|
||||
msg_max = f"(SELECT MAX(_act_m.timestamp) FROM messages _act_m WHERE _act_m.session_id = {session_id_expr})"
|
||||
return (
|
||||
f"COALESCE("
|
||||
f"(SELECT MAX(_act_v.v) FROM ("
|
||||
f"SELECT {activity} AS v "
|
||||
f"UNION ALL "
|
||||
f"SELECT {msg_max}"
|
||||
f") _act_v), "
|
||||
f"{started})"
|
||||
)
|
||||
return (f"COALESCE((SELECT MAX(_act_v.v) FROM (SELECT {activity} AS v UNION ALL SELECT {msg_max}) _act_v), "
|
||||
f"{started})")
|
||||
|
||||
|
||||
def _sql_session_last_active(alias: str = "s") -> str:
|
||||
@@ -248,8 +194,7 @@ def _sql_session_last_active(alias: str = "s") -> str:
|
||||
def _sql_session_last_active_by_id(session_id_expr: str) -> str:
|
||||
"""Same freshest-of expression keyed by a session-id SQL expression."""
|
||||
return _sql_freshest_of(
|
||||
f"(SELECT last_activity_at FROM sessions _act_s WHERE _act_s.id = {session_id_expr})",
|
||||
session_id_expr,
|
||||
f"(SELECT last_activity_at FROM sessions _act_s WHERE _act_s.id = {session_id_expr})", session_id_expr,
|
||||
f"(SELECT started_at FROM sessions _act_s WHERE _act_s.id = {session_id_expr})")
|
||||
|
||||
|
||||
@@ -293,9 +238,8 @@ def _placeholders(items) -> str:
|
||||
return ",".join("?" for _ in range(items if isinstance(items, int) else len(items)))
|
||||
|
||||
|
||||
_FTS_TRIGGERS = (
|
||||
"messages_fts_insert", "messages_fts_delete", "messages_fts_update",
|
||||
"messages_fts_trigram_insert", "messages_fts_trigram_delete", "messages_fts_trigram_update")
|
||||
_FTS_TRIGGERS = ("messages_fts_insert", "messages_fts_delete", "messages_fts_update",
|
||||
"messages_fts_trigram_insert", "messages_fts_trigram_delete", "messages_fts_trigram_update")
|
||||
|
||||
SCHEMA_SQL = """
|
||||
CREATE TABLE IF NOT EXISTS schema_version (
|
||||
@@ -675,8 +619,7 @@ BEGIN
|
||||
END;
|
||||
"""
|
||||
|
||||
_FTS_CJK_TRIGGERS = (
|
||||
"messages_fts_cjk_insert", "messages_fts_cjk_delete", "messages_fts_cjk_update")
|
||||
_FTS_CJK_TRIGGERS = ("messages_fts_cjk_insert", "messages_fts_cjk_delete", "messages_fts_cjk_update")
|
||||
|
||||
# Set when a tokenizer-less process dropped the cjk triggers to keep writes alive: the cjk index is missing
|
||||
# rows and must not serve reads until `hermes sessions optimize-storage` rebuilds it on a capable host.
|
||||
@@ -775,9 +718,7 @@ _LOCK_BREAK_REACQUIRE_SECONDS = 5.0
|
||||
# "Another process holds the lock": flock → EWOULDBLOCK/EAGAIN, msvcrt.locking → EACCES (EDEADLK when its
|
||||
# retry gives up). Anything else (ESTALE, ENOTSUP, ENOLCK, EIO) is a persistent environment failure that
|
||||
# polling cannot fix; treating it as contention burned the full timeout on every attempt.
|
||||
_LOCK_CONTENTION_ERRNOS = {errno.EAGAIN, errno.EACCES, errno.EWOULDBLOCK}
|
||||
if hasattr(errno, "EDEADLK"):
|
||||
_LOCK_CONTENTION_ERRNOS.add(errno.EDEADLK)
|
||||
_LOCK_CONTENTION_ERRNOS = {errno.EAGAIN, errno.EACCES, errno.EWOULDBLOCK, errno.EDEADLK}
|
||||
|
||||
|
||||
def is_advisory_lock_contention(exc: BaseException) -> bool:
|
||||
@@ -811,14 +752,12 @@ def _read_lock_holder_record(handle):
|
||||
|
||||
def _rewrite_lock_file(handle, payload: bytes) -> None:
|
||||
"""Best-effort truncate-and-write of *payload* at offset 0."""
|
||||
try:
|
||||
with contextlib.suppress(OSError, ValueError):
|
||||
handle.seek(0)
|
||||
handle.truncate()
|
||||
if payload:
|
||||
handle.write(payload)
|
||||
handle.flush()
|
||||
except (OSError, ValueError):
|
||||
pass
|
||||
|
||||
|
||||
def _write_lock_holder_record(handle) -> None:
|
||||
@@ -871,7 +810,6 @@ def _acquire_db_flock(lock_path, handle, timeout_seconds, poll_seconds, descript
|
||||
nobody. Every successful acquire verifies its inode still names *lock_path*, so a racer that locked a
|
||||
dead inode retries instead of running alongside the breaker. Indeterminate liveness defers."""
|
||||
import fcntl
|
||||
|
||||
deadline = time.monotonic() + timeout_seconds
|
||||
broke_lock = False
|
||||
while True:
|
||||
@@ -880,10 +818,9 @@ def _acquire_db_flock(lock_path, handle, timeout_seconds, poll_seconds, descript
|
||||
except (BlockingIOError, OSError) as exc:
|
||||
if not is_advisory_lock_contention(exc):
|
||||
# Not a holder and polling cannot fix it: defer NOW.
|
||||
logger.warning(
|
||||
"Could not acquire %s %s (%s) — deferring rather than "
|
||||
"waiting out the %.0fs holder timeout on a non-contention error.",
|
||||
description, lock_path, exc, timeout_seconds)
|
||||
logger.warning("Could not acquire %s %s (%s) — deferring rather than "
|
||||
"waiting out the %.0fs holder timeout on a non-contention error.",
|
||||
description, lock_path, exc, timeout_seconds)
|
||||
return None, handle
|
||||
if time.monotonic() < deadline:
|
||||
time.sleep(poll_seconds)
|
||||
@@ -893,11 +830,9 @@ def _acquire_db_flock(lock_path, handle, timeout_seconds, poll_seconds, descript
|
||||
record = _read_lock_holder_record(handle)
|
||||
if not _lock_holder_provably_dead(record):
|
||||
return False, handle
|
||||
logger.warning(
|
||||
"%s %s is held by an orphaned file descriptor (recorded "
|
||||
"holder pid %s is dead — a forked child inherited the lock "
|
||||
"fd); breaking the stale lock and retaking it on a fresh file.",
|
||||
description, lock_path, (record or {}).get("pid"))
|
||||
logger.warning("%s %s is held by an orphaned file descriptor (recorded holder pid %s is dead — a "
|
||||
"forked child inherited the lock fd); breaking the stale lock and retaking it on a "
|
||||
"fresh file.", description, lock_path, (record or {}).get("pid"))
|
||||
try:
|
||||
os.unlink(lock_path)
|
||||
handle.close()
|
||||
@@ -911,8 +846,7 @@ def _acquire_db_flock(lock_path, handle, timeout_seconds, poll_seconds, descript
|
||||
# Verify the path still names our inode: a breaker may have replaced
|
||||
# the file while we waited, and a lock on a dead inode excludes nobody.
|
||||
try:
|
||||
fd_stat = os.fstat(handle.fileno())
|
||||
path_stat = os.stat(lock_path)
|
||||
fd_stat, path_stat = os.fstat(handle.fileno()), os.stat(lock_path)
|
||||
same_file = fd_stat.st_dev == path_stat.st_dev and fd_stat.st_ino == path_stat.st_ino
|
||||
except OSError:
|
||||
same_file = False
|
||||
@@ -933,18 +867,15 @@ def _describe_lock_holder(record) -> str:
|
||||
if not isinstance(record, dict) or "pid" not in record:
|
||||
return "unknown (no holder record; pre-fix writer or non-Hermes)"
|
||||
age = ""
|
||||
try:
|
||||
with contextlib.suppress(TypeError, ValueError):
|
||||
if record.get("acquired_at") is not None:
|
||||
age = f", acquired {time.time() - float(record['acquired_at']):.0f}s ago"
|
||||
except (TypeError, ValueError):
|
||||
pass
|
||||
return f"pid {record.get('pid')}{age}"
|
||||
|
||||
|
||||
def _acquire_msvcrt_lock(lock_path, handle, timeout):
|
||||
"""Windows counterpart of ``_acquire_db_flock`` (no orphan break); same True / False / None contract."""
|
||||
import msvcrt
|
||||
|
||||
deadline = time.monotonic() + timeout
|
||||
while True:
|
||||
try:
|
||||
@@ -953,9 +884,8 @@ def _acquire_msvcrt_lock(lock_path, handle, timeout):
|
||||
return True
|
||||
except (BlockingIOError, OSError) as exc:
|
||||
if not is_advisory_lock_contention(exc):
|
||||
logger.warning(
|
||||
"Could not acquire FTS rebuild lock %s (%s) — deferring on a non-contention error.",
|
||||
lock_path, exc)
|
||||
logger.warning("Could not acquire FTS rebuild lock %s (%s) — deferring on a non-contention error.",
|
||||
lock_path, exc)
|
||||
return None
|
||||
if time.monotonic() >= deadline:
|
||||
return False
|
||||
@@ -981,10 +911,8 @@ def fts_rebuild_admission(db_path, *, timeout_seconds=None):
|
||||
# may still be rebuilding — yielding True gave every process on a full disk a
|
||||
# concurrent rebuild of the same DB. Deferring costs nothing (the breadcrumb
|
||||
# retries, and the rebuild's own writes could not have committed either).
|
||||
logger.warning(
|
||||
"Could not open FTS rebuild lock %s (%s) — deferring this rebuild "
|
||||
"rather than running it without cross-process authority.",
|
||||
lock_path, exc)
|
||||
logger.warning("Could not open FTS rebuild lock %s (%s) — deferring this rebuild "
|
||||
"rather than running it without cross-process authority.", lock_path, exc)
|
||||
yield False
|
||||
return
|
||||
acquired = False
|
||||
@@ -1001,28 +929,23 @@ def fts_rebuild_admission(db_path, *, timeout_seconds=None):
|
||||
record = None if _IS_WINDOWS else _read_lock_holder_record(handle)
|
||||
if timeout <= 0:
|
||||
# Non-blocking probe from an in-process retry: keep it quiet.
|
||||
logger.info(
|
||||
"FTS rebuild lock %s is busy — deferring this retry "
|
||||
"(the stale-FTS breadcrumb keeps it retryable). Recorded holder: %s.",
|
||||
lock_path, _describe_lock_holder(record))
|
||||
logger.info("FTS rebuild lock %s is busy — deferring this retry "
|
||||
"(the stale-FTS breadcrumb keeps it retryable). Recorded holder: %s.",
|
||||
lock_path, _describe_lock_holder(record))
|
||||
else:
|
||||
logger.warning(
|
||||
"FTS rebuild lock %s held by another process for more than "
|
||||
"%.0fs — deferring this rebuild to avoid racing the holder "
|
||||
"(the stale-FTS breadcrumb keeps it retryable). Recorded holder: %s.",
|
||||
lock_path, timeout, _describe_lock_holder(record))
|
||||
logger.warning("FTS rebuild lock %s held by another process for more than %.0fs — deferring "
|
||||
"this rebuild to avoid racing the holder (the stale-FTS breadcrumb keeps it "
|
||||
"retryable). Recorded holder: %s.", lock_path, timeout, _describe_lock_holder(record))
|
||||
yield acquired
|
||||
finally:
|
||||
try:
|
||||
if acquired:
|
||||
if _IS_WINDOWS:
|
||||
import msvcrt
|
||||
|
||||
handle.seek(0)
|
||||
msvcrt.locking(handle.fileno(), msvcrt.LK_UNLCK, 1)
|
||||
else:
|
||||
import fcntl
|
||||
|
||||
_clear_lock_holder_record(handle)
|
||||
fcntl.flock(handle.fileno(), fcntl.LOCK_UN)
|
||||
except OSError: # pragma: no cover - best effort release
|
||||
|
||||
Reference in New Issue
Block a user