from dataclasses import dataclass, field from .models import * from .ports import Connector, Publisher from .util import now @dataclass class RunResult: run_id: str; status: str; count: int=0; errors: list[str]=field(default_factory=list); cursor: SyncCursor|None=None class Orchestrator: def __init__(self, connectors: dict[str,Connector], publisher: Publisher, max_retries=3): self.connectors,self.publisher,self.max_retries=connectors,publisher,max_retries def run(self, config, cursor=None, mode="incremental", changes=None, run_id=None): run_id=run_id or f"{config.connector_id}:{now()}" for attempt in range(self.max_retries): try: docs,new_cursor=(self.connectors[config.connector_id].full(config) if mode=="full" else self.connectors[config.connector_id].incremental(config,cursor or SyncCursor(config.connector_id,"incremental"),changes)) self.publisher.upsert(docs); return RunResult(run_id,"succeeded",len(docs),cursor=new_cursor) except Exception as e: if attempt==self.max_retries-1: return RunResult(run_id,"dead_letter",errors=[str(e)]) return RunResult(run_id,"dead_letter",errors=["unreachable"]) def webhook(self,config,payload,headers): return self.connectors[config.connector_id].webhook(config,payload,headers)