package input import ( "context" "time" "go.uber.org/zap" "iop/apps/edge/internal/authprojection" edgea2a "iop/apps/edge/internal/input/a2a" edgeopenai "iop/apps/edge/internal/openai" edgeservice "iop/apps/edge/internal/service" "iop/packages/go/config" ) // Manager owns the lifecycle of all Edge inbound input servers (OpenAI-compatible and A2A). type Manager struct { OpenAI *edgeopenai.Server A2A *edgea2a.Server principalProjection *authprojection.Cache } // NewManager creates a Manager wiring both input servers. func NewManager(cfg config.EdgeConfig, svc *edgeservice.Service, logger *zap.Logger) *Manager { openaiServer := edgeopenai.NewServer(cfg.OpenAI, svc, logger.Named("openai")) openaiServer.SetCredentialPlaneManaged(cfg.CredentialPlane.Mode() == config.CredentialPlaneModeManaged) var projection *authprojection.Cache if cfg.CredentialPlane.Mode() == config.CredentialPlaneModeManaged { projection = authprojection.NewCache(authprojection.DefaultLimits(), time.Now) openaiServer.SetPrincipalProjection(projection) } openaiServer.SetEdgeID(cfg.Edge.ID) openaiServer.SetModelCatalog(cfg.Models) openaiServer.SetLongContextThreshold(cfg.LongContextThresholdTokens) a2aServer := edgea2a.NewServer(cfg.A2A, svc, logger.Named("a2a")) return &Manager{OpenAI: openaiServer, A2A: a2aServer, principalProjection: projection} } // PrincipalProjection exposes the shared cache to the authenticated Control // Plane connector. Managed ingress and the pre-send lease fence therefore read // the same immutable generation. func (m *Manager) PrincipalProjection() *authprojection.Cache { if m == nil { return nil } return m.principalProjection } func (m *Manager) SetModelCatalog(catalog []config.ModelCatalogEntry) { if m == nil || m.OpenAI == nil { return } m.OpenAI.SetModelCatalog(catalog) } func (m *Manager) Start(ctx context.Context) error { if err := m.OpenAI.Start(ctx); err != nil { return err } if err := m.A2A.Start(ctx); err != nil { stopCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second) defer cancel() _ = m.OpenAI.Stop(stopCtx) return err } return nil } func (m *Manager) Stop(ctx context.Context) error { a2aErr := m.A2A.Stop(ctx) openaiErr := m.OpenAI.Stop(ctx) if a2aErr != nil { return a2aErr } return openaiErr }