Skip to content
Merged
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
105 changes: 105 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,111 @@ default: off

Disable request mirror until we enable it in the Lua code.

### apisix_stream_metrics_zone \<size\>

context: `stream`

default: off

Reserve a shared memory zone that collects, per stream listening address, the
number of active sessions and the bytes transferred in the four directions
(downstream/upstream × ingress/egress). A session accumulates locally and
merges into the zone once per forwarding pass, so the counters keep moving
during a long-lived connection while the inner read/write loop still costs a
single atomic per direction. Without this directive nothing is collected.

example:

```nginx
stream {
apisix_stream_metrics_zone 1m;
}
```

Only TCP and UDP listening addresses are accounted for. **Every** unix socket
inside `stream{}` is skipped, including one you configured yourself to proxy
on: the filter is on the address family, not on which server owns it. This is
what keeps APISIX's own worker event channel, which lives on a unix socket in
the same block, from being reported as proxied traffic.

Things worth knowing before building alerts on this:

- The byte counters are monotonic from the moment the zone is created. They
restart at zero when the process restarts, when the configured zone size
changes (nginx only reuses a zone whose size is unchanged), and when a
reload removes the directive or the whole `stream{}` block, since the zone
is then released.
- Slots are only ever added. A listening address removed by a reload keeps its
slot, and its accumulated totals, until the process restarts.
- A counter is written once per forwarding pass and again when the session
ends, so a session that has gone quiet has already published everything it
moved.
- A TCP and a UDP listener on the same port share one slot: nginx gives them
the same address text, so their traffic is summed and `$stream_listen_addr`
cannot tell them apart.
- `active` is decremented in the log phase, which nginx runs for every session
it finalizes, including on a graceful shutdown. A worker that dies without
running it (a crash, or `SIGKILL`) leaks its in-flight sessions into the
count, and nothing rebases the zone until the process restarts.
- With `proxy_next_upstream`, the upstream byte counters sum every attempt,
including what was sent to a peer that then failed. `$upstream_bytes_sent`
and `$upstream_bytes_received` instead report one value per attempt, comma
separated, so comparing them against these counters means summing them
first. For a session that reached its upstream on the first try the two
agree directly.

Read the counters from Lua with `resty.apisix.stream.metrics`:

```lua
local metrics = require("resty.apisix.stream.metrics")
local res, err = metrics.dump()
-- res[i] = { listen_addr = "0.0.0.0:9100", active = 3,
-- downstream_ingress = 12, downstream_egress = 34,
-- upstream_egress = 12, upstream_ingress = 34 }
```

## Variable

### $stream_session_reason

context: `stream`

Why the session ended. nginx keeps `$status` at 200 for every failure that
happens after the upstream connection is established, so this variable is what
tells a graceful close apart from a timeout or a reset:

| value | meaning |
|---|---|
| `closed` | one side sent a FIN, the session ended normally |
| `client_rst` | the client reset the connection (`ECONNRESET`) |
| `client_read_error` | reading from the client failed for any other reason, including a TLS protocol failure |
| `client_error` | sending to the client failed |
| `upstream_rst` | the upstream reset the connection (`ECONNRESET`) |
| `upstream_read_error` | reading from the upstream failed for any other reason |
| `upstream_error` | sending to the upstream failed |
| `connect_timeout` | connecting to the upstream timed out and no peer was left |
| `recv_timeout` | `proxy_timeout` expired while waiting for data |
| `send_timeout` | `proxy_timeout` expired with data still buffered towards a peer |
| `upstream_timeout` | a UDP upstream never answered |
| `shutdown` | the worker was shutting down |
| `connect_failed` | no upstream could be reached: refused, no live node, or a failed upstream handshake |
| `-` | no reason applies, for example a session rejected during preread |

A UDP session has no close to observe, so one that ends normally -- including
one that ends on `proxy_timeout` with the expected responses received --
reports `closed`. `closed` therefore does not distinguish the two for UDP.

Available whether or not `apisix_stream_metrics_zone` is configured.

### $stream_listen_addr

context: `stream`

The configured listening address the session came in on, for example
`0.0.0.0:9100`. This is the key the metrics zone uses for its slots; unlike
`$server_addr`, which holds the address the connection was accepted on, it
stays the same on a wildcard listen.

## Block

### lua
Expand Down
79 changes: 79 additions & 0 deletions lib/resty/apisix/stream/metrics.lua
Original file line number Diff line number Diff line change
@@ -0,0 +1,79 @@
local ffi = require("ffi")
local C = ffi.C
local ffi_str = ffi.string
local tonumber = tonumber


-- No allows_subsystem() guard: the zone is process global and the reader
-- touches neither a session nor any stream context, so http can read it too.


ffi.cdef([[
typedef intptr_t ngx_int_t;

typedef struct {
unsigned char addr[128];
uint32_t addr_len;
uint64_t active;
uint64_t bytes[4];
} ngx_stream_apisix_metrics_entry_t;

typedef uintptr_t ngx_uint_t;

ngx_int_t
ngx_stream_apisix_metrics_dump(ngx_stream_apisix_metrics_entry_t *entries, ngx_uint_t max);
]])


-- must stay in sync with ngx_stream_apisix_metrics_module.h
local MAX_ENTRIES = 512
local DOWNSTREAM_INGRESS = 0
local DOWNSTREAM_EGRESS = 1
local UPSTREAM_EGRESS = 2
local UPSTREAM_INGRESS = 3

local entries = ffi.new("ngx_stream_apisix_metrics_entry_t[?]", MAX_ENTRIES)

-- the stream module is a separate addon, so a build can lack it entirely;
-- resolve the symbol once rather than letting every call throw
local has_dump = pcall(function()
return C.ngx_stream_apisix_metrics_dump
end)

local _M = {}


-- Returns an array of per listening address counters:
-- { listen_addr = "0.0.0.0:9100", active = 3,
-- downstream_ingress = 12, downstream_egress = 34,
-- upstream_ingress = 34, upstream_egress = 12 }
-- The byte counters are monotonic totals since the zone was created.
function _M.dump()
if not has_dump then
return nil, "this runtime has no stream metrics support"
end

local n = C.ngx_stream_apisix_metrics_dump(entries, MAX_ENTRIES)
n = tonumber(n)
if n < 0 then
return nil, "stream metrics zone is not configured"
end

local res = {}
for i = 0, n - 1 do
local e = entries[i]
res[i + 1] = {
listen_addr = ffi_str(e.addr, e.addr_len),
active = tonumber(e.active),
downstream_ingress = tonumber(e.bytes[DOWNSTREAM_INGRESS]),
downstream_egress = tonumber(e.bytes[DOWNSTREAM_EGRESS]),
upstream_egress = tonumber(e.bytes[UPSTREAM_EGRESS]),
upstream_ingress = tonumber(e.bytes[UPSTREAM_INGRESS]),
}
end

return res
end


return _M
120 changes: 120 additions & 0 deletions patch/1.21.4/nginx-stream_metrics.patch
Original file line number Diff line number Diff line change
@@ -0,0 +1,120 @@
diff --git src/stream/ngx_stream_proxy_module.c src/stream/ngx_stream_proxy_module.c
--- src/stream/ngx_stream_proxy_module.c
+++ src/stream/ngx_stream_proxy_module.c
@@ -1410,6 +1410,10 @@

if (c->close) {
ngx_log_error(NGX_LOG_INFO, c->log, 0, "shutdown timeout");
+#if (NGX_STREAM_APISIX)
+ ngx_stream_apisix_set_session_reason(s,
+ NGX_STREAM_APISIX_REASON_SHUTDOWN);
+#endif
ngx_stream_proxy_finalize(s, NGX_STREAM_OK);
return;
}
@@ -1473,6 +1477,11 @@

pc->read->error = 1;

+#if (NGX_STREAM_APISIX)
+ ngx_stream_apisix_set_session_reason(s,
+ NGX_STREAM_APISIX_REASON_UPSTREAM_TIMEOUT);
+#endif
+
ngx_stream_proxy_finalize(s, NGX_STREAM_BAD_GATEWAY);

return;
@@ -1480,6 +1489,10 @@

ngx_connection_error(c, NGX_ETIMEDOUT, "connection timed out");

+#if (NGX_STREAM_APISIX)
+ ngx_stream_apisix_set_session_timeout_reason(s);
+#endif
+
ngx_stream_proxy_finalize(s, NGX_STREAM_OK);

return;
@@ -1516,6 +1529,9 @@

if (ev->timedout) {
ngx_log_error(NGX_LOG_ERR, c->log, NGX_ETIMEDOUT, "upstream timed out");
+#if (NGX_STREAM_APISIX)
+ ngx_stream_apisix_set_connect_timeout(s);
+#endif
ngx_stream_proxy_next_upstream(s);
return;
}
@@ -1612,6 +1628,11 @@

c->log->handler = handler;

+#if (NGX_STREAM_APISIX)
+ ngx_stream_apisix_set_session_reason(s,
+ NGX_STREAM_APISIX_REASON_SHUTDOWN);
+#endif
+
ngx_stream_proxy_finalize(s, NGX_STREAM_OK);
return;
}
@@ -1690,6 +1711,16 @@

c->log->action = recv_action;

+#if (NGX_STREAM_APISIX)
+ /*
+ * ngx_ssl_handle_recv() normalises SSL_ERROR_SSL to no error but
+ * leaves errno alone, so clear it here: whatever is readable after
+ * this recv then belongs to this recv, and a leftover ECONNRESET
+ * cannot be mistaken for a real reset.
+ */
+ ngx_set_socket_errno(0);
+#endif
+
n = src->recv(src, b->last, size);

if (n == NGX_AGAIN) {
@@ -1697,6 +1728,10 @@
}

if (n == NGX_ERROR) {
+#if (NGX_STREAM_APISIX)
+ ngx_stream_apisix_set_read_error(s, from_upstream,
+ ngx_socket_errno);
+#endif
src->read->eof = 1;
n = 0;
}
@@ -1751,6 +1786,10 @@

c->log->action = "proxying connection";

+#if (NGX_STREAM_APISIX)
+ ngx_stream_apisix_metrics_update(s);
+#endif
+
if (ngx_stream_proxy_test_finalize(s, from_upstream) == NGX_OK) {
return;
}
@@ -1930,6 +1969,10 @@
ngx_log_debug1(NGX_LOG_DEBUG_STREAM, s->connection->log, 0,
"close proxy upstream connection: %d", pc->fd);

+#if (NGX_STREAM_APISIX)
+ ngx_stream_apisix_metrics_peer_closing(s);
+#endif
+
#if (NGX_STREAM_SSL)
if (pc->ssl) {
pc->ssl->no_wait_shutdown = 1;
@@ -1960,6 +2003,10 @@
ngx_log_debug1(NGX_LOG_DEBUG_STREAM, s->connection->log, 0,
"finalize stream proxy: %i", rc);

+#if (NGX_STREAM_APISIX)
+ ngx_stream_apisix_metrics_finalize(s, rc);
+#endif
+
u = s->upstream;

if (u == NULL) {
Loading
Loading