2020-10-08 10:12:26 +00:00
|
|
|
package connection
|
|
|
|
|
|
|
|
import (
|
2022-06-02 17:57:37 +00:00
|
|
|
"net"
|
2021-07-28 09:02:55 +00:00
|
|
|
"strings"
|
|
|
|
|
2023-04-12 21:41:11 +00:00
|
|
|
"github.com/google/uuid"
|
2020-11-25 06:55:13 +00:00
|
|
|
"github.com/rs/zerolog"
|
2023-04-12 21:41:11 +00:00
|
|
|
|
|
|
|
"github.com/cloudflare/cloudflared/management"
|
2020-10-08 10:12:26 +00:00
|
|
|
)
|
|
|
|
|
2021-01-14 22:33:36 +00:00
|
|
|
const (
|
2023-04-12 21:41:11 +00:00
|
|
|
LogFieldConnectionID = "connection"
|
2021-01-14 22:33:36 +00:00
|
|
|
LogFieldLocation = "location"
|
2022-06-02 17:57:37 +00:00
|
|
|
LogFieldIPAddress = "ip"
|
2023-04-12 21:41:11 +00:00
|
|
|
LogFieldProtocol = "protocol"
|
2021-01-14 22:33:36 +00:00
|
|
|
observerChannelBufferSize = 16
|
|
|
|
)
|
2020-12-28 18:10:01 +00:00
|
|
|
|
2020-10-08 10:12:26 +00:00
|
|
|
type Observer struct {
|
2021-01-14 22:33:36 +00:00
|
|
|
log *zerolog.Logger
|
2021-02-03 18:32:54 +00:00
|
|
|
logTransport *zerolog.Logger
|
2021-01-14 22:33:36 +00:00
|
|
|
metrics *tunnelMetrics
|
|
|
|
tunnelEventChan chan Event
|
|
|
|
addSinkChan chan EventSink
|
|
|
|
}
|
|
|
|
|
|
|
|
type EventSink interface {
|
|
|
|
OnTunnelEvent(event Event)
|
|
|
|
}
|
|
|
|
|
2022-07-20 23:17:29 +00:00
|
|
|
func NewObserver(log, logTransport *zerolog.Logger) *Observer {
|
2021-01-14 22:33:36 +00:00
|
|
|
o := &Observer{
|
|
|
|
log: log,
|
2021-02-03 18:32:54 +00:00
|
|
|
logTransport: logTransport,
|
2021-01-14 22:33:36 +00:00
|
|
|
metrics: newTunnelMetrics(),
|
|
|
|
tunnelEventChan: make(chan Event, observerChannelBufferSize),
|
|
|
|
addSinkChan: make(chan EventSink, observerChannelBufferSize),
|
2020-10-08 10:12:26 +00:00
|
|
|
}
|
2021-01-14 22:33:36 +00:00
|
|
|
go o.dispatchEvents()
|
|
|
|
return o
|
|
|
|
}
|
|
|
|
|
|
|
|
func (o *Observer) RegisterSink(sink EventSink) {
|
|
|
|
o.addSinkChan <- sink
|
2020-10-08 10:12:26 +00:00
|
|
|
}
|
|
|
|
|
2023-04-12 21:41:11 +00:00
|
|
|
func (o *Observer) logConnected(connectionID uuid.UUID, connIndex uint8, location string, address net.IP, protocol Protocol) {
|
2020-12-28 18:10:01 +00:00
|
|
|
o.log.Info().
|
2023-04-12 21:41:11 +00:00
|
|
|
Int(management.EventTypeKey, int(management.Cloudflared)).
|
|
|
|
Str(LogFieldConnectionID, connectionID.String()).
|
2020-12-28 18:10:01 +00:00
|
|
|
Uint8(LogFieldConnIndex, connIndex).
|
|
|
|
Str(LogFieldLocation, location).
|
2022-06-02 17:57:37 +00:00
|
|
|
IPAddr(LogFieldIPAddress, address).
|
2023-04-12 21:41:11 +00:00
|
|
|
Str(LogFieldProtocol, protocol.String()).
|
|
|
|
Msg("Registered tunnel connection")
|
2020-11-09 11:40:48 +00:00
|
|
|
o.metrics.registerServerLocation(uint8ToString(connIndex), location)
|
2020-10-08 10:12:26 +00:00
|
|
|
}
|
|
|
|
|
2021-01-19 12:20:11 +00:00
|
|
|
func (o *Observer) sendRegisteringEvent(connIndex uint8) {
|
|
|
|
o.sendEvent(Event{Index: connIndex, EventType: RegisteringTunnel})
|
2020-10-08 10:12:26 +00:00
|
|
|
}
|
|
|
|
|
2024-11-25 18:43:32 +00:00
|
|
|
func (o *Observer) sendConnectedEvent(connIndex uint8, protocol Protocol, location string, edgeAddress net.IP) {
|
|
|
|
o.sendEvent(Event{Index: connIndex, EventType: Connected, Protocol: protocol, Location: location, EdgeAddress: edgeAddress})
|
2020-10-08 10:12:26 +00:00
|
|
|
}
|
|
|
|
|
2021-07-09 17:52:41 +00:00
|
|
|
func (o *Observer) SendURL(url string) {
|
2020-11-30 20:05:37 +00:00
|
|
|
o.sendEvent(Event{EventType: SetURL, URL: url})
|
2021-07-28 08:27:05 +00:00
|
|
|
|
|
|
|
if !strings.HasPrefix(url, "https://") {
|
|
|
|
// We add https:// in the prefix for backwards compatibility as we used to do that with the old free tunnels
|
|
|
|
// and some tools (like `wrangler tail`) are regexp-ing for that specifically.
|
|
|
|
url = "https://" + url
|
|
|
|
}
|
|
|
|
o.metrics.userHostnamesCounts.WithLabelValues(url).Inc()
|
2020-11-30 20:05:37 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
func (o *Observer) SendReconnect(connIndex uint8) {
|
|
|
|
o.sendEvent(Event{Index: connIndex, EventType: Reconnecting})
|
|
|
|
}
|
|
|
|
|
2021-02-04 21:09:17 +00:00
|
|
|
func (o *Observer) sendUnregisteringEvent(connIndex uint8) {
|
|
|
|
o.sendEvent(Event{Index: connIndex, EventType: Unregistering})
|
|
|
|
}
|
|
|
|
|
2020-11-30 20:05:37 +00:00
|
|
|
func (o *Observer) SendDisconnect(connIndex uint8) {
|
|
|
|
o.sendEvent(Event{Index: connIndex, EventType: Disconnected})
|
|
|
|
}
|
|
|
|
|
|
|
|
func (o *Observer) sendEvent(e Event) {
|
2021-01-14 22:33:36 +00:00
|
|
|
select {
|
|
|
|
case o.tunnelEventChan <- e:
|
|
|
|
break
|
|
|
|
default:
|
|
|
|
o.log.Warn().Msg("observer channel buffer is full")
|
2020-10-08 10:12:26 +00:00
|
|
|
}
|
|
|
|
}
|
2021-01-14 22:33:36 +00:00
|
|
|
|
|
|
|
func (o *Observer) dispatchEvents() {
|
|
|
|
var sinks []EventSink
|
|
|
|
for {
|
|
|
|
select {
|
|
|
|
case sink := <-o.addSinkChan:
|
|
|
|
sinks = append(sinks, sink)
|
|
|
|
case evt := <-o.tunnelEventChan:
|
|
|
|
for _, sink := range sinks {
|
|
|
|
sink.OnTunnelEvent(evt)
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
type EventSinkFunc func(event Event)
|
|
|
|
|
|
|
|
func (f EventSinkFunc) OnTunnelEvent(event Event) {
|
|
|
|
f(event)
|
|
|
|
}
|