From 93b0cad091052f9ac27ad75fc23ee0c067a52bab Mon Sep 17 00:00:00 2001 From: Hiroshi Hatake Date: Mon, 28 Sep 2026 09:44:35 +0900 Subject: [PATCH 1/4] in_windows_exporter_metrics: fix TCP connection state collection Signed-off-by: Hiroshi Hatake --- .../in_windows_exporter_metrics/we_wmi_tcp.c | 60 +++++++++++++------ 1 file changed, 41 insertions(+), 19 deletions(-) diff --git a/plugins/in_windows_exporter_metrics/we_wmi_tcp.c b/plugins/in_windows_exporter_metrics/we_wmi_tcp.c index ef08ee45bfd..ded628bd800 100644 --- a/plugins/in_windows_exporter_metrics/we_wmi_tcp.c +++ b/plugins/in_windows_exporter_metrics/we_wmi_tcp.c @@ -51,7 +51,8 @@ static inline int windows_state_to_index(int state) return 4; case MIB_TCP_STATE_TIME_WAIT: return 5; - /* MIB_TCP_STATE_CLOSED is 1 */ + case MIB_TCP_STATE_CLOSED: + return 6; case MIB_TCP_STATE_CLOSE_WAIT: return 7; case MIB_TCP_STATE_LAST_ACK: @@ -69,12 +70,16 @@ static inline int windows_state_to_index(int state) static int we_tcp_get_state_metrics(struct flb_we *ctx, const char *af_label) { - PMIB_TCPTABLE2 tcp_table = NULL; + void *tcp_table = NULL; + void *new_table; + PMIB_TCPTABLE_OWNER_PID tcp4_table; + PMIB_TCP6TABLE_OWNER_PID tcp6_table; ULONG buffer_size = 0; DWORD result; DWORD idx = 0; int state_index; int i = 0; + int attempt; const char *state_label; uint64_t timestamp = cfl_time_now(); int af_family = (strcmp(af_label, "ipv4") == 0) ? AF_INET : AF_INET6; @@ -87,34 +92,51 @@ static int we_tcp_get_state_metrics(struct flb_we *ctx, const char *af_label) return -1; } - tcp_table = (PMIB_TCPTABLE2)flb_malloc(buffer_size); - if (tcp_table == NULL) { - flb_plg_error(ctx->ins, "TCP state metrics: could not allocate buffer"); - return -1; - } + /* The table can grow between sizing and collection. Bound retries under churn. */ + for (attempt = 0; attempt < 3; attempt++) { + new_table = flb_realloc(tcp_table, buffer_size); + if (new_table == NULL) { + flb_plg_error(ctx->ins, "TCP state metrics: could not allocate buffer"); + flb_free(tcp_table); + return -1; + } + tcp_table = new_table; - result = GetExtendedTcpTable(tcp_table, &buffer_size, FALSE, af_family, TCP_TABLE_OWNER_PID_ALL, 0); + result = GetExtendedTcpTable(tcp_table, &buffer_size, FALSE, af_family, + TCP_TABLE_OWNER_PID_ALL, 0); + if (result != ERROR_INSUFFICIENT_BUFFER) { + break; + } + } if (result != NO_ERROR) { flb_plg_error(ctx->ins, "TCP state metrics: error getting table: %lu", result); flb_free(tcp_table); return -1; } - for (idx = 0; idx < tcp_table->dwNumEntries; idx++) { - state_index = windows_state_to_index(tcp_table->table[idx].dwState); - state_counts[state_index]++; + if (af_family == AF_INET) { + tcp4_table = (PMIB_TCPTABLE_OWNER_PID) tcp_table; + for (idx = 0; idx < tcp4_table->dwNumEntries; idx++) { + state_index = windows_state_to_index(tcp4_table->table[idx].dwState); + state_counts[state_index]++; + } + } + else { + tcp6_table = (PMIB_TCP6TABLE_OWNER_PID) tcp_table; + for (idx = 0; idx < tcp6_table->dwNumEntries; idx++) { + state_index = windows_state_to_index(tcp6_table->table[idx].dwState); + state_counts[state_index]++; + } } flb_free(tcp_table); for (i = 0; i < 13; i++) { - if (state_counts[i] > 0) { - state_label = TCP_STATE_STRINGS[i]; - labels[0] = af_label; - labels[1] = state_label; - cmt_gauge_set(ctx->wmi_tcp->connections_state, timestamp, - (double)state_counts[i], 2, (char **)labels); - } + state_label = TCP_STATE_STRINGS[i]; + labels[0] = af_label; + labels[1] = state_label; + cmt_gauge_set(ctx->wmi_tcp->connections_state, timestamp, + (double)state_counts[i], 2, (char **)labels); } return 0; @@ -358,4 +380,4 @@ int we_wmi_tcp_update(struct flb_we *ctx) we_wmi_cleanup(ctx); return 0; -} \ No newline at end of file +} From e8e448407b1f1c6755259bf8bfbcc454d816dae9 Mon Sep 17 00:00:00 2001 From: Hiroshi Hatake Date: Mon, 28 Sep 2026 09:45:00 +0900 Subject: [PATCH 2/4] tests: cover Windows TCP state collection Signed-off-by: Hiroshi Hatake --- tests/runtime/CMakeLists.txt | 1 + .../runtime/in_windows_exporter_metrics_tcp.c | 270 ++++++++++++++++++ 2 files changed, 271 insertions(+) create mode 100644 tests/runtime/in_windows_exporter_metrics_tcp.c diff --git a/tests/runtime/CMakeLists.txt b/tests/runtime/CMakeLists.txt index f049f52cc15..c7d9fd5b9b3 100644 --- a/tests/runtime/CMakeLists.txt +++ b/tests/runtime/CMakeLists.txt @@ -39,6 +39,7 @@ FLB_RT_TEST(FLB_CHUNK_TRACE "core_chunk_trace.c") FLB_RT_TEST(FLB_IN_EVENT_TEST "in_event_test.c") # FLB_IN_WINEVTLOG is only enabled on Windows builds FLB_RT_TEST(FLB_IN_WINEVTLOG "in_winevtlog.c") +FLB_RT_TEST(FLB_IN_WINDOWS_EXPORTER_METRICS "in_windows_exporter_metrics_tcp.c") if(FLB_OUT_LIB) # These plugins works only on Linux diff --git a/tests/runtime/in_windows_exporter_metrics_tcp.c b/tests/runtime/in_windows_exporter_metrics_tcp.c new file mode 100644 index 00000000000..52c82860dd1 --- /dev/null +++ b/tests/runtime/in_windows_exporter_metrics_tcp.c @@ -0,0 +1,270 @@ +/* -*- Mode: C; tab-width: 4; indent-tabs-mode: nil; c-basic-offset: 4 -*- */ + +/* Fluent Bit + * ========== + * Copyright (C) 2026 The Fluent Bit Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +#include +#include +#include +#include "flb_tests_runtime.h" +#include "../../plugins/in_windows_exporter_metrics/we_wmi_tcp.h" + +static DWORD fixture_rows; +static int growth_failures; +static int table_calls; +static DWORD api_error; +static int live_table; +static void *guard_allocation; + +/* Put the exact requested buffer immediately before an inaccessible page. */ +static void *guard_realloc(void *old_buffer, size_t size) +{ + SYSTEM_INFO info; + size_t committed; + char *allocation; + + GetSystemInfo(&info); + committed = (size + info.dwPageSize - 1) / info.dwPageSize * info.dwPageSize; + allocation = VirtualAlloc(NULL, committed + info.dwPageSize, + MEM_RESERVE, PAGE_NOACCESS); + TEST_ASSERT(allocation != NULL); + TEST_ASSERT(VirtualAlloc(allocation, committed, MEM_COMMIT, PAGE_READWRITE) != NULL); + if (guard_allocation != NULL) { + VirtualFree(guard_allocation, 0, MEM_RELEASE); + } + guard_allocation = allocation; + /* Retries overwrite the table; no previous contents need to be preserved. */ + return allocation + committed - size; +} + +static void guard_free(void *buffer) +{ + if (buffer != NULL) { + TEST_ASSERT(guard_allocation != NULL); + TEST_CHECK(VirtualFree(guard_allocation, 0, MEM_RELEASE) != 0); + guard_allocation = NULL; + } +} + +static DWORD WINAPI fixture_tcp_table(PVOID buffer, PDWORD size, BOOL ordered, + ULONG family, TCP_TABLE_CLASS table_class, + ULONG reserved) +{ + PMIB_TCPTABLE_OWNER_PID table4; + PMIB_TCP6TABLE_OWNER_PID table6; + DWORD required; + DWORD i; + + if (live_table) { + return GetExtendedTcpTable(buffer, size, ordered, family, table_class, reserved); + } + TEST_CHECK(table_class == TCP_TABLE_OWNER_PID_ALL); + TEST_CHECK(family == AF_INET || family == AF_INET6); + if (family == AF_INET) { + required = offsetof(MIB_TCPTABLE_OWNER_PID, table) + + fixture_rows * sizeof(MIB_TCPROW_OWNER_PID); + } + else { + required = offsetof(MIB_TCP6TABLE_OWNER_PID, table) + + fixture_rows * sizeof(MIB_TCP6ROW_OWNER_PID); + } + if (buffer == NULL) { + *size = required; + return ERROR_INSUFFICIENT_BUFFER; + } + table_calls++; + if (growth_failures > 0) { + growth_failures--; + *size += 56; + return ERROR_INSUFFICIENT_BUFFER; + } + if (api_error != NO_ERROR) { + return api_error; + } + TEST_ASSERT(*size >= required); + memset(buffer, 0xa5, required); + if (family == AF_INET) { + table4 = buffer; + table4->dwNumEntries = fixture_rows; + for (i = 0; i < fixture_rows; i++) { + table4->table[i].dwState = i % 13 + 1; + } + } + else { + table6 = buffer; + table6->dwNumEntries = fixture_rows; + for (i = 0; i < fixture_rows; i++) { + table6->table[i].dwState = i % 13 + 1; + } + } + return NO_ERROR; +} + +/* Exercise the private collector with controlled Win32 replies and guard pages. */ +#define GetExtendedTcpTable fixture_tcp_table +#define flb_realloc guard_realloc +#define flb_free guard_free +#define TCP_STATE_STRINGS test_tcp_state_strings +#define we_wmi_tcp_init test_tcp_init +#define we_wmi_tcp_update test_tcp_update +#define we_wmi_tcp_exit test_tcp_exit +#include "../../plugins/in_windows_exporter_metrics/we_wmi_tcp.c" +#undef GetExtendedTcpTable +#undef flb_realloc +#undef flb_free + +static void check_counts(struct flb_we *ctx, char *family, DWORD rows) +{ + char *states[] = { + "CLOSE", "LISTEN", "SYN_SENT", "SYN_RECV", "ESTABLISHED", "FIN_WAIT1", + "FIN_WAIT2", "CLOSE_WAIT", "CLOSING", "LAST_ACK", "TIME_WAIT", + "DELETE_TCB", "UNKNOWN" + }; + char *labels[2]; + double value; + unsigned int expected; + int i; + + labels[0] = family; + for (i = 0; i < 13; i++) { + labels[1] = states[i]; + expected = rows / 13 + (i < rows % 13); + TEST_CHECK(cmt_gauge_get_val(ctx->wmi_tcp->connections_state, 2, labels, &value) == 0); + TEST_CHECK(value == expected); + TEST_MSG("%s %s: expected %u, got %.0f", family, states[i], expected, value); + } +} + +static void run_collector_test(char *family, int live) +{ + struct flb_we ctx = {0}; + struct we_wmi_tcp_counters counters = {0}; + struct flb_input_instance ins = {0}; + char *keys[] = {"af", "state"}; + char *labels[] = {family, "LISTEN"}; + WSADATA wsa; + SOCKET listener = INVALID_SOCKET; + struct sockaddr_storage address = {0}; + struct sockaddr_in *address4; + struct sockaddr_in6 *address6; + int address_size; + int af; + double value; + + ctx.ins = &ins; + ins.log_level = FLB_LOG_OFF; + ctx.wmi_tcp = &counters; + ctx.cmt = cmt_create(); + TEST_ASSERT(ctx.cmt != NULL); + counters.connections_state = cmt_gauge_create(ctx.cmt, "windows", "tcp", + "connections_state", "TCP states", 2, keys); + TEST_ASSERT(counters.connections_state != NULL); + live_table = live; + fixture_rows = 368; + growth_failures = 0; + api_error = NO_ERROR; + table_calls = 0; + + if (live) { + TEST_ASSERT(WSAStartup(MAKEWORD(2, 2), &wsa) == 0); + af = strcmp(family, "ipv4") == 0 ? AF_INET : AF_INET6; + listener = socket(af, SOCK_STREAM, IPPROTO_TCP); + TEST_ASSERT(listener != INVALID_SOCKET); + if (af == AF_INET) { + address4 = (struct sockaddr_in *) &address; + address4->sin_family = AF_INET; + address4->sin_addr.s_addr = htonl(INADDR_LOOPBACK); + address_size = sizeof(*address4); + } + else { + address6 = (struct sockaddr_in6 *) &address; + address6->sin6_family = AF_INET6; + address6->sin6_addr.u.Byte[15] = 1; + address_size = sizeof(*address6); + } + TEST_ASSERT(bind(listener, (struct sockaddr *) &address, address_size) == 0); + TEST_ASSERT(listen(listener, 1) == 0); + } + + TEST_CHECK(we_tcp_get_state_metrics(&ctx, family) == 0); + TEST_CHECK(guard_allocation == NULL); + if (live) { + TEST_CHECK(cmt_gauge_get_val(counters.connections_state, 2, labels, &value) == 0); + TEST_CHECK(value >= 1); + closesocket(listener); + WSACleanup(); + } + else { + check_counts(&ctx, family, fixture_rows); + + /* A growing table succeeds on the final permitted attempt. */ + growth_failures = 2; + table_calls = 0; + TEST_CHECK(we_tcp_get_state_metrics(&ctx, family) == 0); + TEST_CHECK(table_calls == 3); + check_counts(&ctx, family, fixture_rows); + + /* Failed collections must preserve the last successful values. */ + growth_failures = 10; + table_calls = 0; + TEST_CHECK(we_tcp_get_state_metrics(&ctx, family) == -1); + TEST_CHECK(table_calls == 3); + TEST_CHECK(guard_allocation == NULL); + check_counts(&ctx, family, fixture_rows); + + growth_failures = 0; + api_error = ERROR_ACCESS_DENIED; + TEST_CHECK(we_tcp_get_state_metrics(&ctx, family) == -1); + TEST_CHECK(guard_allocation == NULL); + check_counts(&ctx, family, fixture_rows); + api_error = NO_ERROR; + + /* An empty successful snapshot clears every previously nonzero state. */ + fixture_rows = 0; + TEST_CHECK(we_tcp_get_state_metrics(&ctx, family) == 0); + check_counts(&ctx, family, 0); + } + cmt_destroy(ctx.cmt); +} + +static void test_ipv4(void) +{ + run_collector_test("ipv4", 0); +} + +static void test_ipv6(void) +{ + run_collector_test("ipv6", 0); +} + +static void test_live_ipv4(void) +{ + run_collector_test("ipv4", 1); +} + +static void test_live_ipv6(void) +{ + run_collector_test("ipv6", 1); +} + +TEST_LIST = { + {"tcp_ipv4_states", test_ipv4}, + {"tcp_ipv6_states", test_ipv6}, + {"tcp_live_ipv4", test_live_ipv4}, + {"tcp_live_ipv6", test_live_ipv6}, + {NULL, NULL} +}; From 7cd62b83d6c1e266b2ed28497e747544c9c3c1eb Mon Sep 17 00:00:00 2001 From: Hiroshi Hatake Date: Mon, 28 Sep 2026 10:22:44 +0900 Subject: [PATCH 3/4] cmake: enable Windows exporter only in Windows defaults Signed-off-by: Hiroshi Hatake --- cmake/plugins_options.cmake | 2 +- cmake/windows-setup.cmake | 1 + 2 files changed, 2 insertions(+), 1 deletion(-) diff --git a/cmake/plugins_options.cmake b/cmake/plugins_options.cmake index 7fbd0a96bd5..36b01e32885 100644 --- a/cmake/plugins_options.cmake +++ b/cmake/plugins_options.cmake @@ -62,7 +62,7 @@ DEFINE_OPTION(FLB_IN_THERMAL "Enable Thermal plugin" DEFINE_OPTION(FLB_IN_UDP "Enable UDP input plugin" ON) DEFINE_OPTION(FLB_IN_UNIX_SOCKET "Enable Unix socket input plugin" OFF) DEFINE_OPTION(FLB_IN_WINLOG "Enable Windows Log input plugin" OFF) -DEFINE_OPTION(FLB_IN_WINDOWS_EXPORTER_METRICS "Enable windows exporter metrics input plugin" ON) +DEFINE_OPTION(FLB_IN_WINDOWS_EXPORTER_METRICS "Enable windows exporter metrics input plugin" OFF) DEFINE_OPTION(FLB_IN_WINEVTLOG "Enable Windows EvtLog input plugin" OFF) DEFINE_OPTION(FLB_IN_WINSTAT "Enable Windows Stat input plugin" OFF) DEFINE_OPTION(FLB_IN_EBPF "Enable Linux eBPF input plugin" OFF) diff --git a/cmake/windows-setup.cmake b/cmake/windows-setup.cmake index f253382a0cc..a43e31d21f5 100644 --- a/cmake/windows-setup.cmake +++ b/cmake/windows-setup.cmake @@ -74,6 +74,7 @@ if(FLB_WINDOWS_DEFAULTS) set(FLB_IN_WINLOG Yes) set(FLB_IN_WINSTAT Yes) set(FLB_IN_WINEVTLOG Yes) + set(FLB_IN_WINDOWS_EXPORTER_METRICS Yes) set(FLB_IN_COLLECTD No) set(FLB_IN_STATSD Yes) set(FLB_IN_STORAGE_BACKLOG Yes) From 0418e2bf25c963fcb10210da0f380227d0860d4e Mon Sep 17 00:00:00 2001 From: Hiroshi Hatake Date: Mon, 28 Sep 2026 10:22:58 +0900 Subject: [PATCH 4/4] tests: link Windows TCP test to exporter plugin Signed-off-by: Hiroshi Hatake --- tests/runtime/CMakeLists.txt | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/tests/runtime/CMakeLists.txt b/tests/runtime/CMakeLists.txt index c7d9fd5b9b3..7d5c0d5cbfc 100644 --- a/tests/runtime/CMakeLists.txt +++ b/tests/runtime/CMakeLists.txt @@ -348,6 +348,11 @@ foreach(source_file ${CHECK_PROGRAMS}) ${CMAKE_THREAD_LIBS_INIT} ${SYSTEMD_LIB} ) + if(o_source_file_we STREQUAL "in_windows_exporter_metrics_tcp") + target_link_libraries(${source_file_we} + flb-plugin-in_windows_exporter_metrics + ) + endif() if(FLB_AVRO_ENCODER) target_link_libraries(${source_file_we} avro-static jansson) endif()