Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 4 additions & 3 deletions agent/verification_evidence.py
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@
import re
import shlex
import sqlite3
from contextlib import closing
import tempfile
import threading
from dataclasses import dataclass
Expand Down Expand Up @@ -452,7 +453,7 @@ def record_terminal_result(

created_at = _utc_now()
with _DB_LOCK:
with _connect() as conn:
with closing(_connect()) as conn:
cur = conn.execute(
"""
INSERT INTO verification_events(
Expand Down Expand Up @@ -518,7 +519,7 @@ def mark_workspace_edited(
edited_at = _utc_now()

with _DB_LOCK:
with _connect() as conn:
with closing(_connect()) as conn:
row = conn.execute(
"""
SELECT changed_paths_json FROM verification_state
Expand Down Expand Up @@ -568,7 +569,7 @@ def verification_status(
sid = str(session_id or "default")
root = str(facts.get("root") or Path(cwd or ".").resolve())
with _DB_LOCK:
with _connect() as conn:
with closing(_connect()) as conn:
state = conn.execute(
"""
SELECT last_event_id, last_edit_at, changed_paths_json
Expand Down
11 changes: 6 additions & 5 deletions gateway/delivery_ledger.py
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,7 @@
from __future__ import annotations

import hashlib
from contextlib import closing
import json
import logging
import os
Expand Down Expand Up @@ -164,7 +165,7 @@ def record_obligation(
"""Record a final response as owed to the platform (state='pending')."""
now = time.time()
pid, started = _owner_stamp()
with _DB_LOCK, _connect() as conn:
with _DB_LOCK, closing(_connect()) as conn:
conn.execute(
"""INSERT OR REPLACE INTO delivery_obligations
(obligation_id, session_key, platform, chat_id, thread_id,
Expand All @@ -191,7 +192,7 @@ def mark_failed(obligation_id: str, error: str = "") -> None:


def _update_state(obligation_id: str, state: str, error: str = "") -> None:
with _DB_LOCK, _connect() as conn:
with _DB_LOCK, closing(_connect()) as conn:
conn.execute(
"""UPDATE delivery_obligations
SET state=?, updated_at=?, last_error=?
Expand Down Expand Up @@ -224,7 +225,7 @@ def sweep_recoverable(
now = now if now is not None else time.time()
pid, started = _owner_stamp()
claimed: List[Dict[str, Any]] = []
with _DB_LOCK, _connect() as conn:
with _DB_LOCK, closing(_connect()) as conn:
rows = conn.execute(
"""SELECT obligation_id, session_key, platform, chat_id, thread_id,
content, state, attempts, created_at,
Expand Down Expand Up @@ -277,7 +278,7 @@ def _prune(now: Optional[float] = None) -> None:
now = now if now is not None else time.time()
cutoff = now - _RETENTION_SECONDS
try:
with _connect() as conn:
with closing(_connect()) as conn:
conn.execute(
"""DELETE FROM delivery_obligations
WHERE state IN ('delivered', 'abandoned') AND updated_at < ?""",
Expand Down Expand Up @@ -321,7 +322,7 @@ def ledger_enabled(config: Optional[Dict[str, Any]] = None) -> bool:

def debug_rows(limit: int = 20) -> str:
"""Human-readable dump for ad-hoc inspection (sqlite3-free path)."""
with _DB_LOCK, _connect() as conn:
with _DB_LOCK, closing(_connect()) as conn:
rows = conn.execute(
"""SELECT obligation_id, session_key, state, attempts,
created_at, updated_at, last_error
Expand Down
27 changes: 14 additions & 13 deletions tools/async_delegation.py
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,7 @@
import json
import logging
import sqlite3
from contextlib import closing
import threading
import time
import uuid
Expand Down Expand Up @@ -141,7 +142,7 @@ def _persist_dispatch(record: Dict[str, Any]) -> None:
for key in ("goal", "goals", "context", "toolsets", "role", "model", "is_batch")
if key in record
}
with _DB_LOCK, _connect() as conn:
with _DB_LOCK, closing(_connect()) as conn:
conn.execute(
"""INSERT OR REPLACE INTO async_delegations
(delegation_id, origin_session, origin_ui_session_id,
Expand All @@ -158,15 +159,15 @@ def _persist_dispatch(record: Dict[str, Any]) -> None:


def _delete_durable_delegation(delegation_id: str) -> None:
with _DB_LOCK, _connect() as conn:
with _DB_LOCK, closing(_connect()) as conn:
conn.execute("DELETE FROM async_delegations WHERE delegation_id=?", (delegation_id,))


def _prune_durable_records() -> None:
"""Bound terminal history, preferring delivered records for deletion."""
now = time.time()
cutoff = now - _DURABLE_RETENTION_SECONDS
with _DB_LOCK, _connect() as conn:
with _DB_LOCK, closing(_connect()) as conn:
conn.execute(
"DELETE FROM async_delegations WHERE delivery_state='delivered' AND updated_at < ?",
(cutoff,),
Expand Down Expand Up @@ -203,7 +204,7 @@ def _prune_durable_records() -> None:

def _persist_completion(event: Dict[str, Any], result: Dict[str, Any]) -> None:
now = time.time()
with _DB_LOCK, _connect() as conn:
with _DB_LOCK, closing(_connect()) as conn:
conn.execute(
"""UPDATE async_delegations SET state=?, completed_at=?, updated_at=?,
event_json=?, result_json=?, delivery_state='pending'
Expand All @@ -214,7 +215,7 @@ def _persist_completion(event: Dict[str, Any], result: Dict[str, Any]) -> None:


def _note_delivery_attempt(delegation_id: str) -> None:
with _DB_LOCK, _connect() as conn:
with _DB_LOCK, closing(_connect()) as conn:
conn.execute(
"UPDATE async_delegations SET delivery_attempts=delivery_attempts+1, updated_at=? WHERE delegation_id=?",
(time.time(), delegation_id),
Expand All @@ -229,7 +230,7 @@ def recover_abandoned_delegations() -> int:
return 0
now = time.time()
recovered = 0
with _DB_LOCK, _connect() as conn:
with _DB_LOCK, closing(_connect()) as conn:
rows = conn.execute(
"""SELECT delegation_id, origin_session, origin_ui_session_id,
parent_session_id, dispatched_at, owner_pid,
Expand Down Expand Up @@ -281,7 +282,7 @@ def restore_undelivered_completions(target_queue) -> int:
results seconds after boot (#64484).
"""
recover_abandoned_delegations()
with _DB_LOCK, _connect() as conn:
with _DB_LOCK, closing(_connect()) as conn:
rows = conn.execute(
"""SELECT delegation_id, event_json FROM async_delegations
WHERE state != 'running' AND delivery_state='pending' AND event_json IS NOT NULL
Expand All @@ -298,7 +299,7 @@ def restore_undelivered_completions(target_queue) -> int:
def mark_completion_delivered(delegation_id: str) -> bool:
"""Atomically acknowledge successful injection of a durable completion."""
now = time.time()
with _DB_LOCK, _connect() as conn:
with _DB_LOCK, closing(_connect()) as conn:
cur = conn.execute(
"""UPDATE async_delegations SET delivery_state='delivered', delivered_at=?, updated_at=?
WHERE delegation_id=? AND delivery_state!='delivered'""",
Expand All @@ -310,7 +311,7 @@ def mark_completion_delivered(delegation_id: str) -> bool:
def claim_completion_delivery(delegation_id: str, claim_id: str) -> bool:
"""Claim one pending completion across competing consumers/processes."""
now = time.time()
with _DB_LOCK, _connect() as conn:
with _DB_LOCK, closing(_connect()) as conn:
row = conn.execute(
"SELECT delivery_state FROM async_delegations WHERE delegation_id=?",
(delegation_id,),
Expand Down Expand Up @@ -349,7 +350,7 @@ def release_completion_delivery(delegation_id: str, claim_id: str) -> bool:
pending rows).
"""
now = time.time()
with _DB_LOCK, _connect() as conn:
with _DB_LOCK, closing(_connect()) as conn:
capped = conn.execute(
"""UPDATE async_delegations SET delivery_state='dropped',
delivery_claim=NULL, delivery_claimed_at=NULL, updated_at=?
Expand Down Expand Up @@ -384,7 +385,7 @@ def drop_completion_delivery(delegation_id: str, claim_id: str) -> bool:
completion that will be fail-closed dropped again every time.
"""
now = time.time()
with _DB_LOCK, _connect() as conn:
with _DB_LOCK, closing(_connect()) as conn:
cur = conn.execute(
"""UPDATE async_delegations SET delivery_state='dropped',
updated_at=?, delivery_claim=NULL,
Expand All @@ -399,7 +400,7 @@ def drop_completion_delivery(delegation_id: str, claim_id: str) -> bool:
def complete_completion_delivery(delegation_id: str, claim_id: str) -> bool:
"""Acknowledge acceptance for the consumer holding this claim."""
now = time.time()
with _DB_LOCK, _connect() as conn:
with _DB_LOCK, closing(_connect()) as conn:
cur = conn.execute(
"""UPDATE async_delegations SET delivery_state='delivered',
delivered_at=?, updated_at=?, delivery_claim=NULL,
Expand All @@ -422,7 +423,7 @@ def release_event_delivery(evt: Dict[str, Any], claim_id: str) -> None:


def get_durable_delegation(delegation_id: str) -> Optional[Dict[str, Any]]:
with _DB_LOCK, _connect() as conn:
with _DB_LOCK, closing(_connect()) as conn:
row = conn.execute(
"""SELECT origin_session, state, dispatched_at, completed_at,
result_json, delivery_state, delivery_attempts
Expand Down