7a298ede6e
CI / benchmark-publish (pull_request) Has been skipped
CI / lint (pull_request) Successful in 16s
CI / quality (pull_request) Successful in 30s
CI / build (pull_request) Successful in 24s
CI / security (pull_request) Successful in 40s
CI / typecheck (pull_request) Successful in 1m0s
CI / integration_tests (pull_request) Successful in 4m14s
CI / unit_tests (pull_request) Successful in 16m25s
CI / docker (pull_request) Successful in 55s
CI / benchmark-regression (pull_request) Successful in 22m55s
CI / coverage (pull_request) Successful in 38m4s
CI / lint (push) Successful in 13s
CI / build (push) Successful in 15s
CI / quality (push) Successful in 28s
CI / typecheck (push) Successful in 31s
CI / benchmark-regression (push) Has been skipped
CI / security (push) Successful in 31s
CI / integration_tests (push) Successful in 3m19s
CI / benchmark-publish (push) Successful in 14m23s
CI / unit_tests (push) Successful in 14m44s
CI / docker (push) Successful in 38s
CI / coverage (push) Successful in 32m22s
Implemented plan-level and project-level locking with configurable timeouts. Added locks table via Alembic migration storing owner_id, resource_type, resource_id, acquired_at, and expires_at. Locks enforced in PlanLifecycleService transitions. Support for re-entrant acquisition, lock renewal, graceful shutdown release, and startup cleanup of expired locks. Added diagnostics check for stale lock reporting. ISSUES CLOSED: #327
145 lines
3.7 KiB
Python
145 lines
3.7 KiB
Python
"""ASV benchmarks for concurrency lock overhead.
|
|
|
|
Measures the cost of lock acquire, release, renew, and cleanup
|
|
operations against an in-memory SQLite database.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
from datetime import UTC, datetime, timedelta
|
|
from typing import Any
|
|
|
|
from sqlalchemy import create_engine, event
|
|
from sqlalchemy.orm import Session, sessionmaker
|
|
|
|
from cleveragents.application.services.lock_service import LockService
|
|
from cleveragents.infrastructure.database.models import Base, LockModel
|
|
|
|
_bench_counter = 9000
|
|
|
|
|
|
def _next_id() -> str:
|
|
global _bench_counter
|
|
_bench_counter += 1
|
|
return f"bench-res-{_bench_counter}"
|
|
|
|
|
|
def _build_lock_service() -> tuple[LockService, sessionmaker[Session]]:
|
|
engine = create_engine(
|
|
"sqlite:///:memory:",
|
|
echo=False,
|
|
future=True,
|
|
connect_args={"check_same_thread": False},
|
|
)
|
|
|
|
@event.listens_for(engine, "connect")
|
|
def _fk(dbapi_conn: Any, _rec: Any) -> None:
|
|
cursor = dbapi_conn.cursor()
|
|
cursor.execute("PRAGMA foreign_keys=ON")
|
|
cursor.close()
|
|
|
|
Base.metadata.create_all(engine)
|
|
sf: sessionmaker[Session] = sessionmaker(
|
|
bind=engine,
|
|
expire_on_commit=False,
|
|
autoflush=False,
|
|
autocommit=False,
|
|
class_=Session,
|
|
)
|
|
return LockService(session_factory=sf), sf
|
|
|
|
|
|
class TimeAcquireLock:
|
|
"""Benchmark acquiring a new lock."""
|
|
|
|
timeout = 30
|
|
|
|
def setup(self) -> None:
|
|
self.svc, self._sf = _build_lock_service()
|
|
|
|
def time_acquire(self) -> None:
|
|
rid = _next_id()
|
|
self.svc.acquire("owner-bench", "plan", rid)
|
|
|
|
|
|
class TimeReentrantAcquire:
|
|
"""Benchmark re-entrant lock acquisition."""
|
|
|
|
timeout = 30
|
|
|
|
def setup(self) -> None:
|
|
self.svc, self._sf = _build_lock_service()
|
|
self.svc.acquire("owner-bench", "plan", "reentrant-res")
|
|
|
|
def time_reentrant_acquire(self) -> None:
|
|
self.svc.acquire("owner-bench", "plan", "reentrant-res")
|
|
|
|
|
|
class TimeReleaseLock:
|
|
"""Benchmark releasing a lock."""
|
|
|
|
timeout = 30
|
|
|
|
def setup(self) -> None:
|
|
self.svc, self._sf = _build_lock_service()
|
|
self._rid = _next_id()
|
|
self.svc.acquire("owner-bench", "plan", self._rid)
|
|
|
|
def time_release(self) -> None:
|
|
self.svc.release("owner-bench", "plan", self._rid)
|
|
self.svc.acquire("owner-bench", "plan", self._rid)
|
|
|
|
|
|
class TimeRenewLock:
|
|
"""Benchmark renewing a lock."""
|
|
|
|
timeout = 30
|
|
|
|
def setup(self) -> None:
|
|
self.svc, self._sf = _build_lock_service()
|
|
self.svc.acquire("owner-bench", "plan", "renew-res")
|
|
|
|
def time_renew(self) -> None:
|
|
self.svc.renew("owner-bench", "plan", "renew-res")
|
|
|
|
|
|
class TimeCleanupExpired:
|
|
"""Benchmark expired lock cleanup."""
|
|
|
|
timeout = 30
|
|
|
|
def setup(self) -> None:
|
|
self.svc, self._sf = _build_lock_service()
|
|
now = datetime.now(tz=UTC)
|
|
expired = (now - timedelta(seconds=60)).isoformat()
|
|
acquired = (now - timedelta(seconds=120)).isoformat()
|
|
session = self._sf()
|
|
for i in range(50):
|
|
session.add(
|
|
LockModel(
|
|
owner_id=f"expired-{i}",
|
|
resource_type="plan",
|
|
resource_id=f"cleanup-{i}",
|
|
acquired_at=acquired,
|
|
expires_at=expired,
|
|
)
|
|
)
|
|
session.commit()
|
|
session.close()
|
|
|
|
def time_cleanup(self) -> None:
|
|
self.svc.cleanup_expired()
|
|
|
|
|
|
class TimeIsLocked:
|
|
"""Benchmark is_locked check."""
|
|
|
|
timeout = 30
|
|
|
|
def setup(self) -> None:
|
|
self.svc, self._sf = _build_lock_service()
|
|
self.svc.acquire("owner-bench", "plan", "locked-res")
|
|
|
|
def time_is_locked(self) -> None:
|
|
self.svc.is_locked("plan", "locked-res")
|