From 2a5cfd90d76de7a182dfd31bf1712ab1bc85a91a Mon Sep 17 00:00:00 2001 From: HAL9000 Date: Thu, 7 May 2026 20:16:17 +0000 Subject: [PATCH] fix(data-integrity): remove session.rollback() calls from all repositories Remove all session.rollback() calls from repository error handlers. When repositories receive a Session via constructor injection from the UnitOfWork, calling rollback() on that shared session discards ALL pending work from other repositories in the same transaction, causing silent data loss. The UnitOfWork is the sole owner of the session and must manage transaction lifecycle (commit/rollback). Repositories now delegate this responsibility to the UoW by simply propagating exceptions without calling rollback. --- .../infrastructure/database/repositories.py | 66 ------------------- 1 file changed, 66 deletions(-) diff --git a/src/cleveragents/infrastructure/database/repositories.py b/src/cleveragents/infrastructure/database/repositories.py index 0aa0036f6..641a66c33 100644 --- a/src/cleveragents/infrastructure/database/repositories.py +++ b/src/cleveragents/infrastructure/database/repositories.py @@ -174,7 +174,6 @@ class ProjectRepository: project.id = db_project.id # type: ignore return project except (OperationalError, SQLAlchemyDatabaseError) as e: - self.session.rollback() raise DatabaseError(f"Failed to create project: {e}") from e @database_retry @@ -930,7 +929,6 @@ class ActorRepository: self.session.flush() return except IntegrityError: - self.session.rollback() # Another writer raced us — re-fetch and update. row = cast( Any, @@ -1009,13 +1007,11 @@ class ActionRepository(ActionRepositoryProtocol): session.flush() return action except IntegrityError as exc: - session.rollback() # Unique constraint on ``name`` column if "UNIQUE" in str(exc).upper() or "unique" in str(exc).lower(): raise DuplicateActionError(str(action.namespaced_name)) from exc raise DatabaseError(f"Failed to create action: {exc}") from exc except (OperationalError, SQLAlchemyDatabaseError) as exc: - session.rollback() raise DatabaseError(f"Failed to create action: {exc}") from exc @database_retry @@ -1257,7 +1253,6 @@ class ActionRepository(ActionRepositoryProtocol): session.flush() return action except (OperationalError, SQLAlchemyDatabaseError) as exc: - session.rollback() raise DatabaseError( f"Failed to update action {action_name_str}: {exc}" ) from exc @@ -1332,7 +1327,6 @@ class ActionRepository(ActionRepositoryProtocol): except ActionInUseError: raise except (OperationalError, SQLAlchemyDatabaseError) as exc: - session.rollback() raise DatabaseError( f"Failed to delete action {action_name}: {exc}" ) from exc @@ -1402,13 +1396,11 @@ class LifecyclePlanRepository(LifecyclePlanRepositoryProtocol): session.flush() return plan except IntegrityError as exc: - session.rollback() exc_str = str(exc).upper() if "UNIQUE" in exc_str or "PRIMARY" in exc_str: raise DuplicatePlanError(str(plan.identity.plan_id)) from exc raise DatabaseError(f"Failed to create plan: {exc}") from exc except (OperationalError, SQLAlchemyDatabaseError) as exc: - session.rollback() raise DatabaseError(f"Failed to create plan: {exc}") from exc @database_retry @@ -1618,7 +1610,6 @@ class LifecyclePlanRepository(LifecyclePlanRepositoryProtocol): except PlanNotFoundError: raise except (OperationalError, SQLAlchemyDatabaseError) as exc: - session.rollback() raise DatabaseError(f"Failed to update plan {plan_id_str}: {exc}") from exc @database_retry @@ -1662,7 +1653,6 @@ class LifecyclePlanRepository(LifecyclePlanRepositoryProtocol): session.flush() return True except (OperationalError, SQLAlchemyDatabaseError) as exc: - session.rollback() raise DatabaseError(f"Failed to delete plan {plan_id}: {exc}") from exc @database_retry @@ -1943,12 +1933,10 @@ class ResourceTypeRepository: session.flush() return resource_type except IntegrityError as exc: - session.rollback() if "UNIQUE" in str(exc).upper() or "unique" in str(exc).lower(): raise DuplicateResourceTypeError(resource_type.name) from exc raise DatabaseError(f"Failed to create resource type: {exc}") from exc except (OperationalError, SQLAlchemyDatabaseError) as exc: - session.rollback() raise DatabaseError(f"Failed to create resource type: {exc}") from exc @database_retry @@ -2079,7 +2067,6 @@ class ResourceTypeRepository: except ResourceTypeNotFoundError: raise except (OperationalError, SQLAlchemyDatabaseError) as exc: - session.rollback() raise DatabaseError( f"Failed to update resource type '{resource_type.name}': {exc}" ) from exc @@ -2119,7 +2106,6 @@ class ResourceTypeRepository: except ResourceTypeHasResourcesError: raise except (OperationalError, SQLAlchemyDatabaseError) as exc: - session.rollback() raise DatabaseError( f"Failed to delete resource type '{name}': {exc}" ) from exc @@ -2274,14 +2260,12 @@ class ResourceRepository: except (ResourceTypeNotFoundError, DuplicateResourceError): raise except IntegrityError as exc: - session.rollback() if "UNIQUE" in str(exc).upper() or "unique" in str(exc).lower(): raise DuplicateResourceError( resource.name or resource.resource_id ) from exc raise DatabaseError(f"Failed to create resource: {exc}") from exc except (OperationalError, SQLAlchemyDatabaseError) as exc: - session.rollback() raise DatabaseError(f"Failed to create resource: {exc}") from exc @database_retry @@ -2410,7 +2394,6 @@ class ResourceRepository: except ResourceNotFoundRepoError: raise except (OperationalError, SQLAlchemyDatabaseError) as exc: - session.rollback() raise DatabaseError( f"Failed to update resource '{resource.resource_id}': {exc}" ) from exc @@ -2455,7 +2438,6 @@ class ResourceRepository: except ResourceHasEdgesError: raise except (OperationalError, SQLAlchemyDatabaseError) as exc: - session.rollback() raise DatabaseError( f"Failed to delete resource '{resource_id}': {exc}" ) from exc @@ -2547,7 +2529,6 @@ class ResourceRepository: ): raise except IntegrityError as exc: - session.rollback() raise DatabaseError( f"Failed to link {parent_id} -> {child_id}: {exc}" ) from exc @@ -2555,7 +2536,6 @@ class ResourceRepository: OperationalError, SQLAlchemyDatabaseError, ) as exc: - session.rollback() raise DatabaseError( f"Failed to link {parent_id} -> {child_id}: {exc}" ) from exc @@ -2606,7 +2586,6 @@ class ResourceRepository: OperationalError, SQLAlchemyDatabaseError, ) as exc: - session.rollback() raise DatabaseError( f"Failed to unlink {parent_id} -> {child_id}: {exc}" ) from exc @@ -2817,7 +2796,6 @@ class ResourceRepository: OperationalError, SQLAlchemyDatabaseError, ) as exc: - session.rollback() raise DatabaseError( f"Failed to auto-discover children for '{resource_id}': {exc}" ) from exc @@ -3040,14 +3018,12 @@ class NamespacedProjectRepository(ProjectRepositoryProtocol): session.commit() return project except IntegrityError as exc: - session.rollback() if "UNIQUE" in str(exc).upper() or "unique" in str(exc).lower(): raise DatabaseError( f"Project '{project.namespaced_name}' already exists" ) from exc raise DatabaseError(f"Failed to create project: {exc}") from exc except (OperationalError, SQLAlchemyDatabaseError) as exc: - session.rollback() raise DatabaseError(f"Failed to create project: {exc}") from exc finally: session.close() @@ -3192,7 +3168,6 @@ class NamespacedProjectRepository(ProjectRepositoryProtocol): except ProjectNotFoundError: raise except (OperationalError, SQLAlchemyDatabaseError) as exc: - session.rollback() raise DatabaseError(f"Failed to update project '{ns_name}': {exc}") from exc finally: session.close() @@ -3226,7 +3201,6 @@ class NamespacedProjectRepository(ProjectRepositoryProtocol): session.commit() return True except (OperationalError, SQLAlchemyDatabaseError) as exc: - session.rollback() raise DatabaseError( f"Failed to delete project '{namespaced_name}': {exc}" ) from exc @@ -3313,10 +3287,8 @@ class ProjectResourceLinkRepository: except DuplicateLinkError: raise except IntegrityError as exc: - session.rollback() raise DatabaseError(f"Failed to create link: {exc}") from exc except (OperationalError, SQLAlchemyDatabaseError) as exc: - session.rollback() raise DatabaseError(f"Failed to create link: {exc}") from exc finally: session.close() @@ -3390,7 +3362,6 @@ class ProjectResourceLinkRepository: session.commit() return True except (OperationalError, SQLAlchemyDatabaseError) as exc: - session.rollback() raise DatabaseError(f"Failed to remove link '{link_id}': {exc}") from exc finally: session.close() @@ -3503,7 +3474,6 @@ class ToolRegistryRepository: session.commit() return tool except IntegrityError as exc: - session.rollback() name_str = ( tool.get("name", "") if isinstance(tool, dict) @@ -3513,7 +3483,6 @@ class ToolRegistryRepository: raise DuplicateToolError(name_str) from exc raise DatabaseError(f"Failed to create tool: {exc}") from exc except (OperationalError, SQLAlchemyDatabaseError) as exc: - session.rollback() raise DatabaseError(f"Failed to create tool: {exc}") from exc finally: session.close() @@ -3636,7 +3605,6 @@ class ToolRegistryRepository: session.commit() return tool except (OperationalError, SQLAlchemyDatabaseError) as exc: - session.rollback() raise DatabaseError(f"Failed to update tool {name_str}: {exc}") from exc finally: session.close() @@ -3678,7 +3646,6 @@ class ToolRegistryRepository: except ToolInUseError: raise except (OperationalError, SQLAlchemyDatabaseError) as exc: - session.rollback() raise DatabaseError(f"Failed to delete tool {name}: {exc}") from exc finally: session.close() @@ -3977,7 +3944,6 @@ class ValidationAttachmentRepository: "created_at": now_iso, } except IntegrityError as exc: - session.rollback() if "UNIQUE" in str(exc).upper() or "unique" in str(exc).lower(): raise DuplicateValidationAttachmentError( validation_name, @@ -3987,7 +3953,6 @@ class ValidationAttachmentRepository: ) from exc raise DatabaseError(f"Failed to attach validation: {exc}") from exc except (OperationalError, SQLAlchemyDatabaseError) as exc: - session.rollback() raise DatabaseError(f"Failed to attach validation: {exc}") from exc finally: session.close() @@ -4017,7 +3982,6 @@ class ValidationAttachmentRepository: session.commit() return True except (OperationalError, SQLAlchemyDatabaseError) as exc: - session.rollback() raise DatabaseError( f"Failed to detach validation {attachment_id}: {exc}" ) from exc @@ -4165,10 +4129,8 @@ class SessionRepository: db_session.commit() return session except IntegrityError as exc: - db_session.rollback() raise DatabaseError(f"Failed to create session: {exc}") from exc except (OperationalError, SQLAlchemyDatabaseError) as exc: - db_session.rollback() raise DatabaseError(f"Failed to create session: {exc}") from exc finally: if self._auto_commit: @@ -4241,7 +4203,6 @@ class SessionRepository: db_session.commit() return True except (OperationalError, SQLAlchemyDatabaseError) as exc: - db_session.rollback() raise DatabaseError( f"Failed to delete session {session_id}: {exc}" ) from exc @@ -4297,7 +4258,6 @@ class SessionRepository: db_session.commit() return session except (OperationalError, SQLAlchemyDatabaseError) as exc: - db_session.rollback() raise DatabaseError( f"Failed to update session {session.session_id}: {exc}" ) from exc @@ -4360,7 +4320,6 @@ class SessionMessageRepository: db_session.commit() return message except (OperationalError, SQLAlchemyDatabaseError) as exc: - db_session.rollback() raise DatabaseError( f"Failed to append message to session {session_id}: {exc}" ) from exc @@ -4606,7 +4565,6 @@ class AutomationProfileRepository: except AutomationProfileSchemaVersionError: raise except IntegrityError as exc: - session.rollback() raise DuplicateAutomationProfileError( f"Profile '{profile.name}' already exists: {exc}" ) from exc @@ -4614,7 +4572,6 @@ class AutomationProfileRepository: OperationalError, SQLAlchemyDatabaseError, ) as exc: - session.rollback() raise DatabaseError( f"Failed to upsert profile '{profile.name}': {exc}" ) from exc @@ -4648,7 +4605,6 @@ class AutomationProfileRepository: OperationalError, SQLAlchemyDatabaseError, ) as exc: - session.rollback() raise DatabaseError(f"Failed to delete profile '{name}': {exc}") from exc finally: if self._auto_commit: @@ -4950,7 +4906,6 @@ class SkillRepository: session.close() return skill except IntegrityError as exc: - session.rollback() name_str = getattr(skill, "name", "") if "UNIQUE" in str(exc).upper() or "unique" in str(exc).lower(): raise DuplicateSkillError(name_str) from exc @@ -4958,7 +4913,6 @@ class SkillRepository: except DuplicateSkillError: raise except (OperationalError, SQLAlchemyDatabaseError) as exc: - session.rollback() raise DatabaseError(f"Failed to create skill: {exc}") from exc @database_retry @@ -5066,7 +5020,6 @@ class SkillRepository: except SkillNotFoundError: raise except (OperationalError, SQLAlchemyDatabaseError) as exc: - session.rollback() raise DatabaseError(f"Failed to update skill {name_str}: {exc}") from exc @database_retry @@ -5096,7 +5049,6 @@ class SkillRepository: session.close() return True except (OperationalError, SQLAlchemyDatabaseError) as exc: - session.rollback() raise DatabaseError(f"Failed to delete skill {name}: {exc}") from exc # -- Flattened tool cache methods (m4_002) ------------------------------ @@ -5145,7 +5097,6 @@ class SkillRepository: except SkillNotFoundError: raise except (OperationalError, SQLAlchemyDatabaseError) as exc: - session.rollback() raise DatabaseError( f"Failed to update flattened tools for {name}: {exc}", ) from exc @@ -5257,7 +5208,6 @@ class SkillRepository: except SkillNotFoundError: raise except (OperationalError, SQLAlchemyDatabaseError) as exc: - session.rollback() raise DatabaseError( f"Failed to recompute hash for {name}: {exc}", ) from exc @@ -5291,7 +5241,6 @@ class SkillRepository: except SkillNotFoundError: raise except (OperationalError, SQLAlchemyDatabaseError) as exc: - session.rollback() raise DatabaseError( f"Failed to invalidate cached summaries for {name}: {exc}", ) from exc @@ -5360,7 +5309,6 @@ class DecisionRepository(DecisionRepositoryProtocol): session.flush() return decision except IntegrityError as exc: - session.rollback() if "UNIQUE" in str(exc).upper() or "unique" in str(exc).lower(): raise DuplicateDecisionError( str(decision.decision_id), @@ -5369,7 +5317,6 @@ class DecisionRepository(DecisionRepositoryProtocol): f"Failed to create decision: {exc}", ) from exc except (OperationalError, SQLAlchemyDatabaseError) as exc: - session.rollback() raise DatabaseError( f"Failed to create decision: {exc}", ) from exc @@ -5598,7 +5545,6 @@ class DecisionRepository(DecisionRepositoryProtocol): except DecisionNotFoundError: raise except (OperationalError, SQLAlchemyDatabaseError) as exc: - session.rollback() raise DatabaseError( f"Failed to update superseded_by for {decision_id}: {exc}", ) from exc @@ -5666,7 +5612,6 @@ class DecisionRepository(DecisionRepositoryProtocol): session.flush() return True except (OperationalError, SQLAlchemyDatabaseError) as exc: - session.rollback() raise DatabaseError( f"Failed to delete decision {decision_id}: {exc}", ) from exc @@ -5777,12 +5722,10 @@ class CheckpointRepository: session.flush() return model.to_domain() except IntegrityError as exc: - session.rollback() raise DatabaseError( f"Duplicate or constraint violation creating checkpoint: {exc}" ) from exc except (OperationalError, SQLAlchemyDatabaseError) as exc: - session.rollback() raise DatabaseError(f"Failed to create checkpoint: {exc}") from exc @database_retry @@ -5812,7 +5755,6 @@ class CheckpointRepository: except CheckpointNotFoundError: raise except (OperationalError, SQLAlchemyDatabaseError) as exc: - session.rollback() raise DatabaseError( f"Failed to get checkpoint {checkpoint_id}: {exc}" ) from exc @@ -5840,7 +5782,6 @@ class CheckpointRepository: ) return [row.to_domain() for row in rows] except (OperationalError, SQLAlchemyDatabaseError) as exc: - session.rollback() raise DatabaseError( f"Failed to list checkpoints for plan {plan_id}: {exc}" ) from exc @@ -5871,7 +5812,6 @@ class CheckpointRepository: session.flush() return True except (OperationalError, SQLAlchemyDatabaseError) as exc: - session.rollback() raise DatabaseError( f"Failed to delete checkpoint {checkpoint_id}: {exc}" ) from exc @@ -5926,7 +5866,6 @@ class CheckpointRepository: session.flush() return pruned_ids except (OperationalError, SQLAlchemyDatabaseError) as exc: - session.rollback() raise DatabaseError( f"Failed to prune checkpoints for plan {plan_id}: {exc}" ) from exc @@ -5989,7 +5928,6 @@ class CorrectionAttemptRepository: session.flush() return model.to_domain() except IntegrityError as exc: - session.rollback() exc_str = str(exc).upper() if "UNIQUE" in exc_str: raise DuplicateCorrectionAttemptError( @@ -6003,7 +5941,6 @@ class CorrectionAttemptRepository: ) from exc raise DatabaseError(f"Failed to create correction attempt: {exc}") from exc except (OperationalError, SQLAlchemyDatabaseError) as exc: - session.rollback() raise DatabaseError(f"Failed to create correction attempt: {exc}") from exc # --- GET --------------------------------------------------------------- @@ -6177,7 +6114,6 @@ class CorrectionAttemptRepository: except (CorrectionAttemptNotFoundError, InvalidCorrectionStateTransitionError): raise except IntegrityError as exc: - session.rollback() exc_str = str(exc).upper() if "FOREIGN KEY" in exc_str or "FOREIGN_KEY" in exc_str: fk_details = ( @@ -6194,7 +6130,6 @@ class CorrectionAttemptRepository: f"Failed to update correction attempt {correction_attempt_id}: {exc}" ) from exc except (OperationalError, SQLAlchemyDatabaseError) as exc: - session.rollback() raise DatabaseError( f"Failed to update correction attempt {correction_attempt_id}: {exc}" ) from exc @@ -6227,7 +6162,6 @@ class CorrectionAttemptRepository: session.flush() return True except (OperationalError, SQLAlchemyDatabaseError) as exc: - session.rollback() raise DatabaseError( f"Failed to delete correction attempt {correction_attempt_id}: {exc}" ) from exc