2021-08-17 14:30:02 +00:00
|
|
|
package connection
|
|
|
|
|
|
|
|
import (
|
|
|
|
"context"
|
|
|
|
"io"
|
|
|
|
"time"
|
|
|
|
|
|
|
|
"github.com/rs/zerolog"
|
|
|
|
|
|
|
|
tunnelpogs "github.com/cloudflare/cloudflared/tunnelrpc/pogs"
|
|
|
|
)
|
|
|
|
|
|
|
|
// RPCClientFunc derives a named tunnel rpc client that can then be used to register and unregister connections.
|
|
|
|
type RPCClientFunc func(context.Context, io.ReadWriteCloser, *zerolog.Logger) NamedTunnelRPCClient
|
|
|
|
|
|
|
|
type controlStream struct {
|
|
|
|
observer *Observer
|
|
|
|
|
2022-02-07 09:42:07 +00:00
|
|
|
connectedFuse ConnectedFuse
|
|
|
|
namedTunnelProperties *NamedTunnelProperties
|
|
|
|
connIndex uint8
|
2021-08-17 14:30:02 +00:00
|
|
|
|
|
|
|
newRPCClientFunc RPCClientFunc
|
|
|
|
|
|
|
|
gracefulShutdownC <-chan struct{}
|
|
|
|
gracePeriod time.Duration
|
|
|
|
stoppedGracefully bool
|
|
|
|
}
|
|
|
|
|
|
|
|
// ControlStreamHandler registers connections with origintunneld and initiates graceful shutdown.
|
|
|
|
type ControlStreamHandler interface {
|
2022-01-05 16:01:56 +00:00
|
|
|
// ServeControlStream handles the control plane of the transport in the current goroutine calling this
|
2022-04-27 10:51:06 +00:00
|
|
|
ServeControlStream(ctx context.Context, rw io.ReadWriteCloser, connOptions *tunnelpogs.ConnectionOptions, tunnelConfigGetter TunnelConfigJSONGetter) error
|
2022-01-05 16:01:56 +00:00
|
|
|
// IsStopped tells whether the method above has finished
|
2021-08-17 14:30:02 +00:00
|
|
|
IsStopped() bool
|
|
|
|
}
|
|
|
|
|
2022-04-27 10:51:06 +00:00
|
|
|
type TunnelConfigJSONGetter interface {
|
|
|
|
GetConfigJSON() ([]byte, error)
|
|
|
|
}
|
|
|
|
|
2021-08-17 14:30:02 +00:00
|
|
|
// NewControlStream returns a new instance of ControlStreamHandler
|
|
|
|
func NewControlStream(
|
|
|
|
observer *Observer,
|
|
|
|
connectedFuse ConnectedFuse,
|
2022-02-07 09:42:07 +00:00
|
|
|
namedTunnelConfig *NamedTunnelProperties,
|
2021-08-17 14:30:02 +00:00
|
|
|
connIndex uint8,
|
|
|
|
newRPCClientFunc RPCClientFunc,
|
|
|
|
gracefulShutdownC <-chan struct{},
|
|
|
|
gracePeriod time.Duration,
|
|
|
|
) ControlStreamHandler {
|
|
|
|
if newRPCClientFunc == nil {
|
|
|
|
newRPCClientFunc = newRegistrationRPCClient
|
|
|
|
}
|
|
|
|
return &controlStream{
|
2022-02-07 09:42:07 +00:00
|
|
|
observer: observer,
|
|
|
|
connectedFuse: connectedFuse,
|
|
|
|
namedTunnelProperties: namedTunnelConfig,
|
|
|
|
newRPCClientFunc: newRPCClientFunc,
|
|
|
|
connIndex: connIndex,
|
|
|
|
gracefulShutdownC: gracefulShutdownC,
|
|
|
|
gracePeriod: gracePeriod,
|
2021-08-17 14:30:02 +00:00
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
func (c *controlStream) ServeControlStream(
|
|
|
|
ctx context.Context,
|
|
|
|
rw io.ReadWriteCloser,
|
|
|
|
connOptions *tunnelpogs.ConnectionOptions,
|
2022-04-27 10:51:06 +00:00
|
|
|
tunnelConfigGetter TunnelConfigJSONGetter,
|
2021-08-17 14:30:02 +00:00
|
|
|
) error {
|
|
|
|
rpcClient := c.newRPCClientFunc(ctx, rw, c.observer.log)
|
|
|
|
|
2022-04-27 10:51:06 +00:00
|
|
|
registrationDetails, err := rpcClient.RegisterConnection(ctx, c.namedTunnelProperties, connOptions, c.connIndex, c.observer)
|
|
|
|
if err != nil {
|
2021-09-24 11:56:31 +00:00
|
|
|
rpcClient.Close()
|
2021-08-17 14:30:02 +00:00
|
|
|
return err
|
|
|
|
}
|
|
|
|
c.connectedFuse.Connected()
|
|
|
|
|
2022-04-27 10:51:06 +00:00
|
|
|
// if conn index is 0 and tunnel is not remotely managed, then send local ingress rules configuration
|
|
|
|
if c.connIndex == 0 && !registrationDetails.TunnelIsRemotelyManaged {
|
|
|
|
if tunnelConfig, err := tunnelConfigGetter.GetConfigJSON(); err == nil {
|
|
|
|
if err := rpcClient.SendLocalConfiguration(ctx, tunnelConfig, c.observer); err != nil {
|
|
|
|
c.observer.log.Err(err).Msg("unable to send local configuration")
|
|
|
|
}
|
|
|
|
} else {
|
|
|
|
c.observer.log.Err(err).Msg("failed to obtain current configuration")
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2022-01-05 16:01:56 +00:00
|
|
|
c.waitForUnregister(ctx, rpcClient)
|
2021-09-24 09:31:28 +00:00
|
|
|
return nil
|
|
|
|
}
|
|
|
|
|
|
|
|
func (c *controlStream) waitForUnregister(ctx context.Context, rpcClient NamedTunnelRPCClient) {
|
2021-08-17 14:30:02 +00:00
|
|
|
// wait for connection termination or start of graceful shutdown
|
2021-09-24 11:56:31 +00:00
|
|
|
defer rpcClient.Close()
|
2021-08-17 14:30:02 +00:00
|
|
|
select {
|
|
|
|
case <-ctx.Done():
|
|
|
|
break
|
|
|
|
case <-c.gracefulShutdownC:
|
|
|
|
c.stoppedGracefully = true
|
|
|
|
}
|
|
|
|
|
|
|
|
c.observer.sendUnregisteringEvent(c.connIndex)
|
|
|
|
rpcClient.GracefulShutdown(ctx, c.gracePeriod)
|
|
|
|
c.observer.log.Info().Uint8(LogFieldConnIndex, c.connIndex).Msg("Unregistered tunnel connection")
|
|
|
|
}
|
|
|
|
|
|
|
|
func (c *controlStream) IsStopped() bool {
|
|
|
|
return c.stoppedGracefully
|
|
|
|
}
|