fix(data-integrity): remove session.rollback() calls from all repositories
CI / benchmark-publish (pull_request) Has been skipped
CI / lint (pull_request) Successful in 1m17s
CI / helm (pull_request) Successful in 46s
CI / typecheck (pull_request) Successful in 1m45s
CI / security (pull_request) Successful in 1m44s
CI / push-validation (pull_request) Successful in 55s
CI / build (pull_request) Successful in 1m18s
CI / benchmark-regression (pull_request) Failing after 1m44s
CI / quality (pull_request) Successful in 2m33s
CI / integration_tests (pull_request) Failing after 4m4s
CI / e2e_tests (pull_request) Successful in 6m13s
CI / unit_tests (pull_request) Failing after 7m43s
CI / coverage (pull_request) Has been skipped
CI / docker (pull_request) Has been skipped
CI / status-check (pull_request) Failing after 9s

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.
This commit is contained in:
2026-05-07 20:16:17 +00:00
committed by Forgejo
parent 23e9848f95
commit 2a5cfd90d7
@@ -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