fastapi-gsap/gsap_broker/routers/connectors.py
Tyler J King 871541f0eb feat(connectors): add Intune device management connector
Implements ConnectorPlugin for Intune Graph API operations.
Governed invocation: every Intune call requires an active AC
and emits a Chronicle CONNECTOR_INVOKED event.
Operations: list, get, compliance check, sync, lock, retire, wipe.
In-memory compliance cache with configurable TTL.
Conditional registration via intune_enabled setting.

Signed-off-by: Tyler King <tking@guildhouse.dev>
2026-04-14 05:21:47 -04:00

100 lines
3.3 KiB
Python

"""Connectors router — governed connector catalogue and invocation."""
from __future__ import annotations
from typing import Any
from fastapi import APIRouter, HTTPException
from pydantic import BaseModel
from gsap_broker.connectors.base import ConnectorContext
from gsap_broker.connectors.registry import ConnectorRegistry
from gsap_broker.connectors.runtime import ConnectorRuntime
from gsap_broker.connectors.examples.echo_connector import EchoConnector
from gsap_broker.settings import settings
router = APIRouter()
# Module-level registry and runtime
_registry = ConnectorRegistry()
_runtime = ConnectorRuntime(registry=_registry)
# Register built-in connectors
_registry.register(EchoConnector())
# Conditionally register Intune connector
if settings.intune_enabled and settings.entra_client_secret:
from gsap_broker.intune.graph_client import GraphClient
from gsap_broker.intune.device_cache import DeviceComplianceCache
from gsap_broker.connectors.intune import IntuneConnector
_intune_graph = GraphClient(
tenant_id=settings.entra_tenant_id,
client_id=settings.entra_client_id,
client_secret=settings.entra_client_secret,
)
_intune_cache = DeviceComplianceCache(ttl_seconds=settings.intune_compliance_cache_ttl)
_intune_connector = IntuneConnector(graph_client=_intune_graph, cache=_intune_cache)
_registry.register(_intune_connector)
class InvokeRequest(BaseModel):
operation: str
parameters: dict[str, Any] = {}
chronicle_session_id: str = ""
gsap_context_id: str = ""
pipeline_run_id: str = ""
dag_id: str = ""
class InvokeResponse(BaseModel):
success: bool
data: Any = None
error: str | None = None
lineage_cid: str = ""
connector_id: str = ""
operation: str = ""
@router.get("/")
async def catalogue() -> list[dict]:
return _registry.catalogue()
@router.get("/{connector_id}/")
async def get_descriptor(connector_id: str) -> dict:
connector = _registry.get(connector_id)
if connector is None:
raise HTTPException(status_code=404, detail=f"Connector not found: {connector_id}")
return connector.descriptor()
@router.post("/{connector_id}/invoke/")
async def invoke_connector(connector_id: str, body: InvokeRequest) -> InvokeResponse:
connector = _registry.get(connector_id)
if connector is None:
raise HTTPException(status_code=404, detail=f"Connector not found: {connector_id}")
ctx = ConnectorContext(
chronicle_session_id=body.chronicle_session_id,
gsap_context_id=body.gsap_context_id,
pipeline_run_id=body.pipeline_run_id,
dag_id=body.dag_id,
)
result = await _runtime.invoke(connector_id, body.operation, body.parameters, ctx)
return InvokeResponse(
success=result.success,
data=result.data,
error=result.error,
lineage_cid=result.lineage_cid,
connector_id=connector_id,
operation=body.operation,
)
@router.get("/{connector_id}/health/")
async def health_check(connector_id: str) -> dict:
connector = _registry.get(connector_id)
if connector is None:
raise HTTPException(status_code=404, detail=f"Connector not found: {connector_id}")
healthy = connector.health_check()
return {"connector_id": connector_id, "healthy": healthy}