tmux-amq/src/tamq/broker.py

72 lines
2.3 KiB
Python
Raw Normal View History

from __future__ import annotations
import re
from dataclasses import dataclass
from .routing import RoutedMessage, parse_address_line
from .store import Store
from .registry import RegistryError, validate_targets
LEGACY_DELIVERY_SUFFIX = re.compile(
r"\s+\[(m-[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12})\]$",
re.IGNORECASE,
)
@dataclass(frozen=True)
class BrokerIdentity:
endpoint_id: str
source_repo: str
class InputBroker:
"""Convert tap input into durable, identity-bearing outbound messages."""
def __init__(self, store: Store, identity: BrokerIdentity):
self.store = store
self.identity = identity
def inspect_line(self, line: str) -> RoutedMessage | None:
routed = parse_address_line(line)
if routed is None:
return None
receipt = LEGACY_DELIVERY_SUFFIX.search(routed.body)
if line.startswith("#") and receipt:
prior = self.store.message(receipt.group(1))
if (
prior is not None
and prior["sender_repo"] == routed.target_repo
and prior["target_repo"] == self.identity.source_repo
):
return None
try:
validate_targets([routed.target_repo])
except RegistryError:
return None
self.store.add(
self.identity.source_repo,
routed.target_repo,
routed.body,
endpoint=self.identity.endpoint_id,
)
return routed
def deliver_pending(self, control, *, window_for_repo):
"""Inject pending messages through the authoritative control client."""
delivered = 0
for row in self.store.list(state="pending"):
if row["endpoint_id"] not in (None, self.identity.endpoint_id):
continue
target = window_for_repo(row["target_repo"])
lease_id = self.store.claim(row["message_id"], self.identity.endpoint_id)
if lease_id is None:
continue
try:
control.inject(target, f"#{row['sender_repo']}: {row['body']}")
except Exception:
continue
self.store.release(row["message_id"], lease_id, "injected")
delivered += 1
return delivered