iop/apps/edge/internal/bootstrap/runtime.go
toki 5bb8a4bc85 feat: edge local dev config and observability improvements
- Add edge node bootstrap and runtime configuration
- Update observability with test coverage
- Add hostsetup test and template updates
- Add m-edge-local-dev-config-runtime task tracking
2026-05-28 05:44:43 +09:00

143 lines
3.3 KiB
Go

package bootstrap
import (
"context"
"fmt"
"sync"
"time"
"go.uber.org/zap"
edgeevents "iop/apps/edge/internal/events"
edgeinput "iop/apps/edge/internal/input"
edgenode "iop/apps/edge/internal/node"
edgeservice "iop/apps/edge/internal/service"
"iop/apps/edge/internal/transport"
"iop/packages/config"
"iop/packages/observability"
)
// Runtime is the shared edge runtime assembly used by both `serve` and
// `console` entrypoints. NewRuntime wires logger, registry, node store,
// event bus, service, and transport server; Start performs handler wiring
// and binds the server and metrics endpoint.
type Runtime struct {
Cfg *config.EdgeConfig
Logger *zap.Logger
Registry *edgenode.Registry
NodeStore *edgenode.NodeStore
EventBus *edgeevents.Bus
Service *edgeservice.Service
Server *transport.Server
Input *edgeinput.Manager
lifetimeMu sync.Mutex
lifetimeCancel context.CancelFunc
}
func NewRuntime(cfg *config.EdgeConfig) (*Runtime, error) {
logger, err := observability.NewLoggerWithFile(cfg.Logging.Level, cfg.Logging.Pretty, cfg.Logging.Path)
if err != nil {
return nil, err
}
registry := edgenode.NewRegistry()
nodeStore, err := edgenode.LoadFromConfig(cfg.Nodes)
if err != nil {
return nil, fmt.Errorf("edge: seed node store: %w", err)
}
bus := edgeevents.NewBus()
svc := edgeservice.New(registry, bus)
inputManager := edgeinput.NewManager(*cfg, svc, logger.Named("input"))
server, err := transport.NewServer(cfg.Server.Listen, registry, nodeStore, logger)
if err != nil {
return nil, err
}
return &Runtime{
Cfg: cfg,
Logger: logger,
Registry: registry,
NodeStore: nodeStore,
EventBus: bus,
Service: svc,
Server: server,
Input: inputManager,
}, nil
}
func (r *Runtime) wireHandlers() {
r.Server.SetRunEventHandler(r.EventBus.PublishRun)
r.Server.SetNodeEventHandler(r.EventBus.PublishNode)
}
func (r *Runtime) Start(ctx context.Context) error {
if err := ctx.Err(); err != nil {
return err
}
r.wireHandlers()
lifetimeCtx := r.newLifetimeContext()
if err := r.Server.Start(lifetimeCtx); err != nil {
r.cancelLifetime()
return err
}
if err := ctx.Err(); err != nil {
r.cancelLifetime()
_ = r.Server.Stop()
return err
}
if err := r.Input.Start(lifetimeCtx); err != nil {
r.cancelLifetime()
_ = r.Server.Stop()
return err
}
r.startMetrics()
return nil
}
func (r *Runtime) Stop() error {
r.cancelLifetime()
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
inputErr := r.Input.Stop(ctx)
err := r.Server.Stop()
_ = r.Logger.Sync()
if err == nil {
err = inputErr
}
return err
}
func (r *Runtime) newLifetimeContext() context.Context {
r.lifetimeMu.Lock()
defer r.lifetimeMu.Unlock()
if r.lifetimeCancel != nil {
r.lifetimeCancel()
}
lifetimeCtx, cancel := context.WithCancel(context.Background())
r.lifetimeCancel = cancel
return lifetimeCtx
}
func (r *Runtime) cancelLifetime() {
r.lifetimeMu.Lock()
cancel := r.lifetimeCancel
r.lifetimeCancel = nil
r.lifetimeMu.Unlock()
if cancel != nil {
cancel()
}
}
func (r *Runtime) startMetrics() {
if r.Cfg.Metrics.Port <= 0 {
return
}
go func() {
if err := observability.ServeMetrics(r.Cfg.Metrics.Port); err != nil {
r.Logger.Warn("metrics server exited", zap.Error(err))
}
}()
}