Files

22 lines
1.3 KiB
Python

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)