235 lines
14 KiB
Python
235 lines
14 KiB
Python
"""Explicit, fenced recovery of claims left by a proven stopped runtime."""
|
|
from datetime import datetime, timezone
|
|
import hashlib
|
|
import json
|
|
from typing import Any
|
|
|
|
from sqlalchemy.orm import Session
|
|
|
|
from govoplan_core.core.runtime_coordination import (
|
|
DistributedLease, RuntimeNode, acquire_lease, release_lease, process_runtime_identity,
|
|
)
|
|
from govoplan_core.core.recovery import (
|
|
RecoveryOperation, RecoveryStatus, record_recovery_checkpoint,
|
|
transition_recovery_operation, verify_recovery_evidence_chain,
|
|
)
|
|
from govoplan_campaign.backend.db.models import CampaignJob, SendAttempt, ImapAppendAttempt
|
|
from govoplan_campaign.backend.sending.jobs import QueueingError, _update_campaign_after_job
|
|
|
|
|
|
class RecoveryStateConflict(QueueingError):
|
|
pass
|
|
|
|
|
|
def _utc(value: datetime) -> datetime:
|
|
return value.replace(tzinfo=timezone.utc) if value.tzinfo is None else value.astimezone(timezone.utc)
|
|
|
|
|
|
def _key(job: CampaignJob, channel: str) -> str:
|
|
return f"campaign:{'delivery' if channel == 'smtp' else 'imap'}:{job.tenant_id}:{job.id}"
|
|
|
|
|
|
def _claim_metadata(job: CampaignJob, channel: str, lease: DistributedLease | None, node: RuntimeNode | None) -> dict[str, Any]:
|
|
state = job.send_status if channel == "smtp" else job.imap_status
|
|
claim = job.claim_token if channel == "smtp" else job.imap_claim_token
|
|
active = state in ({"claimed", "sending"} if channel == "smtp" else {"appending"})
|
|
reason = "not_active"
|
|
if active:
|
|
if lease is None or not claim or not lease.holder_node_id or node is None:
|
|
reason = "owner_not_confirmed_stopped"
|
|
elif _utc(lease.expires_at) > datetime.now(timezone.utc):
|
|
reason = "live_claim"
|
|
elif node.incarnation == lease.holder_incarnation and node.state != "stopped":
|
|
# Heartbeat age alone is NOT proof that a slow worker is dead.
|
|
reason = "owner_not_confirmed_stopped"
|
|
else:
|
|
reason = "recoverable"
|
|
revision = hashlib.sha256(json.dumps({
|
|
"job": job.id, "channel": channel, "state": state, "claim": claim,
|
|
"attempts": job.attempt_count,
|
|
"lease": [lease.id, lease.fencing_token, str(lease.expires_at), lease.holder_node_id, lease.holder_incarnation] if lease else None,
|
|
"owner": [node.incarnation, node.state] if node else None,
|
|
}, sort_keys=True).encode()).hexdigest()
|
|
return {"eligible": reason == "recoverable", "revision": revision, "reason": reason}
|
|
|
|
|
|
def job_recovery_metadata(session: Session, jobs: list[CampaignJob]) -> dict[str, dict[str, Any]]:
|
|
"""Two bounded metadata reads per loaded page, not per recipient."""
|
|
if not jobs:
|
|
return {}
|
|
installation_id = process_runtime_identity().installation_id
|
|
keys = [_key(job, channel) for job in jobs for channel in ("smtp", "imap")]
|
|
leases = session.query(DistributedLease).filter(
|
|
DistributedLease.installation_id == installation_id,
|
|
DistributedLease.resource_key.in_(keys),
|
|
).all()
|
|
by_key = {lease.resource_key: lease for lease in leases}
|
|
owner_ids = {lease.holder_node_id for lease in leases if lease.holder_node_id}
|
|
nodes = session.query(RuntimeNode).filter(
|
|
RuntimeNode.installation_id == installation_id,
|
|
RuntimeNode.node_id.in_(owner_ids),
|
|
).all() if owner_ids else []
|
|
by_owner = {node.node_id: node for node in nodes}
|
|
result = {}
|
|
for job in jobs:
|
|
channels = {}
|
|
for channel in ("smtp", "imap"):
|
|
lease = by_key.get(_key(job, channel))
|
|
channels[channel] = _claim_metadata(job, channel, lease, by_owner.get(lease.holder_node_id) if lease else None)
|
|
result[job.id] = channels
|
|
return result
|
|
|
|
|
|
def recover_stale_delivery_claim(
|
|
session: Session, *, tenant_id: str, campaign_id: str, job_id: str,
|
|
channel: str, expected_revision: str, note: str,
|
|
) -> dict[str, Any]:
|
|
"""Freeze an abandoned effect as unknown; NEVER infer that it was not sent."""
|
|
if channel not in {"smtp", "imap"} or not note.strip():
|
|
raise QueueingError("Claim recovery requires a channel and an evidence note.")
|
|
job = session.get(CampaignJob, job_id)
|
|
if job is None or job.tenant_id != tenant_id or job.campaign_id != campaign_id:
|
|
raise QueueingError("Campaign job not found or not accessible")
|
|
identity = process_runtime_identity()
|
|
resource_key = _key(job, channel)
|
|
# Same lock order as runtime authority: lease, operation, domain row.
|
|
lease = session.query(DistributedLease).filter(
|
|
DistributedLease.installation_id == identity.installation_id,
|
|
DistributedLease.resource_key == resource_key,
|
|
).with_for_update().populate_existing().one_or_none()
|
|
node = session.query(RuntimeNode).filter(
|
|
RuntimeNode.installation_id == identity.installation_id,
|
|
RuntimeNode.node_id == lease.holder_node_id,
|
|
).with_for_update().populate_existing().one_or_none() if lease and lease.holder_node_id else None
|
|
session.refresh(job)
|
|
metadata = _claim_metadata(job, channel, lease, node)
|
|
if metadata["revision"] != expected_revision:
|
|
raise RecoveryStateConflict("Delivery claim changed; reload its current evidence before recovery.")
|
|
if not metadata["eligible"]:
|
|
raise RecoveryStateConflict("Recovery is blocked until the lease expires and its owning runtime is confirmed stopped or replaced.")
|
|
claim_token = job.claim_token if channel == "smtp" else job.imap_claim_token
|
|
state = job.send_status if channel == "smtp" else job.imap_status
|
|
assert claim_token is not None
|
|
claim_sha = hashlib.sha256(claim_token.encode()).hexdigest()
|
|
operation_key = f"campaign-{'delivery' if channel == 'smtp' else 'imap'}:{job.id}:{claim_sha[:32]}"
|
|
operation = session.query(RecoveryOperation).filter(
|
|
RecoveryOperation.installation_id == identity.installation_id,
|
|
RecoveryOperation.module_id == "campaigns",
|
|
RecoveryOperation.idempotency_key == operation_key,
|
|
RecoveryOperation.lease_resource_key == resource_key,
|
|
).with_for_update().one_or_none()
|
|
if operation is None or operation.status not in {"running", "outcome_unknown"}:
|
|
raise RecoveryStateConflict("The original durable delivery evidence cannot be recovered safely.")
|
|
authority = acquire_lease(
|
|
session, installation_id=identity.installation_id, resource_key=resource_key,
|
|
holder_node_id=identity.node_id, holder_incarnation=identity.incarnation,
|
|
ttl_seconds=300, metadata={"module_id": "campaigns", "recovery_operation_id": operation.id},
|
|
)
|
|
if authority is None:
|
|
raise RecoveryStateConflict("Another runtime acquired this delivery claim.")
|
|
operation.holder_node_id = authority.holder_node_id
|
|
operation.holder_incarnation = authority.holder_incarnation
|
|
operation.fencing_token = authority.fencing_token
|
|
session.add(operation)
|
|
session.flush()
|
|
evidence = {"job_id": job.id, "channel": channel, "previous_state": state, "evidence_note_sha256": hashlib.sha256(note.strip().encode()).hexdigest(), "claim_sha256": claim_sha}
|
|
record_recovery_checkpoint(session, operation, kind="campaign-claim-recovery", summary="An operator fenced a claim owned by a stopped runtime", evidence=evidence, lease_claim=authority)
|
|
if operation.status != "outcome_unknown":
|
|
transition_recovery_operation(session, operation, status=RecoveryStatus.OUTCOME_UNKNOWN, kind="campaign-claim-outcome-unknown", summary="The abandoned provider effect requires explicit reconciliation", evidence=evidence, failure_summary="The original runtime stopped before recording a final provider result", lease_claim=authority)
|
|
if not verify_recovery_evidence_chain(session, operation.id):
|
|
raise RecoveryStateConflict("Durable delivery evidence verification failed.")
|
|
state_column = CampaignJob.send_status if channel == "smtp" else CampaignJob.imap_status
|
|
claim_column = CampaignJob.claim_token if channel == "smtp" else CampaignJob.imap_claim_token
|
|
changes = {state_column: "outcome_unknown", claim_column: None, CampaignJob.last_error: note.strip()}
|
|
if channel == "smtp":
|
|
changes.update({CampaignJob.queue_status: "draft", CampaignJob.outcome_unknown_at: datetime.now(timezone.utc)})
|
|
else:
|
|
changes[CampaignJob.imap_claimed_at] = None
|
|
changed = session.query(CampaignJob).filter(
|
|
CampaignJob.id == job.id, state_column == state, claim_column == claim_token,
|
|
).update(changes, synchronize_session=False)
|
|
if changed != 1:
|
|
raise RecoveryStateConflict("The delivery claim changed before recovery could be recorded.")
|
|
attempt_model = SendAttempt if channel == "smtp" else ImapAppendAttempt
|
|
attempt = session.query(attempt_model).filter(attempt_model.job_id == job.id, attempt_model.claim_token == claim_token).order_by(attempt_model.attempt_number.desc()).first()
|
|
if attempt is not None:
|
|
attempt.status = "outcome_unknown"
|
|
attempt.error_message = note.strip()
|
|
if channel == "smtp":
|
|
attempt.finished_at = datetime.now(timezone.utc)
|
|
session.add(attempt)
|
|
release_lease(session, authority)
|
|
session.expire(job)
|
|
_update_campaign_after_job(session, campaign_id, job.campaign_version_id)
|
|
session.flush()
|
|
return {"campaign_id": campaign_id, "version_id": job.campaign_version_id, "job_id": job.id, "channel": channel, "send_status": job.send_status, "imap_status": job.imap_status, "note": note.strip(), "reconciliation_required": True}
|
|
|
|
|
|
def reconcile_campaign_delivery_operation(
|
|
session: Session, *, job: CampaignJob, channel: str, claim_token: str | None,
|
|
effect_occurred: bool, note: str,
|
|
) -> None:
|
|
"""Resolve only the original Campaign ledger in the caller's audit transaction.
|
|
|
|
Older jobs without a claim-bound operation remain supported. Mail's nested
|
|
provider-effect ledgers are separate evidence and are never rewritten here.
|
|
"""
|
|
# A compound external-channel operation may include Postbox/Print effects.
|
|
# One SMTP decision cannot verify or negate that entire operation.
|
|
if not claim_token or (channel == "smtp" and getattr(job, "delivery_channel_policy", "mail") != "mail"):
|
|
return
|
|
identity = process_runtime_identity()
|
|
key = _key(job, channel)
|
|
claim_sha = hashlib.sha256(claim_token.encode()).hexdigest()
|
|
operation_key = f"campaign-{'delivery' if channel == 'smtp' else 'imap'}:{job.id}:{claim_sha[:32]}"
|
|
# Acquire the same lock order as effect execution and claim recovery.
|
|
lease = session.query(DistributedLease).filter(
|
|
DistributedLease.installation_id == identity.installation_id,
|
|
DistributedLease.resource_key == key,
|
|
).with_for_update().populate_existing().one_or_none()
|
|
operation = session.query(RecoveryOperation).filter(
|
|
RecoveryOperation.installation_id == identity.installation_id,
|
|
RecoveryOperation.module_id == "campaigns",
|
|
RecoveryOperation.idempotency_key == operation_key,
|
|
RecoveryOperation.lease_resource_key == key,
|
|
RecoveryOperation.resource_type == "campaign_job",
|
|
RecoveryOperation.resource_id == job.id,
|
|
).with_for_update().populate_existing().one_or_none()
|
|
if operation is None:
|
|
return
|
|
if operation.status in {"succeeded", "recovered"}:
|
|
if (operation.status == "succeeded") != effect_occurred:
|
|
raise RecoveryStateConflict("The original durable operation already records a different verified outcome.")
|
|
return
|
|
if operation.status != "outcome_unknown" or lease is None:
|
|
raise RecoveryStateConflict("The original durable operation requires guarded claim recovery before reconciliation.")
|
|
# Even the same API process must not borrow another active operation's
|
|
# lease merely because its runtime identity happens to match.
|
|
if lease.holder_node_id is not None:
|
|
raise RecoveryStateConflict("The original operation still has a runtime owner; recover its stopped claim before reconciliation.")
|
|
authority = acquire_lease(session, installation_id=identity.installation_id, resource_key=key,
|
|
holder_node_id=identity.node_id, holder_incarnation=identity.incarnation,
|
|
ttl_seconds=300, metadata={"module_id": "campaigns", "recovery_operation_id": operation.id})
|
|
if authority is None:
|
|
raise RecoveryStateConflict("Another runtime owns the original delivery operation.")
|
|
operation.holder_node_id = authority.holder_node_id
|
|
operation.holder_incarnation = authority.holder_incarnation
|
|
operation.fencing_token = authority.fencing_token
|
|
session.add(operation)
|
|
session.flush()
|
|
evidence = {"verified": True, "checks": {"operator_provider_evidence_recorded": True, "matching_claim_attempt": True},
|
|
"job_id": job.id, "channel": channel, "effect_occurred": effect_occurred,
|
|
"claim_sha256": claim_sha, "evidence_note_sha256": hashlib.sha256(note.strip().encode()).hexdigest()}
|
|
record_recovery_checkpoint(session, operation, kind="campaign-reconciliation-fence", summary="An operator acquired authority for the original Campaign attempt", evidence=evidence, lease_claim=authority)
|
|
if effect_occurred:
|
|
transition_recovery_operation(session, operation, status=RecoveryStatus.SUCCEEDED,
|
|
kind="campaign-reconciled-provider-acceptance", summary="Operator evidence confirms the Campaign effect was accepted", evidence=evidence, lease_claim=authority)
|
|
else:
|
|
for next_status in (RecoveryStatus.RECOVERY_REQUIRED, RecoveryStatus.RECOVERING, RecoveryStatus.RECOVERED):
|
|
transition_recovery_operation(session, operation, status=next_status,
|
|
kind=f"campaign-reconciled-absence-{next_status.value}", summary="Operator evidence confirms the Campaign effect did not occur",
|
|
evidence=evidence, failure_summary="The original external effect was verified absent" if next_status == RecoveryStatus.RECOVERY_REQUIRED else None, lease_claim=authority)
|
|
if not verify_recovery_evidence_chain(session, operation.id):
|
|
raise RecoveryStateConflict("Original Campaign recovery evidence verification failed.")
|
|
release_lease(session, authority)
|