9ff0f04ba3
Enable ANN rules in ruff.toml (flake8-annotations) and resolve all 221 violations: ANN201/ANN202 — return types on 168 public/private functions: - All 28 FastAPI routers: endpoints annotated with dict/list/specific schema/ StreamingResponse/FileResponse/JSONResponse as appropriate - main.py: lifespan→AsyncGenerator[None,None], exception handlers→JSONResponse - database.py: get_db→Generator[Session,None,None], proxy methods→correct types - middleware/request_context.py: dispatch→Response with Callable call_next type ANN001/ANN002/ANN003 — 32 missing argument types: - seed_demo.py: all db parameters typed as Session - domain/unit_of_work.py: __aexit__ exc_type/exc_val/exc_tb typed with TracebackType - services: audit_service user_id→UUID|None, heatmap_service query/model/builder, notification_service test→Test, tempo_service test→Test/user→User, test_workflow_service test_id→UUID, campaign_crud **fields→object, test_crud **fields→object (4 sites) ANN401 — 16 Any usages resolved: - Domain entities (campaign/technique/threat_actor/test_entity): replaced Any with actual ORM types via TYPE_CHECKING guards to avoid circular imports - detection_rule_service: test_id/detection_rule_id/evaluator_id→UUID - score_cache: kept Any with # noqa: ANN401 (genuinely generic cache) - jira_service/tempo_service: kept Any with # noqa: ANN401 (lazy optional deps) - d3fend_import_service: _to_str(v: Any) kept with # noqa: ANN401 ANN204/ANN205/ANN206 — special/static/class methods: - database.py proxy __call__/__getattr__: *args: object/**kwargs: object - schemas/test.py model_validate: obj→object, **kwargs→object - sa_technique_repository._int_type→type All 439 unit tests pass. ruff check app/ → All checks passed!
63 lines
2.0 KiB
Python
63 lines
2.0 KiB
Python
"""Unit of Work — wraps a SQLAlchemy session for explicit transaction control.
|
|
|
|
Usage in routers::
|
|
|
|
with UnitOfWork(db) as uow:
|
|
service_a(db, ...)
|
|
service_b(db, ...)
|
|
uow.commit() # single commit for the entire operation
|
|
|
|
If an exception propagates, ``__exit__`` issues a rollback automatically.
|
|
Services should **never** call ``db.commit()``; they use ``db.add()`` /
|
|
``db.flush()`` to stage work and let the caller decide when to commit.
|
|
|
|
**Documented exceptions** (services that may commit internally):
|
|
- Import services (atomic_import, sigma_import, etc.) — self-contained sync ops.
|
|
- Background jobs (campaign_scheduler, intel_service, stale_detection,
|
|
mitre_sync) — self-contained operations.
|
|
- Self-contained batch ops (e.g. detection_rule_service.auto_associate_rules,
|
|
snapshot_service.create_snapshot, campaign_service.generate_campaign_from_*,
|
|
osint_enrichment_service.enrich_technique_with_cves).
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
from types import TracebackType
|
|
|
|
from sqlalchemy.orm import Session
|
|
|
|
|
|
class UnitOfWork:
|
|
"""Lightweight transaction wrapper around an existing SQLAlchemy session."""
|
|
|
|
def __init__(self, session: Session) -> None:
|
|
self._session = session
|
|
|
|
# -- context manager -----------------------------------------------------
|
|
|
|
def __enter__(self) -> "UnitOfWork":
|
|
return self
|
|
|
|
def __exit__(
|
|
self,
|
|
exc_type: type[BaseException] | None,
|
|
exc_val: BaseException | None,
|
|
exc_tb: TracebackType | None,
|
|
) -> None:
|
|
if exc_type is not None:
|
|
self.rollback()
|
|
|
|
# -- public API ----------------------------------------------------------
|
|
|
|
def commit(self) -> None:
|
|
"""Flush pending changes and commit the transaction."""
|
|
self._session.commit()
|
|
|
|
def rollback(self) -> None:
|
|
"""Roll back the current transaction."""
|
|
self._session.rollback()
|
|
|
|
def flush(self) -> None:
|
|
"""Flush pending changes without committing (useful for getting IDs)."""
|
|
self._session.flush()
|