iop/apps/edge/internal/input/manager.go
toki b4124f0bd6 feat: CLI setup, edge/node transport refactor, and infrastructure updates
- Add CLI core setup for edge and node services
- Refactor edge transport layer (server, integration tests)
- Refactor node transport layer (parser, session, heartbeat, client)
- Add main_test.go files for edge and node commands
- Add input package for edge service
- Add go.work and go.work.sum for workspace support
- Update configs, docs, and project rules
2026-05-20 16:37:42 +09:00

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/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
}