Files
leocrm/app/commands/contact_commands.py
T
Agent Zero abbe7a18fc fix(audit): P0-P3 audit fixes — 838 ruff errors → 0, 30 F821 bugs fixed, 118 files changed
- P0: hooks.py 3-tuple fix, trigger_dispatcher Contract, contacts/plugin unregister_actions_by_owner
- P0: 5 test files — check_permission mocks removed, hardcoded DB credential → env var
- P1: attachment_service DmsFile via Contract helper, restore_registry/history_hooks dedup
- P1: mail/plugin restore unregister, mcp_client datetime.now(UTC), saved_views/filters patterns
- P1: ProtectedRoute fail-closed, 13 test assertion fixes (bcrypt, DB-URLs, SECRET_KEYs)
- P2: deprecated notifications → post_system_message (3 files), forgejo Base, report_generator lazy import
- P2: webhooks permissions, deps.py/roles.py plugin perms removed, import_export default
- P2: address/tags/entity_links patterns removed, worker.py Contract-Umgehungen fixed
- P2: 28 frontend TODOs (hardcoded constants, deprecated notification API)
- P3: dead code, duplicates, deprecated imports, private attr, __import__ inline
- P3: 8 frontend TODOs (LucideIcons, inline styles, XSS, i18n)
- ruff: 838 → 0 (612 auto-fix + 246 manual + 27 F821 regression fix)
- F821: 30 → 0 (AutomationDefinition, DmsFile, user_id, Path, Any, String)
- Contract-Umgehungen: 2 neue gefunden (worker.py:169, worker.py:280) und gefixt
2026-08-16 01:17:18 +02:00

393 lines
13 KiB
Python

"""Contact commands — Create, Update, Delete (soft), Merge.
Each command:
- Checks permissions via ``require_permission`` semantics
- Delegates business logic to existing services
- Creates an AuditLog entry
- Enqueues outbox events
- Validates state transitions via the state machine
Commands do NOT commit — the transaction is managed by ``get_db``.
"""
from __future__ import annotations
import logging
import uuid
from typing import Any
import redis.asyncio as aioredis
from sqlalchemy.ext.asyncio import AsyncSession
from app.commands.base import BaseCommand, CommandResult
from app.core.outbox import enqueue_outbox_event
from app.core.state_machine import StateMachineError, contact_state_machine
from app.models.audit import AuditLog
from app.models.contact import Contact
from app.services import contact_service, dedup_service
logger = logging.getLogger(__name__)
class CreateContactCommand(BaseCommand):
"""Create a new contact (company or person).
Args:
data: Contact fields dict (from ContactCreate schema).
"""
permission = "contacts:write"
def __init__(self, data: dict[str, Any]) -> None:
self.data = data
self._created_contact_id: uuid.UUID | None = None
self._serialized: dict | None = None
async def run(
self,
db: AsyncSession,
redis: aioredis.Redis,
current_user: dict[str, Any],
) -> CommandResult:
tenant_id = self._tenant_id(current_user)
user_id = self._user_id(current_user)
# Set default status for new contacts
if "status" not in self.data:
self.data["status"] = "lead"
# Validate status if provided
status = self.data.get("status", "lead")
if status not in contact_state_machine.transitions:
return CommandResult.fail(f"Invalid contact status: '{status}'")
# Delegate to existing service
try:
serialized = await contact_service.create_contact(
db, tenant_id, user_id, self.data
)
except ValueError as exc:
return CommandResult.fail(str(exc))
self._created_contact_id = uuid.UUID(serialized["id"])
self._serialized = serialized
# Handle company_ids — create ContactPerson links
company_ids = self.data.get("company_ids")
if company_ids:
from app.models.contact import ContactPerson
for cid in company_ids:
cp = ContactPerson(
tenant_id=tenant_id,
contact_id=uuid.UUID(cid),
displayname=serialized.get("displayname", ""),
firstname=self.data.get("firstname"),
lastname=self.data.get("surname"),
email=self.data.get("email_1"),
phone=self.data.get("phone_1"),
created_by=user_id,
updated_by=user_id,
)
db.add(cp)
await db.flush()
# Enqueue outbox events
events: list[dict] = []
await enqueue_outbox_event(db, tenant_id, "contact.created", {
"contact_id": serialized["id"],
"tenant_id": str(tenant_id),
"user_id": str(user_id),
"type": self.data.get("type", "company"),
})
events.append({"event": "contact.created", "contact_id": serialized["id"]})
if self.data.get("type") == "company":
await enqueue_outbox_event(db, tenant_id, "lead.created", {
"contact_id": serialized["id"],
"tenant_id": str(tenant_id),
"user_id": str(user_id),
})
events.append({"event": "lead.created", "contact_id": serialized["id"]})
return CommandResult.ok(data=serialized, events=events)
async def audit(self, db: AsyncSession, current_user: dict[str, Any]) -> None:
tenant_id = self._tenant_id(current_user)
user_id = self._user_id(current_user)
if self._created_contact_id is None:
return
entry = AuditLog(
tenant_id=tenant_id,
user_id=user_id,
action="create",
entity_type="contact",
entity_id=self._created_contact_id,
changes=self._serialized,
)
db.add(entry)
await db.flush()
class UpdateContactCommand(BaseCommand):
"""Update an existing contact.
Args:
contact_id: UUID string of the contact to update.
data: Contact fields to update (from ContactUpdate schema).
"""
permission = "contacts:write"
def __init__(self, contact_id: str, data: dict[str, Any]) -> None:
self.contact_id = contact_id
self.data = data
self._contact_uuid: uuid.UUID | None = None
self._serialized: dict | None = None
self._changes: dict | None = None
async def run(
self,
db: AsyncSession,
redis: aioredis.Redis,
current_user: dict[str, Any],
) -> CommandResult:
tenant_id = self._tenant_id(current_user)
user_id = self._user_id(current_user)
# Validate status transition if status is being updated
if "status" in self.data:
new_status = self.data["status"]
# Query only the status column to avoid lazy-loading issues
from sqlalchemy import select
status_q = select(Contact.status).where(
Contact.id == uuid.UUID(self.contact_id),
Contact.tenant_id == tenant_id,
Contact.deleted_at.is_(None),
)
status_result = await db.execute(status_q)
current_status = status_result.scalar_one_or_none()
if current_status is None:
return CommandResult.fail("Contact not found")
try:
contact_state_machine.transition(current_status, new_status)
except StateMachineError as exc:
return CommandResult.fail(str(exc))
# Delegate to existing service
try:
serialized = await contact_service.update_contact(
db, tenant_id, user_id, self.contact_id, self.data
)
except ValueError as exc:
return CommandResult.fail(str(exc))
self._contact_uuid = uuid.UUID(self.contact_id)
self._serialized = serialized
# Compute changes for audit
self._changes = {k: {"new": v} for k, v in self.data.items()}
# Enqueue outbox event
events: list[dict] = []
await enqueue_outbox_event(db, tenant_id, "contact.updated", {
"contact_id": self.contact_id,
"tenant_id": str(tenant_id),
"user_id": str(user_id),
"type": serialized.get("type"),
})
events.append({"event": "contact.updated", "contact_id": self.contact_id})
return CommandResult.ok(data=serialized, events=events)
async def audit(self, db: AsyncSession, current_user: dict[str, Any]) -> None:
tenant_id = self._tenant_id(current_user)
user_id = self._user_id(current_user)
if self._contact_uuid is None:
return
entry = AuditLog(
tenant_id=tenant_id,
user_id=user_id,
action="update",
entity_type="contact",
entity_id=self._contact_uuid,
changes=self._changes,
)
db.add(entry)
await db.flush()
class DeleteContactCommand(BaseCommand):
"""Soft-delete a contact.
Args:
contact_id: UUID string of the contact to delete.
hard: If True, perform GDPR hard-delete instead of soft-delete.
"""
permission = "contacts:delete"
def __init__(self, contact_id: str, hard: bool = False) -> None:
self.contact_id = contact_id
self.hard = hard
self._contact_uuid: uuid.UUID | None = None
self._snapshot: dict | None = None
async def run(
self,
db: AsyncSession,
redis: aioredis.Redis,
current_user: dict[str, Any],
) -> CommandResult:
tenant_id = self._tenant_id(current_user)
user_id = self._user_id(current_user)
# Fetch contact for snapshot before deletion
from sqlalchemy import select
from sqlalchemy.orm import selectinload
q = (
select(Contact)
.options(selectinload(Contact.contact_persons))
.where(
Contact.id == uuid.UUID(self.contact_id),
Contact.tenant_id == tenant_id,
)
)
result = await db.execute(q)
contact = result.scalar_one_or_none()
if not contact:
return CommandResult.fail("Contact not found")
self._contact_uuid = contact.id
self._snapshot = contact_service._serialize_contact_detail(contact)
if self.hard:
await contact_service.hard_delete_contact(
db, tenant_id, self.contact_id
)
else:
await contact_service.delete_contact(
db, tenant_id, self.contact_id, user_id,
is_system_admin=current_user.get("is_system_admin", False),
)
# Enqueue outbox event
events: list[dict] = []
event_name = "contact.hard_deleted" if self.hard else "contact.deleted"
await enqueue_outbox_event(db, tenant_id, event_name, {
"contact_id": self.contact_id,
"tenant_id": str(tenant_id),
"user_id": str(user_id),
})
events.append({"event": event_name, "contact_id": self.contact_id})
return CommandResult.ok(data=None, events=events)
async def audit(self, db: AsyncSession, current_user: dict[str, Any]) -> None:
tenant_id = self._tenant_id(current_user)
user_id = self._user_id(current_user)
if self._contact_uuid is None:
return
action = "hard_delete" if self.hard else "delete"
entry = AuditLog(
tenant_id=tenant_id,
user_id=user_id,
action=action,
entity_type="contact",
entity_id=self._contact_uuid,
changes={"snapshot": self._snapshot},
)
db.add(entry)
await db.flush()
class MergeContactsCommand(BaseCommand):
"""Merge two contacts (source → target).
Args:
source_contact_id: UUID string of the source contact (will be soft-deleted).
target_contact_id: UUID string of the target contact (will survive).
field_overrides: Optional field overrides to apply to the target.
note: Optional note for the merge history.
"""
permission = "contacts:write"
def __init__(
self,
source_contact_id: str,
target_contact_id: str,
field_overrides: dict[str, Any] | None = None,
note: str | None = None,
) -> None:
self.source_contact_id = source_contact_id
self.target_contact_id = target_contact_id
self.field_overrides = field_overrides
self.note = note
self._result_data: dict | None = None
async def run(
self,
db: AsyncSession,
redis: aioredis.Redis,
current_user: dict[str, Any],
) -> CommandResult:
tenant_id = self._tenant_id(current_user)
user_id = self._user_id(current_user)
if self.source_contact_id == self.target_contact_id:
return CommandResult.fail("Source and target contacts must be different")
try:
result = await dedup_service.merge_contacts(
db, tenant_id, user_id,
source_id=self.source_contact_id,
target_id=self.target_contact_id,
field_overrides=self.field_overrides,
note=self.note,
)
except ValueError as exc:
return CommandResult.fail(str(exc))
self._result_data = result
# Enqueue outbox events
events: list[dict] = []
await enqueue_outbox_event(db, tenant_id, "contact.merged", {
"source_contact_id": self.source_contact_id,
"target_contact_id": self.target_contact_id,
"tenant_id": str(tenant_id),
"user_id": str(user_id),
})
events.append({
"event": "contact.merged",
"source_contact_id": self.source_contact_id,
"target_contact_id": self.target_contact_id,
})
return CommandResult.ok(data=result, events=events)
async def audit(self, db: AsyncSession, current_user: dict[str, Any]) -> None:
tenant_id = self._tenant_id(current_user)
user_id = self._user_id(current_user)
if self._result_data is None:
return
entry = AuditLog(
tenant_id=tenant_id,
user_id=user_id,
action="merge",
entity_type="contact",
entity_id=uuid.UUID(self.target_contact_id),
changes={
"source_contact_id": self.source_contact_id,
"target_contact_id": self.target_contact_id,
"note": self.note,
},
)
db.add(entry)
await db.flush()