From 1a6041c4fed300ac5e8d3b400f8cde8a36fb72a4 Mon Sep 17 00:00:00 2001 From: AlinsRan Date: Mon, 3 Aug 2026 13:26:50 +0800 Subject: [PATCH 1/8] feat(stream): per-listen metrics zone and $stream_session_reason nginx keeps the stream $status at 200 for every failure that happens after the upstream connection is established, and it exposes no way to read the byte counters of a live session, so a TCP proxy cannot report connection outcomes or real-time throughput. Add ngx_stream_apisix_metrics_module: - apisix_stream_metrics_zone reserves a shared memory zone holding, per stream listening address, the active session count and the bytes moved in the four directions. Each session accumulates locally and merges into the zone at most once per second, plus a final flush, so long-lived connections keep the counters moving without an atomic operation per read or write. - $stream_session_reason reports why a session ended: normal close, client or upstream reset, connect/receive/send timeout, worker shutdown. - resty.apisix.stream.metrics exposes the zone to Lua over FFI. The ngx_stream_proxy_module patch only records what nginx destroys on its own: it clears ev->timedout before finalizing, and it frees u->peer.connection before the log phase runs, which is also why the last byte delta towards the upstream has to be taken in the finalize hook. Everything else, including reset detection, is derived from the connection flags. --- README.md | 66 ++ lib/resty/apisix/stream/metrics.lua | 67 ++ patch/1.21.4/nginx-stream_metrics.patch | 69 ++ patch/1.25.3.1/nginx-stream_metrics.patch | 69 ++ patch/1.27.1.1/nginx-stream_metrics.patch | 69 ++ patch/1.29.2.4/nginx-stream_metrics.patch | 69 ++ src/stream/config | 10 +- src/stream/ngx_stream_apisix_metrics_module.c | 786 ++++++++++++++++++ src/stream/ngx_stream_apisix_metrics_module.h | 67 ++ src/stream/ngx_stream_apisix_module.h | 6 + t/stream/metrics.t | 137 +++ 11 files changed, 1414 insertions(+), 1 deletion(-) create mode 100644 lib/resty/apisix/stream/metrics.lua create mode 100644 patch/1.21.4/nginx-stream_metrics.patch create mode 100644 patch/1.25.3.1/nginx-stream_metrics.patch create mode 100644 patch/1.27.1.1/nginx-stream_metrics.patch create mode 100644 patch/1.29.2.4/nginx-stream_metrics.patch create mode 100644 src/stream/ngx_stream_apisix_metrics_module.c create mode 100644 src/stream/ngx_stream_apisix_metrics_module.h create mode 100644 t/stream/metrics.t diff --git a/README.md b/README.md index 5c850b9..5dfe474 100644 --- a/README.md +++ b/README.md @@ -14,6 +14,72 @@ 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). Counters are merged into the zone at +most once per second per session and flushed when the session ends, so they +keep moving during long-lived connections without paying an atomic operation +per read or write. Without this directive nothing is collected. + +example: + +```nginx +stream { + apisix_stream_metrics_zone 1m; +} +``` + +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 | +| `client_error` | sending to the client failed | +| `upstream_rst` | the upstream reset the connection | +| `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 | +| `-` | not applicable, for example a session rejected before `proxy_pass` | + +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..da0347b --- /dev/null +++ b/lib/resty/apisix/stream/metrics.lua @@ -0,0 +1,67 @@ +local ffi = require("ffi") +local base = require("resty.core.base") +local C = ffi.C +local ffi_str = ffi.string +local tonumber = tonumber + + +base.allows_subsystem("stream") + + +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; + +ngx_int_t +ngx_stream_apisix_metrics_dump(ngx_stream_apisix_metrics_entry_t *entries, size_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) + +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() + 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..f6efee4 --- /dev/null +++ b/patch/1.21.4/nginx-stream_metrics.patch @@ -0,0 +1,69 @@ +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; + } +@@ -1751,6 +1767,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; + } +@@ -1960,6 +1980,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..dc89dc0 --- /dev/null +++ b/patch/1.25.3.1/nginx-stream_metrics.patch @@ -0,0 +1,69 @@ +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; + } +@@ -1756,6 +1772,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; + } +@@ -1965,6 +1985,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..dc89dc0 --- /dev/null +++ b/patch/1.27.1.1/nginx-stream_metrics.patch @@ -0,0 +1,69 @@ +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; + } +@@ -1756,6 +1772,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; + } +@@ -1965,6 +1985,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..459fba4 --- /dev/null +++ b/patch/1.29.2.4/nginx-stream_metrics.patch @@ -0,0 +1,69 @@ +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; + } +@@ -1882,6 +1898,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; + } +@@ -2091,6 +2111,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..96ba0b8 --- /dev/null +++ b/src/stream/ngx_stream_apisix_metrics_module.c @@ -0,0 +1,786 @@ +#include +#include +#include +#include "ngx_stream_apisix_metrics_module.h" + + +/* + * Counters are merged into the shared zone at most once per second, so a + * long-lived session keeps the metrics moving without paying an atomic + * operation per read/write. The pending delta is flushed when the session + * ends, which keeps the totals exact. + */ +#define NGX_STREAM_APISIX_METRICS_FLUSH_INTERVAL 1 + +/* + * 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]; + time_t flush_time; + ngx_uint_t reason; + 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, + 0, + 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") +}; + + +/* + * Set once the zone is created, so that the FFI reader does not depend on a + * session or on the stream configuration being reachable. + */ +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 */ + ngx_stream_apisix_metrics_sh = osh; + shm_zone->data = osh; + return NGX_OK; + } + + shpool = (ngx_slab_pool_t *) shm_zone->shm.addr; + + if (shm_zone->shm.exists) { + ngx_stream_apisix_metrics_sh = shpool->data; + 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; + ngx_stream_apisix_metrics_sh = sh; + + return NGX_OK; +} + + +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_stream_apisix_metrics_sh_t *sh; + ngx_stream_apisix_metrics_slot_t *slot; + + sh = ngx_stream_apisix_metrics_sh; + if (sh == NULL) { + return NULL; + } + + for (i = 0; i < sh->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 || sh->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); + + sh->nused++; + + return slot; +} + + +/* + * Claiming happens in the master before the workers are forked, so the slot + * array needs no locking: it is append only and every reader afterwards only + * touches the per-slot atomics. + */ +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; + + 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; + } + + addr = ls[i].addr_text; + + if (addr.len == 0 || addr.len > 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; + + 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, ngx_uint_t force) +{ + 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; + } + + if (!force && ngx_time() - ctx->flush_time + < NGX_STREAM_APISIX_METRICS_FLUSH_INTERVAL) + { + 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]; + + if (delta <= 0) { + continue; + } + + (void) ngx_atomic_fetch_add(&slot->bytes[i], (ngx_atomic_int_t) delta); + ctx->flushed[i] = current[i]; + } + + ctx->flush_time = ngx_time(); +} + + +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, 0); +} + + +void +ngx_stream_apisix_set_session_reason(ngx_stream_session_t *s, + ngx_uint_t reason) +{ + ngx_stream_apisix_metrics_ctx_t *ctx; + + 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 NGX_STREAM_APISIX_REASON_CLIENT_RST; + } + + if (pc && pc->read->error) { + return NGX_STREAM_APISIX_REASON_UPSTREAM_RST; + } + + 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, 1); + + if (ctx->reason == NGX_STREAM_APISIX_REASON_UNSET) { + ctx->reason = ngx_stream_apisix_derive_reason(s, ctx, rc); + } +} + + +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, 1); + + 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; + } + + n = ngx_min(sh->nused, max); + + for (i = 0; i < n; i++) { + slot = &sh->slots[i]; + + entries[i].addr_len = slot->addr_len; + ngx_memcpy(entries[i].addr, slot->addr, slot->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..f2fce3a --- /dev/null +++ b/src/stream/ngx_stream_apisix_metrics_module.h @@ -0,0 +1,67 @@ +#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_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_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..9a1b254 --- /dev/null +++ b/t/stream/metrics.t @@ -0,0 +1,137 @@ +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) + } +--- stream_request eval +"GET /" +--- error_log +status: 502 + + + +=== 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 From 2b67d12af5713206c92c4c6815d9dcf2bee5b5b9 Mon Sep 17 00:00:00 2001 From: AlinsRan Date: Mon, 3 Aug 2026 14:52:38 +0800 Subject: [PATCH 2/8] fix(stream): survive a reload that drops the metrics zone, and account correctly Review follow-ups on the stream metrics module: - The zone pointer was cached when the zone was created and never cleared. A reload that removed the directive (or the whole stream block) left it pointing at shared memory ngx_init_cycle had already unmapped, so the next worker segfaulted on startup. Resolve it from the module main conf of the current cycle instead, which is NULL in exactly those cases. - proxy_next_upstream swaps the peer connection, so pc->sent restarts from zero while the flushed mark still held the previous peer's total. The decrease was read as 'nothing to do' and every retried byte was lost. Treat a decrease as a counter restart. - The interval clock was stamped even when nothing was written. The proxy module calls in once before any data has moved, so the first bytes of a session stayed invisible for a whole second. Only an actual write starts the interval. - A UDP session has no FIN to observe, so one that simply ran to completion reported no reason at all. Fall back to a normal close when the session ended on NGX_STREAM_OK, and record the reason on the UDP shutdown path too. - Unix sockets inside stream{} are internal plumbing, not proxy ports. Counting them reported gateway control traffic as proxied bytes and left a permanent floor of worker connections in the active gauge. - A listening address too long for a slot was dropped silently; warn instead. resty.apisix.stream.metrics no longer restricts itself to the stream subsystem: the zone is process global and the reader touches no session, so http can read it too. Tests now cover what the requirements actually asked for: counters moving while a session is still open (with exact byte counts), UDP accounting and its reason, and the exclusion of internal unix sockets. --- README.md | 20 +++ lib/resty/apisix/stream/metrics.lua | 8 +- patch/1.21.4/nginx-stream_metrics.patch | 16 ++- patch/1.25.3.1/nginx-stream_metrics.patch | 16 ++- patch/1.27.1.1/nginx-stream_metrics.patch | 16 ++- patch/1.29.2.4/nginx-stream_metrics.patch | 16 ++- src/stream/ngx_stream_apisix_metrics_module.c | 93 ++++++++++-- t/stream/metrics.t | 135 ++++++++++++++++++ 8 files changed, 299 insertions(+), 21 deletions(-) diff --git a/README.md b/README.md index 5dfe474..a50fc28 100644 --- a/README.md +++ b/README.md @@ -35,6 +35,26 @@ stream { } ``` +Only TCP and UDP listening addresses are accounted for. Unix sockets inside +`stream{}` are skipped: they are internal plumbing (APISIX puts its worker +event channel there) rather than proxy ports. + +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, or when the configured zone size + changes, because nginx only reuses a zone whose size is unchanged. +- A counter is written at most once per second per session, and again when the + session ends, so the totals are exact but can lag by up to a second. +- `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. + They therefore agree with `$upstream_bytes_sent` / `$upstream_bytes_received` + for a session that reached its upstream on the first try, and are more + accurate than those variables when it did not. + Read the counters from Lua with `resty.apisix.stream.metrics`: ```lua diff --git a/lib/resty/apisix/stream/metrics.lua b/lib/resty/apisix/stream/metrics.lua index da0347b..e0e9bb4 100644 --- a/lib/resty/apisix/stream/metrics.lua +++ b/lib/resty/apisix/stream/metrics.lua @@ -1,11 +1,11 @@ local ffi = require("ffi") -local base = require("resty.core.base") local C = ffi.C local ffi_str = ffi.string local tonumber = tonumber -base.allows_subsystem("stream") +-- 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([[ @@ -18,8 +18,10 @@ typedef struct { 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, size_t max); +ngx_stream_apisix_metrics_dump(ngx_stream_apisix_metrics_entry_t *entries, ngx_uint_t max); ]]) diff --git a/patch/1.21.4/nginx-stream_metrics.patch b/patch/1.21.4/nginx-stream_metrics.patch index f6efee4..f22ab56 100644 --- a/patch/1.21.4/nginx-stream_metrics.patch +++ b/patch/1.21.4/nginx-stream_metrics.patch @@ -45,7 +45,19 @@ diff --git src/stream/ngx_stream_proxy_module.c src/stream/ngx_stream_proxy_modu ngx_stream_proxy_next_upstream(s); return; } -@@ -1751,6 +1767,10 @@ +@@ -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; + } +@@ -1751,6 +1772,10 @@ c->log->action = "proxying connection"; @@ -56,7 +68,7 @@ diff --git src/stream/ngx_stream_proxy_module.c src/stream/ngx_stream_proxy_modu if (ngx_stream_proxy_test_finalize(s, from_upstream) == NGX_OK) { return; } -@@ -1960,6 +1980,10 @@ +@@ -1960,6 +1985,10 @@ ngx_log_debug1(NGX_LOG_DEBUG_STREAM, s->connection->log, 0, "finalize stream proxy: %i", rc); diff --git a/patch/1.25.3.1/nginx-stream_metrics.patch b/patch/1.25.3.1/nginx-stream_metrics.patch index dc89dc0..314cc45 100644 --- a/patch/1.25.3.1/nginx-stream_metrics.patch +++ b/patch/1.25.3.1/nginx-stream_metrics.patch @@ -45,7 +45,19 @@ diff --git src/stream/ngx_stream_proxy_module.c src/stream/ngx_stream_proxy_modu ngx_stream_proxy_next_upstream(s); return; } -@@ -1756,6 +1772,10 @@ +@@ -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; + } +@@ -1756,6 +1777,10 @@ c->log->action = "proxying connection"; @@ -56,7 +68,7 @@ diff --git src/stream/ngx_stream_proxy_module.c src/stream/ngx_stream_proxy_modu if (ngx_stream_proxy_test_finalize(s, from_upstream) == NGX_OK) { return; } -@@ -1965,6 +1985,10 @@ +@@ -1965,6 +1990,10 @@ ngx_log_debug1(NGX_LOG_DEBUG_STREAM, s->connection->log, 0, "finalize stream proxy: %i", rc); diff --git a/patch/1.27.1.1/nginx-stream_metrics.patch b/patch/1.27.1.1/nginx-stream_metrics.patch index dc89dc0..314cc45 100644 --- a/patch/1.27.1.1/nginx-stream_metrics.patch +++ b/patch/1.27.1.1/nginx-stream_metrics.patch @@ -45,7 +45,19 @@ diff --git src/stream/ngx_stream_proxy_module.c src/stream/ngx_stream_proxy_modu ngx_stream_proxy_next_upstream(s); return; } -@@ -1756,6 +1772,10 @@ +@@ -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; + } +@@ -1756,6 +1777,10 @@ c->log->action = "proxying connection"; @@ -56,7 +68,7 @@ diff --git src/stream/ngx_stream_proxy_module.c src/stream/ngx_stream_proxy_modu if (ngx_stream_proxy_test_finalize(s, from_upstream) == NGX_OK) { return; } -@@ -1965,6 +1985,10 @@ +@@ -1965,6 +1990,10 @@ ngx_log_debug1(NGX_LOG_DEBUG_STREAM, s->connection->log, 0, "finalize stream proxy: %i", rc); diff --git a/patch/1.29.2.4/nginx-stream_metrics.patch b/patch/1.29.2.4/nginx-stream_metrics.patch index 459fba4..8e2582e 100644 --- a/patch/1.29.2.4/nginx-stream_metrics.patch +++ b/patch/1.29.2.4/nginx-stream_metrics.patch @@ -45,7 +45,19 @@ diff --git src/stream/ngx_stream_proxy_module.c src/stream/ngx_stream_proxy_modu ngx_stream_proxy_next_upstream(s); return; } -@@ -1882,6 +1898,10 @@ +@@ -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; + } +@@ -1882,6 +1903,10 @@ c->log->action = "proxying connection"; @@ -56,7 +68,7 @@ diff --git src/stream/ngx_stream_proxy_module.c src/stream/ngx_stream_proxy_modu if (ngx_stream_proxy_test_finalize(s, from_upstream) == NGX_OK) { return; } -@@ -2091,6 +2111,10 @@ +@@ -2091,6 +2116,10 @@ ngx_log_debug1(NGX_LOG_DEBUG_STREAM, s->connection->log, 0, "finalize stream proxy: %i", rc); diff --git a/src/stream/ngx_stream_apisix_metrics_module.c b/src/stream/ngx_stream_apisix_metrics_module.c index 96ba0b8..ea3c8ca 100644 --- a/src/stream/ngx_stream_apisix_metrics_module.c +++ b/src/stream/ngx_stream_apisix_metrics_module.c @@ -74,7 +74,7 @@ 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, - 0, + NGX_STREAM_MAIN_CONF_OFFSET, 0, NULL }, @@ -144,8 +144,9 @@ static ngx_str_t ngx_stream_apisix_reasons[] = { /* - * Set once the zone is created, so that the FFI reader does not depend on a - * session or on the stream configuration being reachable. + * 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; @@ -218,7 +219,6 @@ ngx_stream_apisix_metrics_init_zone(ngx_shm_zone_t *shm_zone, void *data) if (osh) { /* reused on reload, the accumulated counters must survive */ - ngx_stream_apisix_metrics_sh = osh; shm_zone->data = osh; return NGX_OK; } @@ -226,7 +226,6 @@ ngx_stream_apisix_metrics_init_zone(ngx_shm_zone_t *shm_zone, void *data) shpool = (ngx_slab_pool_t *) shm_zone->shm.addr; if (shm_zone->shm.exists) { - ngx_stream_apisix_metrics_sh = shpool->data; shm_zone->data = shpool->data; return NGX_OK; } @@ -259,12 +258,33 @@ ngx_stream_apisix_metrics_init_zone(ngx_shm_zone_t *shm_zone, void *data) shpool->data = sh; shm_zone->data = sh; - ngx_stream_apisix_metrics_sh = 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) { @@ -314,6 +334,8 @@ ngx_stream_apisix_metrics_init_module(ngx_cycle_t *cycle) 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; } @@ -325,9 +347,27 @@ ngx_stream_apisix_metrics_init_module(ngx_cycle_t *cycle) 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 (ls[i].sockaddr->sa_family == AF_UNIX) { + continue; + } + addr = ls[i].addr_text; - if (addr.len == 0 || addr.len > NGX_STREAM_APISIX_METRICS_ADDR_LEN) { + 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; } @@ -352,6 +392,8 @@ ngx_stream_apisix_metrics_init_process(ngx_cycle_t *cycle) 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; } @@ -445,7 +487,7 @@ ngx_stream_apisix_metrics_flush(ngx_stream_session_t *s, { off_t current[NGX_STREAM_APISIX_METRICS_DIRECTIONS]; off_t delta; - ngx_uint_t i; + ngx_uint_t i, flushed; ngx_connection_t *pc; ngx_stream_upstream_t *u; ngx_stream_apisix_metrics_slot_t *slot; @@ -469,18 +511,38 @@ ngx_stream_apisix_metrics_flush(ngx_stream_session_t *s, current[NGX_STREAM_APISIX_METRICS_UPSTREAM_EGRESS] = pc ? pc->sent : 0; current[NGX_STREAM_APISIX_METRICS_UPSTREAM_INGRESS] = u ? u->received : 0; + flushed = 0; + for (i = 0; i < NGX_STREAM_APISIX_METRICS_DIRECTIONS; i++) { delta = current[i] - ctx->flushed[i]; - if (delta <= 0) { + /* + * 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]; + flushed = 1; } - ctx->flush_time = ngx_time(); + /* + * Only an actual write starts the interval. The proxy module calls in + * once before any data has moved, and stamping the clock there would + * hide the first bytes of the session for a whole second. + */ + if (flushed) { + ctx->flush_time = ngx_time(); + } } @@ -627,6 +689,17 @@ ngx_stream_apisix_metrics_finalize(ngx_stream_session_t *s, ngx_uint_t rc) 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; + } } diff --git a/t/stream/metrics.t b/t/stream/metrics.t index 9a1b254..a33f3a1 100644 --- a/t/stream/metrics.t +++ b/t/stream/metrics.t @@ -135,3 +135,138 @@ dump: nil err: stream metrics zone is not configured --- 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") + + -- outlive one flush interval without closing anything + ngx.sleep(1.5) + + for _, e in ipairs(metrics.dump()) do + if e.downstream_ingress > 0 then + ngx.say("live active=", e.active, + " di=", e.downstream_ingress, + " ue=", e.upstream_egress) + end + end + + sock:close() + } + } +--- request +GET /probe +--- response_body +live active=1 di=16 ue=16 + + + +=== 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; + return "internal"; +} +--- 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:receive("*a") + sock:close() + + for _, e in ipairs(metrics.dump()) do + if e.listen_addr:find("unix:", 1, true) then + ngx.say("leaked ", e.listen_addr) + end + end + ngx.say("done") + } + } +--- request +GET /probe +--- response_body +done From e3f2479acd4eac0fa424a6db79ee088345ab37bf Mon Sep 17 00:00:00 2001 From: AlinsRan Date: Mon, 3 Aug 2026 15:33:12 +0800 Subject: [PATCH 3/8] fix(stream): bank a failed peer's bytes, publish slots safely, correct the README Second review round: - proxy_next_upstream drops the peer and pc->sent with it. Treating a decrease as a restart only recovered the bytes if the new peer's counter had not yet overtaken the old mark; otherwise the failed attempt was still lost. Flush and reset the upstream mark before the peer goes away instead, which makes the total exact rather than nearly exact. The loss was small in practice -- every next_upstream call site runs before u->connected, so only PROXY protocol and TLS handshake bytes were at stake -- but the requirement is exactness. - Slot publication had no write barrier. Claiming is master-only, but a reload reuses the zone while the previous generation of workers is still reading, so on a weakly ordered architecture a reader could see the incremented count before the address it exposes. The old comment claimed fork made this safe, which is only true for a cold start. - Guard the AF_UNIX check with NGX_HAVE_UNIX_DOMAIN, as nginx does everywhere else, and clamp the address length the FFI reader copies. - resty.apisix.stream.metrics threw on a build without the stream addon rather than returning the documented nil, err. The README promised totals that were 'exact' and upstream counters 'more accurate' than nginx's own variables; neither held on the retry path. It also omitted that removing the directive releases the zone, and that slots are never reclaimed. TEST 9 asserted only that no unix slot exists, which a long enough servroot path would satisfy by tripping the address length check instead. It now asserts the exact slot set. --- README.md | 20 ++++++--- lib/resty/apisix/stream/metrics.lua | 10 +++++ patch/1.21.4/nginx-stream_metrics.patch | 13 +++++- patch/1.25.3.1/nginx-stream_metrics.patch | 13 +++++- patch/1.27.1.1/nginx-stream_metrics.patch | 13 +++++- patch/1.29.2.4/nginx-stream_metrics.patch | 13 +++++- src/stream/ngx_stream_apisix_metrics_module.c | 43 ++++++++++++++++--- src/stream/ngx_stream_apisix_metrics_module.h | 1 + t/stream/metrics.t | 17 +++++--- 9 files changed, 121 insertions(+), 22 deletions(-) diff --git a/README.md b/README.md index a50fc28..adc7614 100644 --- a/README.md +++ b/README.md @@ -42,18 +42,22 @@ event channel there) rather than proxy ports. 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, or when the configured zone size - changes, because nginx only reuses a zone whose size is unchanged. + 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 at most once per second per session, and again when the session ends, so the totals are exact but can lag by up to a second. - `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. - They therefore agree with `$upstream_bytes_sent` / `$upstream_bytes_received` - for a session that reached its upstream on the first try, and are more - accurate than those variables when it did not. +- With `proxy_next_upstream`, the upstream byte counters sum every attempt, + including what was sent to a peer that then failed. They agree with + `$upstream_bytes_sent` / `$upstream_bytes_received` for a session that + reached its upstream on the first try. Read the counters from Lua with `resty.apisix.stream.metrics`: @@ -89,6 +93,10 @@ tells a graceful close apart from a timeout or a reset: | `shutdown` | the worker was shutting down | | `-` | not applicable, for example a session rejected before `proxy_pass` | +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 diff --git a/lib/resty/apisix/stream/metrics.lua b/lib/resty/apisix/stream/metrics.lua index e0e9bb4..0029780 100644 --- a/lib/resty/apisix/stream/metrics.lua +++ b/lib/resty/apisix/stream/metrics.lua @@ -34,6 +34,12 @@ 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 = {} @@ -43,6 +49,10 @@ local _M = {} -- 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 diff --git a/patch/1.21.4/nginx-stream_metrics.patch b/patch/1.21.4/nginx-stream_metrics.patch index f22ab56..6795db6 100644 --- a/patch/1.21.4/nginx-stream_metrics.patch +++ b/patch/1.21.4/nginx-stream_metrics.patch @@ -68,7 +68,18 @@ diff --git src/stream/ngx_stream_proxy_module.c src/stream/ngx_stream_proxy_modu if (ngx_stream_proxy_test_finalize(s, from_upstream) == NGX_OK) { return; } -@@ -1960,6 +1985,10 @@ +@@ -1930,6 +1955,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 +1989,10 @@ ngx_log_debug1(NGX_LOG_DEBUG_STREAM, s->connection->log, 0, "finalize stream proxy: %i", rc); diff --git a/patch/1.25.3.1/nginx-stream_metrics.patch b/patch/1.25.3.1/nginx-stream_metrics.patch index 314cc45..9d3539b 100644 --- a/patch/1.25.3.1/nginx-stream_metrics.patch +++ b/patch/1.25.3.1/nginx-stream_metrics.patch @@ -68,7 +68,18 @@ diff --git src/stream/ngx_stream_proxy_module.c src/stream/ngx_stream_proxy_modu if (ngx_stream_proxy_test_finalize(s, from_upstream) == NGX_OK) { return; } -@@ -1965,6 +1990,10 @@ +@@ -1935,6 +1960,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 +1994,10 @@ ngx_log_debug1(NGX_LOG_DEBUG_STREAM, s->connection->log, 0, "finalize stream proxy: %i", rc); diff --git a/patch/1.27.1.1/nginx-stream_metrics.patch b/patch/1.27.1.1/nginx-stream_metrics.patch index 314cc45..9d3539b 100644 --- a/patch/1.27.1.1/nginx-stream_metrics.patch +++ b/patch/1.27.1.1/nginx-stream_metrics.patch @@ -68,7 +68,18 @@ diff --git src/stream/ngx_stream_proxy_module.c src/stream/ngx_stream_proxy_modu if (ngx_stream_proxy_test_finalize(s, from_upstream) == NGX_OK) { return; } -@@ -1965,6 +1990,10 @@ +@@ -1935,6 +1960,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 +1994,10 @@ ngx_log_debug1(NGX_LOG_DEBUG_STREAM, s->connection->log, 0, "finalize stream proxy: %i", rc); diff --git a/patch/1.29.2.4/nginx-stream_metrics.patch b/patch/1.29.2.4/nginx-stream_metrics.patch index 8e2582e..6be3df4 100644 --- a/patch/1.29.2.4/nginx-stream_metrics.patch +++ b/patch/1.29.2.4/nginx-stream_metrics.patch @@ -68,7 +68,18 @@ diff --git src/stream/ngx_stream_proxy_module.c src/stream/ngx_stream_proxy_modu if (ngx_stream_proxy_test_finalize(s, from_upstream) == NGX_OK) { return; } -@@ -2091,6 +2116,10 @@ +@@ -2061,6 +2086,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 +2120,10 @@ ngx_log_debug1(NGX_LOG_DEBUG_STREAM, s->connection->log, 0, "finalize stream proxy: %i", rc); diff --git a/src/stream/ngx_stream_apisix_metrics_module.c b/src/stream/ngx_stream_apisix_metrics_module.c index ea3c8ca..2ba7aaf 100644 --- a/src/stream/ngx_stream_apisix_metrics_module.c +++ b/src/stream/ngx_stream_apisix_metrics_module.c @@ -316,6 +316,13 @@ ngx_stream_apisix_metrics_lookup(ngx_str_t *addr, ngx_uint_t create) 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; @@ -323,9 +330,10 @@ ngx_stream_apisix_metrics_lookup(ngx_str_t *addr, ngx_uint_t create) /* - * Claiming happens in the master before the workers are forked, so the slot - * array needs no locking: it is append only and every reader afterwards only - * touches the per-slot atomics. + * 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) @@ -353,9 +361,11 @@ ngx_stream_apisix_metrics_init_module(ngx_cycle_t *cycle) * 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; @@ -560,6 +570,28 @@ ngx_stream_apisix_metrics_update(ngx_stream_session_t *s) } +/* + * 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, 1); + + ctx->flushed[NGX_STREAM_APISIX_METRICS_UPSTREAM_EGRESS] = 0; +} + + void ngx_stream_apisix_set_session_reason(ngx_stream_session_t *s, ngx_uint_t reason) @@ -800,8 +832,9 @@ ngx_stream_apisix_metrics_dump(ngx_stream_apisix_metrics_entry_t *entries, for (i = 0; i < n; i++) { slot = &sh->slots[i]; - entries[i].addr_len = slot->addr_len; - ngx_memcpy(entries[i].addr, slot->addr, slot->addr_len); + 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; diff --git a/src/stream/ngx_stream_apisix_metrics_module.h b/src/stream/ngx_stream_apisix_metrics_module.h index f2fce3a..e7e4776 100644 --- a/src/stream/ngx_stream_apisix_metrics_module.h +++ b/src/stream/ngx_stream_apisix_metrics_module.h @@ -56,6 +56,7 @@ void ngx_stream_apisix_set_session_reason(ngx_stream_session_t *s, 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_metrics_finalize(ngx_stream_session_t *s, ngx_uint_t rc); diff --git a/t/stream/metrics.t b/t/stream/metrics.t index a33f3a1..29d8f0c 100644 --- a/t/stream/metrics.t +++ b/t/stream/metrics.t @@ -240,7 +240,7 @@ apisix_stream_metrics_zone 1m; server { listen unix:$TEST_NGINX_HTML_DIR/stream_internal.sock; - return "internal"; + proxy_pass 127.0.0.1:1994; } --- stream_server_config proxy_pass 127.0.0.1:1994; @@ -255,18 +255,21 @@ server { ngx.say("connect: ", err) return end + sock:send("GET / HTTP/1.0\r\nHost: localhost\r\n\r\n") sock:receive("*a") sock:close() - for _, e in ipairs(metrics.dump()) do - if e.listen_addr:find("unix:", 1, true) then - ngx.say("leaked ", e.listen_addr) - end + -- 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("done") + ngx.say("slots=", #res) } } --- request GET /probe --- response_body -done +slot 0.0.0.0:1985 +slots=1 From 18a270ac77ccf2d81aa4e3d338ab14dbe0da380c Mon Sep 17 00:00:00 2001 From: AlinsRan Date: Tue, 4 Aug 2026 08:59:10 +0800 Subject: [PATCH 4/8] fix(stream): bound the reason index, and be precise about scope in the README PR review comments: - ngx_stream_apisix_set_session_reason indexes the reason table with a value that arrives from the ngx_stream_proxy_module patches, which are separate files free to drift from this enum. Refuse an out-of-range value instead of reading past the table. - The unix socket note read as though only internal sockets are skipped. The filter is on the address family, so a unix listener configured for proxying is skipped too; say so. - $upstream_bytes_sent and $upstream_bytes_received report one value per attempt rather than a total, so the retry note now says they have to be summed before they can be compared with these counters. --- README.md | 16 ++++++++++------ src/stream/ngx_stream_apisix_metrics_module.c | 12 ++++++++++++ 2 files changed, 22 insertions(+), 6 deletions(-) diff --git a/README.md b/README.md index adc7614..b338819 100644 --- a/README.md +++ b/README.md @@ -35,9 +35,11 @@ stream { } ``` -Only TCP and UDP listening addresses are accounted for. Unix sockets inside -`stream{}` are skipped: they are internal plumbing (APISIX puts its worker -event channel there) rather than proxy ports. +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: @@ -55,9 +57,11 @@ Things worth knowing before building alerts on this: 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. They agree with - `$upstream_bytes_sent` / `$upstream_bytes_received` for a session that - reached its upstream on the first try. + 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`: diff --git a/src/stream/ngx_stream_apisix_metrics_module.c b/src/stream/ngx_stream_apisix_metrics_module.c index 2ba7aaf..1f414c2 100644 --- a/src/stream/ngx_stream_apisix_metrics_module.c +++ b/src/stream/ngx_stream_apisix_metrics_module.c @@ -142,6 +142,9 @@ static ngx_str_t ngx_stream_apisix_reasons[] = { ngx_string("shutdown") }; +#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 @@ -598,6 +601,15 @@ ngx_stream_apisix_set_session_reason(ngx_stream_session_t *s, { 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; From 96ee0219a21377853a810cba9d5997747c662dfe Mon Sep 17 00:00:00 2001 From: AlinsRan Date: Tue, 4 Aug 2026 09:00:55 +0800 Subject: [PATCH 5/8] fix(stream): pair the slot publication barrier with one on the read side A store barrier alone only orders the writer. A reader that loads nused and then the slot it exposes still lets a weakly ordered CPU hoist the slot load above the count load, which is the very reordering the writer barrier exists to prevent. Snapshot nused, barrier, then walk the slots, in both the lookup and the FFI dump. --- src/stream/ngx_stream_apisix_metrics_module.c | 18 +++++++++++++++--- 1 file changed, 15 insertions(+), 3 deletions(-) diff --git a/src/stream/ngx_stream_apisix_metrics_module.c b/src/stream/ngx_stream_apisix_metrics_module.c index 1f414c2..39178b8 100644 --- a/src/stream/ngx_stream_apisix_metrics_module.c +++ b/src/stream/ngx_stream_apisix_metrics_module.c @@ -292,6 +292,7 @@ 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; @@ -300,7 +301,14 @@ ngx_stream_apisix_metrics_lookup(ngx_str_t *addr, ngx_uint_t create) return NULL; } - for (i = 0; i < sh->nused; i++) { + /* + * 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 @@ -310,7 +318,7 @@ ngx_stream_apisix_metrics_lookup(ngx_str_t *addr, ngx_uint_t create) } } - if (!create || sh->nused == sh->nslots) { + if (!create || nused == sh->nslots) { return NULL; } @@ -839,7 +847,11 @@ ngx_stream_apisix_metrics_dump(ngx_stream_apisix_metrics_entry_t *entries, return NGX_ERROR; } - n = ngx_min(sh->nused, max); + /* 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]; From e8238613e287fc36b3959e4cd3a65bf2294cb4bf Mon Sep 17 00:00:00 2001 From: AlinsRan Date: Tue, 4 Aug 2026 09:21:32 +0800 Subject: [PATCH 6/8] fix(stream): flush every forwarding pass, and name the connect failure Final review round. The once-a-second throttle did not bound staleness, it removed it. A skipped delta has no later trigger: nothing calls back into the forwarding path until more data arrives or the session ends. For request/response traffic on a long lived connection -- the common L4 shape -- a whole response could stay invisible for as long as the connection then stayed idle, which is exactly what keeping the counters in nginx was supposed to avoid, and the README claim of 'up to a second' was simply wrong. The hook already sits where nginx has drained everything readable, so the session-local accumulation still collapses the inner read/write loop into one atomic per direction. The throttle was buying almost nothing for that. Giving up before any peer answers -- connection refused, no live upstream, a failed or timed out upstream handshake -- left no mark on the connection flags and so reported no reason at all. It now reports connect_failed, and the README no longer describes '-' as only meaning a pre-proxy rejection. TEST 7 could not have caught the throttle bug: the flush mark starts at zero, so the first burst always published. It now sends a second burst on the still open session, which is precisely what a deferred flush would hide. TEST 3 gained the reason assertion for the path it was already walking. --- README.md | 19 ++++--- src/stream/ngx_stream_apisix_metrics_module.c | 56 +++++++++---------- src/stream/ngx_stream_apisix_metrics_module.h | 3 +- t/stream/metrics.t | 36 ++++++++---- 4 files changed, 64 insertions(+), 50 deletions(-) diff --git a/README.md b/README.md index b338819..3da2737 100644 --- a/README.md +++ b/README.md @@ -22,10 +22,10 @@ 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). Counters are merged into the zone at -most once per second per session and flushed when the session ends, so they -keep moving during long-lived connections without paying an atomic operation -per read or write. Without this directive nothing is collected. +(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: @@ -50,8 +50,12 @@ Things worth knowing before building alerts on this: 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 at most once per second per session, and again when the - session ends, so the totals are exact but can lag by up to a second. +- 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 @@ -95,7 +99,8 @@ tells a graceful close apart from a timeout or a reset: | `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 | -| `-` | not applicable, for example a session rejected before `proxy_pass` | +| `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 -- diff --git a/src/stream/ngx_stream_apisix_metrics_module.c b/src/stream/ngx_stream_apisix_metrics_module.c index 39178b8..9afdf25 100644 --- a/src/stream/ngx_stream_apisix_metrics_module.c +++ b/src/stream/ngx_stream_apisix_metrics_module.c @@ -5,12 +5,13 @@ /* - * Counters are merged into the shared zone at most once per second, so a - * long-lived session keeps the metrics moving without paying an atomic - * operation per read/write. The pending delta is flushed when the session - * ends, which keeps the totals exact. + * 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. */ -#define NGX_STREAM_APISIX_METRICS_FLUSH_INTERVAL 1 /* * One slot per stream listening address. A deployment has a handful of them, @@ -42,7 +43,6 @@ typedef struct { typedef struct { ngx_stream_apisix_metrics_slot_t *slot; off_t flushed[NGX_STREAM_APISIX_METRICS_DIRECTIONS]; - time_t flush_time; ngx_uint_t reason; unsigned counted:1; unsigned connect_timeout:1; @@ -139,7 +139,8 @@ static ngx_str_t ngx_stream_apisix_reasons[] = { ngx_string("recv_timeout"), ngx_string("send_timeout"), ngx_string("upstream_timeout"), - ngx_string("shutdown") + ngx_string("shutdown"), + ngx_string("connect_failed") }; #define NGX_STREAM_APISIX_REASONS_N \ @@ -504,11 +505,11 @@ ngx_stream_apisix_metrics_post_accept_handler(ngx_stream_session_t *s) static void ngx_stream_apisix_metrics_flush(ngx_stream_session_t *s, - ngx_stream_apisix_metrics_ctx_t *ctx, ngx_uint_t force) + ngx_stream_apisix_metrics_ctx_t *ctx) { off_t current[NGX_STREAM_APISIX_METRICS_DIRECTIONS]; off_t delta; - ngx_uint_t i, flushed; + ngx_uint_t i; ngx_connection_t *pc; ngx_stream_upstream_t *u; ngx_stream_apisix_metrics_slot_t *slot; @@ -518,12 +519,6 @@ ngx_stream_apisix_metrics_flush(ngx_stream_session_t *s, return; } - if (!force && ngx_time() - ctx->flush_time - < NGX_STREAM_APISIX_METRICS_FLUSH_INTERVAL) - { - return; - } - u = s->upstream; pc = (u && u->peer.connection) ? u->peer.connection : NULL; @@ -532,8 +527,6 @@ ngx_stream_apisix_metrics_flush(ngx_stream_session_t *s, current[NGX_STREAM_APISIX_METRICS_UPSTREAM_EGRESS] = pc ? pc->sent : 0; current[NGX_STREAM_APISIX_METRICS_UPSTREAM_INGRESS] = u ? u->received : 0; - flushed = 0; - for (i = 0; i < NGX_STREAM_APISIX_METRICS_DIRECTIONS; i++) { delta = current[i] - ctx->flushed[i]; @@ -553,16 +546,6 @@ ngx_stream_apisix_metrics_flush(ngx_stream_session_t *s, (void) ngx_atomic_fetch_add(&slot->bytes[i], (ngx_atomic_int_t) delta); ctx->flushed[i] = current[i]; - flushed = 1; - } - - /* - * Only an actual write starts the interval. The proxy module calls in - * once before any data has moved, and stamping the clock there would - * hide the first bytes of the session for a whole second. - */ - if (flushed) { - ctx->flush_time = ngx_time(); } } @@ -577,7 +560,7 @@ ngx_stream_apisix_metrics_update(ngx_stream_session_t *s) return; } - ngx_stream_apisix_metrics_flush(s, ctx, 0); + ngx_stream_apisix_metrics_flush(s, ctx); } @@ -597,7 +580,7 @@ ngx_stream_apisix_metrics_peer_closing(ngx_stream_session_t *s) return; } - ngx_stream_apisix_metrics_flush(s, ctx, 1); + ngx_stream_apisix_metrics_flush(s, ctx); ctx->flushed[NGX_STREAM_APISIX_METRICS_UPSTREAM_EGRESS] = 0; } @@ -736,7 +719,7 @@ ngx_stream_apisix_metrics_finalize(ngx_stream_session_t *s, ngx_uint_t rc) ctx->finalized = 1; - ngx_stream_apisix_metrics_flush(s, ctx, 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); @@ -752,6 +735,17 @@ ngx_stream_apisix_metrics_finalize(ngx_stream_session_t *s, ngx_uint_t rc) { 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; + } } @@ -823,7 +817,7 @@ ngx_stream_apisix_metrics_log_handler(ngx_stream_session_t *s) return NGX_OK; } - ngx_stream_apisix_metrics_flush(s, ctx, 1); + ngx_stream_apisix_metrics_flush(s, ctx); if (ctx->counted) { ctx->counted = 0; diff --git a/src/stream/ngx_stream_apisix_metrics_module.h b/src/stream/ngx_stream_apisix_metrics_module.h index e7e4776..c80dde6 100644 --- a/src/stream/ngx_stream_apisix_metrics_module.h +++ b/src/stream/ngx_stream_apisix_metrics_module.h @@ -37,7 +37,8 @@ typedef enum { 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_SHUTDOWN, + NGX_STREAM_APISIX_REASON_CONNECT_FAILED } ngx_stream_apisix_reason_e; diff --git a/t/stream/metrics.t b/t/stream/metrics.t index 29d8f0c..122444f 100644 --- a/t/stream/metrics.t +++ b/t/stream/metrics.t @@ -67,12 +67,13 @@ reason: recv_timeout, status: 200 proxy_connect_timeout 300ms; proxy_pass 127.0.0.1:1979; log_by_lua_block { - ngx.log(ngx.WARN, "status: ", ngx.var.status) + ngx.log(ngx.WARN, "status: ", ngx.var.status, + ", reason: ", ngx.var.stream_session_reason) } --- stream_request eval "GET /" --- error_log -status: 502 +status: 502, reason: connect_failed @@ -159,25 +160,38 @@ apisix_stream_metrics_zone 1m; -- a partial request, so the upstream keeps the session open sock:send("GET / HTTP/1.0\r\n") + ngx.sleep(0.3) - -- outlive one flush interval without closing anything - ngx.sleep(1.5) - - for _, e in ipairs(metrics.dump()) do - if e.downstream_ingress > 0 then - ngx.say("live active=", e.active, - " di=", e.downstream_ingress, - " ue=", e.upstream_egress) + 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.di or second.downstream_ingress, + " ue=", second.upstream_egress) + sock:close() } } --- request GET /probe --- response_body -live active=1 di=16 ue=16 +first active=1 di=16 ue=16 +second di=33 ue=33 From 685b5a633a2227ea26b46f2e1957e4681f507ab0 Mon Sep 17 00:00:00 2001 From: AlinsRan Date: Fri, 7 Aug 2026 15:02:44 +0800 Subject: [PATCH 7/8] fix(stream): only call it a reset when the errno says so Review follow-up from @membphis. nginx raises read->error for every fatal recv(), not only for ECONNRESET, and ngx_ssl_recv adds protocol failures on top of that. Deriving the reason from the flag alone therefore reported client_rst / upstream_rst for aborts and TLS read failures too. The errno is the only thing that separates them and it is gone by the time the session is finalized, so the proxy module now hands it over at the moment the read fails. ECONNRESET keeps client_rst / upstream_rst; everything else gets client_read_error / upstream_read_error. Verified both directions: a real SO_LINGER reset still reports client_rst and upstream_rst, and plaintext sent to a TLS listener -- previously mislabelled as a reset -- now reports client_read_error, which is the new regression test. --- README.md | 6 ++- patch/1.21.4/nginx-stream_metrics.patch | 17 ++++++-- patch/1.25.3.1/nginx-stream_metrics.patch | 17 ++++++-- patch/1.27.1.1/nginx-stream_metrics.patch | 17 ++++++-- patch/1.29.2.4/nginx-stream_metrics.patch | 17 ++++++-- src/stream/ngx_stream_apisix_metrics_module.c | 35 +++++++++++++++-- src/stream/ngx_stream_apisix_metrics_module.h | 6 ++- t/stream/metrics.t | 39 +++++++++++++++++++ 8 files changed, 136 insertions(+), 18 deletions(-) diff --git a/README.md b/README.md index 3da2737..5c09df2 100644 --- a/README.md +++ b/README.md @@ -90,9 +90,11 @@ 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 | +| `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 | +| `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 | diff --git a/patch/1.21.4/nginx-stream_metrics.patch b/patch/1.21.4/nginx-stream_metrics.patch index 6795db6..7e73e0c 100644 --- a/patch/1.21.4/nginx-stream_metrics.patch +++ b/patch/1.21.4/nginx-stream_metrics.patch @@ -57,7 +57,18 @@ diff --git src/stream/ngx_stream_proxy_module.c src/stream/ngx_stream_proxy_modu ngx_stream_proxy_finalize(s, NGX_STREAM_OK); return; } -@@ -1751,6 +1772,10 @@ +@@ -1697,6 +1718,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 +1776,10 @@ c->log->action = "proxying connection"; @@ -68,7 +79,7 @@ diff --git src/stream/ngx_stream_proxy_module.c src/stream/ngx_stream_proxy_modu if (ngx_stream_proxy_test_finalize(s, from_upstream) == NGX_OK) { return; } -@@ -1930,6 +1955,10 @@ +@@ -1930,6 +1959,10 @@ ngx_log_debug1(NGX_LOG_DEBUG_STREAM, s->connection->log, 0, "close proxy upstream connection: %d", pc->fd); @@ -79,7 +90,7 @@ diff --git src/stream/ngx_stream_proxy_module.c src/stream/ngx_stream_proxy_modu #if (NGX_STREAM_SSL) if (pc->ssl) { pc->ssl->no_wait_shutdown = 1; -@@ -1960,6 +1989,10 @@ +@@ -1960,6 +1993,10 @@ ngx_log_debug1(NGX_LOG_DEBUG_STREAM, s->connection->log, 0, "finalize stream proxy: %i", rc); diff --git a/patch/1.25.3.1/nginx-stream_metrics.patch b/patch/1.25.3.1/nginx-stream_metrics.patch index 9d3539b..835d13a 100644 --- a/patch/1.25.3.1/nginx-stream_metrics.patch +++ b/patch/1.25.3.1/nginx-stream_metrics.patch @@ -57,7 +57,18 @@ diff --git src/stream/ngx_stream_proxy_module.c src/stream/ngx_stream_proxy_modu ngx_stream_proxy_finalize(s, NGX_STREAM_OK); return; } -@@ -1756,6 +1777,10 @@ +@@ -1702,6 +1723,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 +1781,10 @@ c->log->action = "proxying connection"; @@ -68,7 +79,7 @@ diff --git src/stream/ngx_stream_proxy_module.c src/stream/ngx_stream_proxy_modu if (ngx_stream_proxy_test_finalize(s, from_upstream) == NGX_OK) { return; } -@@ -1935,6 +1960,10 @@ +@@ -1935,6 +1964,10 @@ ngx_log_debug1(NGX_LOG_DEBUG_STREAM, s->connection->log, 0, "close proxy upstream connection: %d", pc->fd); @@ -79,7 +90,7 @@ diff --git src/stream/ngx_stream_proxy_module.c src/stream/ngx_stream_proxy_modu #if (NGX_STREAM_SSL) if (pc->ssl) { pc->ssl->no_wait_shutdown = 1; -@@ -1965,6 +1994,10 @@ +@@ -1965,6 +1998,10 @@ ngx_log_debug1(NGX_LOG_DEBUG_STREAM, s->connection->log, 0, "finalize stream proxy: %i", rc); diff --git a/patch/1.27.1.1/nginx-stream_metrics.patch b/patch/1.27.1.1/nginx-stream_metrics.patch index 9d3539b..835d13a 100644 --- a/patch/1.27.1.1/nginx-stream_metrics.patch +++ b/patch/1.27.1.1/nginx-stream_metrics.patch @@ -57,7 +57,18 @@ diff --git src/stream/ngx_stream_proxy_module.c src/stream/ngx_stream_proxy_modu ngx_stream_proxy_finalize(s, NGX_STREAM_OK); return; } -@@ -1756,6 +1777,10 @@ +@@ -1702,6 +1723,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 +1781,10 @@ c->log->action = "proxying connection"; @@ -68,7 +79,7 @@ diff --git src/stream/ngx_stream_proxy_module.c src/stream/ngx_stream_proxy_modu if (ngx_stream_proxy_test_finalize(s, from_upstream) == NGX_OK) { return; } -@@ -1935,6 +1960,10 @@ +@@ -1935,6 +1964,10 @@ ngx_log_debug1(NGX_LOG_DEBUG_STREAM, s->connection->log, 0, "close proxy upstream connection: %d", pc->fd); @@ -79,7 +90,7 @@ diff --git src/stream/ngx_stream_proxy_module.c src/stream/ngx_stream_proxy_modu #if (NGX_STREAM_SSL) if (pc->ssl) { pc->ssl->no_wait_shutdown = 1; -@@ -1965,6 +1994,10 @@ +@@ -1965,6 +1998,10 @@ ngx_log_debug1(NGX_LOG_DEBUG_STREAM, s->connection->log, 0, "finalize stream proxy: %i", rc); diff --git a/patch/1.29.2.4/nginx-stream_metrics.patch b/patch/1.29.2.4/nginx-stream_metrics.patch index 6be3df4..c1b828f 100644 --- a/patch/1.29.2.4/nginx-stream_metrics.patch +++ b/patch/1.29.2.4/nginx-stream_metrics.patch @@ -57,7 +57,18 @@ diff --git src/stream/ngx_stream_proxy_module.c src/stream/ngx_stream_proxy_modu ngx_stream_proxy_finalize(s, NGX_STREAM_OK); return; } -@@ -1882,6 +1903,10 @@ +@@ -1828,6 +1849,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 +1907,10 @@ c->log->action = "proxying connection"; @@ -68,7 +79,7 @@ diff --git src/stream/ngx_stream_proxy_module.c src/stream/ngx_stream_proxy_modu if (ngx_stream_proxy_test_finalize(s, from_upstream) == NGX_OK) { return; } -@@ -2061,6 +2086,10 @@ +@@ -2061,6 +2090,10 @@ ngx_log_debug1(NGX_LOG_DEBUG_STREAM, s->connection->log, 0, "close proxy upstream connection: %d", pc->fd); @@ -79,7 +90,7 @@ diff --git src/stream/ngx_stream_proxy_module.c src/stream/ngx_stream_proxy_modu #if (NGX_STREAM_SSL) if (pc->ssl) { pc->ssl->no_wait_shutdown = 1; -@@ -2091,6 +2120,10 @@ +@@ -2091,6 +2124,10 @@ ngx_log_debug1(NGX_LOG_DEBUG_STREAM, s->connection->log, 0, "finalize stream proxy: %i", rc); diff --git a/src/stream/ngx_stream_apisix_metrics_module.c b/src/stream/ngx_stream_apisix_metrics_module.c index 9afdf25..6e3b24d 100644 --- a/src/stream/ngx_stream_apisix_metrics_module.c +++ b/src/stream/ngx_stream_apisix_metrics_module.c @@ -44,6 +44,8 @@ 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; @@ -140,7 +142,9 @@ static ngx_str_t ngx_stream_apisix_reasons[] = { ngx_string("send_timeout"), ngx_string("upstream_timeout"), ngx_string("shutdown"), - ngx_string("connect_failed") + ngx_string("connect_failed"), + ngx_string("client_read_error"), + ngx_string("upstream_read_error") }; #define NGX_STREAM_APISIX_REASONS_N \ @@ -586,6 +590,27 @@ ngx_stream_apisix_metrics_peer_closing(ngx_stream_session_t *s) } +/* + * 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) @@ -679,11 +704,15 @@ ngx_stream_apisix_derive_reason(ngx_stream_session_t *s, pc = (u && u->peer.connection) ? u->peer.connection : NULL; if (c->read->error) { - return NGX_STREAM_APISIX_REASON_CLIENT_RST; + 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 NGX_STREAM_APISIX_REASON_UPSTREAM_RST; + return ctx && ctx->read_err[1] == NGX_ECONNRESET + ? NGX_STREAM_APISIX_REASON_UPSTREAM_RST + : NGX_STREAM_APISIX_REASON_UPSTREAM_READ_ERROR; } if (c->error) { diff --git a/src/stream/ngx_stream_apisix_metrics_module.h b/src/stream/ngx_stream_apisix_metrics_module.h index c80dde6..301ca20 100644 --- a/src/stream/ngx_stream_apisix_metrics_module.h +++ b/src/stream/ngx_stream_apisix_metrics_module.h @@ -38,7 +38,9 @@ typedef enum { 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_CONNECT_FAILED, + NGX_STREAM_APISIX_REASON_CLIENT_READ_ERROR, + NGX_STREAM_APISIX_REASON_UPSTREAM_READ_ERROR } ngx_stream_apisix_reason_e; @@ -58,6 +60,8 @@ 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); diff --git a/t/stream/metrics.t b/t/stream/metrics.t index 122444f..c1fe244 100644 --- a/t/stream/metrics.t +++ b/t/stream/metrics.t @@ -287,3 +287,42 @@ 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 From b62feea5a74aace4fcd0e6e563441e853d3a1643 Mon Sep 17 00:00:00 2001 From: AlinsRan Date: Fri, 7 Aug 2026 15:29:05 +0800 Subject: [PATCH 8/8] fix(stream): do not let a stale errno be read as a reset ngx_ssl_handle_recv() normalises SSL_ERROR_SSL to no error internally but leaves errno untouched, so the value read after a failed TLS read could be a leftover ECONNRESET from an earlier syscall and get reported as a reset that never happened. Clear errno immediately before the recv instead, so whatever is readable afterwards belongs to that recv. Chosen over threading a normalised error out of the receive layer: that means patching ngx_event_openssl.c and ngx_recv.c, which the http subsystem shares, for the same guarantee a single store buys here. Also drop the second.di fallback in TEST 7 -- the field is downstream_ingress, and accepting both would let a rename of the public Lua API pass unnoticed. Verified a real SO_LINGER reset still reports client_rst and upstream_rst. --- patch/1.21.4/nginx-stream_metrics.patch | 25 +++++++++++++++++++---- patch/1.25.3.1/nginx-stream_metrics.patch | 25 +++++++++++++++++++---- patch/1.27.1.1/nginx-stream_metrics.patch | 25 +++++++++++++++++++---- patch/1.29.2.4/nginx-stream_metrics.patch | 25 +++++++++++++++++++---- t/stream/metrics.t | 2 +- 5 files changed, 85 insertions(+), 17 deletions(-) diff --git a/patch/1.21.4/nginx-stream_metrics.patch b/patch/1.21.4/nginx-stream_metrics.patch index 7e73e0c..da086f5 100644 --- a/patch/1.21.4/nginx-stream_metrics.patch +++ b/patch/1.21.4/nginx-stream_metrics.patch @@ -57,7 +57,24 @@ diff --git src/stream/ngx_stream_proxy_module.c src/stream/ngx_stream_proxy_modu ngx_stream_proxy_finalize(s, NGX_STREAM_OK); return; } -@@ -1697,6 +1718,10 @@ +@@ -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) { @@ -68,7 +85,7 @@ diff --git src/stream/ngx_stream_proxy_module.c src/stream/ngx_stream_proxy_modu src->read->eof = 1; n = 0; } -@@ -1751,6 +1776,10 @@ +@@ -1751,6 +1786,10 @@ c->log->action = "proxying connection"; @@ -79,7 +96,7 @@ diff --git src/stream/ngx_stream_proxy_module.c src/stream/ngx_stream_proxy_modu if (ngx_stream_proxy_test_finalize(s, from_upstream) == NGX_OK) { return; } -@@ -1930,6 +1959,10 @@ +@@ -1930,6 +1969,10 @@ ngx_log_debug1(NGX_LOG_DEBUG_STREAM, s->connection->log, 0, "close proxy upstream connection: %d", pc->fd); @@ -90,7 +107,7 @@ diff --git src/stream/ngx_stream_proxy_module.c src/stream/ngx_stream_proxy_modu #if (NGX_STREAM_SSL) if (pc->ssl) { pc->ssl->no_wait_shutdown = 1; -@@ -1960,6 +1993,10 @@ +@@ -1960,6 +2003,10 @@ ngx_log_debug1(NGX_LOG_DEBUG_STREAM, s->connection->log, 0, "finalize stream proxy: %i", rc); diff --git a/patch/1.25.3.1/nginx-stream_metrics.patch b/patch/1.25.3.1/nginx-stream_metrics.patch index 835d13a..027f523 100644 --- a/patch/1.25.3.1/nginx-stream_metrics.patch +++ b/patch/1.25.3.1/nginx-stream_metrics.patch @@ -57,7 +57,24 @@ diff --git src/stream/ngx_stream_proxy_module.c src/stream/ngx_stream_proxy_modu ngx_stream_proxy_finalize(s, NGX_STREAM_OK); return; } -@@ -1702,6 +1723,10 @@ +@@ -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) { @@ -68,7 +85,7 @@ diff --git src/stream/ngx_stream_proxy_module.c src/stream/ngx_stream_proxy_modu src->read->eof = 1; n = 0; } -@@ -1756,6 +1781,10 @@ +@@ -1756,6 +1791,10 @@ c->log->action = "proxying connection"; @@ -79,7 +96,7 @@ diff --git src/stream/ngx_stream_proxy_module.c src/stream/ngx_stream_proxy_modu if (ngx_stream_proxy_test_finalize(s, from_upstream) == NGX_OK) { return; } -@@ -1935,6 +1964,10 @@ +@@ -1935,6 +1974,10 @@ ngx_log_debug1(NGX_LOG_DEBUG_STREAM, s->connection->log, 0, "close proxy upstream connection: %d", pc->fd); @@ -90,7 +107,7 @@ diff --git src/stream/ngx_stream_proxy_module.c src/stream/ngx_stream_proxy_modu #if (NGX_STREAM_SSL) if (pc->ssl) { pc->ssl->no_wait_shutdown = 1; -@@ -1965,6 +1998,10 @@ +@@ -1965,6 +2008,10 @@ ngx_log_debug1(NGX_LOG_DEBUG_STREAM, s->connection->log, 0, "finalize stream proxy: %i", rc); diff --git a/patch/1.27.1.1/nginx-stream_metrics.patch b/patch/1.27.1.1/nginx-stream_metrics.patch index 835d13a..027f523 100644 --- a/patch/1.27.1.1/nginx-stream_metrics.patch +++ b/patch/1.27.1.1/nginx-stream_metrics.patch @@ -57,7 +57,24 @@ diff --git src/stream/ngx_stream_proxy_module.c src/stream/ngx_stream_proxy_modu ngx_stream_proxy_finalize(s, NGX_STREAM_OK); return; } -@@ -1702,6 +1723,10 @@ +@@ -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) { @@ -68,7 +85,7 @@ diff --git src/stream/ngx_stream_proxy_module.c src/stream/ngx_stream_proxy_modu src->read->eof = 1; n = 0; } -@@ -1756,6 +1781,10 @@ +@@ -1756,6 +1791,10 @@ c->log->action = "proxying connection"; @@ -79,7 +96,7 @@ diff --git src/stream/ngx_stream_proxy_module.c src/stream/ngx_stream_proxy_modu if (ngx_stream_proxy_test_finalize(s, from_upstream) == NGX_OK) { return; } -@@ -1935,6 +1964,10 @@ +@@ -1935,6 +1974,10 @@ ngx_log_debug1(NGX_LOG_DEBUG_STREAM, s->connection->log, 0, "close proxy upstream connection: %d", pc->fd); @@ -90,7 +107,7 @@ diff --git src/stream/ngx_stream_proxy_module.c src/stream/ngx_stream_proxy_modu #if (NGX_STREAM_SSL) if (pc->ssl) { pc->ssl->no_wait_shutdown = 1; -@@ -1965,6 +1998,10 @@ +@@ -1965,6 +2008,10 @@ ngx_log_debug1(NGX_LOG_DEBUG_STREAM, s->connection->log, 0, "finalize stream proxy: %i", rc); diff --git a/patch/1.29.2.4/nginx-stream_metrics.patch b/patch/1.29.2.4/nginx-stream_metrics.patch index c1b828f..693b713 100644 --- a/patch/1.29.2.4/nginx-stream_metrics.patch +++ b/patch/1.29.2.4/nginx-stream_metrics.patch @@ -57,7 +57,24 @@ diff --git src/stream/ngx_stream_proxy_module.c src/stream/ngx_stream_proxy_modu ngx_stream_proxy_finalize(s, NGX_STREAM_OK); return; } -@@ -1828,6 +1849,10 @@ +@@ -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) { @@ -68,7 +85,7 @@ diff --git src/stream/ngx_stream_proxy_module.c src/stream/ngx_stream_proxy_modu src->read->eof = 1; n = 0; } -@@ -1882,6 +1907,10 @@ +@@ -1882,6 +1917,10 @@ c->log->action = "proxying connection"; @@ -79,7 +96,7 @@ diff --git src/stream/ngx_stream_proxy_module.c src/stream/ngx_stream_proxy_modu if (ngx_stream_proxy_test_finalize(s, from_upstream) == NGX_OK) { return; } -@@ -2061,6 +2090,10 @@ +@@ -2061,6 +2100,10 @@ ngx_log_debug1(NGX_LOG_DEBUG_STREAM, s->connection->log, 0, "close proxy upstream connection: %d", pc->fd); @@ -90,7 +107,7 @@ diff --git src/stream/ngx_stream_proxy_module.c src/stream/ngx_stream_proxy_modu #if (NGX_STREAM_SSL) if (pc->ssl) { pc->ssl->no_wait_shutdown = 1; -@@ -2091,6 +2124,10 @@ +@@ -2091,6 +2134,10 @@ ngx_log_debug1(NGX_LOG_DEBUG_STREAM, s->connection->log, 0, "finalize stream proxy: %i", rc); diff --git a/t/stream/metrics.t b/t/stream/metrics.t index c1fe244..c1d0af5 100644 --- a/t/stream/metrics.t +++ b/t/stream/metrics.t @@ -181,7 +181,7 @@ apisix_stream_metrics_zone 1m; ngx.sleep(0.3) local second = seen() - ngx.say("second di=", second.di or second.downstream_ingress, + ngx.say("second di=", second.downstream_ingress, " ue=", second.upstream_egress) sock:close()