48 lines
1.2 KiB
Go
48 lines
1.2 KiB
Go
package input
|
|
|
|
import (
|
|
"context"
|
|
"time"
|
|
|
|
"go.uber.org/zap"
|
|
|
|
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
|
|
}
|
|
|
|
// 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"))
|
|
a2aServer := edgea2a.NewServer(cfg.A2A, svc, logger.Named("a2a"))
|
|
return &Manager{OpenAI: openaiServer, A2A: a2aServer}
|
|
}
|
|
|
|
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
|
|
}
|