From d10bf3f45a8abce883b3e6e98d2677f43b3029b9 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micha=C5=82=20Dobranowski?= Date: Thu, 27 Nov 2025 11:39:37 +0100 Subject: [PATCH 01/24] Upgrade prometheus --- rebar.config | 4 ++-- rebar.lock | 30 +++++++++++++++--------------- 2 files changed, 17 insertions(+), 17 deletions(-) diff --git a/rebar.config b/rebar.config index 65a959a38a2..42b7861f38b 100644 --- a/rebar.config +++ b/rebar.config @@ -56,8 +56,8 @@ {exometer_report_statsd, {git, "https://github.com/esl/exometer_report_statsd.git", {branch, "master"}}}, {syslogger, "0.3.0"}, {flatlog, "0.1.2"}, - {prometheus, "5.0.0"}, - {prometheus_cowboy, "0.1.9"}, + {prometheus, "6.1.1"}, + {prometheus_cowboy, "0.2.0"}, %%% Stateless libraries {opuntia, "1.1.2"}, diff --git a/rebar.lock b/rebar.lock index 12d8886864b..bcaed5cdb60 100644 --- a/rebar.lock +++ b/rebar.lock @@ -1,5 +1,5 @@ {"1.2.0", -[{<<"accept">>,{pkg,<<"accept">>,<<"0.3.6">>},2}, +[{<<"accept">>,{pkg,<<"accept">>,<<"0.3.7">>},2}, {<<"amqp_client">>,{pkg,<<"amqp_client">>,<<"4.2.1">>},0}, {<<"backoff">>,{pkg,<<"backoff">>,<<"1.1.6">>},1}, {<<"base16">>,{pkg,<<"base16">>,<<"2.0.1">>},1}, @@ -19,6 +19,7 @@ {<<"credentials_obfuscation">>, {pkg,<<"credentials_obfuscation">>,<<"3.5.0">>}, 1}, + {<<"ddskerl">>,{pkg,<<"ddskerl">>,<<"0.4.2">>},1}, {<<"eini">>,{pkg,<<"eini">>,<<"1.2.9">>},1}, {<<"epgsql">>,{pkg,<<"epgsql">>,<<"4.7.1">>},0}, {<<"eredis">>,{pkg,<<"eredis">>,<<"1.7.1">>},0}, @@ -86,10 +87,9 @@ 0}, {<<"parse_trans">>,{pkg,<<"parse_trans">>,<<"3.4.0">>},1}, {<<"pooler">>,{pkg,<<"pooler">>,<<"1.5.3">>},1}, - {<<"prometheus">>,{pkg,<<"prometheus">>,<<"5.0.0">>},0}, - {<<"prometheus_cowboy">>,{pkg,<<"prometheus_cowboy">>,<<"0.1.9">>},0}, - {<<"prometheus_httpd">>,{pkg,<<"prometheus_httpd">>,<<"2.1.13">>},1}, - {<<"quantile_estimator">>,{pkg,<<"quantile_estimator">>,<<"1.0.2">>},1}, + {<<"prometheus">>,{pkg,<<"prometheus">>,<<"6.1.1">>},0}, + {<<"prometheus_cowboy">>,{pkg,<<"prometheus_cowboy">>,<<"0.2.0">>},0}, + {<<"prometheus_httpd">>,{pkg,<<"prometheus_httpd">>,<<"2.1.15">>},1}, {<<"quickrand">>,{pkg,<<"quickrand">>,<<"2.0.7">>},1}, {<<"rabbit_common">>,{pkg,<<"rabbit_common">>,<<"4.2.1">>},1}, {<<"ranch">>,{pkg,<<"ranch">>,<<"2.2.0">>},0}, @@ -117,7 +117,7 @@ {<<"worker_pool">>,{pkg,<<"worker_pool">>,<<"6.4.0">>},0}]}. [ {pkg_hash,[ - {<<"accept">>, <<"AD44AC7D704BF70EF8FB2E313EF5B978F9D1330BDDAC64509E93AFDA13281215">>}, + {<<"accept">>, <<"CD6E34A2D7E28CA38B2D3CB233734CA0C221EFBC1F171F91FEC5F162CC2D18DA">>}, {<<"amqp_client">>, <<"CFF0CC13186E57457DC5745F1B3A4127C6857717CB8F5920DC457C84D0AD00A2">>}, {<<"backoff">>, <<"83B72ED2108BA1EE8F7D1C22E0B4A00CFE3593A67DBC792799E8CCE9F42F796B">>}, {<<"base16">>, <<"F0549F732E03BE8124ED0D19FD5EE52146CC8BE24C48CBC3F23AB44B157F11A2">>}, @@ -131,6 +131,7 @@ {<<"cowlib">>, <<"623791C56C1CC9DF54A71A9C55147A401549917F00A2E48A6AE12B812C586CED">>}, {<<"cpool">>, <<"76222AA1DAC0F8089707167BD69D221EDA63DB65B8BD67DBF6E133075392EEDC">>}, {<<"credentials_obfuscation">>, <<"61E282ADFB4439486B3994FAAEC69543C7EE6CC7E70C6340E8853FD9DEAF8219">>}, + {<<"ddskerl">>, <<"A51A90BE9AC9B36A94017670BED479C623B10CA9D4BDA1EDF3A0E48CAEEADA2A">>}, {<<"eini">>, <<"FCC3CBD49BBDD9A1D9735C7365DAFFCD84481CCE81E6CB80537883AA44AC4895">>}, {<<"epgsql">>, <<"D4E47CAE46C18C8AFA88E34D59A9B4BAE16368D7CE1EB3DA24FA755EB28393EB">>}, {<<"eredis">>, <<"39E31AA02ADCD651C657F39AAFD4D31A9B2F63C6C700DC9CECE98D4BC3C897AB">>}, @@ -162,10 +163,9 @@ {<<"p1_utils">>, <<"67B0C4AC9FA3BA3EF563B31AA111B0A004439A37FAC85E027F1C3617E1C7EC6C">>}, {<<"parse_trans">>, <<"BB87AC362A03CA674EBB7D9D498F45C03256ADED7214C9101F7035EF44B798C7">>}, {<<"pooler">>, <<"898CD1FA301FC42D4A8ED598CE139B71CA85B54C16AB161152B5CC5FBDCFA1A8">>}, - {<<"prometheus">>, <<"8A37A3216D8DB019D19068602669C9819C099120F8E39994DD1BD3A3F5553376">>}, - {<<"prometheus_cowboy">>, <<"D9D5B300516A61ED5AE31391F8EEEEB202230081D32A1813F2D78772B6F274E1">>}, - {<<"prometheus_httpd">>, <<"F086390B4E4E3F41112889B745BAC53D26437B6139496E6700C2508858F5985B">>}, - {<<"quantile_estimator">>, <<"ECD281D40110FDD9BA62685531E4435E0839A52FD1058DA5564F1763E4642EF7">>}, + {<<"prometheus">>, <<"3C9D8E2D4FCF948450550693F6A82FAE013D3665BA10AA55955B64C2875FADB3">>}, + {<<"prometheus_cowboy">>, <<"526F75D9850A9125496F78BCEECCA0F237BC7B403C976D44508543AE5967DAD9">>}, + {<<"prometheus_httpd">>, <<"8F767D819A5D36275EAB9264AFF40D87279151646776069BF69FBDBBD562BD75">>}, {<<"quickrand">>, <<"D2BD76676A446E6A058D678444B7FDA1387B813710D1AF6D6E29BB92186C8820">>}, {<<"rabbit_common">>, <<"1D64E391E12116B76B1425EB96B7552DE51F0301093EBA669B5334F4759CC1E8">>}, {<<"ranch">>, <<"25528F82BC8D7C6152C57666CA99EC716510FE0925CB188172F41CE93117B1B0">>}, @@ -184,7 +184,7 @@ {<<"uuid">>, <<"B2078D2CC814F53AFA52D36C91E08962C7E7373585C623F4C0EA6DFB04B2AF94">>}, {<<"worker_pool">>, <<"0347B805A8E5804B5676A9885FB3B9B6C1627099C449C3C67C0E8E6AF79E9AA6">>}]}, {pkg_hash_ext,[ - {<<"accept">>, <<"A5167FA1AE90315C3F1DD189446312F8A55D00EFA357E9C569BDA47736B874C3">>}, + {<<"accept">>, <<"CA69388943F5DAD2E7232A5478F16086E3C872F48E32B88B378E1885A59F5649">>}, {<<"amqp_client">>, <<"8AE00B055A58500E0557F73D9C0FFE257487131E603F7F84FE72CBFAAF03838A">>}, {<<"backoff">>, <<"CF0CFFF8995FB20562F822E5CC47D8CCF664C5ECDC26A684CBE85C225F9D7C39">>}, {<<"base16">>, <<"06EA2D48343282E712160BA89F692B471DB8B36ABE8394F3445FF9032251D772">>}, @@ -198,6 +198,7 @@ {<<"cowlib">>, <<"0AF652D1550C8411C3B58EED7A035A7FB088C0B86AFF6BC504B0BC3B7F791AA2">>}, {<<"cpool">>, <<"430E18DF4A9D584EB1ED0D196A87CC02E878AF5B4888BFDC9B65F86A96480E30">>}, {<<"credentials_obfuscation">>, <<"843ADBE3246861CE0F1A0FA3222F384834EB31DEFD8D6B9CBA7AFD2977C957BC">>}, + {<<"ddskerl">>, <<"63F907373D7E548151D584D4DA8A38928FD26EC9477B94C0FFAAD87D7CB69FE7">>}, {<<"eini">>, <<"DA64AE8DB7C2F502E6F20CDF44CD3D9BE364412B87FF49FEBF282540F673DFCB">>}, {<<"epgsql">>, <<"B6D86B7DC42C8555B1D4E20880E5099D6D6D053148000E188E548F98E4E01836">>}, {<<"eredis">>, <<"7C2B54C566FED55FEEF3341CA79B0100A6348FD3F162184B7ED5118D258C3CC1">>}, @@ -229,10 +230,9 @@ {<<"p1_utils">>, <<"D0379E8C1156B98BD64F8129C1DE022FCCA4F2FDB7486CE73BF0ED2C3376B04C">>}, {<<"parse_trans">>, <<"F99E368830BEA44552224E37E04943A54874F08B8590485DE8D13832B63A2DC3">>}, {<<"pooler">>, <<"058D85C5081289B90E97E4DDDBC3BB5A3B4A19A728AB3BC88C689EFCC36A07C7">>}, - {<<"prometheus">>, <<"80D29564A5DC4490B53FD225D752B65FB0DBEBA41497F96D62223338127C5659">>}, - {<<"prometheus_cowboy">>, <<"5F71C039DEB9E9FF9DD6366BC74C907A463872B85286E619EFF0BDA15111695A">>}, - {<<"prometheus_httpd">>, <<"9B5A44D1F6FBB3C3FE6F85F06DAFE680AD9FFD591EC65A10BB51DFF0FBBE45D2">>}, - {<<"quantile_estimator">>, <<"DB404793D6384995A1AC6DD973E2CEE5BE9FCC128765BDBA53D87C564E296B64">>}, + {<<"prometheus">>, <<"C3ECC7A35676948089E28FF4383CCB85DCEF447F0982C01E5630A084008C456A">>}, + {<<"prometheus_cowboy">>, <<"2C7EB12F4B970D91E3B47BAAD0F138F6ADC34E53EEB0AE18068FF0AFAB441B24">>}, + {<<"prometheus_httpd">>, <<"67736D000745184D5013C58A63E947821AB90CB9320BC2E6AE5D3061C6FFE039">>}, {<<"quickrand">>, <<"B8ACBF89A224BC217C3070CA8BEBC6EB236DBE7F9767993B274084EA044D35F0">>}, {<<"rabbit_common">>, <<"FF509B07E639B1784898C28031E5204FEA14260172E4FC339F94405586037E40">>}, {<<"ranch">>, <<"FA0B99A1780C80218A4197A59EA8D3BDAE32FBFF7E88527D7D8A4787EFF4F8E7">>}, From 8f9296629eed2e5ade75d4d95478f4db25f0884c Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micha=C5=82=20Dobranowski?= Date: Fri, 28 Nov 2025 15:47:15 +0100 Subject: [PATCH 02/24] Update prometheus histogram tests --- test/mongoose_instrument_metrics_SUITE.erl | 61 ++++++++++++++++++---- 1 file changed, 50 insertions(+), 11 deletions(-) diff --git a/test/mongoose_instrument_metrics_SUITE.erl b/test/mongoose_instrument_metrics_SUITE.erl index 69b7a7d39a1..eaa28c3a707 100644 --- a/test/mongoose_instrument_metrics_SUITE.erl +++ b/test/mongoose_instrument_metrics_SUITE.erl @@ -10,6 +10,24 @@ -define(HOST_TYPE, <<"localhost">>). -define(HOST_TYPE2, <<"test type">>). +-define(assertRelEqual(Expect, Actual, Tolerance), + (fun() -> + __Exp = (Expect), + __Act = (Actual), + __Tol = (Tolerance), + case __Exp == 0 of + true -> ?assertEqual(__Exp, __Act); + false -> + __RelErr = abs(__Exp - __Act) / __Exp, + __Check = case __RelErr < __Tol of + true -> ok; + false -> {fail, [{expected, __Exp}, {actual, __Act}, + {rel_error, __RelErr}, {tolerance, __Tol}]} + end, + ?assertEqual(ok, __Check) + end + end)()). + -import(mongoose_instrument_exometer, [exometer_metric_name/3]). %% Setup and teardown @@ -30,6 +48,7 @@ groups() -> prometheus_counter_cannot_be_decreased, prometheus_counter_is_updated_separately_for_different_labels, prometheus_histogram_is_created_and_updated, + prometheus_histogram_is_calculated_correctly, prometheus_histogram_is_updated_separately_for_different_labels, multiple_prometheus_metrics_are_updated]}, {exometer, [parallel], [exometer_skips_non_metric_event, @@ -184,13 +203,33 @@ prometheus_histogram_is_created_and_updated(Config) -> ok = mongoose_instrument:set_up(Event, ?LABELS, #{metrics => #{time => histogram}}), %% Prometheus histogram shows no value if there is no data - ?assertEqual(undefined, prometheus_histogram:value(Metric, [?HOST_TYPE])), + ?assertEqual(undefined, prometheus_quantile_summary:value(Metric, [?HOST_TYPE])), ok = mongoose_instrument:execute(Event, ?LABELS, #{time => 1}), - ?assertMatch({[1, 0|_], 1}, prometheus_histogram:value(Metric, [?HOST_TYPE])), + ?assertMatch({1, 1, _}, prometheus_quantile_summary:value(Metric, [?HOST_TYPE])), ok = mongoose_instrument:execute(Event, ?LABELS, #{time => 1}), - ?assertMatch({[2, 0|_], 2}, prometheus_histogram:value(Metric, [?HOST_TYPE])), + ?assertMatch({2, 2, _}, prometheus_quantile_summary:value(Metric, [?HOST_TYPE])), ok = mongoose_instrument:execute(Event, ?LABELS, #{time => 2}), - ?assertMatch({[2, 1|_], 4}, prometheus_histogram:value(Metric, [?HOST_TYPE])). + ?assertMatch({3, 4, _}, prometheus_quantile_summary:value(Metric, [?HOST_TYPE])). + +prometheus_histogram_is_calculated_correctly(Config) -> + Event = ?config(event, Config), + Metric = prom_name(Event, time), + ok = mongoose_instrument:set_up(Event, ?LABELS, #{metrics => #{time => histogram}}), + + %% Prometheus histogram shows no value if there is no data + ?assertEqual(undefined, prometheus_quantile_summary:value(Metric, [?HOST_TYPE])), + Values = [1, 2, 2, 2, 2, 2, 2, 9, 9, 10], + lists:foreach(fun(V) -> ok = mongoose_instrument:execute(Event, ?LABELS, #{time => V}) end, Values), + Sum = lists:sum(Values), + Length = length(Values), + ?assertMatch({Length, Sum, _}, prometheus_quantile_summary:value(Metric, [?HOST_TYPE])), + %% Check some quantiles + {_, _, Quantiles} = prometheus_quantile_summary:value(Metric, [?HOST_TYPE]), + Tolerance = 0.01, + ?assertRelEqual(2, proplists:get_value(0.5, Quantiles), Tolerance), + ?assertRelEqual(9, proplists:get_value(0.9, Quantiles), Tolerance), + ?assertRelEqual(10, proplists:get_value(0.95, Quantiles), Tolerance), + ?assertRelEqual(10, proplists:get_value(0.99, Quantiles), Tolerance). prometheus_histogram_is_updated_separately_for_different_labels(Config) -> Event = ?config(event, Config), @@ -199,8 +238,8 @@ prometheus_histogram_is_updated_separately_for_different_labels(Config) -> ok = mongoose_instrument:set_up(Event, ?LABELS2, #{metrics => #{time => histogram}}), ok = mongoose_instrument:execute(Event, ?LABELS, #{time => 1}), ok = mongoose_instrument:execute(Event, ?LABELS2, #{time => 2}), - ?assertMatch({[1, 0|_], 1}, prometheus_histogram:value(Metric, [?HOST_TYPE])), - ?assertMatch({[0, 1|_], 2}, prometheus_histogram:value(Metric, [?HOST_TYPE2])). + ?assertMatch({1, 1, _}, prometheus_quantile_summary:value(Metric, [?HOST_TYPE])), + ?assertMatch({1, 2, _}, prometheus_quantile_summary:value(Metric, [?HOST_TYPE2])). multiple_prometheus_metrics_are_updated(Config) -> Event = ?config(event, Config), @@ -211,18 +250,18 @@ multiple_prometheus_metrics_are_updated(Config) -> %% Update both metrics ok = mongoose_instrument:execute(Event, ?LABELS, #{count => 1, time => 2}), ?assertEqual(1, prometheus_counter:value(Counter, [?HOST_TYPE])), - HistogramValue = prometheus_histogram:value(Histogram, [?HOST_TYPE]), - ?assertMatch({[0, 1|_], 2}, HistogramValue), + HistogramValue = prometheus_quantile_summary:value(Histogram, [?HOST_TYPE]), + ?assertMatch({1, 2, _}, HistogramValue), %% Update only one metric ok = mongoose_instrument:execute(Event, ?LABELS, #{count => 2}), ?assertEqual(3, prometheus_counter:value(Counter, [?HOST_TYPE])), - ?assertEqual(HistogramValue, prometheus_histogram:value(Histogram, [?HOST_TYPE])), + ?assertEqual(HistogramValue, prometheus_quantile_summary:value(Histogram, [?HOST_TYPE])), %% No update ok = mongoose_instrument:execute(Event, ?LABELS, #{something => irrelevant}), ?assertEqual(3, prometheus_counter:value(Counter, [?HOST_TYPE])), - ?assertEqual(HistogramValue, prometheus_histogram:value(Histogram, [?HOST_TYPE])). + ?assertEqual(HistogramValue, prometheus_quantile_summary:value(Histogram, [?HOST_TYPE])). exometer_skips_non_metric_event(Config) -> Event = ?config(event, Config), @@ -386,7 +425,7 @@ prometheus_and_exometer_metrics_are_updated(Config) -> ?assertEqual({ok, [{count, 1}]}, exometer:get_value([?HOST_TYPE, Event, count], count)), ?assertEqual({ok, [{mean, 2}]}, exometer:get_value([?HOST_TYPE, Event, time], mean)), ?assertEqual(1, prometheus_counter:value(prom_name(Event, count), [?HOST_TYPE])), - ?assertMatch({[0, 1|_], 2}, prometheus_histogram:value(prom_name(Event, time), [?HOST_TYPE])). + ?assertMatch({1, 2, _}, prometheus_quantile_summary:value(prom_name(Event, time), [?HOST_TYPE])). %% Helpers From fd16df3a6a3fc7fcadb9f1bd69688c0199658b9a Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micha=C5=82=20Dobranowski?= Date: Wed, 3 Dec 2025 10:59:33 +0100 Subject: [PATCH 03/24] Use prometheus quantile summaries --- src/instrument/mongoose_instrument_exometer.erl | 2 +- src/instrument/mongoose_instrument_prometheus.erl | 15 +++------------ 2 files changed, 4 insertions(+), 13 deletions(-) diff --git a/src/instrument/mongoose_instrument_exometer.erl b/src/instrument/mongoose_instrument_exometer.erl index eb2109fa151..567a74e59e1 100644 --- a/src/instrument/mongoose_instrument_exometer.erl +++ b/src/instrument/mongoose_instrument_exometer.erl @@ -153,7 +153,7 @@ handle_metric_event(EventName, Labels, MetricName, MetricType, Measurements) -> ok end. --spec update_metric(exometer:name(), spiral | histogram, integer()) -> ok. +-spec update_metric(exometer:name(), mongoose_instrument:metric_type(), integer()) -> ok. update_metric(Name, gauge, Value) when is_integer(Value) -> ok = exometer:update(Name, Value); update_metric(Name, counter, Value) when is_integer(Value) -> diff --git a/src/instrument/mongoose_instrument_prometheus.erl b/src/instrument/mongoose_instrument_prometheus.erl index f0843d37625..df73da95700 100644 --- a/src/instrument/mongoose_instrument_prometheus.erl +++ b/src/instrument/mongoose_instrument_prometheus.erl @@ -57,7 +57,7 @@ declare_metric(MetricSpec, counter) -> declare_metric(MetricSpec, spiral) -> prometheus_counter:declare(MetricSpec); declare_metric(MetricSpec, histogram) -> - prometheus_histogram:declare([{buckets, histogram_buckets()} | MetricSpec]). + prometheus_quantile_summary:declare([{quantiles, [0.1, 0.3, 0.5, 0.9, 0.95, 0.99]} | MetricSpec]). -spec reset_metric(name(), [mongoose_instrument:label_value()], mongoose_instrument:metric_type()) -> boolean(). @@ -68,7 +68,7 @@ reset_metric(Name, LabelValues, counter) -> reset_metric(Name, LabelValues, spiral) -> prometheus_counter:remove(Name, LabelValues); reset_metric(Name, LabelValues, histogram) -> - prometheus_histogram:remove(Name, LabelValues). + prometheus_quantile_summary:remove(Name, LabelValues). -spec initialize_metric(name(), [mongoose_instrument:label_value()], mongoose_instrument:metric_type()) -> ok. @@ -92,15 +92,6 @@ metric_spec(EventName, LabelKeys, MetricName) -> {duration_unit, false} % prevent unwanted implicit conversions, e.g. seconds -> microseconds ]. --spec histogram_buckets() -> [integer()]. -histogram_buckets() -> - histogram_buckets([], 1 bsl 30). % ~1.07 * 10^9 - -histogram_buckets(AccBuckets, Val) when Val > 0 -> - histogram_buckets([Val | AccBuckets], Val bsr 1); -histogram_buckets(AccBuckets, _Val) -> - AccBuckets. - -spec handle_metric_event(mongoose_instrument:event_name(), [mongoose_instrument:label_value()], mongoose_instrument:metric_name(), mongoose_instrument:metric_type(), mongoose_instrument:measurements()) -> ok. @@ -139,4 +130,4 @@ update_metric(Name, Labels, counter, Value) when is_integer(Value) -> update_metric(Name, Labels, spiral, Value) when is_integer(Value), Value >= 0 -> ok = prometheus_counter:inc(Name, Labels, Value); update_metric(Name, Labels, histogram, Value) when is_integer(Value) -> - ok = prometheus_histogram:observe(Name, Labels, Value). + ok = prometheus_quantile_summary:observe(Name, Labels, Value). From b05c8329f3650e5ae46976bddcc1b936ead4f986 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micha=C5=82=20Dobranowski?= Date: Wed, 3 Dec 2025 14:08:42 +0100 Subject: [PATCH 04/24] Fix prometheus_endpoint test --- test/prometheus_endpoint_SUITE.erl | 34 +++++------------------------- 1 file changed, 5 insertions(+), 29 deletions(-) diff --git a/test/prometheus_endpoint_SUITE.erl b/test/prometheus_endpoint_SUITE.erl index 20bbebcf996..b063481e1e7 100644 --- a/test/prometheus_endpoint_SUITE.erl +++ b/test/prometheus_endpoint_SUITE.erl @@ -54,7 +54,7 @@ test_metrics(Event, Labels) -> check_counter(Event, requests, Labels, 1, Scraped), % 'spiral' is a Prometheus counter check_gauge(Event, sessions, Labels, 5, Scraped), % 'counter' is a Prometheus gauge check_gauge(Event, seconds, Labels, 10, Scraped), - check_histogram(Event, time, Labels, #{count => 1, sum => 2, bucket_num => 32}, Scraped). + check_histogram(Event, time, Labels, #{count => 1, sum => 2}, Scraped). %% Checks for the parsed metrics @@ -71,38 +71,14 @@ check_gauge(Event, Metric, Labels, ExpValue, Scraped) -> ?assertEqual({Labels, ExpValue}, Value). check_histogram(Event, Metric, Labels, ExpValues, Scraped) -> - #{count := ExpCount, sum := ExpSum, bucket_num := ExpBucketNum} = ExpValues, - [Type, Help] = get_metric([Event, Metric], Scraped), - ?assertEqual({<<"TYPE">>, <<"histogram">>}, Type), + #{count := ExpCount, sum := ExpSum} = ExpValues, + [Type, Help | _] = get_metric([Event, Metric], Scraped), + ?assertEqual({<<"TYPE">>, <<"summary">>}, Type), ?assertEqual({<<"HELP">>, help(Event, Metric)}, Help), [Count] = get_metric([Event, Metric, count], Scraped), ?assertEqual({Labels, ExpCount}, Count), [Sum] = get_metric([Event, Metric, sum], Scraped), - ?assertEqual({Labels, ExpSum}, Sum), - Buckets = get_metric([Event, Metric, bucket], Scraped), - ?assertEqual(ExpBucketNum, length(Buckets)), - check_buckets(ExpCount, Labels, Buckets). - -%% Check that the histogram buckets have growing thresholds and counts, -%% and that the last bucket has the expected total count (because they are cumulative). -check_buckets(ExpCount, Labels, Buckets) -> - InitState = #{labels => Labels, last_count => 0, last_threshold => 0}, - #{final_count := LastCount} = lists:foldl(fun check_bucket/2, InitState, Buckets), - ?assertEqual(ExpCount, LastCount). - -check_bucket({Labels, Count}, State) -> - #{labels := BaseLabels, last_count := LastCount, last_threshold := LastThreshold} = State, - {ThresholdBin, Labels1} = maps:take(le, Labels), - ?assertEqual(Labels1, BaseLabels), - ?assert(Count >= LastCount), - case ThresholdBin of - <<"+Inf">> -> - #{final_count => Count}; - _ -> - Threshold = binary_to_integer(ThresholdBin), - ?assert(Threshold > LastThreshold), - State#{last_count => Count, last_threshold => Threshold} - end. + ?assertEqual({Labels, ExpSum}, Sum). help(Event, Metric) -> <<"Event: ", (atom_to_binary(Event))/binary, ", Metric: ", (atom_to_binary(Metric))/binary>>. From ef6ed93863f8e1bb8555673963933382bf1a51f9 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micha=C5=82=20Dobranowski?= Date: Thu, 4 Dec 2025 12:37:11 +0100 Subject: [PATCH 05/24] Change predefined quantiles and update docs --- doc/operation-and-maintenance/MongooseIM-metrics.md | 11 ++++++----- src/instrument/mongoose_instrument_prometheus.erl | 2 +- 2 files changed, 7 insertions(+), 6 deletions(-) diff --git a/doc/operation-and-maintenance/MongooseIM-metrics.md b/doc/operation-and-maintenance/MongooseIM-metrics.md index ea17a54d558..98c8c601d10 100644 --- a/doc/operation-and-maintenance/MongooseIM-metrics.md +++ b/doc/operation-and-maintenance/MongooseIM-metrics.md @@ -50,16 +50,17 @@ All metrics are divided into the following groups:

`histogram`

- A histogram collects values and groups them in buckets. + A histogram collects values over a 60-second sliding window and exposes the number of observations and their sum (as counters), along with the 50th, 75th, 90th, 95th, 99th, and 99.9th percentiles. **Example:** ``` - # TYPE xmpp_element_in_byte_size histogram + # TYPE xmpp_element_in_byte_size summary # HELP xmpp_element_in_byte_size Event: xmpp_element_in, Metric: byte_size - xmpp_element_in_byte_size_bucket{connection_type="c2s",host_type="localhost",le="1"} 0 + xmpp_element_in_byte_size_count{connection_type="c2s",host_type="localhost"} 0 + xmpp_element_in_byte_size_sum{connection_type="c2s",host_type="localhost"} 0 + xmpp_element_in_byte_size{connection_type="c2s",host_type="localhost",quantile="0.5"} 0 ... - xmpp_element_in_byte_size_bucket{connection_type="c2s",host_type="localhost",le="1073741824"} 0 - xmpp_element_in_byte_size_bucket{connection_type="c2s",host_type="localhost",le="+Inf"} 0 + xmpp_element_in_byte_size{connection_type="c2s",host_type="localhost",quantile="0.999"} 0 ``` === "Exometer" diff --git a/src/instrument/mongoose_instrument_prometheus.erl b/src/instrument/mongoose_instrument_prometheus.erl index df73da95700..e073867bcc6 100644 --- a/src/instrument/mongoose_instrument_prometheus.erl +++ b/src/instrument/mongoose_instrument_prometheus.erl @@ -57,7 +57,7 @@ declare_metric(MetricSpec, counter) -> declare_metric(MetricSpec, spiral) -> prometheus_counter:declare(MetricSpec); declare_metric(MetricSpec, histogram) -> - prometheus_quantile_summary:declare([{quantiles, [0.1, 0.3, 0.5, 0.9, 0.95, 0.99]} | MetricSpec]). + prometheus_quantile_summary:declare([{quantiles, [0.5, 0.75, 0.90, 0.95, 0.99, 0.999]} | MetricSpec]). -spec reset_metric(name(), [mongoose_instrument:label_value()], mongoose_instrument:metric_type()) -> boolean(). From 42797053f857dd9789ff5eec634d014b4d02b349 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micha=C5=82=20Dobranowski?= Date: Thu, 4 Dec 2025 16:21:39 +0100 Subject: [PATCH 06/24] Set `error` and `bound` in quantile_summary --- doc/operation-and-maintenance/MongooseIM-metrics.md | 2 +- src/instrument/mongoose_instrument_prometheus.erl | 8 +++++++- 2 files changed, 8 insertions(+), 2 deletions(-) diff --git a/doc/operation-and-maintenance/MongooseIM-metrics.md b/doc/operation-and-maintenance/MongooseIM-metrics.md index 98c8c601d10..2c07f607842 100644 --- a/doc/operation-and-maintenance/MongooseIM-metrics.md +++ b/doc/operation-and-maintenance/MongooseIM-metrics.md @@ -50,7 +50,7 @@ All metrics are divided into the following groups:

`histogram`

- A histogram collects values over a 60-second sliding window and exposes the number of observations and their sum (as counters), along with the 50th, 75th, 90th, 95th, 99th, and 99.9th percentiles. + A histogram collects values and exposes the number of observations and their sum (as counters), along with the 50th, 75th, 90th, 95th, 99th, and 99.9th percentiles with 1% accuracy. **Example:** ``` diff --git a/src/instrument/mongoose_instrument_prometheus.erl b/src/instrument/mongoose_instrument_prometheus.erl index e073867bcc6..650b1c642cf 100644 --- a/src/instrument/mongoose_instrument_prometheus.erl +++ b/src/instrument/mongoose_instrument_prometheus.erl @@ -57,7 +57,13 @@ declare_metric(MetricSpec, counter) -> declare_metric(MetricSpec, spiral) -> prometheus_counter:declare(MetricSpec); declare_metric(MetricSpec, histogram) -> - prometheus_quantile_summary:declare([{quantiles, [0.5, 0.75, 0.90, 0.95, 0.99, 0.999]} | MetricSpec]). + prometheus_quantile_summary:declare([ + {quantiles, [0.5, 0.75, 0.90, 0.95, 0.99, 0.999]}, + {error, 0.01}, + %% Measuring in µs suffices for actions lasting up to a day (with 1% accuracy). + %% Measuring in bytes suffices for sizes up to 81 GB (with 1% accuracy). + {bound, 1260} + | MetricSpec]). -spec reset_metric(name(), [mongoose_instrument:label_value()], mongoose_instrument:metric_type()) -> boolean(). From 587a84c901e7cc940c543bd17111febc07fc0a27 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micha=C5=82=20Dobranowski?= Date: Wed, 17 Dec 2025 10:45:15 +0100 Subject: [PATCH 07/24] Update docs --- doc/operation-and-maintenance/MongooseIM-metrics.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/doc/operation-and-maintenance/MongooseIM-metrics.md b/doc/operation-and-maintenance/MongooseIM-metrics.md index 2c07f607842..4887225c599 100644 --- a/doc/operation-and-maintenance/MongooseIM-metrics.md +++ b/doc/operation-and-maintenance/MongooseIM-metrics.md @@ -50,7 +50,7 @@ All metrics are divided into the following groups:

`histogram`

- A histogram collects values and exposes the number of observations and their sum (as counters), along with the 50th, 75th, 90th, 95th, 99th, and 99.9th percentiles with 1% accuracy. + A histogram collects values over a 60-second sliding window and exposes the number of observations and their sum (as counters), along with the 50th, 75th, 90th, 95th, 99th, and 99.9th percentiles with 1% accuracy. **Example:** ``` From f851c51bc24f7781fa0a4031b5f46204b075ba70 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micha=C5=82=20Dobranowski?= Date: Wed, 17 Dec 2025 11:18:04 +0100 Subject: [PATCH 08/24] Add support for sliding window (histogram metrics) --- rebar.config | 1 + rebar.lock | 2 +- .../mongoose_instrument_prometheus.erl | 19 +- .../mongoose_prometheus_sliding_window.erl | 243 ++++++++++++++++++ 4 files changed, 260 insertions(+), 5 deletions(-) create mode 100644 src/instrument/mongoose_prometheus_sliding_window.erl diff --git a/rebar.config b/rebar.config index 42b7861f38b..2993aa3fe21 100644 --- a/rebar.config +++ b/rebar.config @@ -56,6 +56,7 @@ {exometer_report_statsd, {git, "https://github.com/esl/exometer_report_statsd.git", {branch, "master"}}}, {syslogger, "0.3.0"}, {flatlog, "0.1.2"}, + {ddskerl, "0.4.2"}, {prometheus, "6.1.1"}, {prometheus_cowboy, "0.2.0"}, diff --git a/rebar.lock b/rebar.lock index bcaed5cdb60..98dbee6ce43 100644 --- a/rebar.lock +++ b/rebar.lock @@ -19,7 +19,7 @@ {<<"credentials_obfuscation">>, {pkg,<<"credentials_obfuscation">>,<<"3.5.0">>}, 1}, - {<<"ddskerl">>,{pkg,<<"ddskerl">>,<<"0.4.2">>},1}, + {<<"ddskerl">>,{pkg,<<"ddskerl">>,<<"0.4.2">>},0}, {<<"eini">>,{pkg,<<"eini">>,<<"1.2.9">>},1}, {<<"epgsql">>,{pkg,<<"epgsql">>,<<"4.7.1">>},0}, {<<"eredis">>,{pkg,<<"eredis">>,<<"1.7.1">>},0}, diff --git a/src/instrument/mongoose_instrument_prometheus.erl b/src/instrument/mongoose_instrument_prometheus.erl index 650b1c642cf..a125c14fd59 100644 --- a/src/instrument/mongoose_instrument_prometheus.erl +++ b/src/instrument/mongoose_instrument_prometheus.erl @@ -13,6 +13,13 @@ start(#{}) -> Apps = [prometheus, prometheus_httpd, prometheus_cowboy], {ok, _} = application:ensure_all_started(Apps, permanent), + %% Start sliding window manager for histogram metrics + case whereis(mongoose_prometheus_sliding_window) of + undefined -> + {ok, _} = mongoose_prometheus_sliding_window:start_link(); + _ -> + ok + end, ok. -spec set_up(mongoose_instrument:event_name(), mongoose_instrument:labels(), @@ -57,13 +64,15 @@ declare_metric(MetricSpec, counter) -> declare_metric(MetricSpec, spiral) -> prometheus_counter:declare(MetricSpec); declare_metric(MetricSpec, histogram) -> - prometheus_quantile_summary:declare([ - {quantiles, [0.5, 0.75, 0.90, 0.95, 0.99, 0.999]}, - {error, 0.01}, + SWResult = mongoose_prometheus_sliding_window:declare(MetricSpec), + PromResult = prometheus_quantile_summary:declare([ + {quantiles, mongoose_prometheus_sliding_window:default_quantiles()}, + {error, mongoose_prometheus_sliding_window:default_error()}, %% Measuring in µs suffices for actions lasting up to a day (with 1% accuracy). %% Measuring in bytes suffices for sizes up to 81 GB (with 1% accuracy). {bound, 1260} - | MetricSpec]). + | MetricSpec]), + SWResult or PromResult. -spec reset_metric(name(), [mongoose_instrument:label_value()], mongoose_instrument:metric_type()) -> boolean(). @@ -74,6 +83,7 @@ reset_metric(Name, LabelValues, counter) -> reset_metric(Name, LabelValues, spiral) -> prometheus_counter:remove(Name, LabelValues); reset_metric(Name, LabelValues, histogram) -> + mongoose_prometheus_sliding_window:remove(Name, LabelValues), prometheus_quantile_summary:remove(Name, LabelValues). -spec initialize_metric(name(), [mongoose_instrument:label_value()], @@ -136,4 +146,5 @@ update_metric(Name, Labels, counter, Value) when is_integer(Value) -> update_metric(Name, Labels, spiral, Value) when is_integer(Value), Value >= 0 -> ok = prometheus_counter:inc(Name, Labels, Value); update_metric(Name, Labels, histogram, Value) when is_integer(Value) -> + ok = mongoose_prometheus_sliding_window:observe(Name, Labels, Value), ok = prometheus_quantile_summary:observe(Name, Labels, Value). diff --git a/src/instrument/mongoose_prometheus_sliding_window.erl b/src/instrument/mongoose_prometheus_sliding_window.erl new file mode 100644 index 00000000000..d317baeefdb --- /dev/null +++ b/src/instrument/mongoose_prometheus_sliding_window.erl @@ -0,0 +1,243 @@ +-module(mongoose_prometheus_sliding_window). + +%% Manages a sliding window of 20 DDSketch instances (3 seconds each = 60 seconds total) +%% Each window uses ddskerl to store quantile summaries + +-export([declare/1, + observe/3, + value/2, + remove/2]). + +-export([start_link/0, + init/1, + handle_call/3, + handle_cast/2, + handle_info/2, + terminate/2, + code_change/3]). + +-export([default_quantiles/0, default_error/0]). + +-behaviour(gen_server). + +-define(WINDOW_COUNT, 20). +-define(WINDOW_DURATION_MS, 3000). % 3 seconds per window +-define(TOTAL_WINDOW_MS, (?WINDOW_COUNT * ?WINDOW_DURATION_MS)). % 60 seconds total + +-spec default_quantiles() -> [number()]. +default_quantiles() -> [0.5, 0.75, 0.90, 0.95, 0.99, 0.999]. + +-spec default_error() -> number(). +default_error() -> 0.01. + +-record(state, { + timer_ref :: reference() | undefined, + metrics = #{} :: #{name() => #{label_values() => metric_state()}} +}). + +-type name() :: string(). +-type label_values() :: [mongoose_instrument:label_value()]. +-type metric_state() :: #{windows => [{ddskerl_ets:ddsketch(), non_neg_integer()}], + current_index => non_neg_integer()}. + +%% Public API + +-spec declare(proplists:proplist()) -> boolean(). +declare(MetricSpec) -> + Name = proplists:get_value(name, MetricSpec), + gen_server:call(?MODULE, {declare, Name, MetricSpec}). + +-spec observe(name(), label_values(), number()) -> ok. +observe(Name, LabelValues, Value) -> + gen_server:cast(?MODULE, {observe, Name, LabelValues, Value}). + +-spec value(name(), label_values()) -> + undefined | {non_neg_integer(), number(), [{number(), number()}]}. +value(Name, LabelValues) -> + gen_server:call(?MODULE, {value, Name, LabelValues}). + +-spec remove(name(), label_values()) -> boolean(). +remove(Name, LabelValues) -> + gen_server:call(?MODULE, {remove, Name, LabelValues}). + +%% gen_server callbacks + +start_link() -> + gen_server:start_link({local, ?MODULE}, ?MODULE, [], []). + +init([]) -> + TimerRef = erlang:start_timer(?WINDOW_DURATION_MS, self(), rotate), + {ok, #state{timer_ref = TimerRef}}. + +handle_call({declare, Name, _MetricSpec}, _From, State) -> + case maps:is_key(Name, State#state.metrics) of + true -> + {reply, false, State}; + false -> + NewMetrics = maps:put(Name, #{}, State#state.metrics), + {reply, true, State#state{metrics = NewMetrics}} + end; + +handle_call({observe, Name, LabelValues, Value}, _From, State) -> + NewState = do_observe(Name, LabelValues, Value, State), + {reply, ok, NewState}; + +handle_call({value, Name, LabelValues}, _From, State) -> + {Result, NewState} = get_value(Name, LabelValues, State), + {reply, Result, NewState}; + +handle_call({remove, Name, LabelValues}, _From, State) -> + NewState = do_remove(Name, LabelValues, State), + {reply, true, NewState}; + +handle_call(_Request, _From, State) -> + {reply, ok, State}. + +handle_cast({observe, Name, LabelValues, Value}, State) -> + NewState = do_observe(Name, LabelValues, Value, State), + {noreply, NewState}; + +handle_cast(_Msg, State) -> + {noreply, State}. + +handle_info({timeout, _TimerRef, rotate}, State) -> + %% The rotation is handled in ensure_windows_rotated, so we just reschedule + TimerRef = erlang:start_timer(?WINDOW_DURATION_MS, self(), rotate), + {noreply, State#state{timer_ref = TimerRef}}; + +handle_info(_Info, State) -> + {noreply, State}. + +terminate(_Reason, _State) -> + ok. + +code_change(_OldVsn, State, _Extra) -> + {ok, State}. + +%% Internal functions + +-spec do_observe(name(), label_values(), number(), #state{}) -> #state{}. +do_observe(Name, LabelValues, Value, State) -> + case maps:find(Name, State#state.metrics) of + error -> + State; % metric not declared, ignore + {ok, MetricState} -> + Key = {Name, LabelValues}, + CurrentTime = erlang:monotonic_time(millisecond), + {Windows, CurrentIndex} = + ensure_windows_rotated(Key, CurrentTime, State), + + %% Get current window and add observation + {CurrentWindow, WindowStartTime} = lists:nth(CurrentIndex + 1, Windows), + UpdatedWindow = ddskerl_std:insert(CurrentWindow, Value), + + %% Update windows list and metric state + UpdatedWindows = set_nth(CurrentIndex + 1, {UpdatedWindow, WindowStartTime}, Windows), + UpdatedMetricState = maps:put(Key, + #{windows => UpdatedWindows, + current_index => CurrentIndex}, + MetricState), + + UpdatedMetrics = maps:put(Name, UpdatedMetricState, State#state.metrics), + State#state{metrics = UpdatedMetrics} + end. + +-spec ensure_windows_rotated({name(), label_values()}, non_neg_integer(), #state{}) -> + {[{ddskerl_std:ddsketch(), non_neg_integer()}], non_neg_integer()}. +ensure_windows_rotated(Key, CurrentTime, State) -> + case get_metric_state(Key, State) of + undefined -> + %% Initialize new windows, all starting at current time (they'll be rotated as needed) + Windows = [{ddskerl_std:new(#{error => default_error()}), CurrentTime} || _ <- lists:seq(1, ?WINDOW_COUNT)], + {Windows, 0}; + #{windows := Windows, current_index := CurrentIndex} -> + {_CurrentWindow, WindowStartTime} = lists:nth(CurrentIndex + 1, Windows), + Elapsed = CurrentTime - WindowStartTime, + if + Elapsed >= ?WINDOW_DURATION_MS -> + NewIndex = (CurrentIndex + 1) rem ?WINDOW_COUNT, + NewWindow = {ddskerl_std:new(#{error => default_error()}), CurrentTime}, + UpdatedWindows = set_nth(CurrentIndex + 1, NewWindow, Windows), + {UpdatedWindows, NewIndex}; + true -> + {Windows, CurrentIndex} + end + end. + +-spec get_value(name(), label_values(), #state{}) -> + {undefined | {non_neg_integer(), number(), [{number(), number()}]}, #state{}}. +get_value(Name, LabelValues, State) -> + Key = {Name, LabelValues}, + case get_metric_state(Key, State) of + undefined -> + {undefined, State}; + #{} -> + CurrentTime = erlang:monotonic_time(millisecond), + {UpdatedWindows, NewIndex} = ensure_windows_rotated(Key, CurrentTime, State), + + %% Update metric state with rotated windows + {ok, MetricState} = maps:find(Name, State#state.metrics), + UpdatedMetricState = maps:put(Key, + #{windows => UpdatedWindows, + current_index => NewIndex}, + MetricState), + UpdatedMetrics = maps:put(Name, UpdatedMetricState, State#state.metrics), + NewState = State#state{metrics = UpdatedMetrics}, + + %% Filter windows to only include those within the last 60 seconds + CutoffTime = CurrentTime - ?TOTAL_WINDOW_MS, + ActiveWindows = [Window || {Window, StartTime} <- UpdatedWindows, + StartTime >= CutoffTime], + + %% Merge all active windows + Merged = lists:foldl(fun(Window, Acc) -> + case Acc of + undefined -> Window; + _ -> ddskerl_std:merge(Acc, Window) + end + end, undefined, ActiveWindows), + + Result = case Merged of + undefined -> + undefined; + _ -> + Count = ddskerl_std:total(Merged), + Sum = ddskerl_std:sum(Merged), + Quantiles = [{Q, ddskerl_std:quantile(Merged, Q)} || Q <- default_quantiles()], + {Count, Sum, Quantiles} + end, + {Result, NewState} + end. + +-spec do_remove(name(), label_values(), #state{}) -> #state{}. +do_remove(Name, LabelValues, State) -> + Key = {Name, LabelValues}, + case maps:find(Name, State#state.metrics) of + error -> + State; + {ok, MetricState} -> + UpdatedMetricState = maps:remove(Key, MetricState), + if + map_size(UpdatedMetricState) == 0 -> + UpdatedMetrics = maps:remove(Name, State#state.metrics), + State#state{metrics = UpdatedMetrics}; + true -> + UpdatedMetrics = maps:put(Name, UpdatedMetricState, State#state.metrics), + State#state{metrics = UpdatedMetrics} + end + end. + +-spec get_metric_state({name(), label_values()}, #state{}) -> undefined | metric_state(). +get_metric_state({Name, LabelValues}, State) -> + case State#state.metrics of + #{Name := MetricState} -> + maps:get({Name, LabelValues}, MetricState, undefined); + #{} -> + undefined + end. + +-spec set_nth(pos_integer(), T, [T]) -> [T]. +set_nth(1, Value, [_ | Rest]) -> + [Value | Rest]; +set_nth(N, Value, [H | Rest]) when N > 1 -> + [H | set_nth(N - 1, Value, Rest)]. From 657a0fafe762f122c6537846624d1e1f8edbc906 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micha=C5=82=20Dobranowski?= Date: Wed, 17 Dec 2025 12:17:49 +0100 Subject: [PATCH 09/24] Add ddskerl to dialyzer `plt_extra_apps` --- rebar.config | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/rebar.config b/rebar.config index 2993aa3fe21..13b971ed25c 100644 --- a/rebar.config +++ b/rebar.config @@ -219,7 +219,7 @@ {plt_extra_apps, [jid, cowboy, cowlib, lasse, p1_utils, ranch, gen_fsm_compat, epgsql, cqerl, common_test, tools, amqp_client, jiffy, erl_csv, inets, compiler, jsx, rabbit_common, mysql, eredis, erlcloud, telemetry, - nksip, nklib, nkservice, nkpacket, prometheus, exometer_core, + nksip, nklib, nkservice, nkpacket, ddskerl, prometheus, exometer_core, mnesia, cets, cpool, tirerl, erlang_doctor]}]}. {cover_print_enabled, true}. From 36bff41b759331684f478e1541f2da63ae4726fd Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micha=C5=82=20Dobranowski?= Date: Thu, 18 Dec 2025 13:22:39 +0100 Subject: [PATCH 10/24] Create and register new prometheus collector --- .../mongoose_instrument_prometheus.erl | 1 + .../mongoose_prometheus_sliding_window.erl | 60 +++++++++++++++---- ...se_prometheus_sliding_window_collector.erl | 57 ++++++++++++++++++ 3 files changed, 107 insertions(+), 11 deletions(-) create mode 100644 src/instrument/mongoose_prometheus_sliding_window_collector.erl diff --git a/src/instrument/mongoose_instrument_prometheus.erl b/src/instrument/mongoose_instrument_prometheus.erl index a125c14fd59..002bc62a3b7 100644 --- a/src/instrument/mongoose_instrument_prometheus.erl +++ b/src/instrument/mongoose_instrument_prometheus.erl @@ -20,6 +20,7 @@ start(#{}) -> _ -> ok end, + prometheus_registry:register_collector(mongoose_prometheus_sliding_window_collector), ok. -spec set_up(mongoose_instrument:event_name(), mongoose_instrument:labels(), diff --git a/src/instrument/mongoose_prometheus_sliding_window.erl b/src/instrument/mongoose_prometheus_sliding_window.erl index d317baeefdb..4205c6cb4ad 100644 --- a/src/instrument/mongoose_prometheus_sliding_window.erl +++ b/src/instrument/mongoose_prometheus_sliding_window.erl @@ -5,8 +5,10 @@ -export([declare/1, observe/3, - value/2, - remove/2]). + values/1, + remove/2, + get_all_metric_names/0, + get_metric_spec/1]). -export([start_link/0, init/1, @@ -32,7 +34,8 @@ default_error() -> 0.01. -record(state, { timer_ref :: reference() | undefined, - metrics = #{} :: #{name() => #{label_values() => metric_state()}} + metrics = #{} :: #{name() => #{label_values() => metric_state()}}, + metric_specs = #{} :: #{name() => proplists:proplist()} }). -type name() :: string(). @@ -51,15 +54,22 @@ declare(MetricSpec) -> observe(Name, LabelValues, Value) -> gen_server:cast(?MODULE, {observe, Name, LabelValues, Value}). --spec value(name(), label_values()) -> - undefined | {non_neg_integer(), number(), [{number(), number()}]}. -value(Name, LabelValues) -> - gen_server:call(?MODULE, {value, Name, LabelValues}). +-spec values(name()) -> [{label_values(), {non_neg_integer(), number(), [{number(), number()}]}}]. +values(Name) -> + gen_server:call(?MODULE, {values, Name}). -spec remove(name(), label_values()) -> boolean(). remove(Name, LabelValues) -> gen_server:call(?MODULE, {remove, Name, LabelValues}). +-spec get_all_metric_names() -> [name()]. +get_all_metric_names() -> + gen_server:call(?MODULE, get_all_metric_names). + +-spec get_metric_spec(name()) -> proplists:proplist() | undefined. +get_metric_spec(Name) -> + gen_server:call(?MODULE, {get_metric_spec, Name}). + %% gen_server callbacks start_link() -> @@ -69,27 +79,36 @@ init([]) -> TimerRef = erlang:start_timer(?WINDOW_DURATION_MS, self(), rotate), {ok, #state{timer_ref = TimerRef}}. -handle_call({declare, Name, _MetricSpec}, _From, State) -> +handle_call({declare, Name, MetricSpec}, _From, State) -> case maps:is_key(Name, State#state.metrics) of true -> {reply, false, State}; false -> NewMetrics = maps:put(Name, #{}, State#state.metrics), - {reply, true, State#state{metrics = NewMetrics}} + NewSpecs = maps:put(Name, MetricSpec, State#state.metric_specs), + {reply, true, State#state{metrics = NewMetrics, metric_specs = NewSpecs}} end; handle_call({observe, Name, LabelValues, Value}, _From, State) -> NewState = do_observe(Name, LabelValues, Value, State), {reply, ok, NewState}; -handle_call({value, Name, LabelValues}, _From, State) -> - {Result, NewState} = get_value(Name, LabelValues, State), +handle_call({values, Name}, _From, State) -> + {Result, NewState} = get_values(Name, State), {reply, Result, NewState}; handle_call({remove, Name, LabelValues}, _From, State) -> NewState = do_remove(Name, LabelValues, State), {reply, true, NewState}; +handle_call(get_all_metric_names, _From, State) -> + Names = maps:keys(State#state.metrics), + {reply, Names, State}; + +handle_call({get_metric_spec, Name}, _From, State) -> + Spec = maps:get(Name, State#state.metric_specs, undefined), + {reply, Spec, State}; + handle_call(_Request, _From, State) -> {reply, ok, State}. @@ -209,6 +228,25 @@ get_value(Name, LabelValues, State) -> {Result, NewState} end. +-spec get_values(name(), #state{}) -> {[{label_values(), {non_neg_integer(), number(), [{number(), number()}]}}], #state{}}. +get_values(Name, State) -> + case maps:find(Name, State#state.metrics) of + error -> + {[], State}; + {ok, MetricState} -> + Keys = maps:keys(MetricState), + {ResultsRev, FinalState} = lists:foldl( + fun({_, LabelValues}, {AccRes, S}) -> + {Res, S2} = get_value(Name, LabelValues, S), + NewAcc = case Res of + undefined -> AccRes; + _ -> [{LabelValues, Res} | AccRes] + end, + {NewAcc, S2} + end, {[], State}, Keys), + {lists:reverse(ResultsRev), FinalState} + end. + -spec do_remove(name(), label_values(), #state{}) -> #state{}. do_remove(Name, LabelValues, State) -> Key = {Name, LabelValues}, diff --git a/src/instrument/mongoose_prometheus_sliding_window_collector.erl b/src/instrument/mongoose_prometheus_sliding_window_collector.erl new file mode 100644 index 00000000000..47d747027b5 --- /dev/null +++ b/src/instrument/mongoose_prometheus_sliding_window_collector.erl @@ -0,0 +1,57 @@ +-module(mongoose_prometheus_sliding_window_collector). + +-behaviour(prometheus_collector). + +-export([deregister_cleanup/1, + collect_mf/2]). + +%% Collector behavior callbacks + +deregister_cleanup(_) -> + ok. + +collect_mf(_Registry, Callback) -> + %% Get all metric names from the sliding window manager + MetricNames = mongoose_prometheus_sliding_window:get_all_metric_names(), + lists:foreach(fun(Name) -> + collect_metric_family(Name, Callback) + end, MetricNames), + ok. + +%% Internal functions + +collect_metric_family(Name, Callback) -> + Values = mongoose_prometheus_sliding_window:values(Name), + case Values of + [] -> + ok; + _ -> + SWName = sliding_window_metric_name(Name), + MetricSpec = mongoose_prometheus_sliding_window:get_metric_spec(Name), + Help = case MetricSpec of + undefined -> + "60-second sliding window quantile summary for " ++ Name; + Spec -> + BaseHelp = proplists:get_value(help, Spec, ""), + BaseHelp ++ " (60s sliding window)" + end, + LabelKeys = case MetricSpec of + undefined -> []; + Spec2 -> proplists:get_value(labels, Spec2, []) + end, + %% Convert to Prometheus format + Metrics = [create_summary_metric(LabelKeys, LabelValues, Count, Sum, Quantiles) + || {LabelValues, {Count, Sum, Quantiles}} <- Values], + %% Create and send the metric family + MF = prometheus_model_helpers:create_mf(SWName, Help, summary, Metrics), + Callback(MF) + end. + +sliding_window_metric_name(Name) -> + Name ++ "_60s". + +create_summary_metric(LabelKeys, LabelValues, Count, Sum, Quantiles) -> + %% Convert label keys and values to label pairs format + LabelPairs = lists:zip(LabelKeys, LabelValues), + QuantilePairs = [{Q, V} || {Q, V} <- Quantiles], + prometheus_model_helpers:summary_metric(LabelPairs, Count, Sum, QuantilePairs). From 6f69268f751c62d6c61c29b64ffcc63fed0c0540 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micha=C5=82=20Dobranowski?= Date: Fri, 19 Dec 2025 13:18:00 +0100 Subject: [PATCH 11/24] Make sure that sliding_window is started --- .../mongoose_prometheus_sliding_window.erl | 21 +++++++++++++++++++ 1 file changed, 21 insertions(+) diff --git a/src/instrument/mongoose_prometheus_sliding_window.erl b/src/instrument/mongoose_prometheus_sliding_window.erl index 4205c6cb4ad..c4c6e4805a0 100644 --- a/src/instrument/mongoose_prometheus_sliding_window.erl +++ b/src/instrument/mongoose_prometheus_sliding_window.erl @@ -47,27 +47,33 @@ default_error() -> 0.01. -spec declare(proplists:proplist()) -> boolean(). declare(MetricSpec) -> + ok = ensure_started(), Name = proplists:get_value(name, MetricSpec), gen_server:call(?MODULE, {declare, Name, MetricSpec}). -spec observe(name(), label_values(), number()) -> ok. observe(Name, LabelValues, Value) -> + ok = ensure_started(), gen_server:cast(?MODULE, {observe, Name, LabelValues, Value}). -spec values(name()) -> [{label_values(), {non_neg_integer(), number(), [{number(), number()}]}}]. values(Name) -> + ok = ensure_started(), gen_server:call(?MODULE, {values, Name}). -spec remove(name(), label_values()) -> boolean(). remove(Name, LabelValues) -> + ok = ensure_started(), gen_server:call(?MODULE, {remove, Name, LabelValues}). -spec get_all_metric_names() -> [name()]. get_all_metric_names() -> + ok = ensure_started(), gen_server:call(?MODULE, get_all_metric_names). -spec get_metric_spec(name()) -> proplists:proplist() | undefined. get_metric_spec(Name) -> + ok = ensure_started(), gen_server:call(?MODULE, {get_metric_spec, Name}). %% gen_server callbacks @@ -133,6 +139,21 @@ terminate(_Reason, _State) -> code_change(_OldVsn, State, _Extra) -> {ok, State}. +%% Internal helpers + +-spec ensure_started() -> ok. +ensure_started() -> + case whereis(?MODULE) of + undefined -> + case start_link() of + {ok, _Pid} -> ok; + {error, {already_started, _Pid}} -> ok; + {error, Reason} -> exit(Reason) + end; + _Pid -> + ok + end. + %% Internal functions -spec do_observe(name(), label_values(), number(), #state{}) -> #state{}. From bdd9b9423e6a4f30bd90e863fcbe9f598b67a303 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micha=C5=82=20Dobranowski?= Date: Mon, 22 Dec 2025 13:57:11 +0100 Subject: [PATCH 12/24] Change `ddskerl_std` -> `ddskerl_ets` --- .../mongoose_instrument_prometheus.erl | 2 +- .../mongoose_prometheus_sliding_window.erl | 118 +++++++++++++++--- 2 files changed, 99 insertions(+), 21 deletions(-) diff --git a/src/instrument/mongoose_instrument_prometheus.erl b/src/instrument/mongoose_instrument_prometheus.erl index 002bc62a3b7..a982f32270e 100644 --- a/src/instrument/mongoose_instrument_prometheus.erl +++ b/src/instrument/mongoose_instrument_prometheus.erl @@ -71,7 +71,7 @@ declare_metric(MetricSpec, histogram) -> {error, mongoose_prometheus_sliding_window:default_error()}, %% Measuring in µs suffices for actions lasting up to a day (with 1% accuracy). %% Measuring in bytes suffices for sizes up to 81 GB (with 1% accuracy). - {bound, 1260} + {bound, mongoose_prometheus_sliding_window:default_bound()} | MetricSpec]), SWResult or PromResult. diff --git a/src/instrument/mongoose_prometheus_sliding_window.erl b/src/instrument/mongoose_prometheus_sliding_window.erl index c4c6e4805a0..e9ecc84e949 100644 --- a/src/instrument/mongoose_prometheus_sliding_window.erl +++ b/src/instrument/mongoose_prometheus_sliding_window.erl @@ -18,7 +18,7 @@ terminate/2, code_change/3]). --export([default_quantiles/0, default_error/0]). +-export([default_quantiles/0, default_error/0, default_bound/0]). -behaviour(gen_server). @@ -26,21 +26,29 @@ -define(WINDOW_DURATION_MS, 3000). % 3 seconds per window -define(TOTAL_WINDOW_MS, (?WINDOW_COUNT * ?WINDOW_DURATION_MS)). % 60 seconds total + -spec default_quantiles() -> [number()]. default_quantiles() -> [0.5, 0.75, 0.90, 0.95, 0.99, 0.999]. -spec default_error() -> number(). default_error() -> 0.01. +-spec default_bound() -> non_neg_integer(). +default_bound() -> 1260. + -record(state, { timer_ref :: reference() | undefined, - metrics = #{} :: #{name() => #{label_values() => metric_state()}}, + ets_table :: ets:tab() | undefined, + metrics = #{} :: #{name() => #{{name(), label_values()} => metric_state()}}, metric_specs = #{} :: #{name() => proplists:proplist()} }). -type name() :: string(). -type label_values() :: [mongoose_instrument:label_value()]. --type metric_state() :: #{windows => [{ddskerl_ets:ddsketch(), non_neg_integer()}], +-type window_data() :: #{sketch := ddskerl_ets:ddsketch(), + ref := ets:tab(), + name := term()}. +-type metric_state() :: #{windows => [{window_data(), non_neg_integer()}], current_index => non_neg_integer()}. %% Public API @@ -82,8 +90,9 @@ start_link() -> gen_server:start_link({local, ?MODULE}, ?MODULE, [], []). init([]) -> + EtsTable = ets:new(?MODULE, [set, private, {read_concurrency, true}, {write_concurrency, true}]), TimerRef = erlang:start_timer(?WINDOW_DURATION_MS, self(), rotate), - {ok, #state{timer_ref = TimerRef}}. + {ok, #state{timer_ref = TimerRef, ets_table = EtsTable}}. handle_call({declare, Name, MetricSpec}, _From, State) -> case maps:is_key(Name, State#state.metrics) of @@ -169,7 +178,9 @@ do_observe(Name, LabelValues, Value, State) -> %% Get current window and add observation {CurrentWindow, WindowStartTime} = lists:nth(CurrentIndex + 1, Windows), - UpdatedWindow = ddskerl_std:insert(CurrentWindow, Value), + Sketch = maps:get(sketch, CurrentWindow), + UpdatedSketch = ddskerl_ets:insert(Sketch, Value), + UpdatedWindow = CurrentWindow#{sketch := UpdatedSketch}, %% Update windows list and metric state UpdatedWindows = set_nth(CurrentIndex + 1, {UpdatedWindow, WindowStartTime}, Windows), @@ -183,12 +194,12 @@ do_observe(Name, LabelValues, Value, State) -> end. -spec ensure_windows_rotated({name(), label_values()}, non_neg_integer(), #state{}) -> - {[{ddskerl_std:ddsketch(), non_neg_integer()}], non_neg_integer()}. + {[{window_data(), non_neg_integer()}], non_neg_integer()}. ensure_windows_rotated(Key, CurrentTime, State) -> case get_metric_state(Key, State) of undefined -> %% Initialize new windows, all starting at current time (they'll be rotated as needed) - Windows = [{ddskerl_std:new(#{error => default_error()}), CurrentTime} || _ <- lists:seq(1, ?WINDOW_COUNT)], + Windows = [new_window(Key, Index, CurrentTime, State) || Index <- lists:seq(0, ?WINDOW_COUNT - 1)], {Windows, 0}; #{windows := Windows, current_index := CurrentIndex} -> {_CurrentWindow, WindowStartTime} = lists:nth(CurrentIndex + 1, Windows), @@ -196,8 +207,9 @@ ensure_windows_rotated(Key, CurrentTime, State) -> if Elapsed >= ?WINDOW_DURATION_MS -> NewIndex = (CurrentIndex + 1) rem ?WINDOW_COUNT, - NewWindow = {ddskerl_std:new(#{error => default_error()}), CurrentTime}, - UpdatedWindows = set_nth(CurrentIndex + 1, NewWindow, Windows), + {CurrentWindow, _} = lists:nth(CurrentIndex + 1, Windows), + ResetWindow = reset_window(CurrentWindow), + UpdatedWindows = set_nth(CurrentIndex + 1, {ResetWindow, CurrentTime}, Windows), {UpdatedWindows, NewIndex}; true -> {Windows, CurrentIndex} @@ -230,21 +242,26 @@ get_value(Name, LabelValues, State) -> StartTime >= CutoffTime], %% Merge all active windows - Merged = lists:foldl(fun(Window, Acc) -> - case Acc of - undefined -> Window; - _ -> ddskerl_std:merge(Acc, Window) - end - end, undefined, ActiveWindows), + WindowTuples = [window_tuple(Window) || Window <- ActiveWindows], + Merged = lists:foldl(fun + (undefined, Acc) -> Acc; + (Val, undefined) -> Val; + (Val, Acc) -> ddskerl_ets:merge_tuples(Acc, Val) + end, undefined, WindowTuples), Result = case Merged of undefined -> undefined; _ -> - Count = ddskerl_std:total(Merged), - Sum = ddskerl_std:sum(Merged), - Quantiles = [{Q, ddskerl_std:quantile(Merged, Q)} || Q <- default_quantiles()], - {Count, Sum, Quantiles} + Count = ddskerl_ets:total_tuple(Merged), + case Count of + 0 -> + undefined; + _ -> + Sum = ddskerl_ets:sum_tuple(Merged), + Quantiles = [{Q, ddskerl_ets:quantile_tuple(Merged, Q)} || Q <- default_quantiles()], + {Count, Sum, Quantiles} + end end, {Result, NewState} end. @@ -275,7 +292,7 @@ do_remove(Name, LabelValues, State) -> error -> State; {ok, MetricState} -> - UpdatedMetricState = maps:remove(Key, MetricState), + UpdatedMetricState = remove_metric_state(Key, MetricState), if map_size(UpdatedMetricState) == 0 -> UpdatedMetrics = maps:remove(Name, State#state.metrics), @@ -300,3 +317,64 @@ set_nth(1, Value, [_ | Rest]) -> [Value | Rest]; set_nth(N, Value, [H | Rest]) when N > 1 -> [H | set_nth(N - 1, Value, Rest)]. + +%% Window helpers + +-spec metric_options(name(), #state{}) -> {number(), non_neg_integer()}. +metric_options(Name, State) -> + MetricSpec = maps:get(Name, State#state.metric_specs, []), + Error = proplists:get_value(error, MetricSpec, default_error()), + Bound = proplists:get_value(bound, MetricSpec, default_bound()), + {Error, Bound}. + +-spec window_name({name(), label_values()}, non_neg_integer()) -> term(). +window_name({Name, LabelValues}, Index) -> + {Name, LabelValues, Index}. + +-spec new_window({name(), label_values()}, non_neg_integer(), non_neg_integer(), #state{}) -> + {window_data(), non_neg_integer()}. +new_window(Key = {Name, _LabelValues}, Index, StartTime, State) -> + {Error, Bound} = metric_options(Name, State), + WindowName = window_name(Key, Index), + Sketch = ddskerl_ets:new(#{ets_table => State#state.ets_table, + name => WindowName, + error => Error, + bound => Bound}), + WindowData = #{sketch => Sketch, + ref => State#state.ets_table, + name => WindowName}, + {WindowData, StartTime}. + +-spec reset_window(window_data()) -> window_data(). +reset_window(WindowData) -> + Sketch = maps:get(sketch, WindowData), + ResetSketch = ddskerl_ets:reset(Sketch), + WindowData#{sketch := ResetSketch}. + +-spec window_tuple(window_data()) -> ddskerl_ets:object() | undefined. +window_tuple(WindowData) -> + Ref = maps:get(ref, WindowData), + Name = maps:get(name, WindowData), + case ets:lookup(Ref, Name) of + [Val] -> Val; + [] -> undefined + end. + +-spec remove_metric_state({name(), label_values()}, #{}) -> #{}. +remove_metric_state(Key, MetricState) -> + UpdatedMetricState = + case maps:find(Key, MetricState) of + error -> + MetricState; + {ok, #{windows := Windows}} -> + lists:foreach(fun({Window, _}) -> delete_window(Window) end, Windows), + maps:remove(Key, MetricState) + end, + UpdatedMetricState. + +-spec delete_window(window_data()) -> ok. +delete_window(WindowData) -> + Ref = maps:get(ref, WindowData), + Name = maps:get(name, WindowData), + ets:delete(Ref, Name), + ok. From 3400ad4d7861fb85c90e6c7d5267e0e3437e1388 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micha=C5=82=20Dobranowski?= Date: Tue, 23 Dec 2025 13:54:21 +0100 Subject: [PATCH 13/24] Upgrade ddskerl (and prometheus) --- rebar.config | 4 ++-- rebar.lock | 12 ++++++------ 2 files changed, 8 insertions(+), 8 deletions(-) diff --git a/rebar.config b/rebar.config index 13b971ed25c..c597110540e 100644 --- a/rebar.config +++ b/rebar.config @@ -56,8 +56,8 @@ {exometer_report_statsd, {git, "https://github.com/esl/exometer_report_statsd.git", {branch, "master"}}}, {syslogger, "0.3.0"}, {flatlog, "0.1.2"}, - {ddskerl, "0.4.2"}, - {prometheus, "6.1.1"}, + {ddskerl, "0.4.3"}, + {prometheus, "6.1.2"}, {prometheus_cowboy, "0.2.0"}, %%% Stateless libraries diff --git a/rebar.lock b/rebar.lock index 98dbee6ce43..68bac0aae92 100644 --- a/rebar.lock +++ b/rebar.lock @@ -19,7 +19,7 @@ {<<"credentials_obfuscation">>, {pkg,<<"credentials_obfuscation">>,<<"3.5.0">>}, 1}, - {<<"ddskerl">>,{pkg,<<"ddskerl">>,<<"0.4.2">>},0}, + {<<"ddskerl">>,{pkg,<<"ddskerl">>,<<"0.4.3">>},0}, {<<"eini">>,{pkg,<<"eini">>,<<"1.2.9">>},1}, {<<"epgsql">>,{pkg,<<"epgsql">>,<<"4.7.1">>},0}, {<<"eredis">>,{pkg,<<"eredis">>,<<"1.7.1">>},0}, @@ -87,7 +87,7 @@ 0}, {<<"parse_trans">>,{pkg,<<"parse_trans">>,<<"3.4.0">>},1}, {<<"pooler">>,{pkg,<<"pooler">>,<<"1.5.3">>},1}, - {<<"prometheus">>,{pkg,<<"prometheus">>,<<"6.1.1">>},0}, + {<<"prometheus">>,{pkg,<<"prometheus">>,<<"6.1.2">>},0}, {<<"prometheus_cowboy">>,{pkg,<<"prometheus_cowboy">>,<<"0.2.0">>},0}, {<<"prometheus_httpd">>,{pkg,<<"prometheus_httpd">>,<<"2.1.15">>},1}, {<<"quickrand">>,{pkg,<<"quickrand">>,<<"2.0.7">>},1}, @@ -131,7 +131,7 @@ {<<"cowlib">>, <<"623791C56C1CC9DF54A71A9C55147A401549917F00A2E48A6AE12B812C586CED">>}, {<<"cpool">>, <<"76222AA1DAC0F8089707167BD69D221EDA63DB65B8BD67DBF6E133075392EEDC">>}, {<<"credentials_obfuscation">>, <<"61E282ADFB4439486B3994FAAEC69543C7EE6CC7E70C6340E8853FD9DEAF8219">>}, - {<<"ddskerl">>, <<"A51A90BE9AC9B36A94017670BED479C623B10CA9D4BDA1EDF3A0E48CAEEADA2A">>}, + {<<"ddskerl">>, <<"BB97B90EEF1B5906520CEBDD38DA29B18A770B5567F6297927F53566EAE9BEBB">>}, {<<"eini">>, <<"FCC3CBD49BBDD9A1D9735C7365DAFFCD84481CCE81E6CB80537883AA44AC4895">>}, {<<"epgsql">>, <<"D4E47CAE46C18C8AFA88E34D59A9B4BAE16368D7CE1EB3DA24FA755EB28393EB">>}, {<<"eredis">>, <<"39E31AA02ADCD651C657F39AAFD4D31A9B2F63C6C700DC9CECE98D4BC3C897AB">>}, @@ -163,7 +163,7 @@ {<<"p1_utils">>, <<"67B0C4AC9FA3BA3EF563B31AA111B0A004439A37FAC85E027F1C3617E1C7EC6C">>}, {<<"parse_trans">>, <<"BB87AC362A03CA674EBB7D9D498F45C03256ADED7214C9101F7035EF44B798C7">>}, {<<"pooler">>, <<"898CD1FA301FC42D4A8ED598CE139B71CA85B54C16AB161152B5CC5FBDCFA1A8">>}, - {<<"prometheus">>, <<"3C9D8E2D4FCF948450550693F6A82FAE013D3665BA10AA55955B64C2875FADB3">>}, + {<<"prometheus">>, <<"E5C3C567CD8B0994425920763405635211183C15BDC47D0CD524807DDE20CA1D">>}, {<<"prometheus_cowboy">>, <<"526F75D9850A9125496F78BCEECCA0F237BC7B403C976D44508543AE5967DAD9">>}, {<<"prometheus_httpd">>, <<"8F767D819A5D36275EAB9264AFF40D87279151646776069BF69FBDBBD562BD75">>}, {<<"quickrand">>, <<"D2BD76676A446E6A058D678444B7FDA1387B813710D1AF6D6E29BB92186C8820">>}, @@ -198,7 +198,7 @@ {<<"cowlib">>, <<"0AF652D1550C8411C3B58EED7A035A7FB088C0B86AFF6BC504B0BC3B7F791AA2">>}, {<<"cpool">>, <<"430E18DF4A9D584EB1ED0D196A87CC02E878AF5B4888BFDC9B65F86A96480E30">>}, {<<"credentials_obfuscation">>, <<"843ADBE3246861CE0F1A0FA3222F384834EB31DEFD8D6B9CBA7AFD2977C957BC">>}, - {<<"ddskerl">>, <<"63F907373D7E548151D584D4DA8A38928FD26EC9477B94C0FFAAD87D7CB69FE7">>}, + {<<"ddskerl">>, <<"4E0F6047C6A002CE38F6EA155276DD918CD635FD0A5EDB9E0B46AEB6F7AFF2C2">>}, {<<"eini">>, <<"DA64AE8DB7C2F502E6F20CDF44CD3D9BE364412B87FF49FEBF282540F673DFCB">>}, {<<"epgsql">>, <<"B6D86B7DC42C8555B1D4E20880E5099D6D6D053148000E188E548F98E4E01836">>}, {<<"eredis">>, <<"7C2B54C566FED55FEEF3341CA79B0100A6348FD3F162184B7ED5118D258C3CC1">>}, @@ -230,7 +230,7 @@ {<<"p1_utils">>, <<"D0379E8C1156B98BD64F8129C1DE022FCCA4F2FDB7486CE73BF0ED2C3376B04C">>}, {<<"parse_trans">>, <<"F99E368830BEA44552224E37E04943A54874F08B8590485DE8D13832B63A2DC3">>}, {<<"pooler">>, <<"058D85C5081289B90E97E4DDDBC3BB5A3B4A19A728AB3BC88C689EFCC36A07C7">>}, - {<<"prometheus">>, <<"C3ECC7A35676948089E28FF4383CCB85DCEF447F0982C01E5630A084008C456A">>}, + {<<"prometheus">>, <<"109662001328112B860DE112D49B575310BE8FF33FE04BF0482EEC4B1E5A1278">>}, {<<"prometheus_cowboy">>, <<"2C7EB12F4B970D91E3B47BAAD0F138F6ADC34E53EEB0AE18068FF0AFAB441B24">>}, {<<"prometheus_httpd">>, <<"67736D000745184D5013C58A63E947821AB90CB9320BC2E6AE5D3061C6FFE039">>}, {<<"quickrand">>, <<"B8ACBF89A224BC217C3070CA8BEBC6EB236DBE7F9767993B274084EA044D35F0">>}, From 92b5da291a7911746682e758feda90b65eef3c13 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micha=C5=82=20Dobranowski?= Date: Fri, 16 Jan 2026 10:34:08 +0100 Subject: [PATCH 14/24] Fix format and remove unused function --- src/instrument/mongoose_prometheus_sliding_window.erl | 11 ----------- 1 file changed, 11 deletions(-) diff --git a/src/instrument/mongoose_prometheus_sliding_window.erl b/src/instrument/mongoose_prometheus_sliding_window.erl index e9ecc84e949..1b5ae252e80 100644 --- a/src/instrument/mongoose_prometheus_sliding_window.erl +++ b/src/instrument/mongoose_prometheus_sliding_window.erl @@ -103,34 +103,24 @@ handle_call({declare, Name, MetricSpec}, _From, State) -> NewSpecs = maps:put(Name, MetricSpec, State#state.metric_specs), {reply, true, State#state{metrics = NewMetrics, metric_specs = NewSpecs}} end; - -handle_call({observe, Name, LabelValues, Value}, _From, State) -> - NewState = do_observe(Name, LabelValues, Value, State), - {reply, ok, NewState}; - handle_call({values, Name}, _From, State) -> {Result, NewState} = get_values(Name, State), {reply, Result, NewState}; - handle_call({remove, Name, LabelValues}, _From, State) -> NewState = do_remove(Name, LabelValues, State), {reply, true, NewState}; - handle_call(get_all_metric_names, _From, State) -> Names = maps:keys(State#state.metrics), {reply, Names, State}; - handle_call({get_metric_spec, Name}, _From, State) -> Spec = maps:get(Name, State#state.metric_specs, undefined), {reply, Spec, State}; - handle_call(_Request, _From, State) -> {reply, ok, State}. handle_cast({observe, Name, LabelValues, Value}, State) -> NewState = do_observe(Name, LabelValues, Value, State), {noreply, NewState}; - handle_cast(_Msg, State) -> {noreply, State}. @@ -138,7 +128,6 @@ handle_info({timeout, _TimerRef, rotate}, State) -> %% The rotation is handled in ensure_windows_rotated, so we just reschedule TimerRef = erlang:start_timer(?WINDOW_DURATION_MS, self(), rotate), {noreply, State#state{timer_ref = TimerRef}}; - handle_info(_Info, State) -> {noreply, State}. From 8c7223bcfc44a623900a37b350ad89c5b21198e1 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micha=C5=82=20Dobranowski?= Date: Wed, 21 Jan 2026 11:31:40 +0100 Subject: [PATCH 15/24] Add sliding window to the supervision tree --- src/ejabberd_sup.erl | 2 ++ .../mongoose_instrument_prometheus.erl | 7 ---- .../mongoose_prometheus_sliding_window.erl | 33 ++++++++----------- test/mongoose_instrument_metrics_SUITE.erl | 7 +++- test/prometheus_endpoint_SUITE.erl | 3 +- 5 files changed, 23 insertions(+), 29 deletions(-) diff --git a/src/ejabberd_sup.erl b/src/ejabberd_sup.erl index 7750ed1103f..87662fedb51 100644 --- a/src/ejabberd_sup.erl +++ b/src/ejabberd_sup.erl @@ -45,6 +45,7 @@ start_link() -> init(noargs) -> Hooks = worker_spec(gen_hook), Instrument = worker_spec(mongoose_instrument), + PrometheusSlidingWindow = mongoose_prometheus_sliding_window:child_spec(), Cleaner = worker_spec(mongoose_cleaner), Router = worker_spec(ejabberd_router), S2S = worker_spec(ejabberd_s2s), @@ -70,6 +71,7 @@ init(noargs) -> PG, Hooks, Instrument, + PrometheusSlidingWindow, Cleaner, SMBackendSupervisor, OutgoingPoolsSupervisor diff --git a/src/instrument/mongoose_instrument_prometheus.erl b/src/instrument/mongoose_instrument_prometheus.erl index a982f32270e..7d29073903c 100644 --- a/src/instrument/mongoose_instrument_prometheus.erl +++ b/src/instrument/mongoose_instrument_prometheus.erl @@ -13,13 +13,6 @@ start(#{}) -> Apps = [prometheus, prometheus_httpd, prometheus_cowboy], {ok, _} = application:ensure_all_started(Apps, permanent), - %% Start sliding window manager for histogram metrics - case whereis(mongoose_prometheus_sliding_window) of - undefined -> - {ok, _} = mongoose_prometheus_sliding_window:start_link(); - _ -> - ok - end, prometheus_registry:register_collector(mongoose_prometheus_sliding_window_collector), ok. diff --git a/src/instrument/mongoose_prometheus_sliding_window.erl b/src/instrument/mongoose_prometheus_sliding_window.erl index 1b5ae252e80..fdb68b66c55 100644 --- a/src/instrument/mongoose_prometheus_sliding_window.erl +++ b/src/instrument/mongoose_prometheus_sliding_window.erl @@ -3,7 +3,8 @@ %% Manages a sliding window of 20 DDSketch instances (3 seconds each = 60 seconds total) %% Each window uses ddskerl to store quantile summaries --export([declare/1, +-export([child_spec/0, + declare/1, observe/3, values/1, remove/2, @@ -20,6 +21,8 @@ -export([default_quantiles/0, default_error/0, default_bound/0]). +-ignore_xref([start_link/0]). + -behaviour(gen_server). -define(WINDOW_COUNT, 20). @@ -55,33 +58,27 @@ default_bound() -> 1260. -spec declare(proplists:proplist()) -> boolean(). declare(MetricSpec) -> - ok = ensure_started(), Name = proplists:get_value(name, MetricSpec), gen_server:call(?MODULE, {declare, Name, MetricSpec}). -spec observe(name(), label_values(), number()) -> ok. observe(Name, LabelValues, Value) -> - ok = ensure_started(), gen_server:cast(?MODULE, {observe, Name, LabelValues, Value}). -spec values(name()) -> [{label_values(), {non_neg_integer(), number(), [{number(), number()}]}}]. values(Name) -> - ok = ensure_started(), gen_server:call(?MODULE, {values, Name}). -spec remove(name(), label_values()) -> boolean(). remove(Name, LabelValues) -> - ok = ensure_started(), gen_server:call(?MODULE, {remove, Name, LabelValues}). -spec get_all_metric_names() -> [name()]. get_all_metric_names() -> - ok = ensure_started(), gen_server:call(?MODULE, get_all_metric_names). -spec get_metric_spec(name()) -> proplists:proplist() | undefined. get_metric_spec(Name) -> - ok = ensure_started(), gen_server:call(?MODULE, {get_metric_spec, Name}). %% gen_server callbacks @@ -137,20 +134,16 @@ terminate(_Reason, _State) -> code_change(_OldVsn, State, _Extra) -> {ok, State}. -%% Internal helpers +%% Child spec for supervision tree --spec ensure_started() -> ok. -ensure_started() -> - case whereis(?MODULE) of - undefined -> - case start_link() of - {ok, _Pid} -> ok; - {error, {already_started, _Pid}} -> ok; - {error, Reason} -> exit(Reason) - end; - _Pid -> - ok - end. +-spec child_spec() -> supervisor:child_spec(). +child_spec() -> + #{id => ?MODULE, + start => {?MODULE, start_link, []}, + restart => permanent, + shutdown => timer:seconds(5), + type => worker, + modules => [?MODULE]}. %% Internal functions diff --git a/test/mongoose_instrument_metrics_SUITE.erl b/test/mongoose_instrument_metrics_SUITE.erl index eaa28c3a707..bd8931972c1 100644 --- a/test/mongoose_instrument_metrics_SUITE.erl +++ b/test/mongoose_instrument_metrics_SUITE.erl @@ -80,10 +80,15 @@ init_per_group(Group, Config) -> mongoose_config:set_opts(#{hosts => [?HOST_TYPE], host_types => [?HOST_TYPE2], instrumentation => opts(Group)}), - Config1 = async_helper:start(Config, mongoose_instrument, start_link, []), + Config1 = async_helper:start(Config, extra_processes(Group) ++ + [{mongoose_instrument, start_link, []}]), mongoose_instrument:persist(), Config1 ++ extra_config(Group). +extra_processes(prometheus) -> [{mongoose_prometheus_sliding_window, start_link, []}]; +extra_processes(prometheus_and_exometer) -> [{mongoose_prometheus_sliding_window, start_link, []}]; +extra_processes(_) -> []. + end_per_group(Group, Config) -> async_helper:stop_all(Config), mongoose_config:erase_opts(), diff --git a/test/prometheus_endpoint_SUITE.erl b/test/prometheus_endpoint_SUITE.erl index b063481e1e7..1e589a8ed19 100644 --- a/test/prometheus_endpoint_SUITE.erl +++ b/test/prometheus_endpoint_SUITE.erl @@ -23,7 +23,8 @@ init_per_suite(Config) -> mongoose_config:set_opts(opts()), Config1 = async_helper:start(Config, [{mongoose_instrument, start_link, []}, {mim_ct_sup, start_link, [ejabberd_sup]}, - {mongoose_listener_sup, start_link, []}]), + {mongoose_listener_sup, start_link, []}, + {mongoose_prometheus_sliding_window, start_link, []}]), mongoose_listener:start(), Config1. From 185e948674510c9f64687130ba9b7be93c8cb147 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micha=C5=82=20Dobranowski?= Date: Wed, 21 Jan 2026 15:25:05 +0100 Subject: [PATCH 16/24] Make window parameters configurable --- .../mongoose_prometheus_sliding_window.erl | 26 +++++++++++-------- ...se_prometheus_sliding_window_collector.erl | 9 ++++--- 2 files changed, 21 insertions(+), 14 deletions(-) diff --git a/src/instrument/mongoose_prometheus_sliding_window.erl b/src/instrument/mongoose_prometheus_sliding_window.erl index fdb68b66c55..59c04ea3604 100644 --- a/src/instrument/mongoose_prometheus_sliding_window.erl +++ b/src/instrument/mongoose_prometheus_sliding_window.erl @@ -1,7 +1,7 @@ -module(mongoose_prometheus_sliding_window). -%% Manages a sliding window of 20 DDSketch instances (3 seconds each = 60 seconds total) -%% Each window uses ddskerl to store quantile summaries +%% Manages a sliding window using DDSketch instances for quantile summaries. +%% Window parameters are configurable via WINDOW_SIZE_MS and WINDOW_STEP_MS constants. -export([child_spec/0, declare/1, @@ -19,15 +19,16 @@ terminate/2, code_change/3]). --export([default_quantiles/0, default_error/0, default_bound/0]). +-export([default_quantiles/0, default_error/0, default_bound/0, window_size_s/0]). -ignore_xref([start_link/0]). -behaviour(gen_server). --define(WINDOW_COUNT, 20). --define(WINDOW_DURATION_MS, 3000). % 3 seconds per window --define(TOTAL_WINDOW_MS, (?WINDOW_COUNT * ?WINDOW_DURATION_MS)). % 60 seconds total +%% Configurable window parameters +-define(WINDOW_SIZE_MS, 60000). % Total window size: 60 seconds +-define(WINDOW_STEP_MS, 3000). % Step size: 3 seconds per sub-window +-define(WINDOW_COUNT, (?WINDOW_SIZE_MS div ?WINDOW_STEP_MS)). % Number of sub-windows -spec default_quantiles() -> [number()]. @@ -39,6 +40,9 @@ default_error() -> 0.01. -spec default_bound() -> non_neg_integer(). default_bound() -> 1260. +-spec window_size_s() -> pos_integer(). +window_size_s() -> ?WINDOW_SIZE_MS div 1000. + -record(state, { timer_ref :: reference() | undefined, ets_table :: ets:tab() | undefined, @@ -88,7 +92,7 @@ start_link() -> init([]) -> EtsTable = ets:new(?MODULE, [set, private, {read_concurrency, true}, {write_concurrency, true}]), - TimerRef = erlang:start_timer(?WINDOW_DURATION_MS, self(), rotate), + TimerRef = erlang:start_timer(?WINDOW_STEP_MS, self(), rotate), {ok, #state{timer_ref = TimerRef, ets_table = EtsTable}}. handle_call({declare, Name, MetricSpec}, _From, State) -> @@ -123,7 +127,7 @@ handle_cast(_Msg, State) -> handle_info({timeout, _TimerRef, rotate}, State) -> %% The rotation is handled in ensure_windows_rotated, so we just reschedule - TimerRef = erlang:start_timer(?WINDOW_DURATION_MS, self(), rotate), + TimerRef = erlang:start_timer(?WINDOW_STEP_MS, self(), rotate), {noreply, State#state{timer_ref = TimerRef}}; handle_info(_Info, State) -> {noreply, State}. @@ -187,7 +191,7 @@ ensure_windows_rotated(Key, CurrentTime, State) -> {_CurrentWindow, WindowStartTime} = lists:nth(CurrentIndex + 1, Windows), Elapsed = CurrentTime - WindowStartTime, if - Elapsed >= ?WINDOW_DURATION_MS -> + Elapsed >= ?WINDOW_STEP_MS -> NewIndex = (CurrentIndex + 1) rem ?WINDOW_COUNT, {CurrentWindow, _} = lists:nth(CurrentIndex + 1, Windows), ResetWindow = reset_window(CurrentWindow), @@ -218,8 +222,8 @@ get_value(Name, LabelValues, State) -> UpdatedMetrics = maps:put(Name, UpdatedMetricState, State#state.metrics), NewState = State#state{metrics = UpdatedMetrics}, - %% Filter windows to only include those within the last 60 seconds - CutoffTime = CurrentTime - ?TOTAL_WINDOW_MS, + %% Filter windows to only include those within the last window size + CutoffTime = CurrentTime - ?WINDOW_SIZE_MS, ActiveWindows = [Window || {Window, StartTime} <- UpdatedWindows, StartTime >= CutoffTime], diff --git a/src/instrument/mongoose_prometheus_sliding_window_collector.erl b/src/instrument/mongoose_prometheus_sliding_window_collector.erl index 47d747027b5..fc759c2eafb 100644 --- a/src/instrument/mongoose_prometheus_sliding_window_collector.erl +++ b/src/instrument/mongoose_prometheus_sliding_window_collector.erl @@ -28,12 +28,14 @@ collect_metric_family(Name, Callback) -> _ -> SWName = sliding_window_metric_name(Name), MetricSpec = mongoose_prometheus_sliding_window:get_metric_spec(Name), + WindowSizeS = mongoose_prometheus_sliding_window:window_size_s(), + WindowSizeSStr = integer_to_list(WindowSizeS), Help = case MetricSpec of undefined -> - "60-second sliding window quantile summary for " ++ Name; + WindowSizeSStr ++ "-second sliding window quantile summary for " ++ Name; Spec -> BaseHelp = proplists:get_value(help, Spec, ""), - BaseHelp ++ " (60s sliding window)" + BaseHelp ++ " (" ++ WindowSizeSStr ++ "s sliding window)" end, LabelKeys = case MetricSpec of undefined -> []; @@ -48,7 +50,8 @@ collect_metric_family(Name, Callback) -> end. sliding_window_metric_name(Name) -> - Name ++ "_60s". + WindowSizeS = mongoose_prometheus_sliding_window:window_size_s(), + Name ++ "_" ++ integer_to_list(WindowSizeS) ++ "s". create_summary_metric(LabelKeys, LabelValues, Count, Sum, Quantiles) -> %% Convert label keys and values to label pairs format From 1141781c56fe930c34500449a4636f0ba6d07c3e Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micha=C5=82=20Dobranowski?= Date: Thu, 22 Jan 2026 14:15:28 +0100 Subject: [PATCH 17/24] Remove old quantile_summary metrics use only sliding_window; fix `do_remove` function --- .../mongoose_instrument_prometheus.erl | 16 ++----- .../mongoose_prometheus_sliding_window.erl | 43 ++++++++++++------- ...se_prometheus_sliding_window_collector.erl | 13 ++---- test/mongoose_instrument_metrics_SUITE.erl | 30 ++++++------- 4 files changed, 49 insertions(+), 53 deletions(-) diff --git a/src/instrument/mongoose_instrument_prometheus.erl b/src/instrument/mongoose_instrument_prometheus.erl index 7d29073903c..6e74991cebf 100644 --- a/src/instrument/mongoose_instrument_prometheus.erl +++ b/src/instrument/mongoose_instrument_prometheus.erl @@ -58,15 +58,7 @@ declare_metric(MetricSpec, counter) -> declare_metric(MetricSpec, spiral) -> prometheus_counter:declare(MetricSpec); declare_metric(MetricSpec, histogram) -> - SWResult = mongoose_prometheus_sliding_window:declare(MetricSpec), - PromResult = prometheus_quantile_summary:declare([ - {quantiles, mongoose_prometheus_sliding_window:default_quantiles()}, - {error, mongoose_prometheus_sliding_window:default_error()}, - %% Measuring in µs suffices for actions lasting up to a day (with 1% accuracy). - %% Measuring in bytes suffices for sizes up to 81 GB (with 1% accuracy). - {bound, mongoose_prometheus_sliding_window:default_bound()} - | MetricSpec]), - SWResult or PromResult. + mongoose_prometheus_sliding_window:declare(MetricSpec). -spec reset_metric(name(), [mongoose_instrument:label_value()], mongoose_instrument:metric_type()) -> boolean(). @@ -77,8 +69,7 @@ reset_metric(Name, LabelValues, counter) -> reset_metric(Name, LabelValues, spiral) -> prometheus_counter:remove(Name, LabelValues); reset_metric(Name, LabelValues, histogram) -> - mongoose_prometheus_sliding_window:remove(Name, LabelValues), - prometheus_quantile_summary:remove(Name, LabelValues). + mongoose_prometheus_sliding_window:remove(Name, LabelValues). -spec initialize_metric(name(), [mongoose_instrument:label_value()], mongoose_instrument:metric_type()) -> ok. @@ -140,5 +131,4 @@ update_metric(Name, Labels, counter, Value) when is_integer(Value) -> update_metric(Name, Labels, spiral, Value) when is_integer(Value), Value >= 0 -> ok = prometheus_counter:inc(Name, Labels, Value); update_metric(Name, Labels, histogram, Value) when is_integer(Value) -> - ok = mongoose_prometheus_sliding_window:observe(Name, Labels, Value), - ok = prometheus_quantile_summary:observe(Name, Labels, Value). + ok = mongoose_prometheus_sliding_window:observe(Name, Labels, Value). diff --git a/src/instrument/mongoose_prometheus_sliding_window.erl b/src/instrument/mongoose_prometheus_sliding_window.erl index 59c04ea3604..943bf84c822 100644 --- a/src/instrument/mongoose_prometheus_sliding_window.erl +++ b/src/instrument/mongoose_prometheus_sliding_window.erl @@ -6,6 +6,7 @@ -export([child_spec/0, declare/1, observe/3, + value/2, values/1, remove/2, get_all_metric_names/0, @@ -19,9 +20,7 @@ terminate/2, code_change/3]). --export([default_quantiles/0, default_error/0, default_bound/0, window_size_s/0]). - --ignore_xref([start_link/0]). +-ignore_xref([start_link/0, value/2]). -behaviour(gen_server). @@ -40,9 +39,6 @@ default_error() -> 0.01. -spec default_bound() -> non_neg_integer(). default_bound() -> 1260. --spec window_size_s() -> pos_integer(). -window_size_s() -> ?WINDOW_SIZE_MS div 1000. - -record(state, { timer_ref :: reference() | undefined, ets_table :: ets:tab() | undefined, @@ -69,9 +65,17 @@ declare(MetricSpec) -> observe(Name, LabelValues, Value) -> gen_server:cast(?MODULE, {observe, Name, LabelValues, Value}). +-spec value(name(), label_values()) -> undefined | {non_neg_integer(), number(), [{number(), number()}]}. +value(Name, LabelValues) -> + gen_server:call(?MODULE, {value, Name, LabelValues}). + -spec values(name()) -> [{label_values(), {non_neg_integer(), number(), [{number(), number()}]}}]. values(Name) -> - gen_server:call(?MODULE, {values, Name}). + [{LV, V} || LV <- get_label_values(Name), V <- [value(Name, LV)], V =/= undefined]. + +-spec get_label_values(name()) -> [label_values()]. +get_label_values(Name) -> + gen_server:call(?MODULE, {get_label_values, Name}). -spec remove(name(), label_values()) -> boolean(). remove(Name, LabelValues) -> @@ -104,6 +108,12 @@ handle_call({declare, Name, MetricSpec}, _From, State) -> NewSpecs = maps:put(Name, MetricSpec, State#state.metric_specs), {reply, true, State#state{metrics = NewMetrics, metric_specs = NewSpecs}} end; +handle_call({value, Name, LabelValues}, _From, State) -> + {Result, NewState} = get_value(Name, LabelValues, State), + {reply, Result, NewState}; +handle_call({get_label_values, Name}, _From, State) -> + Result = do_get_label_values(Name, State), + {reply, Result, State}; handle_call({values, Name}, _From, State) -> {Result, NewState} = get_values(Name, State), {reply, Result, NewState}; @@ -271,6 +281,15 @@ get_values(Name, State) -> {lists:reverse(ResultsRev), FinalState} end. +-spec do_get_label_values(name(), #state{}) -> [label_values()]. +do_get_label_values(Name, State) -> + case State#state.metrics of + #{Name := MetricState} -> + [LabelValues || {_, LabelValues} <- maps:keys(MetricState)]; + #{} -> + [] + end. + -spec do_remove(name(), label_values(), #state{}) -> #state{}. do_remove(Name, LabelValues, State) -> Key = {Name, LabelValues}, @@ -279,14 +298,8 @@ do_remove(Name, LabelValues, State) -> State; {ok, MetricState} -> UpdatedMetricState = remove_metric_state(Key, MetricState), - if - map_size(UpdatedMetricState) == 0 -> - UpdatedMetrics = maps:remove(Name, State#state.metrics), - State#state{metrics = UpdatedMetrics}; - true -> - UpdatedMetrics = maps:put(Name, UpdatedMetricState, State#state.metrics), - State#state{metrics = UpdatedMetrics} - end + UpdatedMetrics = maps:put(Name, UpdatedMetricState, State#state.metrics), + State#state{metrics = UpdatedMetrics} end. -spec get_metric_state({name(), label_values()}, #state{}) -> undefined | metric_state(). diff --git a/src/instrument/mongoose_prometheus_sliding_window_collector.erl b/src/instrument/mongoose_prometheus_sliding_window_collector.erl index fc759c2eafb..7dd46613338 100644 --- a/src/instrument/mongoose_prometheus_sliding_window_collector.erl +++ b/src/instrument/mongoose_prometheus_sliding_window_collector.erl @@ -28,14 +28,9 @@ collect_metric_family(Name, Callback) -> _ -> SWName = sliding_window_metric_name(Name), MetricSpec = mongoose_prometheus_sliding_window:get_metric_spec(Name), - WindowSizeS = mongoose_prometheus_sliding_window:window_size_s(), - WindowSizeSStr = integer_to_list(WindowSizeS), Help = case MetricSpec of - undefined -> - WindowSizeSStr ++ "-second sliding window quantile summary for " ++ Name; - Spec -> - BaseHelp = proplists:get_value(help, Spec, ""), - BaseHelp ++ " (" ++ WindowSizeSStr ++ "s sliding window)" + undefined -> ""; + Spec -> proplists:get_value(help, Spec, "") end, LabelKeys = case MetricSpec of undefined -> []; @@ -49,9 +44,7 @@ collect_metric_family(Name, Callback) -> Callback(MF) end. -sliding_window_metric_name(Name) -> - WindowSizeS = mongoose_prometheus_sliding_window:window_size_s(), - Name ++ "_" ++ integer_to_list(WindowSizeS) ++ "s". +sliding_window_metric_name(Name) -> Name. create_summary_metric(LabelKeys, LabelValues, Count, Sum, Quantiles) -> %% Convert label keys and values to label pairs format diff --git a/test/mongoose_instrument_metrics_SUITE.erl b/test/mongoose_instrument_metrics_SUITE.erl index bd8931972c1..2f5b1182457 100644 --- a/test/mongoose_instrument_metrics_SUITE.erl +++ b/test/mongoose_instrument_metrics_SUITE.erl @@ -11,7 +11,7 @@ -define(HOST_TYPE2, <<"test type">>). -define(assertRelEqual(Expect, Actual, Tolerance), - (fun() -> + (fun() -> __Exp = (Expect), __Act = (Actual), __Tol = (Tolerance), @@ -21,7 +21,7 @@ __RelErr = abs(__Exp - __Act) / __Exp, __Check = case __RelErr < __Tol of true -> ok; - false -> {fail, [{expected, __Exp}, {actual, __Act}, + false -> {fail, [{expected, __Exp}, {actual, __Act}, {rel_error, __RelErr}, {tolerance, __Tol}]} end, ?assertEqual(ok, __Check) @@ -208,13 +208,13 @@ prometheus_histogram_is_created_and_updated(Config) -> ok = mongoose_instrument:set_up(Event, ?LABELS, #{metrics => #{time => histogram}}), %% Prometheus histogram shows no value if there is no data - ?assertEqual(undefined, prometheus_quantile_summary:value(Metric, [?HOST_TYPE])), + ?assertEqual(undefined, mongoose_prometheus_sliding_window:value(Metric, [?HOST_TYPE])), ok = mongoose_instrument:execute(Event, ?LABELS, #{time => 1}), - ?assertMatch({1, 1, _}, prometheus_quantile_summary:value(Metric, [?HOST_TYPE])), + ?assertMatch({1, 1, _}, mongoose_prometheus_sliding_window:value(Metric, [?HOST_TYPE])), ok = mongoose_instrument:execute(Event, ?LABELS, #{time => 1}), - ?assertMatch({2, 2, _}, prometheus_quantile_summary:value(Metric, [?HOST_TYPE])), + ?assertMatch({2, 2, _}, mongoose_prometheus_sliding_window:value(Metric, [?HOST_TYPE])), ok = mongoose_instrument:execute(Event, ?LABELS, #{time => 2}), - ?assertMatch({3, 4, _}, prometheus_quantile_summary:value(Metric, [?HOST_TYPE])). + ?assertMatch({3, 4, _}, mongoose_prometheus_sliding_window:value(Metric, [?HOST_TYPE])). prometheus_histogram_is_calculated_correctly(Config) -> Event = ?config(event, Config), @@ -222,14 +222,14 @@ prometheus_histogram_is_calculated_correctly(Config) -> ok = mongoose_instrument:set_up(Event, ?LABELS, #{metrics => #{time => histogram}}), %% Prometheus histogram shows no value if there is no data - ?assertEqual(undefined, prometheus_quantile_summary:value(Metric, [?HOST_TYPE])), + ?assertEqual(undefined, mongoose_prometheus_sliding_window:value(Metric, [?HOST_TYPE])), Values = [1, 2, 2, 2, 2, 2, 2, 9, 9, 10], lists:foreach(fun(V) -> ok = mongoose_instrument:execute(Event, ?LABELS, #{time => V}) end, Values), Sum = lists:sum(Values), Length = length(Values), - ?assertMatch({Length, Sum, _}, prometheus_quantile_summary:value(Metric, [?HOST_TYPE])), + ?assertMatch({Length, Sum, _}, mongoose_prometheus_sliding_window:value(Metric, [?HOST_TYPE])), %% Check some quantiles - {_, _, Quantiles} = prometheus_quantile_summary:value(Metric, [?HOST_TYPE]), + {_, _, Quantiles} = mongoose_prometheus_sliding_window:value(Metric, [?HOST_TYPE]), Tolerance = 0.01, ?assertRelEqual(2, proplists:get_value(0.5, Quantiles), Tolerance), ?assertRelEqual(9, proplists:get_value(0.9, Quantiles), Tolerance), @@ -243,8 +243,8 @@ prometheus_histogram_is_updated_separately_for_different_labels(Config) -> ok = mongoose_instrument:set_up(Event, ?LABELS2, #{metrics => #{time => histogram}}), ok = mongoose_instrument:execute(Event, ?LABELS, #{time => 1}), ok = mongoose_instrument:execute(Event, ?LABELS2, #{time => 2}), - ?assertMatch({1, 1, _}, prometheus_quantile_summary:value(Metric, [?HOST_TYPE])), - ?assertMatch({1, 2, _}, prometheus_quantile_summary:value(Metric, [?HOST_TYPE2])). + ?assertMatch({1, 1, _}, mongoose_prometheus_sliding_window:value(Metric, [?HOST_TYPE])), + ?assertMatch({1, 2, _}, mongoose_prometheus_sliding_window:value(Metric, [?HOST_TYPE2])). multiple_prometheus_metrics_are_updated(Config) -> Event = ?config(event, Config), @@ -255,18 +255,18 @@ multiple_prometheus_metrics_are_updated(Config) -> %% Update both metrics ok = mongoose_instrument:execute(Event, ?LABELS, #{count => 1, time => 2}), ?assertEqual(1, prometheus_counter:value(Counter, [?HOST_TYPE])), - HistogramValue = prometheus_quantile_summary:value(Histogram, [?HOST_TYPE]), + HistogramValue = mongoose_prometheus_sliding_window:value(Histogram, [?HOST_TYPE]), ?assertMatch({1, 2, _}, HistogramValue), %% Update only one metric ok = mongoose_instrument:execute(Event, ?LABELS, #{count => 2}), ?assertEqual(3, prometheus_counter:value(Counter, [?HOST_TYPE])), - ?assertEqual(HistogramValue, prometheus_quantile_summary:value(Histogram, [?HOST_TYPE])), + ?assertEqual(HistogramValue, mongoose_prometheus_sliding_window:value(Histogram, [?HOST_TYPE])), %% No update ok = mongoose_instrument:execute(Event, ?LABELS, #{something => irrelevant}), ?assertEqual(3, prometheus_counter:value(Counter, [?HOST_TYPE])), - ?assertEqual(HistogramValue, prometheus_quantile_summary:value(Histogram, [?HOST_TYPE])). + ?assertEqual(HistogramValue, mongoose_prometheus_sliding_window:value(Histogram, [?HOST_TYPE])). exometer_skips_non_metric_event(Config) -> Event = ?config(event, Config), @@ -430,7 +430,7 @@ prometheus_and_exometer_metrics_are_updated(Config) -> ?assertEqual({ok, [{count, 1}]}, exometer:get_value([?HOST_TYPE, Event, count], count)), ?assertEqual({ok, [{mean, 2}]}, exometer:get_value([?HOST_TYPE, Event, time], mean)), ?assertEqual(1, prometheus_counter:value(prom_name(Event, count), [?HOST_TYPE])), - ?assertMatch({1, 2, _}, prometheus_quantile_summary:value(prom_name(Event, time), [?HOST_TYPE])). + ?assertMatch({1, 2, _}, mongoose_prometheus_sliding_window:value(prom_name(Event, time), [?HOST_TYPE])). %% Helpers From 4e4d34c1d8dfb0d06c2e51489b30f3a3a13b00ba Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micha=C5=82=20Dobranowski?= Date: Thu, 22 Jan 2026 15:48:13 +0100 Subject: [PATCH 18/24] Raise and error when metric is unknown & clean up --- .../mongoose_prometheus_sliding_window.erl | 24 +------------------ 1 file changed, 1 insertion(+), 23 deletions(-) diff --git a/src/instrument/mongoose_prometheus_sliding_window.erl b/src/instrument/mongoose_prometheus_sliding_window.erl index 943bf84c822..91160636eb3 100644 --- a/src/instrument/mongoose_prometheus_sliding_window.erl +++ b/src/instrument/mongoose_prometheus_sliding_window.erl @@ -114,9 +114,6 @@ handle_call({value, Name, LabelValues}, _From, State) -> handle_call({get_label_values, Name}, _From, State) -> Result = do_get_label_values(Name, State), {reply, Result, State}; -handle_call({values, Name}, _From, State) -> - {Result, NewState} = get_values(Name, State), - {reply, Result, NewState}; handle_call({remove, Name, LabelValues}, _From, State) -> NewState = do_remove(Name, LabelValues, State), {reply, true, NewState}; @@ -165,7 +162,7 @@ child_spec() -> do_observe(Name, LabelValues, Value, State) -> case maps:find(Name, State#state.metrics) of error -> - State; % metric not declared, ignore + erlang:error({unknown_metric, Name}); {ok, MetricState} -> Key = {Name, LabelValues}, CurrentTime = erlang:monotonic_time(millisecond), @@ -262,25 +259,6 @@ get_value(Name, LabelValues, State) -> {Result, NewState} end. --spec get_values(name(), #state{}) -> {[{label_values(), {non_neg_integer(), number(), [{number(), number()}]}}], #state{}}. -get_values(Name, State) -> - case maps:find(Name, State#state.metrics) of - error -> - {[], State}; - {ok, MetricState} -> - Keys = maps:keys(MetricState), - {ResultsRev, FinalState} = lists:foldl( - fun({_, LabelValues}, {AccRes, S}) -> - {Res, S2} = get_value(Name, LabelValues, S), - NewAcc = case Res of - undefined -> AccRes; - _ -> [{LabelValues, Res} | AccRes] - end, - {NewAcc, S2} - end, {[], State}, Keys), - {lists:reverse(ResultsRev), FinalState} - end. - -spec do_get_label_values(name(), #state{}) -> [label_values()]. do_get_label_values(Name, State) -> case State#state.metrics of From 7d52e9b2a9d832fdc9b31e764a10a3faccc5c524 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micha=C5=82=20Dobranowski?= Date: Fri, 30 Jan 2026 15:25:57 +0100 Subject: [PATCH 19/24] Init sliding_window only if prometheus is enabled --- src/ejabberd_sup.erl | 11 ++++++++--- src/instrument/mongoose_prometheus_sliding_window.erl | 4 ++++ 2 files changed, 12 insertions(+), 3 deletions(-) diff --git a/src/ejabberd_sup.erl b/src/ejabberd_sup.erl index 87662fedb51..a3d00db9924 100644 --- a/src/ejabberd_sup.erl +++ b/src/ejabberd_sup.erl @@ -45,7 +45,6 @@ start_link() -> init(noargs) -> Hooks = worker_spec(gen_hook), Instrument = worker_spec(mongoose_instrument), - PrometheusSlidingWindow = mongoose_prometheus_sliding_window:child_spec(), Cleaner = worker_spec(mongoose_cleaner), Router = worker_spec(ejabberd_router), S2S = worker_spec(ejabberd_s2s), @@ -70,8 +69,8 @@ init(noargs) -> [StartIdServer, PG, Hooks, - Instrument, - PrometheusSlidingWindow, + Instrument + ] ++ prometheus_sliding_window_spec() ++ [ Cleaner, SMBackendSupervisor, OutgoingPoolsSupervisor @@ -124,6 +123,12 @@ template_supervisor_spec(Name, Module) -> supervisor_spec(Mod) -> {Mod, {Mod, start_link, []}, permanent, infinity, supervisor, [Mod]}. +prometheus_sliding_window_spec() -> + case mongoose_config:get_opt([instrumentation, prometheus], undefined) of + undefined -> []; + #{} -> [mongoose_prometheus_sliding_window:child_spec()] + end. + worker_spec(Mod) -> worker_spec(Mod, []). diff --git a/src/instrument/mongoose_prometheus_sliding_window.erl b/src/instrument/mongoose_prometheus_sliding_window.erl index 91160636eb3..5c7ca9912c0 100644 --- a/src/instrument/mongoose_prometheus_sliding_window.erl +++ b/src/instrument/mongoose_prometheus_sliding_window.erl @@ -36,6 +36,10 @@ default_quantiles() -> [0.5, 0.75, 0.90, 0.95, 0.99, 0.999]. -spec default_error() -> number(). default_error() -> 0.01. +%% Setting the default bound to 1260 means that: +%% * measuring in µs suffices for actions lasting up to a day (with 1% accuracy), +%% * measuring in bytes suffices for sizes up to 81 GB (with 1% accuracy). +%% See: https://hexdocs.pm/ddskerl/ddskerl_ets.html -spec default_bound() -> non_neg_integer(). default_bound() -> 1260. From 10d9817ca2788830ee134803eb34717b3af0b948 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micha=C5=82=20Dobranowski?= Date: Thu, 5 Feb 2026 11:42:43 +0100 Subject: [PATCH 20/24] Fix rotating windows, add tests --- .../mongoose_prometheus_sliding_window.erl | 89 ++++++++++--------- test/mongoose_instrument_metrics_SUITE.erl | 30 +++++++ 2 files changed, 78 insertions(+), 41 deletions(-) diff --git a/src/instrument/mongoose_prometheus_sliding_window.erl b/src/instrument/mongoose_prometheus_sliding_window.erl index 5c7ca9912c0..b75b96159c5 100644 --- a/src/instrument/mongoose_prometheus_sliding_window.erl +++ b/src/instrument/mongoose_prometheus_sliding_window.erl @@ -1,7 +1,7 @@ -module(mongoose_prometheus_sliding_window). %% Manages a sliding window using DDSketch instances for quantile summaries. -%% Window parameters are configurable via WINDOW_SIZE_MS and WINDOW_STEP_MS constants. +%% Window parameters are configurable via window_size_ms and window_step_ms constants. -export([child_spec/0, declare/1, @@ -20,15 +20,24 @@ terminate/2, code_change/3]). --ignore_xref([start_link/0, value/2]). +-export([window_size_ms/0, window_step_ms/0, window_count/0]). % Exported for mocking in tests + +-ignore_xref([start_link/0, value/2, window_size_ms/0, window_step_ms/0, window_count/0]). -behaviour(gen_server). -%% Configurable window parameters --define(WINDOW_SIZE_MS, 60000). % Total window size: 60 seconds --define(WINDOW_STEP_MS, 3000). % Step size: 3 seconds per sub-window --define(WINDOW_COUNT, (?WINDOW_SIZE_MS div ?WINDOW_STEP_MS)). % Number of sub-windows +%% Window parameters + +-spec window_size_ms() -> pos_integer(). +window_size_ms() -> 60000. % Total window size: 60 seconds + +-spec window_step_ms() -> pos_integer(). +window_step_ms() -> 3000. % Step size: 3 seconds per sub-window +-spec window_count() -> pos_integer(). +window_count() -> ?MODULE:window_size_ms() div ?MODULE:window_step_ms(). + +%% DDSketch parameters -spec default_quantiles() -> [number()]. default_quantiles() -> [0.5, 0.75, 0.90, 0.95, 0.99, 0.999]. @@ -100,7 +109,7 @@ start_link() -> init([]) -> EtsTable = ets:new(?MODULE, [set, private, {read_concurrency, true}, {write_concurrency, true}]), - TimerRef = erlang:start_timer(?WINDOW_STEP_MS, self(), rotate), + TimerRef = erlang:start_timer(?MODULE:window_step_ms(), self(), rotate), {ok, #state{timer_ref = TimerRef, ets_table = EtsTable}}. handle_call({declare, Name, MetricSpec}, _From, State) -> @@ -137,9 +146,10 @@ handle_cast(_Msg, State) -> {noreply, State}. handle_info({timeout, _TimerRef, rotate}, State) -> - %% The rotation is handled in ensure_windows_rotated, so we just reschedule - TimerRef = erlang:start_timer(?WINDOW_STEP_MS, self(), rotate), - {noreply, State#state{timer_ref = TimerRef}}; + CurrentTime = erlang:monotonic_time(millisecond), + NewState = rotate_all_windows(CurrentTime, State), + TimerRef = erlang:start_timer(?MODULE:window_step_ms(), self(), rotate), + {noreply, NewState#state{timer_ref = TimerRef}}; handle_info(_Info, State) -> {noreply, State}. @@ -171,7 +181,7 @@ do_observe(Name, LabelValues, Value, State) -> Key = {Name, LabelValues}, CurrentTime = erlang:monotonic_time(millisecond), {Windows, CurrentIndex} = - ensure_windows_rotated(Key, CurrentTime, State), + ensure_windows_initialized(Key, CurrentTime, State), %% Get current window and add observation {CurrentWindow, WindowStartTime} = lists:nth(CurrentIndex + 1, Windows), @@ -190,27 +200,17 @@ do_observe(Name, LabelValues, Value, State) -> State#state{metrics = UpdatedMetrics} end. --spec ensure_windows_rotated({name(), label_values()}, non_neg_integer(), #state{}) -> +-spec ensure_windows_initialized({name(), label_values()}, non_neg_integer(), #state{}) -> {[{window_data(), non_neg_integer()}], non_neg_integer()}. -ensure_windows_rotated(Key, CurrentTime, State) -> +ensure_windows_initialized(Key, CurrentTime, State) -> + WindowCount = ?MODULE:window_count(), case get_metric_state(Key, State) of undefined -> - %% Initialize new windows, all starting at current time (they'll be rotated as needed) - Windows = [new_window(Key, Index, CurrentTime, State) || Index <- lists:seq(0, ?WINDOW_COUNT - 1)], + %% Initialize new windows, all starting at current time + Windows = [new_window(Key, Index, CurrentTime, State) || Index <- lists:seq(0, WindowCount - 1)], {Windows, 0}; #{windows := Windows, current_index := CurrentIndex} -> - {_CurrentWindow, WindowStartTime} = lists:nth(CurrentIndex + 1, Windows), - Elapsed = CurrentTime - WindowStartTime, - if - Elapsed >= ?WINDOW_STEP_MS -> - NewIndex = (CurrentIndex + 1) rem ?WINDOW_COUNT, - {CurrentWindow, _} = lists:nth(CurrentIndex + 1, Windows), - ResetWindow = reset_window(CurrentWindow), - UpdatedWindows = set_nth(CurrentIndex + 1, {ResetWindow, CurrentTime}, Windows), - {UpdatedWindows, NewIndex}; - true -> - {Windows, CurrentIndex} - end + {Windows, CurrentIndex} end. -spec get_value(name(), label_values(), #state{}) -> @@ -220,22 +220,12 @@ get_value(Name, LabelValues, State) -> case get_metric_state(Key, State) of undefined -> {undefined, State}; - #{} -> + #{windows := Windows} -> CurrentTime = erlang:monotonic_time(millisecond), - {UpdatedWindows, NewIndex} = ensure_windows_rotated(Key, CurrentTime, State), - - %% Update metric state with rotated windows - {ok, MetricState} = maps:find(Name, State#state.metrics), - UpdatedMetricState = maps:put(Key, - #{windows => UpdatedWindows, - current_index => NewIndex}, - MetricState), - UpdatedMetrics = maps:put(Name, UpdatedMetricState, State#state.metrics), - NewState = State#state{metrics = UpdatedMetrics}, %% Filter windows to only include those within the last window size - CutoffTime = CurrentTime - ?WINDOW_SIZE_MS, - ActiveWindows = [Window || {Window, StartTime} <- UpdatedWindows, + CutoffTime = CurrentTime - ?MODULE:window_size_ms(), + ActiveWindows = [Window || {Window, StartTime} <- Windows, StartTime >= CutoffTime], %% Merge all active windows @@ -260,7 +250,7 @@ get_value(Name, LabelValues, State) -> {Count, Sum, Quantiles} end end, - {Result, NewState} + {Result, State} end. -spec do_get_label_values(name(), #state{}) -> [label_values()]. @@ -326,6 +316,23 @@ new_window(Key = {Name, _LabelValues}, Index, StartTime, State) -> name => WindowName}, {WindowData, StartTime}. +-spec rotate_all_windows(non_neg_integer(), #state{}) -> #state{}. +rotate_all_windows(CurrentTime, State) -> + WindowCount = ?MODULE:window_count(), + UpdatedMetrics = maps:map( + fun(_Name, MetricState) -> + maps:map( + fun(_Key, #{windows := Windows, current_index := CurrentIndex}) -> + %% Move to next window and reset it + NewIndex = (CurrentIndex + 1) rem WindowCount, + {WindowToReset, _OldTime} = lists:nth(NewIndex + 1, Windows), + ResetWindow = reset_window(WindowToReset), + UpdatedWindows = set_nth(NewIndex + 1, {ResetWindow, CurrentTime}, Windows), + #{windows => UpdatedWindows, current_index => NewIndex} + end, MetricState) + end, State#state.metrics), + State#state{metrics = UpdatedMetrics}. + -spec reset_window(window_data()) -> window_data(). reset_window(WindowData) -> Sketch = maps:get(sketch, WindowData), diff --git a/test/mongoose_instrument_metrics_SUITE.erl b/test/mongoose_instrument_metrics_SUITE.erl index 2f5b1182457..99d8ab98232 100644 --- a/test/mongoose_instrument_metrics_SUITE.erl +++ b/test/mongoose_instrument_metrics_SUITE.erl @@ -50,6 +50,7 @@ groups() -> prometheus_histogram_is_created_and_updated, prometheus_histogram_is_calculated_correctly, prometheus_histogram_is_updated_separately_for_different_labels, + prometheus_histogram_sliding_window_expires_data, multiple_prometheus_metrics_are_updated]}, {exometer, [parallel], [exometer_skips_non_metric_event, exometer_gauge_is_created_and_updated, @@ -246,6 +247,35 @@ prometheus_histogram_is_updated_separately_for_different_labels(Config) -> ?assertMatch({1, 1, _}, mongoose_prometheus_sliding_window:value(Metric, [?HOST_TYPE])), ?assertMatch({1, 2, _}, mongoose_prometheus_sliding_window:value(Metric, [?HOST_TYPE2])). +prometheus_histogram_sliding_window_expires_data(Config) -> + Event = ?config(event, Config), + Metric = prom_name(Event, time), + + OriginalWindowStep = mongoose_prometheus_sliding_window:window_step_ms(), + + %% We mock the window parameters to use small values and wait a short time. + meck:new(mongoose_prometheus_sliding_window, [passthrough]), + meck:expect(mongoose_prometheus_sliding_window, window_size_ms, fun() -> 1000 end), + meck:expect(mongoose_prometheus_sliding_window, window_step_ms, fun() -> 100 end), + + %% Wait a bit for the mocked module to be used + timer:sleep(OriginalWindowStep), + + try + ok = mongoose_instrument:set_up(Event, ?LABELS, #{metrics => #{time => histogram}}), + ok = mongoose_instrument:execute(Event, ?LABELS, #{time => 10}), + ?assertMatch({1, 10, _}, mongoose_prometheus_sliding_window:value(Metric, [?HOST_TYPE])), + timer:sleep(500), + ok = mongoose_instrument:execute(Event, ?LABELS, #{time => 20}), + ?assertMatch({2, 30, _}, mongoose_prometheus_sliding_window:value(Metric, [?HOST_TYPE])), + timer:sleep(700), % Wait enough for the first value to expire + ?assertMatch({1, 20, _}, mongoose_prometheus_sliding_window:value(Metric, [?HOST_TYPE])), + timer:sleep(500), % Wait enough for the second value to expire + ?assertEqual(undefined, mongoose_prometheus_sliding_window:value(Metric, [?HOST_TYPE])) + after + meck:unload(mongoose_prometheus_sliding_window) + end. + multiple_prometheus_metrics_are_updated(Config) -> Event = ?config(event, Config), Counter = prom_name(Event, count), From 604909c0ee179400b42a9e38f1c037d03579c679 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micha=C5=82=20Dobranowski?= Date: Wed, 11 Feb 2026 15:17:55 +0100 Subject: [PATCH 21/24] Check quantile endpoint too --- test/prometheus_endpoint_SUITE.erl | 14 +++++++++++--- 1 file changed, 11 insertions(+), 3 deletions(-) diff --git a/test/prometheus_endpoint_SUITE.erl b/test/prometheus_endpoint_SUITE.erl index 1e589a8ed19..ac601d78f07 100644 --- a/test/prometheus_endpoint_SUITE.erl +++ b/test/prometheus_endpoint_SUITE.erl @@ -55,7 +55,7 @@ test_metrics(Event, Labels) -> check_counter(Event, requests, Labels, 1, Scraped), % 'spiral' is a Prometheus counter check_gauge(Event, sessions, Labels, 5, Scraped), % 'counter' is a Prometheus gauge check_gauge(Event, seconds, Labels, 10, Scraped), - check_histogram(Event, time, Labels, #{count => 1, sum => 2}, Scraped). + check_histogram(Event, time, Labels, #{count => 1, sum => 2, median => 2}, Scraped). %% Checks for the parsed metrics @@ -72,14 +72,17 @@ check_gauge(Event, Metric, Labels, ExpValue, Scraped) -> ?assertEqual({Labels, ExpValue}, Value). check_histogram(Event, Metric, Labels, ExpValues, Scraped) -> - #{count := ExpCount, sum := ExpSum} = ExpValues, + #{count := ExpCount, sum := ExpSum, median := ExpMedian} = ExpValues, [Type, Help | _] = get_metric([Event, Metric], Scraped), ?assertEqual({<<"TYPE">>, <<"summary">>}, Type), ?assertEqual({<<"HELP">>, help(Event, Metric)}, Help), [Count] = get_metric([Event, Metric, count], Scraped), ?assertEqual({Labels, ExpCount}, Count), [Sum] = get_metric([Event, Metric, sum], Scraped), - ?assertEqual({Labels, ExpSum}, Sum). + ?assertEqual({Labels, ExpSum}, Sum), + MedianLabels = Labels#{quantile => <<"0.5">>}, + [{MedianLabels, Median}] = get_metric([Event, Metric], MedianLabels, Scraped), + ?assert(abs(Median - ExpMedian) < 0.01). help(Event, Metric) -> <<"Event: ", (atom_to_binary(Event))/binary, ", Metric: ", (atom_to_binary(Metric))/binary>>. @@ -89,6 +92,11 @@ get_metric(Parts, Scraped) when is_list(Parts) -> get_metric(Name, Scraped) when is_binary(Name) -> [{LabelsOrKey, Value} || {N, LabelsOrKey, Value} <- Scraped, N =:= Name]. +get_metric(Parts, Labels, Scraped) when is_list(Parts) -> + get_metric(binary_name(Parts), Labels, Scraped); +get_metric(Name, Labels, Scraped) when is_binary(Name) -> + [{L, Value} || {N, L, Value} <- Scraped, N =:= Name, is_map(L), maps:merge(L, Labels) =:= L]. + binary_name(Atoms) -> list_to_binary(string:join([atom_to_list(Atom) || Atom <- Atoms], "_")). From e018a13e2fccb39309cfd61c6b3e93af05ace087 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micha=C5=82=20Dobranowski?= Date: Wed, 11 Feb 2026 16:52:39 +0100 Subject: [PATCH 22/24] Use `timer:send_interval` instead of `start_timer` and restart gen_server instead of waiting during tests --- .../mongoose_prometheus_sliding_window.erl | 12 +++++++----- test/mongoose_instrument_metrics_SUITE.erl | 7 +++---- 2 files changed, 10 insertions(+), 9 deletions(-) diff --git a/src/instrument/mongoose_prometheus_sliding_window.erl b/src/instrument/mongoose_prometheus_sliding_window.erl index b75b96159c5..f44628d6807 100644 --- a/src/instrument/mongoose_prometheus_sliding_window.erl +++ b/src/instrument/mongoose_prometheus_sliding_window.erl @@ -53,7 +53,7 @@ default_error() -> 0.01. default_bound() -> 1260. -record(state, { - timer_ref :: reference() | undefined, + timer_ref :: {ok, timer:tref()} | {error, term()} | undefined, ets_table :: ets:tab() | undefined, metrics = #{} :: #{name() => #{{name(), label_values()} => metric_state()}}, metric_specs = #{} :: #{name() => proplists:proplist()} @@ -109,7 +109,7 @@ start_link() -> init([]) -> EtsTable = ets:new(?MODULE, [set, private, {read_concurrency, true}, {write_concurrency, true}]), - TimerRef = erlang:start_timer(?MODULE:window_step_ms(), self(), rotate), + TimerRef = timer:send_interval(?MODULE:window_step_ms(), rotate), {ok, #state{timer_ref = TimerRef, ets_table = EtsTable}}. handle_call({declare, Name, MetricSpec}, _From, State) -> @@ -145,14 +145,16 @@ handle_cast({observe, Name, LabelValues, Value}, State) -> handle_cast(_Msg, State) -> {noreply, State}. -handle_info({timeout, _TimerRef, rotate}, State) -> +handle_info(rotate, State) -> CurrentTime = erlang:monotonic_time(millisecond), NewState = rotate_all_windows(CurrentTime, State), - TimerRef = erlang:start_timer(?MODULE:window_step_ms(), self(), rotate), - {noreply, NewState#state{timer_ref = TimerRef}}; + {noreply, NewState}; handle_info(_Info, State) -> {noreply, State}. +terminate(_Reason, #state{timer_ref = {ok, TRef}}) -> + timer:cancel(TRef), + ok; terminate(_Reason, _State) -> ok. diff --git a/test/mongoose_instrument_metrics_SUITE.erl b/test/mongoose_instrument_metrics_SUITE.erl index 99d8ab98232..47ab245ba32 100644 --- a/test/mongoose_instrument_metrics_SUITE.erl +++ b/test/mongoose_instrument_metrics_SUITE.erl @@ -251,15 +251,14 @@ prometheus_histogram_sliding_window_expires_data(Config) -> Event = ?config(event, Config), Metric = prom_name(Event, time), - OriginalWindowStep = mongoose_prometheus_sliding_window:window_step_ms(), - %% We mock the window parameters to use small values and wait a short time. meck:new(mongoose_prometheus_sliding_window, [passthrough]), meck:expect(mongoose_prometheus_sliding_window, window_size_ms, fun() -> 1000 end), meck:expect(mongoose_prometheus_sliding_window, window_step_ms, fun() -> 100 end), - %% Wait a bit for the mocked module to be used - timer:sleep(OriginalWindowStep), + %% Restart the gen_server to make sure that the new window parameters are applied + ok = gen_server:stop(mongoose_prometheus_sliding_window), + {ok, _} = mongoose_prometheus_sliding_window:start_link(), try ok = mongoose_instrument:set_up(Event, ?LABELS, #{metrics => #{time => histogram}}), From 1b64dc88ef0fa5667fbd48dc0b857ef887bc688b Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micha=C5=82=20Dobranowski?= Date: Wed, 11 Feb 2026 17:17:32 +0100 Subject: [PATCH 23/24] Remove some redundant code --- .../mongoose_prometheus_sliding_window.erl | 81 +++++++------------ 1 file changed, 31 insertions(+), 50 deletions(-) diff --git a/src/instrument/mongoose_prometheus_sliding_window.erl b/src/instrument/mongoose_prometheus_sliding_window.erl index f44628d6807..16505d90759 100644 --- a/src/instrument/mongoose_prometheus_sliding_window.erl +++ b/src/instrument/mongoose_prometheus_sliding_window.erl @@ -64,7 +64,7 @@ default_bound() -> 1260. -type window_data() :: #{sketch := ddskerl_ets:ddsketch(), ref := ets:tab(), name := term()}. --type metric_state() :: #{windows => [{window_data(), non_neg_integer()}], +-type metric_state() :: #{windows => [window_data()], current_index => non_neg_integer()}. %% Public API @@ -122,8 +122,8 @@ handle_call({declare, Name, MetricSpec}, _From, State) -> {reply, true, State#state{metrics = NewMetrics, metric_specs = NewSpecs}} end; handle_call({value, Name, LabelValues}, _From, State) -> - {Result, NewState} = get_value(Name, LabelValues, State), - {reply, Result, NewState}; + Result = get_value(Name, LabelValues, State), + {reply, Result, State}; handle_call({get_label_values, Name}, _From, State) -> Result = do_get_label_values(Name, State), {reply, Result, State}; @@ -146,8 +146,7 @@ handle_cast(_Msg, State) -> {noreply, State}. handle_info(rotate, State) -> - CurrentTime = erlang:monotonic_time(millisecond), - NewState = rotate_all_windows(CurrentTime, State), + NewState = rotate_all_windows(State), {noreply, NewState}; handle_info(_Info, State) -> {noreply, State}. @@ -181,18 +180,17 @@ do_observe(Name, LabelValues, Value, State) -> erlang:error({unknown_metric, Name}); {ok, MetricState} -> Key = {Name, LabelValues}, - CurrentTime = erlang:monotonic_time(millisecond), {Windows, CurrentIndex} = - ensure_windows_initialized(Key, CurrentTime, State), + ensure_windows_initialized(Key, State), %% Get current window and add observation - {CurrentWindow, WindowStartTime} = lists:nth(CurrentIndex + 1, Windows), + CurrentWindow = lists:nth(CurrentIndex + 1, Windows), Sketch = maps:get(sketch, CurrentWindow), UpdatedSketch = ddskerl_ets:insert(Sketch, Value), UpdatedWindow = CurrentWindow#{sketch := UpdatedSketch}, %% Update windows list and metric state - UpdatedWindows = set_nth(CurrentIndex + 1, {UpdatedWindow, WindowStartTime}, Windows), + UpdatedWindows = set_nth(CurrentIndex + 1, UpdatedWindow, Windows), UpdatedMetricState = maps:put(Key, #{windows => UpdatedWindows, current_index => CurrentIndex}, @@ -202,43 +200,31 @@ do_observe(Name, LabelValues, Value, State) -> State#state{metrics = UpdatedMetrics} end. --spec ensure_windows_initialized({name(), label_values()}, non_neg_integer(), #state{}) -> - {[{window_data(), non_neg_integer()}], non_neg_integer()}. -ensure_windows_initialized(Key, CurrentTime, State) -> +-spec ensure_windows_initialized({name(), label_values()}, #state{}) -> {[window_data()], non_neg_integer()}. +ensure_windows_initialized(Key, State) -> WindowCount = ?MODULE:window_count(), case get_metric_state(Key, State) of undefined -> - %% Initialize new windows, all starting at current time - Windows = [new_window(Key, Index, CurrentTime, State) || Index <- lists:seq(0, WindowCount - 1)], + Windows = [new_window(Key, Index, State) || Index <- lists:seq(0, WindowCount - 1)], {Windows, 0}; #{windows := Windows, current_index := CurrentIndex} -> {Windows, CurrentIndex} end. --spec get_value(name(), label_values(), #state{}) -> - {undefined | {non_neg_integer(), number(), [{number(), number()}]}, #state{}}. +-spec get_value(name(), label_values(), #state{}) -> undefined | {non_neg_integer(), number(), [{number(), number()}]}. get_value(Name, LabelValues, State) -> Key = {Name, LabelValues}, case get_metric_state(Key, State) of undefined -> - {undefined, State}; + undefined; #{windows := Windows} -> - CurrentTime = erlang:monotonic_time(millisecond), - - %% Filter windows to only include those within the last window size - CutoffTime = CurrentTime - ?MODULE:window_size_ms(), - ActiveWindows = [Window || {Window, StartTime} <- Windows, - StartTime >= CutoffTime], - - %% Merge all active windows - WindowTuples = [window_tuple(Window) || Window <- ActiveWindows], + WindowTuples = [window_tuple(Window) || Window <- Windows], Merged = lists:foldl(fun (undefined, Acc) -> Acc; (Val, undefined) -> Val; (Val, Acc) -> ddskerl_ets:merge_tuples(Acc, Val) end, undefined, WindowTuples), - - Result = case Merged of + case Merged of undefined -> undefined; _ -> @@ -251,8 +237,7 @@ get_value(Name, LabelValues, State) -> Quantiles = [{Q, ddskerl_ets:quantile_tuple(Merged, Q)} || Q <- default_quantiles()], {Count, Sum, Quantiles} end - end, - {Result, State} + end end. -spec do_get_label_values(name(), #state{}) -> [label_values()]. @@ -304,22 +289,20 @@ metric_options(Name, State) -> window_name({Name, LabelValues}, Index) -> {Name, LabelValues, Index}. --spec new_window({name(), label_values()}, non_neg_integer(), non_neg_integer(), #state{}) -> - {window_data(), non_neg_integer()}. -new_window(Key = {Name, _LabelValues}, Index, StartTime, State) -> +-spec new_window({name(), label_values()}, non_neg_integer(), #state{}) -> window_data(). +new_window(Key = {Name, _LabelValues}, Index, State) -> {Error, Bound} = metric_options(Name, State), WindowName = window_name(Key, Index), Sketch = ddskerl_ets:new(#{ets_table => State#state.ets_table, name => WindowName, error => Error, bound => Bound}), - WindowData = #{sketch => Sketch, - ref => State#state.ets_table, - name => WindowName}, - {WindowData, StartTime}. + #{sketch => Sketch, + ref => State#state.ets_table, + name => WindowName}. --spec rotate_all_windows(non_neg_integer(), #state{}) -> #state{}. -rotate_all_windows(CurrentTime, State) -> +-spec rotate_all_windows(#state{}) -> #state{}. +rotate_all_windows(State) -> WindowCount = ?MODULE:window_count(), UpdatedMetrics = maps:map( fun(_Name, MetricState) -> @@ -327,9 +310,9 @@ rotate_all_windows(CurrentTime, State) -> fun(_Key, #{windows := Windows, current_index := CurrentIndex}) -> %% Move to next window and reset it NewIndex = (CurrentIndex + 1) rem WindowCount, - {WindowToReset, _OldTime} = lists:nth(NewIndex + 1, Windows), + WindowToReset = lists:nth(NewIndex + 1, Windows), ResetWindow = reset_window(WindowToReset), - UpdatedWindows = set_nth(NewIndex + 1, {ResetWindow, CurrentTime}, Windows), + UpdatedWindows = set_nth(NewIndex + 1, ResetWindow, Windows), #{windows => UpdatedWindows, current_index => NewIndex} end, MetricState) end, State#state.metrics), @@ -352,15 +335,13 @@ window_tuple(WindowData) -> -spec remove_metric_state({name(), label_values()}, #{}) -> #{}. remove_metric_state(Key, MetricState) -> - UpdatedMetricState = - case maps:find(Key, MetricState) of - error -> - MetricState; - {ok, #{windows := Windows}} -> - lists:foreach(fun({Window, _}) -> delete_window(Window) end, Windows), - maps:remove(Key, MetricState) - end, - UpdatedMetricState. + case maps:find(Key, MetricState) of + error -> + MetricState; + {ok, #{windows := Windows}} -> + lists:foreach(fun(Window) -> delete_window(Window) end, Windows), + maps:remove(Key, MetricState) + end. -spec delete_window(window_data()) -> ok. delete_window(WindowData) -> From 1b68f82c72aeee91fd01a29cffc990f56bebc594 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micha=C5=82=20Dobranowski?= Date: Wed, 11 Feb 2026 17:51:01 +0100 Subject: [PATCH 24/24] Create a new sequential test group --- test/mongoose_instrument_metrics_SUITE.erl | 7 ++++++- 1 file changed, 6 insertions(+), 1 deletion(-) diff --git a/test/mongoose_instrument_metrics_SUITE.erl b/test/mongoose_instrument_metrics_SUITE.erl index 47ab245ba32..b3293226cd8 100644 --- a/test/mongoose_instrument_metrics_SUITE.erl +++ b/test/mongoose_instrument_metrics_SUITE.erl @@ -33,6 +33,7 @@ all() -> [{group, prometheus}, + {group, prometheus_sliding_window}, % separate, because it can't be run in parallel {group, exometer}, {group, exometer_global}, {group, prometheus_and_exometer} @@ -50,8 +51,8 @@ groups() -> prometheus_histogram_is_created_and_updated, prometheus_histogram_is_calculated_correctly, prometheus_histogram_is_updated_separately_for_different_labels, - prometheus_histogram_sliding_window_expires_data, multiple_prometheus_metrics_are_updated]}, + {prometheus_sliding_window, [], [prometheus_histogram_sliding_window_expires_data]}, {exometer, [parallel], [exometer_skips_non_metric_event, exometer_gauge_is_created_and_updated, exometer_gauge_is_updated_separately_for_different_labels, @@ -87,6 +88,7 @@ init_per_group(Group, Config) -> Config1 ++ extra_config(Group). extra_processes(prometheus) -> [{mongoose_prometheus_sliding_window, start_link, []}]; +extra_processes(prometheus_sliding_window) -> [{mongoose_prometheus_sliding_window, start_link, []}]; extra_processes(prometheus_and_exometer) -> [{mongoose_prometheus_sliding_window, start_link, []}]; extra_processes(_) -> []. @@ -103,12 +105,15 @@ end_per_testcase(_Case, _Config) -> log_helper:unsubscribe(). apps(prometheus) -> [prometheus, prometheus_httpd, prometheus_cowboy]; +apps(prometheus_sliding_window) -> apps(prometheus); apps(exometer) -> [exometer_core]; apps(exometer_global) -> [exometer_core]; apps(prometheus_and_exometer) -> apps(prometheus) ++ apps(exometer). opts(prometheus) -> #{prometheus => #{}}; +opts(prometheus_sliding_window) -> + opts(prometheus); opts(exometer) -> #{exometer => #{all_metrics_are_global => false, report => #{}}}; opts(exometer_global) ->