Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
27 changes: 25 additions & 2 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@ The following documentation is available:
* [Installation](#installation)
- [Supported environment variables](#supported-environment-variables)
- [Unix Domain Sockets Client](#unix-domain-sockets-client)
- [Vsock Client (experimental)](#vsock-client-experimental)
* [Usage](#usage)
- [Metrics](#metrics)
- [Events](#events)
Expand Down Expand Up @@ -85,12 +86,13 @@ Find a list of all the available options for your DogStatsD Client in the [Datad
### Supported environment variables

* If the `addr` parameter is empty, the client will:
* First use the `DD_DOGSTATSD_URL` environment variables to build a target address. This must be a URL that start with either `udp://` (to connect using UDP) or with `unix://` (to use a Unix Domain Socket).
* First use the `DD_DOGSTATSD_URL` environment variables to build a target address. This must be a URL that start with either `udp://` (to connect using UDP), with `unix://` (to use a Unix Domain Socket) or with `vsock://` (to use a vsock socket).
Example for UDP url: `DD_DOGSTATSD_URL=udp://localhost:8125`
Example for UDS: `DD_DOGSTATSD_URL=unix:///var/run/datadog/dsd.socket`
Example for vsock: `DD_DOGSTATSD_URL=vsock://host:8125`
Example for Windows named pipe`DD_AGENT_HOST=\\.\pipe\my_windows_pipe`
* Fallback to the `DD_AGENT_HOST` environment variables to build a target address.
Example: `DD_AGENT_HOST=127.0.0.1:8125` for UDP, `DD_AGENT_HOST=unix:///path/to/socket` for UDS and `DD_AGENT_HOST=\\.\pipe\my_windows_pipe` for Windows named pipe.
Example: `DD_AGENT_HOST=127.0.0.1:8125` for UDP, `DD_AGENT_HOST=unix:///path/to/socket` for UDS, `DD_AGENT_HOST=vsock://host:8125` for vsock and `DD_AGENT_HOST=\\.\pipe\my_windows_pipe` for Windows named pipe.
* If `DD_AGENT_HOST` has no port it will default the port to `8125`
* You can use `DD_AGENT_PORT` to set the port if `DD_AGENT_HOST` does not have a port set for UDP
Example: `DD_AGENT_HOST=127.0.0.1` and `DD_AGENT_PORT=1234` will create a UDP connection to `127.0.0.1:1234`.
Expand All @@ -112,6 +114,27 @@ env:

Agent v6+ accepts packets through a Unix Socket datagram connection. Details about the advantages of using UDS over UDP are available in the [DogStatsD Unix Socket documentation](https://docs.datadoghq.com/developers/dogstatsd/unix_socket/). You can use this protocol by giving a `unix:///path/to/dsd.socket` address argument to the `New` constructor.

### Vsock Client (experimental)

VM Sockets (vsock) are a Linux-only transport available for allowing hypervisors and guest virtual machines
to communicate with each other in a fast and secure way, similar to Unix Domain Sockets.

You can use this protocol, on Linux only, by giving a `vsock://<CID>:<port>` address argument to the
`New` constructor, where `<CID>` is either a context ID or one of the following shorthands:

| Shorthand | Context ID | Destination |
|--------------|------------|------------------------------------------|
| `hypervisor` | 0 | The hypervisor process |
| `local` | 1 | The local machine, for loopback purposes |
| `host` | 2 | Any process running on the host |

For example, `vsock://host:8125` sends to port `8125` of the host running the virtual machine. Like
Unix Domain Socket streams, payloads are prefixed with their length so that the Agent can tell them
apart. Other CIDs can be passed in their raw numerical form, which is required for non-standard CIDs,
such as those utilized by [AWS Nitro Enclaves](https://docs.aws.amazon.com/enclaves/latest/user/nitro-enclave-concepts.html#term-socket).

This feature is experimental, and depends on experimental support in the Agent.

## Usage

In order to use DogStatsD metrics, events, and Service Checks, the Agent must be [running and available](https://docs.datadoghq.com/developers/dogstatsd/?code-lang=go).
Expand Down
1 change: 1 addition & 0 deletions go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ require (
github.com/golang/mock v1.6.0
github.com/stretchr/testify v1.8.1
golang.org/x/net v0.0.0-20210405180319-a5a99cb37ef4
golang.org/x/sys v0.0.0-20220715151400-c0bba94af5f8
)

replace github.com/sirupsen/logrus v1.7.0 => github.com/sirupsen/logrus v1.9.3
189 changes: 189 additions & 0 deletions statsd/conn.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,189 @@
//go:build !windows
// +build !windows

package statsd

import (
"encoding/binary"
"net"
"strings"
"sync"
"time"
)

// connDialer establishes the connections used by a connWriter.
type connDialer interface {
// dial connects to the Agent, giving up after connectTimeout.
dial(connectTimeout time.Duration) (net.Conn, error)

// transportName returns the name of the transport. It can depend on the connection that
// dial ultimately established, as UDS guesses between datagram and stream sockets.
transportName() string
}

// connWriter is an internal class wrapping around management of a connection to the Agent. The
// connection is established on the first write and re-established whenever the Agent disconnects.
type connWriter struct {
// Dialer used to establish new connections
dialer connDialer
// Established connection object, or nil if not connected yet
conn net.Conn
// write timeout
writeTimeout time.Duration
// connect timeout
connectTimeout time.Duration
sync.RWMutex // used to lock conn / writer can replace it
}

// newConnWriter returns a pointer to a new connWriter using the given dialer.
func newConnWriter(dialer connDialer, writeTimeout time.Duration, connectTimeout time.Duration) *connWriter {
// Defer connection to first Write
return &connWriter{dialer: dialer, conn: nil, writeTimeout: writeTimeout, connectTimeout: connectTimeout}
}

// GetTransportName returns the transport used by the writer
func (w *connWriter) GetTransportName() string {
w.RLock()
defer w.RUnlock()

return w.dialer.transportName()
}

// isStreamConn reports whether conn needs length-delimited framing: datagram transports preserve
// message boundaries, stream transports do not.
func isStreamConn(conn net.Conn) bool {
return conn.LocalAddr().Network() != "unixgram"
}

func (w *connWriter) shouldCloseConnection(err error, partialWrite bool) bool {
if err != nil && partialWrite {
// We can't recover from a partial write
return true
}
if err, isNetworkErr := err.(net.Error); err != nil && (!isNetworkErr || !err.Timeout()) {
// Statsd server disconnected, retry connecting at next packet
return true
}
return false
}

// Write data to the connection with write timeout and minimal error handling:
// create the connection if nil, and destroy it if the statsd server has disconnected
func (w *connWriter) Write(data []byte) (int, error) {
var n int
partialWrite := false
conn, err := w.ensureConnection()
if err != nil {
return 0, err
}
stream := isStreamConn(conn)

// When using streams the deadline will only make us drop the packet if we can't write it at all,
// once we've started writing we need to finish.
conn.SetWriteDeadline(time.Now().Add(w.writeTimeout))

// When using streams, we append the length of the packet to the data
if stream {
bs := []byte{0, 0, 0, 0}
binary.LittleEndian.PutUint32(bs, uint32(len(data)))
_, err = conn.Write(bs)

partialWrite = true

// W need to be able to finish to write partially written packets once we have started.
// But we will reset the connection if we can't write anything at all for a long time.
conn.SetWriteDeadline(time.Now().Add(w.connectTimeout))

// Continue writing only if we've written the length of the packet
if err == nil {
n, err = conn.Write(data)
if err == nil {
partialWrite = false
}
}
} else {
n, err = conn.Write(data)
}

if w.shouldCloseConnection(err, partialWrite) {
w.unsetConnection()
}
return n, err
}

func (w *connWriter) Close() error {
if w.conn != nil {
return w.conn.Close()
}
return nil
}

func (w *connWriter) ensureConnection() (net.Conn, error) {
// Check if we've already got a socket we can use
w.RLock()
currentConn := w.conn
w.RUnlock()

if currentConn != nil {
return currentConn, nil
}

// Looks like we might need to connect - try again with write locking.
w.Lock()
defer w.Unlock()
if w.conn != nil {
return w.conn, nil
}

newConn, err := w.dialer.dial(w.connectTimeout)
if err != nil {
return nil, err
}
w.conn = newConn
return newConn, nil
}

func (w *connWriter) unsetConnection() {
w.Lock()
defer w.Unlock()
_ = w.conn.Close()
w.conn = nil
}

// isConnectionRefused reports whether err means that nothing is listening on the other end. The
// error message is matched, rather than the error itself, because errors.Is is not available in the
// oldest versions of Go this library supports.
func isConnectionRefused(err error) bool {
return strings.HasSuffix(err.Error(), "connection refused")
}

// dialWithRetry calls dial until it succeeds, the connect timeout expires, or dial fails with an
// error that isRetryable rejects. Errors meaning that nothing is listening are worth retrying: it's
// likely that the Agent is restarting in that case, and that it will be back shortly.
func dialWithRetry(connectTimeout time.Duration, isRetryable func(error) bool, dial func(timeout time.Duration) (net.Conn, error)) (net.Conn, error) {
connectAttemptsLeft := 3
connectDeadline := time.Now().Add(connectTimeout)

// Calculate the backoff time for connection refused errors, but don't exceed one second: this means we won't waste
// longer than 1 seconds worth of time if the socket becomes available immediately after our last connect attempt
connRefusedBackoff := connectTimeout / time.Duration(connectAttemptsLeft+1)
if connRefusedBackoff > time.Second {
connRefusedBackoff = time.Second
}

for {
connectAttemptsLeft--

perCallTimeout := time.Until(connectDeadline)
newConn, err := dial(perCallTimeout)
if err != nil {
if isRetryable(err) && connectAttemptsLeft > 0 {
// If we get a retryable error, we need to wait a bit before trying again.
time.Sleep(connRefusedBackoff)
continue
}
return nil, err
}
return newConn, nil
}
}
Loading
Loading