refactor(state): AST-neutral line packing; derive v22 session_model_usage DDL from the heal DDL
This commit is contained in:
@@ -309,12 +309,8 @@ def stat_db_file_identity(path) -> "tuple[int, int] | None":
|
||||
|
||||
|
||||
_FTS_TRIGGERS = (
|
||||
"messages_fts_insert",
|
||||
"messages_fts_delete",
|
||||
"messages_fts_update",
|
||||
"messages_fts_trigram_insert",
|
||||
"messages_fts_trigram_delete",
|
||||
"messages_fts_trigram_update",
|
||||
"messages_fts_insert", "messages_fts_delete", "messages_fts_update",
|
||||
"messages_fts_trigram_insert", "messages_fts_trigram_delete", "messages_fts_trigram_update",
|
||||
)
|
||||
|
||||
|
||||
@@ -706,9 +702,7 @@ END;
|
||||
|
||||
|
||||
_FTS_CJK_TRIGGERS = (
|
||||
"messages_fts_cjk_insert",
|
||||
"messages_fts_cjk_delete",
|
||||
"messages_fts_cjk_update",
|
||||
"messages_fts_cjk_insert", "messages_fts_cjk_delete", "messages_fts_cjk_update",
|
||||
)
|
||||
|
||||
|
||||
|
||||
@@ -224,10 +224,8 @@ class SessionPortabilityMixin:
|
||||
return None
|
||||
messages = [msg for seg in segments for msg in (seg.get("messages") or [])]
|
||||
return {
|
||||
**segments[-1],
|
||||
"segments": segments,
|
||||
"lineage_session_ids": [seg["id"] for seg in segments],
|
||||
"message_count": len(messages),
|
||||
**segments[-1], "segments": segments,
|
||||
"lineage_session_ids": [seg["id"] for seg in segments], "message_count": len(messages),
|
||||
"messages": messages,
|
||||
}
|
||||
|
||||
@@ -236,7 +234,7 @@ class SessionPortabilityMixin:
|
||||
return [self._with_messages(s) for s in self.search_sessions(source=source, limit=100000)]
|
||||
|
||||
def adopt_session_lineage_from(
|
||||
self, donor_db: Any, session_id: str, *, retire_donor: bool = True,
|
||||
self, donor_db: Any, session_id: str, *, retire_donor: bool = True
|
||||
) -> Dict[str, Any]:
|
||||
"""Adopt *session_id*'s full compression lineage from *donor_db* (a
|
||||
full SessionDB — the mixin cannot import it).
|
||||
@@ -486,8 +484,7 @@ class SessionPortabilityMixin:
|
||||
"""INSERT one normalized session + its messages; counts fixed up after."""
|
||||
started_at = self._float_or_none(raw.get("started_at"))
|
||||
params = {
|
||||
"id": session_id,
|
||||
"source": str(raw.get("source") or "import"),
|
||||
"id": session_id, "source": str(raw.get("source") or "import"),
|
||||
"system_prompt_hash": self._store_system_prompt(conn, raw.get("system_prompt")),
|
||||
"started_at": time.time() if started_at is None else started_at,
|
||||
"archived": 1 if raw.get("archived") else 0,
|
||||
@@ -581,12 +578,8 @@ class SessionPortabilityMixin:
|
||||
imported_ids.append(session_id)
|
||||
detached = self._attach_import_parents(conn, parent_updates)
|
||||
return {
|
||||
"ok": True,
|
||||
"imported": len(imported_ids),
|
||||
"skipped": len(skipped_ids),
|
||||
"detached": detached,
|
||||
"imported_ids": imported_ids,
|
||||
"skipped_ids": skipped_ids,
|
||||
"ok": True, "imported": len(imported_ids), "skipped": len(skipped_ids),
|
||||
"detached": detached, "imported_ids": imported_ids, "skipped_ids": skipped_ids,
|
||||
"errors": [],
|
||||
}
|
||||
|
||||
|
||||
+9
-32
@@ -61,11 +61,6 @@ _SESSION_MODEL_USAGE_INDEX_SQL = (
|
||||
"CREATE INDEX IF NOT EXISTS idx_session_model_usage_session ON session_model_usage(session_id)",
|
||||
"CREATE INDEX IF NOT EXISTS idx_session_model_usage_model ON session_model_usage(model)",
|
||||
)
|
||||
_SESSION_MODEL_USAGE_COLS = """session_id, model, billing_provider, billing_base_url,
|
||||
billing_mode, task, api_call_count, input_tokens,
|
||||
output_tokens, cache_read_tokens, cache_write_tokens,
|
||||
reasoning_tokens, estimated_cost_usd, actual_cost_usd,
|
||||
cost_status, cost_source, first_seen, last_seen"""
|
||||
_SESSION_MODEL_USAGE_HEAL_DDL = """CREATE TABLE session_model_usage (
|
||||
session_id TEXT NOT NULL REFERENCES sessions(id) ON DELETE CASCADE,
|
||||
model TEXT NOT NULL,
|
||||
@@ -87,28 +82,13 @@ _SESSION_MODEL_USAGE_HEAL_DDL = """CREATE TABLE session_model_usage (
|
||||
last_seen REAL,
|
||||
PRIMARY KEY (session_id, model, billing_provider, billing_base_url, billing_mode, task)
|
||||
)"""
|
||||
# Same table, v22-migration indentation (statement text is pinned by the SQL trace harness).
|
||||
_SESSION_MODEL_USAGE_V22_DDL = """CREATE TABLE session_model_usage (
|
||||
session_id TEXT NOT NULL REFERENCES sessions(id) ON DELETE CASCADE,
|
||||
model TEXT NOT NULL,
|
||||
billing_provider TEXT NOT NULL DEFAULT '',
|
||||
billing_base_url TEXT NOT NULL DEFAULT '',
|
||||
billing_mode TEXT NOT NULL DEFAULT '',
|
||||
task TEXT NOT NULL DEFAULT '',
|
||||
api_call_count INTEGER NOT NULL DEFAULT 0,
|
||||
input_tokens INTEGER NOT NULL DEFAULT 0,
|
||||
output_tokens INTEGER NOT NULL DEFAULT 0,
|
||||
cache_read_tokens INTEGER NOT NULL DEFAULT 0,
|
||||
cache_write_tokens INTEGER NOT NULL DEFAULT 0,
|
||||
reasoning_tokens INTEGER NOT NULL DEFAULT 0,
|
||||
estimated_cost_usd REAL NOT NULL DEFAULT 0,
|
||||
actual_cost_usd REAL NOT NULL DEFAULT 0,
|
||||
cost_status TEXT,
|
||||
cost_source TEXT,
|
||||
first_seen REAL,
|
||||
last_seen REAL,
|
||||
PRIMARY KEY (session_id, model, billing_provider, billing_base_url, billing_mode, task)
|
||||
)"""
|
||||
# Same table as emitted by the v22 migration: column lines at 35 spaces, closing
|
||||
# paren at 31 (statement text is pinned by the SQL trace harness).
|
||||
_SESSION_MODEL_USAGE_V22_DDL = "\n".join(
|
||||
[_SESSION_MODEL_USAGE_HEAL_DDL.splitlines()[0]]
|
||||
+ [" " * 35 + ln.strip() for ln in _SESSION_MODEL_USAGE_HEAL_DDL.splitlines()[1:-1]]
|
||||
+ [" " * 31 + ")"]
|
||||
)
|
||||
# Statement text pinned by the SQL trace harness (whitespace included).
|
||||
_SESSION_MODEL_USAGE_V20_SEED_SQL = """INSERT OR IGNORE INTO session_model_usage (
|
||||
session_id, model, billing_provider,
|
||||
@@ -309,8 +289,7 @@ class SessionSchemaMixin:
|
||||
"UPDATE OF migration; marked stale and unavailable"
|
||||
)
|
||||
logger.info(
|
||||
"Migrated %d broad FTS UPDATE trigger(s) to AFTER UPDATE OF "
|
||||
"(no rebuild required)",
|
||||
"Migrated %d broad FTS UPDATE trigger(s) to AFTER UPDATE OF " "(no rebuild required)",
|
||||
len(to_drop),
|
||||
)
|
||||
return len(to_drop)
|
||||
@@ -429,9 +408,7 @@ class SessionSchemaMixin:
|
||||
if first_seen > now or first_seen < 0:
|
||||
first_seen = now
|
||||
diagnostic = {
|
||||
"first_seen": first_seen,
|
||||
"last_seen": now,
|
||||
"attempts": attempts,
|
||||
"first_seen": first_seen, "last_seen": now, "attempts": attempts,
|
||||
"holder_pids": sorted({pid for pid, _path in foreign_holders if pid > 0}),
|
||||
}
|
||||
cursor.execute(
|
||||
|
||||
@@ -532,7 +532,7 @@ class SessionSearchMixin:
|
||||
self._conn.commit()
|
||||
|
||||
def optimize_fts_storage(
|
||||
self, *, progress_cb: Optional[Callable[[Dict[str, Any]], None]] = None, vacuum: bool = True,
|
||||
self, *, progress_cb: Optional[Callable[[Dict[str, Any]], None]] = None, vacuum: bool = True
|
||||
) -> Dict[str, Any]:
|
||||
"""Migrate a legacy v22 inline-FTS DB to the v23 external-content
|
||||
schema, foreground and to completion; re-running resumes an
|
||||
@@ -722,15 +722,14 @@ class SessionSearchMixin:
|
||||
return msg
|
||||
|
||||
return {
|
||||
"window": filtered_window,
|
||||
"messages_before": primitive["messages_before"],
|
||||
"window": filtered_window, "messages_before": primitive["messages_before"],
|
||||
"messages_after": primitive["messages_after"],
|
||||
"bookend_start": [_hydrate(r) for r in bookend_start_rows],
|
||||
"bookend_end": [_hydrate(r) for r in bookend_end_rows],
|
||||
}
|
||||
|
||||
def list_recent_user_messages(
|
||||
self, session_id: str, limit: int = 20, include_inactive: bool = False,
|
||||
self, session_id: str, limit: int = 20, include_inactive: bool = False
|
||||
) -> List[Dict[str, Any]]:
|
||||
"""The *limit* most-recent real user turns, newest first, as
|
||||
``{id, timestamp, preview}`` (preview = first 80 chars, whitespace
|
||||
@@ -1020,7 +1019,7 @@ class SessionSearchMixin:
|
||||
return " OR ".join(compiled_groups), params, snippet_term
|
||||
|
||||
def _search_messages_like_fallback(
|
||||
self, query: str, *, limit: int, offset: int, sort: Optional[str], **filters,
|
||||
self, query: str, *, limit: int, offset: int, sort: Optional[str], **filters
|
||||
) -> List[Dict[str, Any]]:
|
||||
"""Search canonical messages while derived FTS state is stale."""
|
||||
predicate, params, snippet_term = self._compile_like_boolean_query(query)
|
||||
@@ -1049,7 +1048,7 @@ class SessionSearchMixin:
|
||||
self._fts_cjk_available = False
|
||||
|
||||
def _finalize_search_matches(
|
||||
self, matches: List[Dict[str, Any]], result_fields: Optional[Collection[str]] = None,
|
||||
self, matches: List[Dict[str, Any]], result_fields: Optional[Collection[str]] = None
|
||||
) -> List[Dict[str, Any]]:
|
||||
"""Attach neighboring messages (1 before + after, only when the
|
||||
projection consumes ``context``) and trim full content. Each context
|
||||
|
||||
Reference in New Issue
Block a user