diff --git a/README.md b/README.md index 5c850b9..5c09df2 100644 --- a/README.md +++ b/README.md @@ -14,6 +14,111 @@ default: off Disable request mirror until we enable it in the Lua code. +### apisix_stream_metrics_zone \ + +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 diff --git a/lib/resty/apisix/stream/metrics.lua b/lib/resty/apisix/stream/metrics.lua new file mode 100644 index 0000000..0029780 --- /dev/null +++ b/lib/resty/apisix/stream/metrics.lua @@ -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 diff --git a/patch/1.21.4/nginx-stream_metrics.patch b/patch/1.21.4/nginx-stream_metrics.patch new file mode 100644 index 0000000..da086f5 --- /dev/null +++ b/patch/1.21.4/nginx-stream_metrics.patch @@ -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) { diff --git a/patch/1.25.3.1/nginx-stream_metrics.patch b/patch/1.25.3.1/nginx-stream_metrics.patch new file mode 100644 index 0000000..027f523 --- /dev/null +++ b/patch/1.25.3.1/nginx-stream_metrics.patch @@ -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 +@@ -1416,6 +1416,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; + } +@@ -1479,6 +1483,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; +@@ -1486,6 +1495,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; +@@ -1522,6 +1535,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; + } +@@ -1618,6 +1634,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; + } +@@ -1695,6 +1716,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) { +@@ -1702,6 +1733,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; + } +@@ -1756,6 +1791,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; + } +@@ -1935,6 +1974,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; +@@ -1965,6 +2008,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) { diff --git a/patch/1.27.1.1/nginx-stream_metrics.patch b/patch/1.27.1.1/nginx-stream_metrics.patch new file mode 100644 index 0000000..027f523 --- /dev/null +++ b/patch/1.27.1.1/nginx-stream_metrics.patch @@ -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 +@@ -1416,6 +1416,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; + } +@@ -1479,6 +1483,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; +@@ -1486,6 +1495,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; +@@ -1522,6 +1535,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; + } +@@ -1618,6 +1634,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; + } +@@ -1695,6 +1716,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) { +@@ -1702,6 +1733,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; + } +@@ -1756,6 +1791,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; + } +@@ -1935,6 +1974,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; +@@ -1965,6 +2008,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) { diff --git a/patch/1.29.2.4/nginx-stream_metrics.patch b/patch/1.29.2.4/nginx-stream_metrics.patch new file mode 100644 index 0000000..693b713 --- /dev/null +++ b/patch/1.29.2.4/nginx-stream_metrics.patch @@ -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 +@@ -1542,6 +1542,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; + } +@@ -1605,6 +1609,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; +@@ -1612,6 +1621,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; +@@ -1648,6 +1661,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; + } +@@ -1744,6 +1760,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; + } +@@ -1821,6 +1842,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) { +@@ -1828,6 +1859,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; + } +@@ -1882,6 +1917,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; + } +@@ -2061,6 +2100,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; +@@ -2091,6 +2134,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) { diff --git a/src/stream/config b/src/stream/config index efdfad8..6c483f0 100644 --- a/src/stream/config +++ b/src/stream/config @@ -6,6 +6,14 @@ ngx_module_incs="$ngx_addon_dir/" . auto/module -ngx_addon_name=$ngx_module_name +ngx_module_type=STREAM +ngx_module_name=ngx_stream_apisix_metrics_module +ngx_module_srcs="$ngx_addon_dir/ngx_stream_apisix_metrics_module.c" +ngx_module_deps=$ngx_addon_dir/ngx_stream_apisix_metrics_module.h +ngx_module_incs="$ngx_addon_dir/" + +. auto/module + +ngx_addon_name=ngx_stream_apisix_module have=NGX_STREAM_APISIX . auto/have diff --git a/src/stream/ngx_stream_apisix_metrics_module.c b/src/stream/ngx_stream_apisix_metrics_module.c new file mode 100644 index 0000000..6e3b24d --- /dev/null +++ b/src/stream/ngx_stream_apisix_metrics_module.c @@ -0,0 +1,939 @@ +#include +#include +#include +#include "ngx_stream_apisix_metrics_module.h" + + +/* + * A session accumulates locally and merges into the shared zone once per + * proxy_process pass, which is the point where nginx has drained everything + * currently readable. That coalesces the inner read/write loop into a single + * atomic per direction, without letting a quiesced session hold its last + * bytes back: there is no later event to flush them on, so anything deferred + * here would stay invisible until the session ended. + */ + +/* + * One slot per stream listening address. A deployment has a handful of them, + * the cap only keeps a huge zone from reserving a pointless amount of memory. + */ +#define NGX_STREAM_APISIX_METRICS_MAX_SLOTS 512 + + +typedef struct { + ngx_atomic_t active; + ngx_atomic_t bytes[NGX_STREAM_APISIX_METRICS_DIRECTIONS]; + uint32_t addr_len; + u_char addr[NGX_STREAM_APISIX_METRICS_ADDR_LEN]; +} ngx_stream_apisix_metrics_slot_t; + + +typedef struct { + ngx_uint_t nslots; + ngx_uint_t nused; + ngx_stream_apisix_metrics_slot_t *slots; +} ngx_stream_apisix_metrics_sh_t; + + +typedef struct { + ngx_shm_zone_t *shm_zone; +} ngx_stream_apisix_metrics_main_conf_t; + + +typedef struct { + ngx_stream_apisix_metrics_slot_t *slot; + off_t flushed[NGX_STREAM_APISIX_METRICS_DIRECTIONS]; + ngx_uint_t reason; + /* errno of the fatal recv(), indexed by from_upstream */ + ngx_err_t read_err[2]; + unsigned counted:1; + unsigned connect_timeout:1; + unsigned finalized:1; +} ngx_stream_apisix_metrics_ctx_t; + + +static ngx_int_t ngx_stream_apisix_metrics_preconf(ngx_conf_t *cf); +static ngx_int_t ngx_stream_apisix_metrics_postconf(ngx_conf_t *cf); +static void *ngx_stream_apisix_metrics_create_main_conf(ngx_conf_t *cf); +static char *ngx_stream_apisix_metrics_zone(ngx_conf_t *cf, ngx_command_t *cmd, + void *conf); +static ngx_int_t ngx_stream_apisix_metrics_init_zone(ngx_shm_zone_t *shm_zone, + void *data); +static ngx_int_t ngx_stream_apisix_metrics_init_module(ngx_cycle_t *cycle); +static ngx_int_t ngx_stream_apisix_metrics_init_process(ngx_cycle_t *cycle); +static ngx_int_t ngx_stream_apisix_metrics_post_accept_handler( + ngx_stream_session_t *s); +static ngx_int_t ngx_stream_apisix_metrics_log_handler( + ngx_stream_session_t *s); +static ngx_int_t ngx_stream_apisix_session_reason_variable( + ngx_stream_session_t *s, ngx_stream_variable_value_t *v, uintptr_t data); +static ngx_int_t ngx_stream_apisix_listen_addr_variable( + ngx_stream_session_t *s, ngx_stream_variable_value_t *v, uintptr_t data); + + +static ngx_command_t ngx_stream_apisix_metrics_cmds[] = { + + { ngx_string("apisix_stream_metrics_zone"), + NGX_STREAM_MAIN_CONF|NGX_CONF_TAKE1, + ngx_stream_apisix_metrics_zone, + NGX_STREAM_MAIN_CONF_OFFSET, + 0, + NULL }, + + ngx_null_command +}; + + +static ngx_stream_module_t ngx_stream_apisix_metrics_module_ctx = { + ngx_stream_apisix_metrics_preconf, /* preconfiguration */ + ngx_stream_apisix_metrics_postconf, /* postconfiguration */ + + ngx_stream_apisix_metrics_create_main_conf, /* create main configuration */ + NULL, /* init main configuration */ + + NULL, /* create server configuration */ + NULL, /* merge server configuration */ +}; + + +ngx_module_t ngx_stream_apisix_metrics_module = { + NGX_MODULE_V1, + &ngx_stream_apisix_metrics_module_ctx, /* module context */ + ngx_stream_apisix_metrics_cmds, /* module directives */ + NGX_STREAM_MODULE, /* module type */ + NULL, /* init master */ + ngx_stream_apisix_metrics_init_module, /* init module */ + ngx_stream_apisix_metrics_init_process, /* init process */ + NULL, /* init thread */ + NULL, /* exit thread */ + NULL, /* exit process */ + NULL, /* exit master */ + NGX_MODULE_V1_PADDING +}; + + +static ngx_str_t ngx_stream_apisix_metrics_zone_name = + ngx_string("apisix_stream_metrics"); + + +static ngx_stream_variable_t ngx_stream_apisix_metrics_vars[] = { + + { ngx_string("stream_session_reason"), NULL, + ngx_stream_apisix_session_reason_variable, 0, + NGX_STREAM_VAR_NOCACHEABLE, 0 }, + + { ngx_string("stream_listen_addr"), NULL, + ngx_stream_apisix_listen_addr_variable, 0, + NGX_STREAM_VAR_NOCACHEABLE, 0 }, + + ngx_stream_null_variable +}; + + +static ngx_str_t ngx_stream_apisix_reasons[] = { + ngx_string("-"), + ngx_string("closed"), + ngx_string("client_rst"), + ngx_string("client_error"), + ngx_string("upstream_rst"), + ngx_string("upstream_error"), + ngx_string("connect_timeout"), + ngx_string("recv_timeout"), + ngx_string("send_timeout"), + ngx_string("upstream_timeout"), + ngx_string("shutdown"), + ngx_string("connect_failed"), + ngx_string("client_read_error"), + ngx_string("upstream_read_error") +}; + +#define NGX_STREAM_APISIX_REASONS_N \ + (sizeof(ngx_stream_apisix_reasons) / sizeof(ngx_stream_apisix_reasons[0])) + + +/* + * Rebound per cycle by ngx_stream_apisix_metrics_bind_zone(), so that the FFI + * reader does not depend on a session, and so that a reload dropping the zone + * cannot leave this pointing into unmapped shared memory. + */ +static ngx_stream_apisix_metrics_sh_t *ngx_stream_apisix_metrics_sh = NULL; + +/* + * Maps a listening socket to its slot without a lookup on the hot path. + * Rebuilt per cycle in init_process, so a reload picks up the new listen set. + */ +static ngx_stream_apisix_metrics_slot_t **ngx_stream_apisix_metrics_map = NULL; +static ngx_listening_t *ngx_stream_apisix_metrics_ls = NULL; +static ngx_uint_t ngx_stream_apisix_metrics_nls = 0; + + +static void * +ngx_stream_apisix_metrics_create_main_conf(ngx_conf_t *cf) +{ + return ngx_pcalloc(cf->pool, sizeof(ngx_stream_apisix_metrics_main_conf_t)); +} + + +static char * +ngx_stream_apisix_metrics_zone(ngx_conf_t *cf, ngx_command_t *cmd, void *conf) +{ + ngx_stream_apisix_metrics_main_conf_t *mcf = conf; + + ssize_t size; + ngx_str_t *value; + + if (mcf->shm_zone) { + return "is duplicate"; + } + + value = cf->args->elts; + + size = ngx_parse_size(&value[1]); + if (size == NGX_ERROR) { + ngx_conf_log_error(NGX_LOG_EMERG, cf, 0, + "invalid zone size \"%V\"", &value[1]); + return NGX_CONF_ERROR; + } + + if (size < (ssize_t) (8 * ngx_pagesize)) { + ngx_conf_log_error(NGX_LOG_EMERG, cf, 0, + "zone \"%V\" is too small", + &ngx_stream_apisix_metrics_zone_name); + return NGX_CONF_ERROR; + } + + mcf->shm_zone = ngx_shared_memory_add(cf, + &ngx_stream_apisix_metrics_zone_name, + size, + &ngx_stream_apisix_metrics_module); + if (mcf->shm_zone == NULL) { + return NGX_CONF_ERROR; + } + + mcf->shm_zone->init = ngx_stream_apisix_metrics_init_zone; + + return NGX_CONF_OK; +} + + +static ngx_int_t +ngx_stream_apisix_metrics_init_zone(ngx_shm_zone_t *shm_zone, void *data) +{ + ngx_stream_apisix_metrics_sh_t *osh = data; + + ngx_uint_t nslots; + ngx_slab_pool_t *shpool; + ngx_stream_apisix_metrics_sh_t *sh; + + if (osh) { + /* reused on reload, the accumulated counters must survive */ + shm_zone->data = osh; + return NGX_OK; + } + + shpool = (ngx_slab_pool_t *) shm_zone->shm.addr; + + if (shm_zone->shm.exists) { + shm_zone->data = shpool->data; + return NGX_OK; + } + + /* + * The slot array is sized once and never grows: the set of stream + * listening addresses is fixed at configuration time. Only half of the + * zone is handed out so that the slab bookkeeping always fits. + */ + nslots = (shm_zone->shm.size - sizeof(ngx_slab_pool_t)) / 2 + / sizeof(ngx_stream_apisix_metrics_slot_t); + + if (nslots > NGX_STREAM_APISIX_METRICS_MAX_SLOTS) { + nslots = NGX_STREAM_APISIX_METRICS_MAX_SLOTS; + } + + sh = ngx_slab_calloc(shpool, sizeof(ngx_stream_apisix_metrics_sh_t)); + if (sh == NULL) { + return NGX_ERROR; + } + + sh->slots = ngx_slab_calloc(shpool, + nslots * sizeof(ngx_stream_apisix_metrics_slot_t)); + if (sh->slots == NULL) { + return NGX_ERROR; + } + + sh->nslots = nslots; + sh->nused = 0; + + shpool->data = sh; + shm_zone->data = sh; + + return NGX_OK; +} + + +/* + * Resolved per cycle rather than cached when the zone is created: a reload + * that drops the directive (or the whole stream block) leaves no zone in the + * new cycle, and ngx_init_cycle then unmaps the old one. A pointer kept from + * the previous cycle would dangle into freed shared memory. + */ +static void +ngx_stream_apisix_metrics_bind_zone(ngx_cycle_t *cycle) +{ + ngx_stream_apisix_metrics_main_conf_t *mcf; + + ngx_stream_apisix_metrics_sh = NULL; + + mcf = ngx_stream_cycle_get_module_main_conf(cycle, + ngx_stream_apisix_metrics_module); + + if (mcf != NULL && mcf->shm_zone != NULL) { + ngx_stream_apisix_metrics_sh = mcf->shm_zone->data; + } +} + + +static ngx_stream_apisix_metrics_slot_t * +ngx_stream_apisix_metrics_lookup(ngx_str_t *addr, ngx_uint_t create) +{ + ngx_uint_t i; + ngx_uint_t nused; + ngx_stream_apisix_metrics_sh_t *sh; + ngx_stream_apisix_metrics_slot_t *slot; + + sh = ngx_stream_apisix_metrics_sh; + if (sh == NULL) { + return NULL; + } + + /* + * Pair with the barrier the writer takes before publishing the count: the + * slot contents must not be read before the count that exposes them. + */ + nused = sh->nused; + ngx_memory_barrier(); + + for (i = 0; i < nused; i++) { + slot = &sh->slots[i]; + + if (slot->addr_len == addr->len + && ngx_memcmp(slot->addr, addr->data, addr->len) == 0) + { + return slot; + } + } + + if (!create || nused == sh->nslots) { + return NULL; + } + + slot = &sh->slots[sh->nused]; + + slot->addr_len = (uint32_t) addr->len; + ngx_memcpy(slot->addr, addr->data, addr->len); + + /* + * A reload reuses the zone, so the master appends here while the previous + * generation of workers is still reading. Publish the contents before the + * count that exposes them. + */ + ngx_memory_barrier(); + + sh->nused++; + + return slot; +} + + +/* + * Only the master ever claims, and only by appending, so no lock is needed: + * readers either see a slot fully or do not see it at all. On a reload the + * previous generation of workers is still reading, which is what the barrier + * in the lookup above is for. + */ +static ngx_int_t +ngx_stream_apisix_metrics_init_module(ngx_cycle_t *cycle) +{ + ngx_str_t addr; + ngx_uint_t i; + ngx_listening_t *ls; + + ngx_stream_apisix_metrics_bind_zone(cycle); + + if (ngx_stream_apisix_metrics_sh == NULL) { + return NGX_OK; + } + + ls = cycle->listening.elts; + + for (i = 0; i < cycle->listening.nelts; i++) { + if (ls[i].handler != ngx_stream_init_connection) { + continue; + } + + /* + * Unix sockets inside stream{} are internal plumbing, not proxy + * ports: APISIX puts its worker event channel there, and counting it + * would report gateway control traffic as proxied bytes and leave a + * permanent floor of worker connections in the active gauge. + */ +#if (NGX_HAVE_UNIX_DOMAIN) + if (ls[i].sockaddr->sa_family == AF_UNIX) { + continue; + } +#endif + + addr = ls[i].addr_text; + + if (addr.len == 0) { + continue; + } + + if (addr.len > NGX_STREAM_APISIX_METRICS_ADDR_LEN) { + ngx_log_error(NGX_LOG_WARN, cycle->log, 0, + "apisix stream metrics: listening address \"%V\" is " + "longer than %d bytes and is not accounted for", + &addr, NGX_STREAM_APISIX_METRICS_ADDR_LEN); + continue; + } + + if (ngx_stream_apisix_metrics_lookup(&addr, 1) == NULL) { + ngx_log_error(NGX_LOG_WARN, cycle->log, 0, + "apisix stream metrics zone is too small to hold " + "listening address \"%V\"", &addr); + } + } + + return NGX_OK; +} + + +static ngx_int_t +ngx_stream_apisix_metrics_init_process(ngx_cycle_t *cycle) +{ + ngx_uint_t i; + ngx_listening_t *ls; + + ngx_stream_apisix_metrics_map = NULL; + ngx_stream_apisix_metrics_ls = NULL; + ngx_stream_apisix_metrics_nls = 0; + + ngx_stream_apisix_metrics_bind_zone(cycle); + + if (ngx_stream_apisix_metrics_sh == NULL || cycle->listening.nelts == 0) { + return NGX_OK; + } + + ls = cycle->listening.elts; + + ngx_stream_apisix_metrics_map = ngx_pcalloc(cycle->pool, + cycle->listening.nelts * sizeof(ngx_stream_apisix_metrics_slot_t *)); + if (ngx_stream_apisix_metrics_map == NULL) { + return NGX_ERROR; + } + + for (i = 0; i < cycle->listening.nelts; i++) { + if (ls[i].handler != ngx_stream_init_connection) { + continue; + } + + ngx_stream_apisix_metrics_map[i] = + ngx_stream_apisix_metrics_lookup(&ls[i].addr_text, 0); + } + + ngx_stream_apisix_metrics_ls = ls; + ngx_stream_apisix_metrics_nls = cycle->listening.nelts; + + return NGX_OK; +} + + +static ngx_stream_apisix_metrics_slot_t * +ngx_stream_apisix_metrics_slot(ngx_connection_t *c) +{ + ngx_uint_t i; + + if (ngx_stream_apisix_metrics_map == NULL || c->listening == NULL) { + return NULL; + } + + i = (ngx_uint_t) (c->listening - ngx_stream_apisix_metrics_ls); + + if (i >= ngx_stream_apisix_metrics_nls) { + return NULL; + } + + return ngx_stream_apisix_metrics_map[i]; +} + + +static ngx_stream_apisix_metrics_ctx_t * +ngx_stream_apisix_metrics_get_ctx(ngx_stream_session_t *s) +{ + return ngx_stream_get_module_ctx(s, ngx_stream_apisix_metrics_module); +} + + +static ngx_int_t +ngx_stream_apisix_metrics_post_accept_handler(ngx_stream_session_t *s) +{ + ngx_stream_apisix_metrics_ctx_t *ctx; + ngx_stream_apisix_metrics_slot_t *slot; + + /* + * The context is always allocated: $stream_session_reason needs it to + * carry the reasons recorded by the proxy module even when no metrics + * zone is configured. + */ + ctx = ngx_pcalloc(s->connection->pool, + sizeof(ngx_stream_apisix_metrics_ctx_t)); + if (ctx == NULL) { + return NGX_ERROR; + } + + ngx_stream_set_ctx(s, ctx, ngx_stream_apisix_metrics_module); + + slot = ngx_stream_apisix_metrics_slot(s->connection); + if (slot == NULL) { + return NGX_DECLINED; + } + + ctx->slot = slot; + ctx->counted = 1; + + (void) ngx_atomic_fetch_add(&slot->active, 1); + + return NGX_DECLINED; +} + + +static void +ngx_stream_apisix_metrics_flush(ngx_stream_session_t *s, + ngx_stream_apisix_metrics_ctx_t *ctx) +{ + off_t current[NGX_STREAM_APISIX_METRICS_DIRECTIONS]; + off_t delta; + ngx_uint_t i; + ngx_connection_t *pc; + ngx_stream_upstream_t *u; + ngx_stream_apisix_metrics_slot_t *slot; + + slot = ctx->slot; + if (slot == NULL) { + return; + } + + u = s->upstream; + pc = (u && u->peer.connection) ? u->peer.connection : NULL; + + current[NGX_STREAM_APISIX_METRICS_DOWNSTREAM_INGRESS] = s->received; + current[NGX_STREAM_APISIX_METRICS_DOWNSTREAM_EGRESS] = s->connection->sent; + current[NGX_STREAM_APISIX_METRICS_UPSTREAM_EGRESS] = pc ? pc->sent : 0; + current[NGX_STREAM_APISIX_METRICS_UPSTREAM_INGRESS] = u ? u->received : 0; + + for (i = 0; i < NGX_STREAM_APISIX_METRICS_DIRECTIONS; i++) { + delta = current[i] - ctx->flushed[i]; + + /* + * proxy_next_upstream replaces the peer connection, so pc->sent + * restarts from zero while the flushed mark still holds the previous + * peer's total. Treat any decrease as a restart and count what the + * new counter holds, otherwise the retried bytes are lost. + */ + if (delta < 0) { + delta = current[i]; + } + + if (delta == 0) { + continue; + } + + (void) ngx_atomic_fetch_add(&slot->bytes[i], (ngx_atomic_int_t) delta); + ctx->flushed[i] = current[i]; + } +} + + +void +ngx_stream_apisix_metrics_update(ngx_stream_session_t *s) +{ + ngx_stream_apisix_metrics_ctx_t *ctx; + + ctx = ngx_stream_apisix_metrics_get_ctx(s); + if (ctx == NULL) { + return; + } + + ngx_stream_apisix_metrics_flush(s, ctx); +} + + +/* + * proxy_next_upstream is about to drop the peer, taking pc->sent with it. + * Bank what it sent and reset the mark, otherwise the next peer starts from + * zero below the old mark and the retried bytes are only recovered if that + * counter happens to overtake it before the next sample. + */ +void +ngx_stream_apisix_metrics_peer_closing(ngx_stream_session_t *s) +{ + ngx_stream_apisix_metrics_ctx_t *ctx; + + ctx = ngx_stream_apisix_metrics_get_ctx(s); + if (ctx == NULL) { + return; + } + + ngx_stream_apisix_metrics_flush(s, ctx); + + ctx->flushed[NGX_STREAM_APISIX_METRICS_UPSTREAM_EGRESS] = 0; +} + + +/* + * nginx raises read->error for every fatal recv(), not only for a reset, and + * ngx_ssl_recv adds protocol failures on top. The errno is the only thing + * that tells them apart, and it is gone by the time the session is finalized, + * so the proxy module hands it over the moment the read fails. + */ +void +ngx_stream_apisix_set_read_error(ngx_stream_session_t *s, + ngx_uint_t from_upstream, ngx_err_t err) +{ + ngx_stream_apisix_metrics_ctx_t *ctx; + + ctx = ngx_stream_apisix_metrics_get_ctx(s); + if (ctx == NULL) { + return; + } + + ctx->read_err[from_upstream ? 1 : 0] = err; +} + + +void +ngx_stream_apisix_set_session_reason(ngx_stream_session_t *s, + ngx_uint_t reason) +{ + ngx_stream_apisix_metrics_ctx_t *ctx; + + /* + * The callers live in the ngx_stream_proxy_module patches, which are + * separate files that can drift from this enum. The value indexes the + * reason table, so refuse anything out of range rather than read past it. + */ + if (reason >= NGX_STREAM_APISIX_REASONS_N) { + return; + } + + ctx = ngx_stream_apisix_metrics_get_ctx(s); + if (ctx == NULL) { + return; + } + + /* the first terminating event wins, later ones are consequences of it */ + if (ctx->reason == NGX_STREAM_APISIX_REASON_UNSET) { + ctx->reason = reason; + } +} + + +/* + * A connect timeout is not terminal on its own: nginx may still reach another + * peer. It is only reported when the session really ended on a bad gateway. + */ +void +ngx_stream_apisix_set_connect_timeout(ngx_stream_session_t *s) +{ + ngx_stream_apisix_metrics_ctx_t *ctx; + + ctx = ngx_stream_apisix_metrics_get_ctx(s); + if (ctx == NULL) { + return; + } + + ctx->connect_timeout = 1; +} + + +/* + * proxy_timeout is a single idle timeout shared by both directions. Data + * still buffered towards a peer means that peer stopped reading, which is + * what the product calls a send timeout; anything else is a receive timeout. + */ +void +ngx_stream_apisix_set_session_timeout_reason(ngx_stream_session_t *s) +{ + ngx_uint_t reason; + ngx_connection_t *pc; + ngx_stream_upstream_t *u; + + u = s->upstream; + pc = (u && u->peer.connection) ? u->peer.connection : NULL; + + if (s->connection->buffered || (pc && pc->buffered)) { + reason = NGX_STREAM_APISIX_REASON_SEND_TIMEOUT; + + } else { + reason = NGX_STREAM_APISIX_REASON_RECV_TIMEOUT; + } + + ngx_stream_apisix_set_session_reason(s, reason); +} + + +/* + * Reasons that nginx does not record itself are derived from the connection + * flags: a reset leaves read->error set even though the proxy module turns it + * into an EOF afterwards, and a failed write leaves error set on the + * destination connection. + */ +static ngx_uint_t +ngx_stream_apisix_derive_reason(ngx_stream_session_t *s, + ngx_stream_apisix_metrics_ctx_t *ctx, ngx_uint_t status) +{ + ngx_connection_t *c, *pc; + ngx_stream_upstream_t *u; + + if (ctx && ctx->connect_timeout && status == NGX_STREAM_BAD_GATEWAY) { + return NGX_STREAM_APISIX_REASON_CONNECT_TIMEOUT; + } + + c = s->connection; + u = s->upstream; + pc = (u && u->peer.connection) ? u->peer.connection : NULL; + + if (c->read->error) { + return ctx && ctx->read_err[0] == NGX_ECONNRESET + ? NGX_STREAM_APISIX_REASON_CLIENT_RST + : NGX_STREAM_APISIX_REASON_CLIENT_READ_ERROR; + } + + if (pc && pc->read->error) { + return ctx && ctx->read_err[1] == NGX_ECONNRESET + ? NGX_STREAM_APISIX_REASON_UPSTREAM_RST + : NGX_STREAM_APISIX_REASON_UPSTREAM_READ_ERROR; + } + + if (c->error) { + return NGX_STREAM_APISIX_REASON_CLIENT_ERROR; + } + + if (pc && pc->error) { + return NGX_STREAM_APISIX_REASON_UPSTREAM_ERROR; + } + + if (c->read->eof || (pc && pc->read->eof)) { + return NGX_STREAM_APISIX_REASON_CLOSED; + } + + return NGX_STREAM_APISIX_REASON_UNSET; +} + + +/* + * Called from the proxy module before it tears the upstream connection down: + * by log time u->peer.connection is already NULL, so both the reason and the + * last byte delta towards the upstream have to be taken here. + */ +void +ngx_stream_apisix_metrics_finalize(ngx_stream_session_t *s, ngx_uint_t rc) +{ + ngx_stream_apisix_metrics_ctx_t *ctx; + + ctx = ngx_stream_apisix_metrics_get_ctx(s); + if (ctx == NULL || ctx->finalized) { + return; + } + + ctx->finalized = 1; + + ngx_stream_apisix_metrics_flush(s, ctx); + + if (ctx->reason == NGX_STREAM_APISIX_REASON_UNSET) { + ctx->reason = ngx_stream_apisix_derive_reason(s, ctx, rc); + } + + /* + * A UDP session has no FIN to observe: ngx_udp_shared_recv never sets + * read->eof, so a session that simply ran to completion would otherwise + * report no reason at all. + */ + if (ctx->reason == NGX_STREAM_APISIX_REASON_UNSET + && rc == NGX_STREAM_OK) + { + ctx->reason = NGX_STREAM_APISIX_REASON_CLOSED; + } + + /* + * Everything that gives up before a peer answers lands here: connection + * refused, no live upstream, a failed or timed out upstream handshake. + * None of them leave a mark on the connection flags. + */ + if (ctx->reason == NGX_STREAM_APISIX_REASON_UNSET + && rc == NGX_STREAM_BAD_GATEWAY) + { + ctx->reason = NGX_STREAM_APISIX_REASON_CONNECT_FAILED; + } +} + + +static ngx_int_t +ngx_stream_apisix_session_reason_variable(ngx_stream_session_t *s, + ngx_stream_variable_value_t *v, uintptr_t data) +{ + ngx_uint_t reason_index; + ngx_str_t *reason; + ngx_stream_apisix_metrics_ctx_t *ctx; + + ctx = ngx_stream_apisix_metrics_get_ctx(s); + + if (ctx && ctx->reason != NGX_STREAM_APISIX_REASON_UNSET) { + reason_index = ctx->reason; + + } else { + /* sessions rejected before proxy_pass never reach the finalize hook */ + reason_index = ngx_stream_apisix_derive_reason(s, ctx, s->status); + } + + reason = &ngx_stream_apisix_reasons[reason_index]; + + v->len = reason->len; + v->data = reason->data; + v->valid = 1; + v->no_cacheable = 1; + v->not_found = 0; + + return NGX_OK; +} + + +/* + * The configured listening address, which is how the metrics zone keys its + * slots. $server_addr holds the address the connection was accepted on, so on + * a wildcard listen the two would not agree. + */ +static ngx_int_t +ngx_stream_apisix_listen_addr_variable(ngx_stream_session_t *s, + ngx_stream_variable_value_t *v, uintptr_t data) +{ + ngx_listening_t *ls; + + ls = s->connection->listening; + + if (ls == NULL) { + v->not_found = 1; + return NGX_OK; + } + + v->len = ls->addr_text.len; + v->data = ls->addr_text.data; + v->valid = 1; + v->no_cacheable = 1; + v->not_found = 0; + + return NGX_OK; +} + + +static ngx_int_t +ngx_stream_apisix_metrics_log_handler(ngx_stream_session_t *s) +{ + ngx_stream_apisix_metrics_ctx_t *ctx; + + ctx = ngx_stream_apisix_metrics_get_ctx(s); + if (ctx == NULL) { + return NGX_OK; + } + + ngx_stream_apisix_metrics_flush(s, ctx); + + if (ctx->counted) { + ctx->counted = 0; + (void) ngx_atomic_fetch_add(&ctx->slot->active, -1); + } + + return NGX_OK; +} + + +ngx_int_t +ngx_stream_apisix_metrics_dump(ngx_stream_apisix_metrics_entry_t *entries, + ngx_uint_t max) +{ + ngx_uint_t i, j, n; + ngx_stream_apisix_metrics_sh_t *sh; + ngx_stream_apisix_metrics_slot_t *slot; + + sh = ngx_stream_apisix_metrics_sh; + if (sh == NULL) { + return NGX_ERROR; + } + + /* see the matching barrier in ngx_stream_apisix_metrics_lookup() */ + n = sh->nused; + ngx_memory_barrier(); + + n = ngx_min(n, max); + + for (i = 0; i < n; i++) { + slot = &sh->slots[i]; + + entries[i].addr_len = ngx_min(slot->addr_len, + NGX_STREAM_APISIX_METRICS_ADDR_LEN); + ngx_memcpy(entries[i].addr, slot->addr, entries[i].addr_len); + + entries[i].active = (uint64_t) slot->active; + + for (j = 0; j < NGX_STREAM_APISIX_METRICS_DIRECTIONS; j++) { + entries[i].bytes[j] = (uint64_t) slot->bytes[j]; + } + } + + return (ngx_int_t) n; +} + + +static ngx_int_t +ngx_stream_apisix_metrics_preconf(ngx_conf_t *cf) +{ + ngx_stream_variable_t *var, *v; + + for (v = ngx_stream_apisix_metrics_vars; v->name.len; v++) { + var = ngx_stream_add_variable(cf, &v->name, v->flags); + if (var == NULL) { + return NGX_ERROR; + } + + var->get_handler = v->get_handler; + var->data = v->data; + } + + return NGX_OK; +} + + +static ngx_int_t +ngx_stream_apisix_metrics_postconf(ngx_conf_t *cf) +{ + ngx_stream_handler_pt *h; + ngx_stream_core_main_conf_t *cmcf; + + cmcf = ngx_stream_conf_get_module_main_conf(cf, ngx_stream_core_module); + + h = ngx_array_push(&cmcf->phases[NGX_STREAM_POST_ACCEPT_PHASE].handlers); + if (h == NULL) { + return NGX_ERROR; + } + + *h = ngx_stream_apisix_metrics_post_accept_handler; + + h = ngx_array_push(&cmcf->phases[NGX_STREAM_LOG_PHASE].handlers); + if (h == NULL) { + return NGX_ERROR; + } + + *h = ngx_stream_apisix_metrics_log_handler; + + return NGX_OK; +} diff --git a/src/stream/ngx_stream_apisix_metrics_module.h b/src/stream/ngx_stream_apisix_metrics_module.h new file mode 100644 index 0000000..301ca20 --- /dev/null +++ b/src/stream/ngx_stream_apisix_metrics_module.h @@ -0,0 +1,73 @@ +#ifndef _NGX_STREAM_APISIX_METRICS_H_INCLUDED_ +#define _NGX_STREAM_APISIX_METRICS_H_INCLUDED_ + + +#include + + +/* + * Byte counters kept per listening address. The direction is always named + * from the gateway's point of view, matching the `type` label of the HTTP + * side `apisix_bandwidth` metric: ingress is what the gateway receives, + * egress is what the gateway sends. + */ +#define NGX_STREAM_APISIX_METRICS_DOWNSTREAM_INGRESS 0 +#define NGX_STREAM_APISIX_METRICS_DOWNSTREAM_EGRESS 1 +#define NGX_STREAM_APISIX_METRICS_UPSTREAM_EGRESS 2 +#define NGX_STREAM_APISIX_METRICS_UPSTREAM_INGRESS 3 +#define NGX_STREAM_APISIX_METRICS_DIRECTIONS 4 + +#define NGX_STREAM_APISIX_METRICS_ADDR_LEN 128 + + +/* + * Why a session ended. nginx keeps `s->status` at 200 for every failure that + * happens after the upstream connection is established, so the reasons below + * are what the Lua side needs to tell a graceful close apart from a timeout + * or a reset. Exposed as $stream_session_reason. + */ +typedef enum { + NGX_STREAM_APISIX_REASON_UNSET = 0, + NGX_STREAM_APISIX_REASON_CLOSED, + NGX_STREAM_APISIX_REASON_CLIENT_RST, + NGX_STREAM_APISIX_REASON_CLIENT_ERROR, + NGX_STREAM_APISIX_REASON_UPSTREAM_RST, + NGX_STREAM_APISIX_REASON_UPSTREAM_ERROR, + NGX_STREAM_APISIX_REASON_CONNECT_TIMEOUT, + NGX_STREAM_APISIX_REASON_RECV_TIMEOUT, + NGX_STREAM_APISIX_REASON_SEND_TIMEOUT, + NGX_STREAM_APISIX_REASON_UPSTREAM_TIMEOUT, + NGX_STREAM_APISIX_REASON_SHUTDOWN, + NGX_STREAM_APISIX_REASON_CONNECT_FAILED, + NGX_STREAM_APISIX_REASON_CLIENT_READ_ERROR, + NGX_STREAM_APISIX_REASON_UPSTREAM_READ_ERROR +} ngx_stream_apisix_reason_e; + + +/* the layout Lua reads through FFI, one entry per listening address */ +typedef struct { + u_char addr[NGX_STREAM_APISIX_METRICS_ADDR_LEN]; + uint32_t addr_len; + uint64_t active; + uint64_t bytes[NGX_STREAM_APISIX_METRICS_DIRECTIONS]; +} ngx_stream_apisix_metrics_entry_t; + + +/* called from the patched ngx_stream_proxy_module */ +void ngx_stream_apisix_set_session_reason(ngx_stream_session_t *s, + ngx_uint_t reason); +void ngx_stream_apisix_set_session_timeout_reason(ngx_stream_session_t *s); +void ngx_stream_apisix_set_connect_timeout(ngx_stream_session_t *s); +void ngx_stream_apisix_metrics_update(ngx_stream_session_t *s); +void ngx_stream_apisix_metrics_peer_closing(ngx_stream_session_t *s); +void ngx_stream_apisix_set_read_error(ngx_stream_session_t *s, + ngx_uint_t from_upstream, ngx_err_t err); +void ngx_stream_apisix_metrics_finalize(ngx_stream_session_t *s, + ngx_uint_t rc); + +/* called from Lua through FFI */ +ngx_int_t ngx_stream_apisix_metrics_dump( + ngx_stream_apisix_metrics_entry_t *entries, ngx_uint_t max); + + +#endif /* _NGX_STREAM_APISIX_METRICS_H_INCLUDED_ */ diff --git a/src/stream/ngx_stream_apisix_module.h b/src/stream/ngx_stream_apisix_module.h index 9ded301..45d838c 100644 --- a/src/stream/ngx_stream_apisix_module.h +++ b/src/stream/ngx_stream_apisix_module.h @@ -4,6 +4,12 @@ #include +/* + * This header is the single entry point the ngx_stream_proxy_module patches + * include, so everything they call has to be reachable from here. + */ +#include + ngx_int_t ngx_stream_apisix_is_proxy_ssl_enabled(ngx_stream_session_t *s); diff --git a/t/stream/metrics.t b/t/stream/metrics.t new file mode 100644 index 0000000..c1d0af5 --- /dev/null +++ b/t/stream/metrics.t @@ -0,0 +1,328 @@ +use t::APISIX_NGINX 'no_plan'; + +add_block_preprocessor(sub { + my ($block) = @_; + + if (!$block->http_config) { + my $http_config = <<'_EOC_'; + server { + listen 1994; + + location / { + content_by_lua_block { + ngx.print("hello") + } + } + } + +_EOC_ + + $block->set_value("http_config", $http_config); + } + + # blocks that need to run without a zone set stream_config themselves + if (defined $block->stream_server_config && !defined $block->stream_config) { + $block->set_value("stream_config", "apisix_stream_metrics_zone 1m;\n"); + } +}); + +run_tests(); + +__DATA__ + +=== TEST 1: a session closed by the upstream is reported as a normal close +--- stream_server_config + proxy_pass 127.0.0.1:1994; + log_by_lua_block { + ngx.log(ngx.WARN, "reason: ", ngx.var.stream_session_reason, + ", listen: ", ngx.var.stream_listen_addr) + } +--- stream_request eval +"GET / HTTP/1.0\r\nHost: localhost\r\n\r\n" +--- stream_response_like: hello +--- error_log eval +qr/reason: closed, listen: [\d.]+:\d+/ + + + +=== TEST 2: an idle session is reported as a receive timeout, not as 200/closed +--- stream_server_config + proxy_timeout 100ms; + proxy_pass 127.0.0.1:1994; + log_by_lua_block { + ngx.log(ngx.WARN, "reason: ", ngx.var.stream_session_reason, + ", status: ", ngx.var.status) + } +--- stream_request eval +"GET /" +--- timeout: 5 +--- wait: 0.5 +--- error_log +reason: recv_timeout, status: 200 + + + +=== TEST 3: an unreachable upstream keeps the native 502 +--- stream_server_config + proxy_connect_timeout 300ms; + proxy_pass 127.0.0.1:1979; + log_by_lua_block { + ngx.log(ngx.WARN, "status: ", ngx.var.status, + ", reason: ", ngx.var.stream_session_reason) + } +--- stream_request eval +"GET /" +--- error_log +status: 502, reason: connect_failed + + + +=== TEST 4: bytes are accounted per listening address in both directions +--- stream_server_config + proxy_pass 127.0.0.1:1994; + log_by_lua_block { + local metrics = require("resty.apisix.stream.metrics") + local res, err = metrics.dump() + if not res then + ngx.log(ngx.ERR, "dump failed: ", err) + return + end + + for _, e in ipairs(res) do + ngx.log(ngx.WARN, "slot ", e.listen_addr, + " active=", e.active, + " di=", e.downstream_ingress > 0, + " de=", e.downstream_egress > 0, + " ue=", e.upstream_egress > 0, + " ui=", e.upstream_ingress > 0) + end + } +--- stream_request eval +"GET / HTTP/1.0\r\nHost: localhost\r\n\r\n" +--- stream_response_like: hello +--- error_log +active=1 di=true de=true ue=true ui=true + + + +=== TEST 5: the FFI reader reports a missing zone instead of returning garbage +--- stream_config +# intentionally no apisix_stream_metrics_zone here +--- stream_server_config + proxy_pass 127.0.0.1:1994; + log_by_lua_block { + local metrics = require("resty.apisix.stream.metrics") + local res, err = metrics.dump() + ngx.log(ngx.WARN, "dump: ", res, " err: ", err) + } +--- stream_request eval +"GET / HTTP/1.0\r\nHost: localhost\r\n\r\n" +--- stream_response_like: hello +--- error_log +dump: nil err: stream metrics zone is not configured + + + +=== TEST 6: $stream_session_reason still works without a metrics zone +--- stream_config +# intentionally no apisix_stream_metrics_zone here +--- stream_server_config + proxy_pass 127.0.0.1:1994; + log_by_lua_block { + ngx.log(ngx.WARN, "reason: ", ngx.var.stream_session_reason) + } +--- stream_request eval +"GET / HTTP/1.0\r\nHost: localhost\r\n\r\n" +--- stream_response_like: hello +--- error_log +reason: closed + + + +=== TEST 7: counters move while the session is still open +This is the point of keeping the counters in nginx: a long-lived connection +must not stay invisible until it ends. +--- stream_config +apisix_stream_metrics_zone 1m; +--- stream_server_config + proxy_pass 127.0.0.1:1994; +--- config + location /probe { + content_by_lua_block { + local metrics = require("resty.apisix.stream.metrics") + + local sock = ngx.socket.tcp() + local ok, err = sock:connect("127.0.0.1", $TEST_NGINX_SERVER_PORT + 1) + if not ok then + ngx.say("connect: ", err) + return + end + + -- a partial request, so the upstream keeps the session open + sock:send("GET / HTTP/1.0\r\n") + ngx.sleep(0.3) + + local function seen() + for _, e in ipairs(metrics.dump()) do + if e.downstream_ingress > 0 then + return e + end + end + end + + local first = seen() + ngx.say("first active=", first.active, + " di=", first.downstream_ingress, + " ue=", first.upstream_egress) + + -- a second burst on the same still-open session: this is what a + -- deferred flush would hide + sock:send("Host: localhost\r\n") + ngx.sleep(0.3) + + local second = seen() + ngx.say("second di=", second.downstream_ingress, + " ue=", second.upstream_egress) + + sock:close() + } + } +--- request +GET /probe +--- response_body +first active=1 di=16 ue=16 +second di=33 ue=33 + + + +=== TEST 8: UDP sessions are accounted for and report a reason +--- stream_config +apisix_stream_metrics_zone 1m; + +server { + listen 127.0.0.1:1988 udp; + return "pong"; +} + +server { + listen 127.0.0.1:1987 udp; + proxy_responses 1; + proxy_pass 127.0.0.1:1988; + log_by_lua_block { + ngx.log(ngx.WARN, "udp reason: ", ngx.var.stream_session_reason, + ", listen: ", ngx.var.stream_listen_addr) + } +} +--- stream_server_config + proxy_pass 127.0.0.1:1994; +--- config + location /probe { + content_by_lua_block { + local metrics = require("resty.apisix.stream.metrics") + + local sock = ngx.socket.udp() + sock:setpeername("127.0.0.1", 1987) + sock:send("ping") + local data = sock:receive() + sock:close() + + ngx.sleep(0.2) + + for _, e in ipairs(metrics.dump()) do + if e.listen_addr == "127.0.0.1:1987" then + ngx.say("udp ", e.listen_addr, + " di=", e.downstream_ingress, + " de=", e.downstream_egress, + " ue=", e.upstream_egress, + " ui=", e.upstream_ingress) + end + end + } + } +--- request +GET /probe +--- response_body +udp 127.0.0.1:1987 di=4 de=4 ue=4 ui=4 +--- error_log +udp reason: closed, listen: 127.0.0.1:1987 + + + +=== TEST 9: the internal unix socket of a stream block is not accounted for +--- stream_config +apisix_stream_metrics_zone 1m; + +server { + listen unix:$TEST_NGINX_HTML_DIR/stream_internal.sock; + proxy_pass 127.0.0.1:1994; +} +--- stream_server_config + proxy_pass 127.0.0.1:1994; +--- config + location /probe { + content_by_lua_block { + local metrics = require("resty.apisix.stream.metrics") + + local sock = ngx.socket.tcp() + local ok, err = sock:connect("unix:$TEST_NGINX_HTML_DIR/stream_internal.sock") + if not ok then + ngx.say("connect: ", err) + return + end + sock:send("GET / HTTP/1.0\r\nHost: localhost\r\n\r\n") + sock:receive("*a") + sock:close() + + -- the tcp listener must be the only slot: proving the unix one was + -- skipped for being a unix socket, not for having a long address + local res = metrics.dump() + for _, e in ipairs(res) do + ngx.say("slot ", e.listen_addr) + end + ngx.say("slots=", #res) + } + } +--- request +GET /probe +--- response_body +slot 0.0.0.0:1985 +slots=1 + + + +=== TEST 10: a read failure that is not a reset is not reported as one +nginx raises read->error for every fatal recv(), so the errno is the only +thing separating a reset from a protocol failure. Plaintext sent to a TLS +listener is the latter. +--- stream_config +apisix_stream_metrics_zone 1m; +--- stream_server_config + listen 12346 ssl; + ssl_certificate ../../certs/mtls_server.crt; + ssl_certificate_key ../../certs/mtls_server.key; + proxy_pass 127.0.0.1:1994; + log_by_lua_block { + ngx.log(ngx.WARN, "tls reason: ", ngx.var.stream_session_reason) + } +--- config + location /probe { + content_by_lua_block { + local sock = ngx.socket.tcp() + local ok, err = sock:connect("127.0.0.1", 12346) + if not ok then + ngx.say("connect: ", err) + return + end + sock:send("not-tls-at-all\r\n\r\n") + sock:receive("*a") + sock:close() + ngx.sleep(0.2) + ngx.say("done") + } + } +--- request +GET /probe +--- response_body +done +--- error_log +tls reason: client_read_error