iop/apps/node/internal/bootstrap/module.go
toki a2cee38cfa feat: console events, session handling, and transport improvements
- Add console events tracking in edge node
- Improve session management in node transport
- Add integration tests for edge and node transport
- Update events package with new event types
- Update README with current status
2026-05-17 06:26:15 +09:00

141 lines
3.8 KiB
Go

// Package bootstrap wires the IOP Node application using go.uber.org/fx.
package bootstrap
import (
"context"
"fmt"
"io"
"os"
"path/filepath"
"go.uber.org/fx"
"go.uber.org/zap"
"iop/apps/node/internal/adapters"
"iop/apps/node/internal/node"
"iop/apps/node/internal/router"
"iop/apps/node/internal/store"
"iop/apps/node/internal/transport"
"iop/packages/config"
"iop/packages/events"
"iop/packages/observability"
iop "iop/proto/gen/iop"
)
// Module returns the fx options that wire the node application.
func Module(cfg *config.NodeConfig) fx.Option {
return fx.Options(
fx.Provide(
func() *config.NodeConfig { return cfg },
func(cfg *config.NodeConfig) (*zap.Logger, error) {
return observability.NewLogger(cfg.Logging.Level, cfg.Logging.Pretty)
},
),
fx.Invoke(func(lc fx.Lifecycle, cfg *config.NodeConfig, logger *zap.Logger) {
var sess *transport.Session
var st *store.Store
var reg *adapters.Registry
lc.Append(fx.Hook{
OnStart: func(ctx context.Context) error {
result, err := transport.DialEdge(ctx, cfg.Transport.EdgeAddr, cfg.Transport.Token, logger)
if err != nil {
return fmt.Errorf("bootstrap: dial edge: %w", err)
}
reg, err = adapters.BuildFromPayload(result.Config, logger)
if err != nil {
_ = result.Session.Close()
return fmt.Errorf("bootstrap: build adapters: %w", err)
}
dsn, err := storeDSN(result.Config.GetRuntime().GetWorkspaceRoot())
if err != nil {
_ = result.Session.Close()
return err
}
st, err = store.New(dsn, logger)
if err != nil {
_ = result.Session.Close()
return fmt.Errorf("bootstrap: store: %w", err)
}
if err := reg.Start(ctx); err != nil {
_ = reg.Stop(context.Background())
_ = result.Session.Close()
_ = st.Close()
return fmt.Errorf("bootstrap: start adapters: %w", err)
}
rtr := router.New(reg, logger)
n := node.New(result.NodeID, rtr, st, os.Stdout, logger)
result.Session.SetEventHandler(func(event *iop.EdgeNodeEvent) {
printEdgeEvent(os.Stdout, event)
})
result.Session.SetHandler(n)
sess = result.Session
go func() {
if err := observability.ServeMetrics(cfg.Metrics.Port); err != nil {
logger.Warn("metrics server exited", zap.Error(err))
}
}()
return nil
},
OnStop: func(_ context.Context) error {
if reg != nil {
if err := reg.Stop(context.Background()); err != nil {
return err
}
}
if sess != nil {
if err := sess.Close(); err != nil {
return err
}
}
if st != nil {
return st.Close()
}
return nil
},
})
}),
)
}
func storeDSN(workspaceRoot string) (string, error) {
if workspaceRoot == "" {
return "file:iop.db?cache=shared&mode=rwc", nil
}
if err := os.MkdirAll(workspaceRoot, 0o755); err != nil {
return "", fmt.Errorf("bootstrap: workspace: %w", err)
}
return "file:" + filepath.Join(workspaceRoot, "iop.db") + "?cache=shared&mode=rwc", nil
}
func printEdgeEvent(out io.Writer, event *iop.EdgeNodeEvent) {
if out == nil {
return
}
switch event.GetType() {
case events.TypeEdgeDisconnected:
fmt.Fprintf(out, "[edge-event] disconnected reason=%q%s\n", event.GetReason(), transportCloseDetail(event.GetMetadata()))
default:
fmt.Fprintf(out, "[edge-event] %s reason=%q%s\n", event.GetType(), event.GetReason(), transportCloseDetail(event.GetMetadata()))
}
}
func transportCloseDetail(metadata map[string]string) string {
if metadata == nil {
return ""
}
detail := ""
if reason := metadata[events.MetadataTransportCloseReason]; reason != "" {
detail += fmt.Sprintf(" transport_close_reason=%q", reason)
}
if err := metadata[events.MetadataTransportCloseError]; err != "" {
detail += fmt.Sprintf(" transport_close_error=%q", err)
}
return detail
}