Refactor calendar outbox delivery paths

This commit is contained in:
2026-07-21 12:30:03 +02:00
parent 760a87ab99
commit 30f6d79f9e

View File

@@ -858,15 +858,18 @@ def _fail_operation(
session.flush() session.flush()
def execute_calendar_outbox_operation( def _lock_leased_outbox_operation(
session: Session, session: Session,
*, *,
operation_id: str, operation_id: str,
lease_token: str, lease_token: str,
client_factory: Callable[[Session, CalendarSyncSource], object] | None = None, ) -> tuple[CalendarOutboxOperation, CalendarSyncSource | None]:
) -> CalendarOutboxOperation:
operation = session.get(CalendarOutboxOperation, operation_id) operation = session.get(CalendarOutboxOperation, operation_id)
if operation is None or operation.status != "in_progress" or operation.lease_token != lease_token: if (
operation is None
or operation.status != "in_progress"
or operation.lease_token != lease_token
):
raise ValueError("Calendar outbox lease is no longer valid") raise ValueError("Calendar outbox lease is no longer valid")
source = ( source = (
session.query(CalendarSyncSource) session.query(CalendarSyncSource)
@@ -882,36 +885,86 @@ def execute_calendar_outbox_operation(
) )
if operation.status != "in_progress" or operation.lease_token != lease_token: if operation.status != "in_progress" or operation.lease_token != lease_token:
raise ValueError("Calendar outbox lease is no longer valid") raise ValueError("Calendar outbox lease is no longer valid")
operation_metadata = operation.metadata_ or {} return operation, source
unavailable_reason = _operation_source_unavailable_reason(operation, source)
if unavailable_reason:
def _cancel_undeliverable_operation(
session: Session,
*,
operation: CalendarOutboxOperation,
source: CalendarSyncSource | None,
reason: str,
) -> None:
if source is not None and source.tenant_id == operation.tenant_id: if source is not None and source.tenant_id == operation.tenant_id:
_fail_operation( _fail_operation(
session, session,
operation=operation, operation=operation,
source=source, source=source,
error=unavailable_reason, error=reason,
terminal_status="cancelled", terminal_status="cancelled",
) )
else: return
operation.status = "cancelled" operation.status = "cancelled"
operation.completed_at = utcnow() operation.completed_at = utcnow()
operation.last_error = unavailable_reason operation.last_error = reason
operation.lease_token = None operation.lease_token = None
operation.lease_expires_at = None operation.lease_expires_at = None
session.flush() session.flush()
return operation
overwrite = bool(operation_metadata.get("overwrite")) if isinstance(operation_metadata, dict) else False
client: object | None = None def _calendar_outbox_client(
try: session: Session,
if client_factory is None: *,
source: CalendarSyncSource,
client_factory: Callable[[Session, CalendarSyncSource], object] | None,
) -> object:
if client_factory is not None:
return client_factory(session, source)
from govoplan_calendar.backend.service import caldav_client_for_source from govoplan_calendar.backend.service import caldav_client_for_source
client = caldav_client_for_source(session, source) return caldav_client_for_source(session, source)
else:
client = client_factory(session, source)
if operation.operation_kind == "put": def _try_remote_match(
client: object,
operation: CalendarOutboxOperation,
) -> tuple[bool, str | None] | None:
try:
return _remote_matches_operation(client, operation)
except Exception:
return None
def _complete_if_remote_match_is_usable(
session: Session,
*,
operation: CalendarOutboxOperation,
source: CalendarSyncSource,
remote_match: tuple[bool, str | None] | None,
) -> bool:
if remote_match is None:
return False
matched, remote_etag = remote_match
if not matched or (operation.operation_kind != "delete" and not remote_etag):
return False
_complete_operation(
session,
operation=operation,
source=source,
remote_etag=remote_etag,
reconciled=True,
)
return True
def _execute_calendar_outbox_put(
session: Session,
*,
operation: CalendarOutboxOperation,
source: CalendarSyncSource,
client: object,
overwrite: bool,
) -> None:
result = client.put_object( result = client.put_object(
operation.resource_href, operation.resource_href,
operation.payload_ics or "", operation.payload_ics or "",
@@ -920,41 +973,48 @@ def execute_calendar_outbox_operation(
overwrite=overwrite, overwrite=overwrite,
) )
remote_etag = result.etag remote_etag = result.etag
reconciled = False if remote_etag:
if not remote_etag: _complete_operation(
try: session,
matched, remote_etag = _remote_matches_operation(client, operation) operation=operation,
except Exception: source=source,
matched, remote_etag = False, None remote_etag=remote_etag,
if not matched or not remote_etag: reconciled=False,
)
return
remote_match = _try_remote_match(client, operation)
if _complete_if_remote_match_is_usable(
session,
operation=operation,
source=source,
remote_match=remote_match,
):
return
_fail_operation( _fail_operation(
session, session,
operation=operation, operation=operation,
source=source, source=source,
error="CalDAV PUT succeeded without an ETag and could not be reconciled safely", error="CalDAV PUT succeeded without an ETag and could not be reconciled safely",
) )
return operation
reconciled = True
_complete_operation(
session,
operation=operation,
source=source,
remote_etag=remote_etag,
reconciled=reconciled,
)
return operation
def _execute_calendar_outbox_delete(
session: Session,
*,
operation: CalendarOutboxOperation,
source: CalendarSyncSource,
client: object,
overwrite: bool,
) -> None:
if not operation.expected_etag and not overwrite: if not operation.expected_etag and not overwrite:
matched, remote_etag = _remote_matches_operation(client, operation) remote_match = _remote_matches_operation(client, operation)
if matched: if _complete_if_remote_match_is_usable(
_complete_operation(
session, session,
operation=operation, operation=operation,
source=source, source=source,
remote_etag=remote_etag, remote_match=remote_match,
reconciled=True, ):
) return
return operation
_fail_operation( _fail_operation(
session, session,
operation=operation, operation=operation,
@@ -965,7 +1025,7 @@ def execute_calendar_outbox_operation(
), ),
terminal_status="conflict", terminal_status="conflict",
) )
return operation return
result = client.delete_object( result = client.delete_object(
operation.resource_href, operation.resource_href,
etag=operation.expected_etag, etag=operation.expected_etag,
@@ -978,25 +1038,28 @@ def execute_calendar_outbox_operation(
remote_etag=result.etag, remote_etag=result.etag,
reconciled=result.status == 404, reconciled=result.status == 404,
) )
return operation
except CalDAVPreconditionFailed as exc:
if client is None: def _handle_outbox_precondition_failure(
_fail_operation(session, operation=operation, source=source, error=exc) session: Session,
return operation *,
try: operation: CalendarOutboxOperation,
matched, remote_etag = _remote_matches_operation(client, operation) source: CalendarSyncSource,
except Exception: client: object | None,
_fail_operation(session, operation=operation, source=source, error=exc) error: CalDAVPreconditionFailed,
else: ) -> None:
if matched and (operation.operation_kind == "delete" or remote_etag): remote_match = _try_remote_match(client, operation) if client is not None else None
_complete_operation( if _complete_if_remote_match_is_usable(
session, session,
operation=operation, operation=operation,
source=source, source=source,
remote_etag=remote_etag, remote_match=remote_match,
reconciled=True, ):
) return
else: if remote_match is None:
_fail_operation(session, operation=operation, source=source, error=error)
return
matched, _remote_etag = remote_match
_fail_operation( _fail_operation(
session, session,
operation=operation, operation=operation,
@@ -1008,24 +1071,95 @@ def execute_calendar_outbox_operation(
), ),
terminal_status=None if matched else "conflict", terminal_status=None if matched else "conflict",
) )
return operation
except Exception as exc:
matched, remote_etag = False, None def _handle_outbox_delivery_failure(
if client is not None: session: Session,
try: *,
matched, remote_etag = _remote_matches_operation(client, operation) operation: CalendarOutboxOperation,
except Exception: source: CalendarSyncSource,
matched, remote_etag = False, None client: object | None,
if matched and (operation.operation_kind == "delete" or remote_etag): error: Exception,
_complete_operation( ) -> None:
remote_match = _try_remote_match(client, operation) if client is not None else None
if _complete_if_remote_match_is_usable(
session, session,
operation=operation, operation=operation,
source=source, source=source,
remote_etag=remote_etag, remote_match=remote_match,
reconciled=True, ):
return
_fail_operation(session, operation=operation, source=source, error=error)
def execute_calendar_outbox_operation(
session: Session,
*,
operation_id: str,
lease_token: str,
client_factory: Callable[[Session, CalendarSyncSource], object] | None = None,
) -> CalendarOutboxOperation:
operation, source = _lock_leased_outbox_operation(
session,
operation_id=operation_id,
lease_token=lease_token,
)
unavailable_reason = _operation_source_unavailable_reason(operation, source)
if unavailable_reason:
_cancel_undeliverable_operation(
session,
operation=operation,
source=source,
reason=unavailable_reason,
)
return operation
assert source is not None
operation_metadata = operation.metadata_ or {}
overwrite = (
bool(operation_metadata.get("overwrite"))
if isinstance(operation_metadata, dict)
else False
)
client: object | None = None
try:
client = _calendar_outbox_client(
session,
source=source,
client_factory=client_factory,
)
if operation.operation_kind == "put":
_execute_calendar_outbox_put(
session,
operation=operation,
source=source,
client=client,
overwrite=overwrite,
) )
else: else:
_fail_operation(session, operation=operation, source=source, error=exc) _execute_calendar_outbox_delete(
session,
operation=operation,
source=source,
client=client,
overwrite=overwrite,
)
except CalDAVPreconditionFailed as exc:
_handle_outbox_precondition_failure(
session,
operation=operation,
source=source,
client=client,
error=exc,
)
except Exception as exc:
_handle_outbox_delivery_failure(
session,
operation=operation,
source=source,
client=client,
error=exc,
)
return operation return operation