- Update SDD and milestone documentation - Add node transport client test - Update edge local dev guide - Add field docs smoke tests
134 lines
3.7 KiB
Go
134 lines
3.7 KiB
Go
package transport
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
"net"
|
|
"strconv"
|
|
"time"
|
|
|
|
toki "git.toki-labs.com/toki/proto-socket/go"
|
|
"go.uber.org/zap"
|
|
|
|
iop "iop/proto/gen/iop"
|
|
)
|
|
|
|
const (
|
|
heartbeatIntervalSec = 30
|
|
// heartbeatWaitSec is kept above heartbeatIntervalSec as defence in depth:
|
|
// the library wait-timer callback already self-heals stale state (see
|
|
// proto-socket go/base_client.go sendHeartBeat), but a larger wait window
|
|
// gives the peer's next heartbeat an extra chance to overwrite any stray
|
|
// timer on slow or jittery links before it fires.
|
|
heartbeatWaitSec = 45
|
|
tcpWriteTimeout = 10 * time.Second
|
|
registerTimeout = 10 * time.Second
|
|
registerInitialWait = 100 * time.Millisecond
|
|
registerRetryWait = 250 * time.Millisecond
|
|
registerAttempts = 3
|
|
)
|
|
|
|
type writeDeadlineConn struct {
|
|
net.Conn
|
|
timeout time.Duration
|
|
}
|
|
|
|
func (c *writeDeadlineConn) Write(b []byte) (int, error) {
|
|
if c.timeout > 0 {
|
|
if err := c.Conn.SetWriteDeadline(time.Now().Add(c.timeout)); err != nil {
|
|
return 0, err
|
|
}
|
|
defer c.Conn.SetWriteDeadline(time.Time{})
|
|
}
|
|
return c.Conn.Write(b)
|
|
}
|
|
|
|
// RegisterResult is returned by DialEdge after successful registration.
|
|
type RegisterResult struct {
|
|
Session *Session
|
|
NodeID string
|
|
Alias string
|
|
Config *iop.NodeConfigPayload
|
|
}
|
|
|
|
// DialEdge connects to edge, performs the registration handshake, and returns
|
|
// a RegisterResult. Call result.Session.SetHandler after creating node.Node.
|
|
func DialEdge(ctx context.Context, addr, token string, logger *zap.Logger) (*RegisterResult, error) {
|
|
host, portStr, err := net.SplitHostPort(addr)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("transport: invalid addr %q: %w", addr, err)
|
|
}
|
|
port, err := strconv.Atoi(portStr)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("transport: invalid port %q: %w", portStr, err)
|
|
}
|
|
|
|
dialer := net.Dialer{KeepAlive: 15 * time.Second}
|
|
conn, err := dialer.DialContext(ctx, "tcp", net.JoinHostPort(host, strconv.Itoa(port)))
|
|
if err != nil {
|
|
return nil, fmt.Errorf("transport: dial edge %s: %w", addr, err)
|
|
}
|
|
client := toki.NewTcpClient(&writeDeadlineConn{Conn: conn, timeout: tcpWriteTimeout}, heartbeatIntervalSec, heartbeatWaitSec, nodeParserMap())
|
|
|
|
resp, err := registerWithEdge(ctx, client, token, logger)
|
|
if err != nil {
|
|
_ = client.Close()
|
|
return nil, fmt.Errorf("transport: register: %w", err)
|
|
}
|
|
if !resp.GetAccepted() {
|
|
_ = client.Close()
|
|
return nil, fmt.Errorf("transport: register rejected: %s", resp.GetReason())
|
|
}
|
|
|
|
sess := newSession(client, logger, resp.GetNodeId(), resp.GetAlias())
|
|
logger.Info("registered with edge",
|
|
zap.String("node_id", resp.GetNodeId()),
|
|
zap.String("alias", resp.GetAlias()),
|
|
)
|
|
return &RegisterResult{
|
|
Session: sess,
|
|
NodeID: resp.GetNodeId(),
|
|
Alias: resp.GetAlias(),
|
|
Config: resp.GetConfig(),
|
|
}, nil
|
|
}
|
|
|
|
func registerWithEdge(ctx context.Context, client *toki.TcpClient, token string, logger *zap.Logger) (*iop.RegisterResponse, error) {
|
|
timer := time.NewTimer(registerInitialWait)
|
|
select {
|
|
case <-ctx.Done():
|
|
timer.Stop()
|
|
return nil, ctx.Err()
|
|
case <-timer.C:
|
|
}
|
|
|
|
var lastErr error
|
|
for attempt := 1; attempt <= registerAttempts; attempt++ {
|
|
resp, err := toki.SendRequestTyped[*iop.RegisterRequest, *iop.RegisterResponse](
|
|
&client.Communicator,
|
|
&iop.RegisterRequest{Token: token},
|
|
registerTimeout,
|
|
)
|
|
if err == nil {
|
|
return resp, nil
|
|
}
|
|
lastErr = err
|
|
if attempt == registerAttempts || !client.IsAlive() {
|
|
break
|
|
}
|
|
if logger != nil {
|
|
logger.Warn("register request failed, retrying",
|
|
zap.Int("attempt", attempt),
|
|
zap.Error(err),
|
|
)
|
|
}
|
|
timer := time.NewTimer(registerRetryWait)
|
|
select {
|
|
case <-ctx.Done():
|
|
timer.Stop()
|
|
return nil, ctx.Err()
|
|
case <-timer.C:
|
|
}
|
|
}
|
|
return nil, lastErr
|
|
}
|