iop/apps/node/internal/transport/session.go
toki cf3760f394 feat: console emitter interface, persistent status, and profile proto message updates
- Add console emitter interface for event-driven console output
- Implement persistent cancel reason and explicit completion tracking
- Add profile proto message definitions
- Update edge and node transport layers with adapter execution terminology
- Add new packages/events module
- Update architecture documentation and README files
2026-05-16 21:16:48 +09:00

154 lines
3.7 KiB
Go

package transport
import (
"context"
"sync"
toki "git.toki-labs.com/toki/common-proto-socket/go"
"go.uber.org/zap"
"google.golang.org/protobuf/proto"
"iop/packages/events"
iop "iop/proto/gen/iop"
)
// Handler processes IOP messages received from edge.
type Handler interface {
OnRunRequest(ctx context.Context, sess *Session, req *iop.RunRequest) error
OnCancel(ctx context.Context, sess *Session, req *iop.CancelRequest) error
OnCommandRequest(ctx context.Context, sess *Session, req *iop.NodeCommandRequest) (*iop.NodeCommandResponse, error)
}
// Session represents the node's persistent connection to edge.
type Session struct {
client *toki.TcpClient
logger *zap.Logger
nodeID string
alias string
mu sync.RWMutex
handler Handler
eventHandler func(*iop.EdgeNodeEvent)
closeReason string
}
func newSession(client *toki.TcpClient, logger *zap.Logger, nodeID, alias string) *Session {
s := &Session{client: client, logger: logger, nodeID: nodeID, alias: alias}
toki.AddListenerTyped[*iop.RunRequest](&client.Communicator, func(req *iop.RunRequest) {
go func() {
s.mu.RLock()
h := s.handler
s.mu.RUnlock()
if h == nil {
return
}
if err := h.OnRunRequest(context.Background(), s, req); err != nil {
logger.Warn("run request error",
zap.String("run_id", req.GetRunId()),
zap.Error(err),
)
}
}()
})
toki.AddListenerTyped[*iop.CancelRequest](&client.Communicator, func(req *iop.CancelRequest) {
s.mu.RLock()
h := s.handler
s.mu.RUnlock()
if h == nil {
return
}
if err := h.OnCancel(context.Background(), s, req); err != nil {
logger.Warn("cancel error",
zap.String("run_id", req.GetRunId()),
zap.Error(err),
)
}
})
toki.AddRequestListenerTyped[*iop.NodeCommandRequest, *iop.NodeCommandResponse](&client.Communicator, func(req *iop.NodeCommandRequest) (*iop.NodeCommandResponse, error) {
s.mu.RLock()
h := s.handler
s.mu.RUnlock()
if h == nil {
return &iop.NodeCommandResponse{Error: "handler not ready"}, nil
}
resp, err := h.OnCommandRequest(context.Background(), s, req)
if err != nil {
return &iop.NodeCommandResponse{Error: err.Error()}, nil
}
return resp, nil
})
toki.AddListenerTyped[*iop.EdgeNodeEvent](&client.Communicator, func(event *iop.EdgeNodeEvent) {
s.emitEvent(event)
})
client.AddDisconnectListener(func(_ *toki.TcpClient) {
logger.Info("disconnected from edge")
s.emitEvent(events.NewEdgeNodeEvent(
events.SourceNode,
events.TypeEdgeDisconnected,
nodeID,
alias,
s.disconnectReason(),
nil,
))
})
return s
}
// SetHandler attaches the message handler. Called after registration completes.
func (s *Session) SetHandler(h Handler) {
s.mu.Lock()
s.handler = h
s.mu.Unlock()
}
func (s *Session) SetEventHandler(handler func(*iop.EdgeNodeEvent)) {
s.mu.Lock()
s.eventHandler = handler
s.mu.Unlock()
}
// Send transmits a proto message to edge.
func (s *Session) Send(m proto.Message) error {
return s.client.Send(m)
}
// IsAlive reports whether the connection is active.
func (s *Session) IsAlive() bool { return s.client.IsAlive() }
// Close terminates the connection to edge.
func (s *Session) Close() error {
s.setCloseReason(events.ReasonLocalShutdown)
return s.client.Close()
}
func (s *Session) emitEvent(event *iop.EdgeNodeEvent) {
s.mu.RLock()
handler := s.eventHandler
s.mu.RUnlock()
if handler != nil {
handler(event)
}
}
func (s *Session) setCloseReason(reason string) {
s.mu.Lock()
if s.closeReason == "" {
s.closeReason = reason
}
s.mu.Unlock()
}
func (s *Session) disconnectReason() string {
s.mu.RLock()
reason := s.closeReason
s.mu.RUnlock()
if reason == "" {
return events.ReasonTransportClosed
}
return reason
}