22 lines
1.3 KiB
Python
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)
|