generated from kgod/ai-review-template
feat: improve resume optimization and import reliability
This commit is contained in:
@@ -7,6 +7,7 @@ from uuid import uuid4
|
||||
|
||||
from sqlalchemy import Connection, Engine, create_engine, delete, func, insert, select, update
|
||||
|
||||
from .database import SessionRevisionConflict
|
||||
from .db.schema import build_session_tables
|
||||
from .models import BusinessResume, ComponentBlock, ConversationTurn, SessionView
|
||||
from .resume_document_core import attach_gap_report_staleness
|
||||
@@ -19,15 +20,46 @@ def _now() -> datetime:
|
||||
class PostgresDatabase:
|
||||
"""PostgreSQL implementation of the Resume Agent persistence contract."""
|
||||
|
||||
def __init__(self, database_url: str, *, schema: str = "resume_agent") -> None:
|
||||
self.engine: Engine = create_engine(database_url, pool_pre_ping=True)
|
||||
def __init__(
|
||||
self,
|
||||
database_url: str,
|
||||
*,
|
||||
schema: str = "resume_agent",
|
||||
pool_size: int = 10,
|
||||
max_overflow: int = 10,
|
||||
pool_timeout_seconds: float = 5.0,
|
||||
statement_timeout_ms: int = 10_000,
|
||||
lock_timeout_ms: int = 3_000,
|
||||
idle_transaction_timeout_ms: int = 15_000,
|
||||
) -> None:
|
||||
self.engine: Engine = create_engine(
|
||||
database_url,
|
||||
pool_pre_ping=True,
|
||||
pool_size=pool_size,
|
||||
max_overflow=max_overflow,
|
||||
pool_timeout=pool_timeout_seconds,
|
||||
)
|
||||
self.schema = schema
|
||||
self.statement_timeout_ms = statement_timeout_ms
|
||||
self.lock_timeout_ms = lock_timeout_ms
|
||||
self.idle_transaction_timeout_ms = idle_transaction_timeout_ms
|
||||
self.metadata, self.tables = build_session_tables(schema)
|
||||
|
||||
@contextmanager
|
||||
def transaction(self, *, immediate: bool = False) -> Iterator[Connection]:
|
||||
del immediate
|
||||
with self.engine.begin() as connection:
|
||||
# These only bound database work. LLM and document processing must
|
||||
# run before this context is entered, so a slow remote call cannot
|
||||
# consume a pool connection or leave a long transaction open.
|
||||
connection.exec_driver_sql(
|
||||
f"SET LOCAL statement_timeout = {self.statement_timeout_ms}"
|
||||
)
|
||||
connection.exec_driver_sql(f"SET LOCAL lock_timeout = {self.lock_timeout_ms}")
|
||||
connection.exec_driver_sql(
|
||||
"SET LOCAL idle_in_transaction_session_timeout = "
|
||||
f"{self.idle_transaction_timeout_ms}"
|
||||
)
|
||||
yield connection
|
||||
|
||||
def initialize(self) -> None:
|
||||
@@ -113,6 +145,7 @@ class PostgresDatabase:
|
||||
self, connection: Connection, session_id: str, *, stage: str,
|
||||
profile: dict[str, Any], draft_id: str | None = None,
|
||||
resume_id: str | None = None, increment_revision: bool = True,
|
||||
expected_revision: int | None = None,
|
||||
) -> dict[str, Any]:
|
||||
sessions = self.tables["sessions"]
|
||||
current = connection.execute(
|
||||
@@ -128,7 +161,12 @@ class PostgresDatabase:
|
||||
"resume_id": resume_id if resume_id is not None else current["resume_id"],
|
||||
"updated_at": _now(),
|
||||
}
|
||||
connection.execute(update(sessions).where(sessions.c.id == session_id).values(**values))
|
||||
statement = update(sessions).where(sessions.c.id == session_id)
|
||||
if expected_revision is not None:
|
||||
statement = statement.where(sessions.c.revision == expected_revision)
|
||||
result = connection.execute(statement.values(**values))
|
||||
if result.rowcount != 1:
|
||||
raise SessionRevisionConflict(session_id)
|
||||
return self.fetch_session(connection, session_id) # type: ignore[return-value]
|
||||
|
||||
def insert_turn(
|
||||
@@ -167,18 +205,22 @@ class PostgresDatabase:
|
||||
|
||||
def update_block(
|
||||
self, connection: Connection, block_id: str, *, lifecycle: str,
|
||||
data: dict[str, Any] | None = None,
|
||||
data: dict[str, Any] | None = None, expected_version: int | None = None,
|
||||
) -> None:
|
||||
blocks = self.tables["blocks"]
|
||||
current = connection.execute(
|
||||
select(blocks.c.data).where(blocks.c.id == block_id).with_for_update()
|
||||
select(blocks.c.data, blocks.c.version).where(blocks.c.id == block_id).with_for_update()
|
||||
).first()
|
||||
if current is None:
|
||||
raise KeyError(block_id)
|
||||
connection.execute(update(blocks).where(blocks.c.id == block_id).values(
|
||||
statement = update(blocks).where(blocks.c.id == block_id)
|
||||
if expected_version is not None:
|
||||
statement = statement.where(blocks.c.version == expected_version)
|
||||
if connection.execute(statement.values(
|
||||
lifecycle=lifecycle, data=current._mapping["data"] if data is None else data,
|
||||
version=blocks.c.version + 1, updated_at=_now(),
|
||||
))
|
||||
)).rowcount != 1:
|
||||
raise SessionRevisionConflict(block_id)
|
||||
|
||||
def supersede_active_components(self, connection: Connection, session_id: str) -> None:
|
||||
blocks = self.tables["blocks"]
|
||||
@@ -234,11 +276,14 @@ class PostgresDatabase:
|
||||
created_at=session["created_at"], updated_at=session["updated_at"],
|
||||
)
|
||||
|
||||
def fetch_resume(self, connection: Connection, session_id: str) -> dict[str, Any] | None:
|
||||
def fetch_resume(
|
||||
self, connection: Connection, session_id: str, *, for_update: bool = False
|
||||
) -> dict[str, Any] | None:
|
||||
resumes = self.tables["resumes"]
|
||||
row = connection.execute(select(resumes).where(
|
||||
resumes.c.session_id == session_id
|
||||
)).mappings().first()
|
||||
statement = select(resumes).where(resumes.c.session_id == session_id)
|
||||
if for_update:
|
||||
statement = statement.with_for_update()
|
||||
row = connection.execute(statement).mappings().first()
|
||||
return dict(row) if row else None
|
||||
|
||||
def insert_resume(
|
||||
@@ -253,13 +298,23 @@ class PostgresDatabase:
|
||||
return self.fetch_resume(connection, session_id) # type: ignore[return-value]
|
||||
|
||||
def update_resume(
|
||||
self, connection: Connection, session_id: str, content: dict[str, Any]
|
||||
self,
|
||||
connection: Connection,
|
||||
session_id: str,
|
||||
content: dict[str, Any],
|
||||
*,
|
||||
expected_revision: int | None = None,
|
||||
) -> dict[str, Any]:
|
||||
resumes = self.tables["resumes"]
|
||||
result = connection.execute(update(resumes).where(
|
||||
resumes.c.session_id == session_id
|
||||
).values(content=content, revision=resumes.c.revision + 1, updated_at=_now()))
|
||||
statement = update(resumes).where(resumes.c.session_id == session_id)
|
||||
if expected_revision is not None:
|
||||
statement = statement.where(resumes.c.revision == expected_revision)
|
||||
result = connection.execute(statement.values(
|
||||
content=content, revision=resumes.c.revision + 1, updated_at=_now()
|
||||
))
|
||||
if result.rowcount != 1:
|
||||
if expected_revision is not None:
|
||||
raise SessionRevisionConflict(session_id)
|
||||
raise KeyError(session_id)
|
||||
return self.fetch_resume(connection, session_id) # type: ignore[return-value]
|
||||
|
||||
|
||||
Reference in New Issue
Block a user