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)
def get_instance(self, orchestrator_id: str) -> Optional[Any]:
66    def get_instance(self, orchestrator_id: str) -> Optional[Any]:
67        return self._instances.get(orchestrator_id)
def list_instances(self) -> List[Any]:
69    def list_instances(self) -> List[Any]:
70        return list(self._instances.values())
async def get_entry( self, orchestrator_id: str) -> Optional[gmf_forge_ai_orchestration.api.schemas.OrchestratorRegistryEntry]:
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)
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 heartbeat(self, orchestrator_id: str) -> None:
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)
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.