gmf_forge_ai_orchestration.api
Public API exports for orchestrator observability routes and registry.
1"""Public API exports for orchestrator observability routes and registry.""" 2 3from typing import Any 4 5from gmf_forge_ai_orchestration.api.orchestrator_registry import OrchestratorRegistry 6 7 8def build_orchestrator_api_router(*args: Any, **kwargs: Any) -> Any: 9 """Lazily import FastAPI router builder so FastAPI stays an app-level dependency.""" 10 from gmf_forge_ai_orchestration.api.routes import build_orchestrator_api_router as _builder 11 12 return _builder(*args, **kwargs) 13 14 15def build_agent_checkpoint_router(*args: Any, **kwargs: Any) -> Any: 16 """Lazily import agent checkpoint router so FastAPI stays an app-level dependency.""" 17 from gmf_forge_ai_orchestration.api.routes import build_agent_checkpoint_router as _builder 18 19 return _builder(*args, **kwargs) 20 21 22def build_agent_metrics_router(*args: Any, **kwargs: Any) -> Any: 23 """Lazily import agent metrics router so FastAPI stays an app-level dependency.""" 24 from gmf_forge_ai_orchestration.api.routes import build_agent_metrics_router as _builder 25 26 return _builder(*args, **kwargs) 27 28 29def build_agent_metadata_router(*args: Any, **kwargs: Any) -> Any: 30 """Lazily import agent metadata router so FastAPI stays an app-level dependency.""" 31 from gmf_forge_ai_orchestration.api.agent_metadata_routes import build_agent_metadata_router as _builder 32 33 return _builder(*args, **kwargs) 34 35 36__all__ = [ 37 "OrchestratorRegistry", 38 "build_orchestrator_api_router", 39 "build_agent_checkpoint_router", 40 "build_agent_metrics_router", 41 "build_agent_metadata_router", 42]
class
OrchestratorRegistry:
16class OrchestratorRegistry: 17 """Tracks live orchestrator objects and persists discoverable metadata.""" 18 19 def __init__( 20 self, 21 store: BaseStateStore, 22 default_ttl: Optional[int] = None, 23 run_ttl: int = 86400 * 7, 24 max_runs: int = 200, 25 ) -> None: 26 self._store = store 27 self._default_ttl = default_ttl 28 self._run_ttl = run_ttl 29 self._max_runs = max_runs 30 self._instances: Dict[str, Any] = {} 31 32 def _entry_key(self, orchestrator_id: str) -> str: 33 return f"{_REGISTRY_DATA_KEY_PREFIX}{orchestrator_id}" 34 35 async def register(self, orchestrator: Any) -> None: 36 orchestrator_id = orchestrator.get_orchestrator_id() 37 self._instances[orchestrator_id] = orchestrator 38 now = datetime.now(timezone.utc) 39 40 entry = OrchestratorRegistryEntry( 41 orchestrator_id=orchestrator_id, 42 orchestrator_type=orchestrator.get_orchestrator_type(), 43 base_url=None, 44 status="online", 45 registered_at=now, 46 last_heartbeat=now, 47 tenant_tags={}, 48 ) 49 await self._upsert_entry(entry) 50 51 async def unregister(self, orchestrator_id: str, remove_entry: bool = False) -> None: 52 self._instances.pop(orchestrator_id, None) 53 if remove_entry: 54 await self._store.delete(self._entry_key(orchestrator_id)) 55 index: List[str] = await self._store.get(_REGISTRY_INDEX_KEY) or [] 56 updated = [oid for oid in index if oid != orchestrator_id] 57 await self._store.set(_REGISTRY_INDEX_KEY, updated, ttl=self._default_ttl) 58 return 59 60 entry = await self.get_entry(orchestrator_id) 61 if entry is not None: 62 entry.status = "offline" 63 entry.last_heartbeat = datetime.now(timezone.utc) 64 await self._upsert_entry(entry) 65 66 def get_instance(self, orchestrator_id: str) -> Optional[Any]: 67 return self._instances.get(orchestrator_id) 68 69 def list_instances(self) -> List[Any]: 70 return list(self._instances.values()) 71 72 async def get_entry(self, orchestrator_id: str) -> Optional[OrchestratorRegistryEntry]: 73 raw = await self._store.get(self._entry_key(orchestrator_id)) 74 if raw is None: 75 return None 76 return OrchestratorRegistryEntry(**raw) 77 78 async def list_entries(self) -> List[OrchestratorRegistryEntry]: 79 index: List[str] = await self._store.get(_REGISTRY_INDEX_KEY) or [] 80 entries: List[OrchestratorRegistryEntry] = [] 81 for orchestrator_id in index: 82 entry = await self.get_entry(orchestrator_id) 83 if entry is not None: 84 entries.append(entry) 85 return entries 86 87 async def heartbeat(self, orchestrator_id: str) -> None: 88 entry = await self.get_entry(orchestrator_id) 89 if not entry: 90 return 91 entry.last_heartbeat = datetime.now(timezone.utc) 92 if orchestrator_id in self._instances: 93 entry.status = "online" 94 await self._upsert_entry(entry) 95 96 async def register_url(self, base_url: str, tags: Optional[Dict[str, str]] = None) -> Dict[str, Any]: 97 tags = tags or {} 98 url = base_url.rstrip("/") + "/.well-known/agent.json" 99 async with httpx.AsyncClient(timeout=8.0) as client: 100 response = await client.get(url) 101 response.raise_for_status() 102 data = response.json() 103 104 orchestrator_id = str(data.get("orchestrator_id") or data.get("id") or "") 105 if not orchestrator_id: 106 raise ValueError("Missing orchestrator_id in .well-known/agent.json response") 107 108 existing = await self.get_entry(orchestrator_id) 109 already_existed = existing is not None 110 now = datetime.now(timezone.utc) 111 112 entry = OrchestratorRegistryEntry( 113 orchestrator_id=orchestrator_id, 114 orchestrator_type=str(data.get("orchestrator_type") or data.get("type") or "unknown"), 115 base_url=base_url, 116 status="online", 117 registered_at=(existing.registered_at if existing else now), 118 last_heartbeat=now, 119 tenant_tags=tags, 120 ) 121 await self._upsert_entry(entry) 122 123 return { 124 "orchestrator_id": orchestrator_id, 125 "entry": entry, 126 "already_existed": already_existed, 127 } 128 129 async def _upsert_entry(self, entry: OrchestratorRegistryEntry) -> None: 130 key = self._entry_key(entry.orchestrator_id) 131 await self._store.set(key, entry.model_dump(mode="json"), ttl=self._default_ttl) 132 133 index: List[str] = await self._store.get(_REGISTRY_INDEX_KEY) or [] 134 if entry.orchestrator_id not in index: 135 index.append(entry.orchestrator_id) 136 await self._store.set(_REGISTRY_INDEX_KEY, index, ttl=self._default_ttl) 137 138 # ── Run persistence ──────────────────────────────────────────────────────── 139 140 def _run_key(self, orchestrator_id: str, run_id: str) -> str: 141 return f"__run__{orchestrator_id}__{run_id}" 142 143 def _run_index_key(self, orchestrator_id: str) -> str: 144 return f"__run_index__{orchestrator_id}" 145 146 async def save_run(self, orchestrator_id: str, payload: Dict[str, Any]) -> None: 147 """Persist a completed/failed run record via the configured state store.""" 148 run_id = payload["run_id"] 149 await self._store.set(self._run_key(orchestrator_id, run_id), payload, ttl=self._run_ttl) 150 index: List[str] = await self._store.get(self._run_index_key(orchestrator_id)) or [] 151 if run_id not in index: 152 index.insert(0, run_id) 153 if len(index) > self._max_runs: 154 evicted = index[self._max_runs:] 155 index = index[:self._max_runs] 156 for evicted_id in evicted: 157 await self._store.delete(self._run_key(orchestrator_id, evicted_id)) 158 await self._store.set(self._run_index_key(orchestrator_id), index) 159 160 async def list_persisted_runs(self, orchestrator_id: str) -> List[Dict[str, Any]]: 161 """Return all persisted run records for an orchestrator (newest first).""" 162 index: List[str] = await self._store.get(self._run_index_key(orchestrator_id)) or [] 163 runs = [] 164 for run_id in index: 165 raw = await self._store.get(self._run_key(orchestrator_id, run_id)) 166 if raw: 167 runs.append(raw) 168 return runs 169 170 async def get_persisted_run(self, orchestrator_id: str, run_id: str) -> Optional[Dict[str, Any]]: 171 """Return a single persisted run record, or None if not found.""" 172 return await self._store.get(self._run_key(orchestrator_id, run_id))
Tracks live orchestrator objects and persists discoverable metadata.
OrchestratorRegistry( store: gmf_forge_ai_orchestration.BaseStateStore, default_ttl: Optional[int] = None, run_ttl: int = 604800, max_runs: int = 200)
19 def __init__( 20 self, 21 store: BaseStateStore, 22 default_ttl: Optional[int] = None, 23 run_ttl: int = 86400 * 7, 24 max_runs: int = 200, 25 ) -> None: 26 self._store = store 27 self._default_ttl = default_ttl 28 self._run_ttl = run_ttl 29 self._max_runs = max_runs 30 self._instances: Dict[str, Any] = {}
async def
register(self, orchestrator: Any) -> None:
35 async def register(self, orchestrator: Any) -> None: 36 orchestrator_id = orchestrator.get_orchestrator_id() 37 self._instances[orchestrator_id] = orchestrator 38 now = datetime.now(timezone.utc) 39 40 entry = OrchestratorRegistryEntry( 41 orchestrator_id=orchestrator_id, 42 orchestrator_type=orchestrator.get_orchestrator_type(), 43 base_url=None, 44 status="online", 45 registered_at=now, 46 last_heartbeat=now, 47 tenant_tags={}, 48 ) 49 await self._upsert_entry(entry)
async def
unregister(self, orchestrator_id: str, remove_entry: bool = False) -> None:
51 async def unregister(self, orchestrator_id: str, remove_entry: bool = False) -> None: 52 self._instances.pop(orchestrator_id, None) 53 if remove_entry: 54 await self._store.delete(self._entry_key(orchestrator_id)) 55 index: List[str] = await self._store.get(_REGISTRY_INDEX_KEY) or [] 56 updated = [oid for oid in index if oid != orchestrator_id] 57 await self._store.set(_REGISTRY_INDEX_KEY, updated, ttl=self._default_ttl) 58 return 59 60 entry = await self.get_entry(orchestrator_id) 61 if entry is not None: 62 entry.status = "offline" 63 entry.last_heartbeat = datetime.now(timezone.utc) 64 await self._upsert_entry(entry)
async def
get_entry( self, orchestrator_id: str) -> Optional[gmf_forge_ai_orchestration.api.schemas.OrchestratorRegistryEntry]:
async def
list_entries( self) -> List[gmf_forge_ai_orchestration.api.schemas.OrchestratorRegistryEntry]:
78 async def list_entries(self) -> List[OrchestratorRegistryEntry]: 79 index: List[str] = await self._store.get(_REGISTRY_INDEX_KEY) or [] 80 entries: List[OrchestratorRegistryEntry] = [] 81 for orchestrator_id in index: 82 entry = await self.get_entry(orchestrator_id) 83 if entry is not None: 84 entries.append(entry) 85 return entries
async def
register_url( self, base_url: str, tags: Optional[Dict[str, str]] = None) -> Dict[str, Any]:
96 async def register_url(self, base_url: str, tags: Optional[Dict[str, str]] = None) -> Dict[str, Any]: 97 tags = tags or {} 98 url = base_url.rstrip("/") + "/.well-known/agent.json" 99 async with httpx.AsyncClient(timeout=8.0) as client: 100 response = await client.get(url) 101 response.raise_for_status() 102 data = response.json() 103 104 orchestrator_id = str(data.get("orchestrator_id") or data.get("id") or "") 105 if not orchestrator_id: 106 raise ValueError("Missing orchestrator_id in .well-known/agent.json response") 107 108 existing = await self.get_entry(orchestrator_id) 109 already_existed = existing is not None 110 now = datetime.now(timezone.utc) 111 112 entry = OrchestratorRegistryEntry( 113 orchestrator_id=orchestrator_id, 114 orchestrator_type=str(data.get("orchestrator_type") or data.get("type") or "unknown"), 115 base_url=base_url, 116 status="online", 117 registered_at=(existing.registered_at if existing else now), 118 last_heartbeat=now, 119 tenant_tags=tags, 120 ) 121 await self._upsert_entry(entry) 122 123 return { 124 "orchestrator_id": orchestrator_id, 125 "entry": entry, 126 "already_existed": already_existed, 127 }
async def
save_run(self, orchestrator_id: str, payload: Dict[str, Any]) -> None:
146 async def save_run(self, orchestrator_id: str, payload: Dict[str, Any]) -> None: 147 """Persist a completed/failed run record via the configured state store.""" 148 run_id = payload["run_id"] 149 await self._store.set(self._run_key(orchestrator_id, run_id), payload, ttl=self._run_ttl) 150 index: List[str] = await self._store.get(self._run_index_key(orchestrator_id)) or [] 151 if run_id not in index: 152 index.insert(0, run_id) 153 if len(index) > self._max_runs: 154 evicted = index[self._max_runs:] 155 index = index[:self._max_runs] 156 for evicted_id in evicted: 157 await self._store.delete(self._run_key(orchestrator_id, evicted_id)) 158 await self._store.set(self._run_index_key(orchestrator_id), index)
Persist a completed/failed run record via the configured state store.
async def
list_persisted_runs(self, orchestrator_id: str) -> List[Dict[str, Any]]:
160 async def list_persisted_runs(self, orchestrator_id: str) -> List[Dict[str, Any]]: 161 """Return all persisted run records for an orchestrator (newest first).""" 162 index: List[str] = await self._store.get(self._run_index_key(orchestrator_id)) or [] 163 runs = [] 164 for run_id in index: 165 raw = await self._store.get(self._run_key(orchestrator_id, run_id)) 166 if raw: 167 runs.append(raw) 168 return runs
Return all persisted run records for an orchestrator (newest first).
async def
get_persisted_run(self, orchestrator_id: str, run_id: str) -> Optional[Dict[str, Any]]:
170 async def get_persisted_run(self, orchestrator_id: str, run_id: str) -> Optional[Dict[str, Any]]: 171 """Return a single persisted run record, or None if not found.""" 172 return await self._store.get(self._run_key(orchestrator_id, run_id))
Return a single persisted run record, or None if not found.
def
build_orchestrator_api_router(*args: Any, **kwargs: Any) -> Any:
9def build_orchestrator_api_router(*args: Any, **kwargs: Any) -> Any: 10 """Lazily import FastAPI router builder so FastAPI stays an app-level dependency.""" 11 from gmf_forge_ai_orchestration.api.routes import build_orchestrator_api_router as _builder 12 13 return _builder(*args, **kwargs)
Lazily import FastAPI router builder so FastAPI stays an app-level dependency.
def
build_agent_checkpoint_router(*args: Any, **kwargs: Any) -> Any:
16def build_agent_checkpoint_router(*args: Any, **kwargs: Any) -> Any: 17 """Lazily import agent checkpoint router so FastAPI stays an app-level dependency.""" 18 from gmf_forge_ai_orchestration.api.routes import build_agent_checkpoint_router as _builder 19 20 return _builder(*args, **kwargs)
Lazily import agent checkpoint router so FastAPI stays an app-level dependency.
def
build_agent_metrics_router(*args: Any, **kwargs: Any) -> Any:
23def build_agent_metrics_router(*args: Any, **kwargs: Any) -> Any: 24 """Lazily import agent metrics router so FastAPI stays an app-level dependency.""" 25 from gmf_forge_ai_orchestration.api.routes import build_agent_metrics_router as _builder 26 27 return _builder(*args, **kwargs)
Lazily import agent metrics router so FastAPI stays an app-level dependency.
def
build_agent_metadata_router(*args: Any, **kwargs: Any) -> Any:
30def build_agent_metadata_router(*args: Any, **kwargs: Any) -> Any: 31 """Lazily import agent metadata router so FastAPI stays an app-level dependency.""" 32 from gmf_forge_ai_orchestration.api.agent_metadata_routes import build_agent_metadata_router as _builder 33 34 return _builder(*args, **kwargs)
Lazily import agent metadata router so FastAPI stays an app-level dependency.