Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
0cf551cb28 |
+26
-20
@@ -4,6 +4,8 @@ JSONL persistence layer for the database.
|
|||||||
|
|
||||||
import copy
|
import copy
|
||||||
import logging
|
import logging
|
||||||
|
import os
|
||||||
|
import signal
|
||||||
from collections import deque
|
from collections import deque
|
||||||
from contextlib import contextmanager
|
from contextlib import contextmanager
|
||||||
from datetime import UTC, datetime
|
from datetime import UTC, datetime
|
||||||
@@ -69,22 +71,25 @@ def create_change_record(
|
|||||||
# Actions that are allowed to create a new database file
|
# Actions that are allowed to create a new database file
|
||||||
_BOOTSTRAP_ACTIONS = frozenset({"bootstrap"})
|
_BOOTSTRAP_ACTIONS = frozenset({"bootstrap"})
|
||||||
|
|
||||||
|
# Flag to prevent duplicate error messages on fatal flush failure
|
||||||
|
_flush_failed = False
|
||||||
|
|
||||||
|
|
||||||
async def flush_changes(
|
async def flush_changes(
|
||||||
db_path: Path,
|
db_path: Path,
|
||||||
pending_changes: deque[_ChangeRecord],
|
pending_changes: deque[_ChangeRecord],
|
||||||
) -> bool:
|
) -> None:
|
||||||
"""Write all pending changes to disk.
|
"""Write all pending changes to disk.
|
||||||
|
|
||||||
Args:
|
Args:
|
||||||
db_path: Path to the JSONL database file
|
db_path: Path to the JSONL database file
|
||||||
pending_changes: Queue of pending change records (will be cleared on success)
|
pending_changes: Queue of pending change records (will be cleared on success)
|
||||||
|
|
||||||
Returns:
|
On failure, logs an error and sends SIGTERM to trigger graceful shutdown.
|
||||||
True if flush succeeded, False otherwise
|
|
||||||
"""
|
"""
|
||||||
if not pending_changes:
|
global _flush_failed
|
||||||
return True
|
if _flush_failed or not pending_changes:
|
||||||
|
return
|
||||||
|
|
||||||
if not db_path.exists():
|
if not db_path.exists():
|
||||||
first_action = pending_changes[0].a
|
first_action = pending_changes[0].a
|
||||||
@@ -94,26 +99,25 @@ async def flush_changes(
|
|||||||
"only bootstrap can create a new database",
|
"only bootstrap can create a new database",
|
||||||
first_action,
|
first_action,
|
||||||
)
|
)
|
||||||
pending_changes.clear()
|
_flush_failed = True
|
||||||
return False
|
os.kill(os.getpid(), signal.SIGTERM)
|
||||||
|
return
|
||||||
|
|
||||||
changes_to_write = list(pending_changes)
|
changes_to_write = list(pending_changes)
|
||||||
pending_changes.clear()
|
|
||||||
|
|
||||||
try:
|
try:
|
||||||
lines = [_change_encoder.encode(change) for change in changes_to_write]
|
lines = [_change_encoder.encode(change) for change in changes_to_write]
|
||||||
if not lines:
|
if not lines:
|
||||||
return True
|
pending_changes.clear()
|
||||||
|
return
|
||||||
|
|
||||||
async with aiofiles.open(db_path, "ab") as f:
|
async with aiofiles.open(db_path, "ab") as f:
|
||||||
await f.write(b"\n".join(lines) + b"\n")
|
await f.write(b"\n".join(lines) + b"\n")
|
||||||
return True
|
pending_changes.clear()
|
||||||
except OSError:
|
except OSError as e:
|
||||||
_logger.exception("Failed to flush database changes")
|
_logger.error("Failed to flush database: %s", e)
|
||||||
# Re-queue the changes on failure
|
_flush_failed = True
|
||||||
for change in reversed(changes_to_write):
|
os.kill(os.getpid(), signal.SIGTERM)
|
||||||
pending_changes.appendleft(change)
|
|
||||||
return False
|
|
||||||
|
|
||||||
|
|
||||||
class JsonlStore:
|
class JsonlStore:
|
||||||
@@ -155,8 +159,10 @@ class JsonlStore:
|
|||||||
self._current_version = change.get("v", 0)
|
self._current_version = change.get("v", 0)
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
raise ValueError(f"Error parsing line {line_num}: {e}")
|
raise ValueError(f"Error parsing line {line_num}: {e}")
|
||||||
except (OSError, ValueError, msgspec.DecodeError) as e:
|
except OSError as e:
|
||||||
raise ValueError(f"Failed to load database: {e}")
|
raise SystemExit(f"Failed to load database: {e}")
|
||||||
|
except (ValueError, msgspec.DecodeError) as e:
|
||||||
|
raise SystemExit(f"Failed to load database: {e}")
|
||||||
|
|
||||||
if not data_dict:
|
if not data_dict:
|
||||||
return
|
return
|
||||||
@@ -282,6 +288,6 @@ class JsonlStore:
|
|||||||
self._in_transaction = False
|
self._in_transaction = False
|
||||||
self._transaction_snapshot = None
|
self._transaction_snapshot = None
|
||||||
|
|
||||||
async def flush(self) -> bool:
|
async def flush(self) -> None:
|
||||||
"""Write all pending changes to disk."""
|
"""Write all pending changes to disk."""
|
||||||
return await flush_changes(self.db_path, self._pending_changes)
|
await flush_changes(self.db_path, self._pending_changes)
|
||||||
|
|||||||
@@ -836,5 +836,5 @@ def get_config() -> Config:
|
|||||||
|
|
||||||
async def set_config(config: Config) -> None:
|
async def set_config(config: Config) -> None:
|
||||||
"""Update the stored configuration."""
|
"""Update the stored configuration."""
|
||||||
async with _db.transaction("update_config"):
|
with _db.transaction("update_config"):
|
||||||
_db.config = config
|
_db.config = config
|
||||||
|
|||||||
Reference in New Issue
Block a user