Copilot commented on code in PR #13754:
URL: https://github.com/apache/trafficserver/pull/13754#discussion_r4169514920
##########
plugins/stats_over_http/stats_over_http.cc:
##########
@@ -1508,6 +1289,911 @@ config_handler(TSCont cont, TSEvent /* event ATS_UNUSED
*/, void * /* edata ATS_
return 0;
}
+//
+// Requests for the stats. A render on a task thread answers the requests to
the global plugin and to remap rules.
+//
+
+static constexpr std::string_view STATS_FORMAT_FIELD = "X-Stats-Format";
+static constexpr std::string_view ALLOWED_METHODS = "GET, HEAD";
+
+static const char REMAP_USAGE[] =
"[--format=json|csv|prometheus|prometheus_v2] [--integer-counters]
[--wrap-counters] "
+ "[--no-prometheus-help] [--max-age-ms=N]
[--wait-timeout-ms=N] [--config=FILE] "
+ "[--on-config-error=fail|503]";
+
+static std::string_view
+format_name(output_format_t format)
+{
+ switch (format) {
+ case output_format_t::JSON_OUTPUT:
+ return "json";
+ case output_format_t::CSV_OUTPUT:
+ return "csv";
+ case output_format_t::PROMETHEUS_OUTPUT:
+ return "prometheus";
+ case output_format_t::PROMETHEUS_V2_OUTPUT:
+ return "prometheus_v2";
+ }
+ return "json";
+}
+
+static std::string_view
+format_content_type(output_format_t format)
+{
+ switch (format) {
+ case output_format_t::JSON_OUTPUT:
+ return "text/json";
+ case output_format_t::CSV_OUTPUT:
+ return "text/csv";
+ case output_format_t::PROMETHEUS_OUTPUT:
+ return "text/plain; version=0.0.4; charset=utf-8";
+ case output_format_t::PROMETHEUS_V2_OUTPUT:
+ return "text/plain; version=2.0.0; charset=utf-8";
+ }
+ return "text/json";
+}
+
+static std::string_view
+encoding_name(encoding_format_t encoding)
+{
+ switch (encoding) {
+ case encoding_format_t::DEFLATE:
+ return "deflate";
+ case encoding_format_t::GZIP:
+ return "gzip";
+ case encoding_format_t::BR:
+ return "br";
+ case encoding_format_t::NONE:
+ break;
+ }
+ return {};
+}
+
+static bool
+parse_format(std::string_view name, output_format_t &format)
+{
+ for (auto candidate : {output_format_t::JSON_OUTPUT,
output_format_t::CSV_OUTPUT, output_format_t::PROMETHEUS_OUTPUT,
+ output_format_t::PROMETHEUS_V2_OUTPUT}) {
+ if (name == format_name(candidate)) {
+ format = candidate;
+ return true;
+ }
+ }
+ return false;
+}
+
+// A rendered body in one encoding. It does not change after the render
publishes it, so requests on any thread share its
+// blocks.
+struct stats_snapshot {
+ stats_snapshot() : body(TSIOBufferCreate()),
reader(TSIOBufferReaderAlloc(body)) {}
+ ~stats_snapshot() { TSIOBufferDestroy(body); }
+ stats_snapshot(const stats_snapshot &) = delete;
+ stats_snapshot &operator=(const stats_snapshot &) = delete;
+
+ TSIOBuffer body;
+ TSIOBufferReader reader;
+ int64_t bytes = 0;
+ encoding_format_t encoding = encoding_format_t::NONE;
+ std::chrono::steady_clock::time_point rendered;
+};
+
+using snapshot_ptr = std::shared_ptr<const stats_snapshot>;
+using snapshot_table = snapshot_ptr[FORMAT_COUNT][ENCODING_COUNT];
+
+// One request for the stats. The waiter list, the transaction hooks and the
intercept share it.
+struct stats_scrape {
+ stats_scrape(TSHttpTxn txn, output_format_t fmt, encoding_format_t enc, bool
rule)
+ : txnp(txn), format(fmt), encoding(enc), remap(rule)
+ {
+ }
+
+ TSHttpTxn txnp;
+ output_format_t format;
+ encoding_format_t encoding;
+ bool remap; // A remap rule serves the request, rather than the
global plugin.
+ // When the request started to wait for a render.
+ std::chrono::steady_clock::time_point arrived;
+
+ // TSRemapDoRemap, scrape_wait, the render or the watchdog sets these before
the transaction continues, and the intercept
+ // reads them only after that. The snapshot is null for a HEAD request to a
remap rule and for a 503.
+ snapshot_ptr snapshot;
+ TSHttpStatus status = TS_HTTP_STATUS_OK;
+};
+
+struct stats_watchdog;
+
+// The stats of one remap rule, or of the global plugin. The remap rule and
each continuation that works for the instance
+// hold a reference, so that a render or a watchdog can outlive the rule.
+struct stats_instance {
+ explicit stats_instance(const stats_options &opts) : options(opts) {}
+ ~stats_instance() { Dbg(dbg_ctl, "Freeing stats instance %p", this); }
+ stats_instance(const stats_instance &) = delete;
+ stats_instance &operator=(const stats_instance &) = delete;
+
+ PrometheusRenderer *
+ prometheus(output_format_t format)
+ {
+ std::unique_ptr<PrometheusRenderer> *renderer = nullptr;
+
+ if (format == output_format_t::PROMETHEUS_OUTPUT) {
+ renderer = &prometheus_v1;
+ } else if (format == output_format_t::PROMETHEUS_V2_OUTPUT) {
+ renderer = &prometheus_v2;
+ } else {
+ return nullptr;
+ }
+ if (*renderer == nullptr) {
+ *renderer = make_prometheus_renderer(format, options);
+ }
+ return renderer->get();
+ }
+
+ const stats_options options;
+ // Only the render in flight uses these, so the mutex does not guard them.
+ std::unique_ptr<PrometheusRenderer> prometheus_v1;
+ std::unique_ptr<PrometheusRenderer> prometheus_v2;
+ zlib_streams zlib;
+
+ // The mutex guards the members below it. No thread holds it while it calls
an API that can run a continuation.
+ std::mutex mutex;
+ snapshot_table snapshots;
+ std::vector<std::shared_ptr<stats_scrape>> waiters;
+ bool rendering = false;
+ stats_watchdog *watchdog = nullptr; // Set while
a request waits.
+};
+
+static std::shared_ptr<stats_instance>
+make_stats_instance(const stats_options &options)
+{
+ return std::make_shared<stats_instance>(options);
+}
+
+// Renders each wanted format once and compresses it for each wanted encoding.
+static void
+render_snapshots(stats_instance &instance, const bool
(&wanted)[FORMAT_COUNT][ENCODING_COUNT], snapshot_table &fresh)
+{
+ auto const now = std::chrono::steady_clock::now();
+
+ for (size_t f = 0; f < FORMAT_COUNT; ++f) {
+ if (std::find(std::begin(wanted[f]), std::end(wanted[f]), true) ==
std::end(wanted[f])) {
+ continue;
+ }
+
+ auto const format = static_cast<output_format_t>(f);
+ auto body = std::make_shared<stats_snapshot>();
+ render_state render{body->body, &instance.options};
+
+ render_stats(format, &render, instance.prometheus(format));
+ body->bytes = TSIOBufferReaderAvail(body->reader);
+ body->rendered = now;
+
+ for (size_t e = 1; e < ENCODING_COUNT; ++e) {
+ if (!wanted[f][e]) {
+ continue;
+ }
+
+ auto const encoding = static_cast<encoding_format_t>(e);
+ auto compressed = std::make_shared<stats_snapshot>();
+
+ if (compress_body(encoding, instance.zlib, body->reader,
compressed->body)) {
+ compressed->bytes = TSIOBufferReaderAvail(compressed->reader);
+ compressed->encoding = encoding;
+ compressed->rendered = now;
+ fresh[f][e] = std::move(compressed);
+ } else {
+ TSError("[%s] Cannot compress the stats, sending them uncompressed",
PLUGIN_NAME);
+ fresh[f][e] = body;
+ }
+ }
+ fresh[f][static_cast<size_t>(encoding_format_t::NONE)] = std::move(body);
+ }
+}
+
+static int watchdog_handler(TSCont contp, TSEvent event, void *edata);
+
+// Answers the waiting requests with a 503 when the render takes longer than
wait_timeout_ms.
+struct stats_watchdog {
+ explicit stats_watchdog(std::shared_ptr<stats_instance> inst)
+ : instance(std::move(inst)), cont(TSContCreate(watchdog_handler,
TSMutexCreate()))
+ {
+ TSContDataSet(cont, this);
+ }
+
+ std::shared_ptr<stats_instance> instance;
+ TSCont cont;
+
+ // The mutex of the continuation guards these.
+ TSAction action = nullptr;
+ bool fired = false;
+};
+
+static int
+watchdog_handler(TSCont contp, TSEvent /* event ATS_UNUSED */, void * /* edata
ATS_UNUSED */)
+{
+ auto *watchdog =
static_cast<stats_watchdog *>(TSContDataGet(contp));
+ stats_instance &instance = *watchdog->instance;
+ std::vector<std::shared_ptr<stats_scrape>> expired;
+
+ watchdog->fired = true;
+ {
+ std::lock_guard lock{instance.mutex};
+
+ // A render took this watchdog or replaced it, and that render destroys it.
+ if (instance.watchdog != watchdog) {
+ return 0;
+ }
+ instance.watchdog = nullptr;
+ expired.swap(instance.waiters);
+ }
+
+ Dbg(dbg_ctl, "Answering %zu requests with a 503 after %" PRId64 " ms without
a render of stats instance %p", expired.size(),
+ instance.options.wait_timeout_ms, &instance);
+ count(metrics().waiter_timeouts, expired.size());
+ for (auto const &scrape : expired) {
+ scrape->status = TS_HTTP_STATUS_SERVICE_UNAVAILABLE;
+ TSHttpTxnReenable(scrape->txnp, TS_EVENT_HTTP_CONTINUE);
+ }
+ TSContDestroy(contp);
+ delete watchdog;
+ return 0;
+}
+
+// Stops and destroys a watchdog that a render took from its instance.
+static void
+stop_watchdog(stats_watchdog *watchdog)
+{
+ TSMutex mutex = TSContMutexGet(watchdog->cont);
+
+ // Only the holder of its continuation's mutex can cancel an action.
+ TSMutexLock(mutex);
+ if (!watchdog->fired) {
+ TSActionCancel(watchdog->action);
+ }
+ TSMutexUnlock(mutex);
+ TSContDestroy(watchdog->cont);
+ delete watchdog;
+}
+
+// Creates a watchdog for the caller to publish under the instance mutex. The
calling thread holds the mutex of the watchdog
+// until arm_watchdog sets its action, so that a render that finishes first
can cancel that action.
+static stats_watchdog *
+new_watchdog(std::shared_ptr<stats_instance> instance)
+{
+ auto *watchdog = new stats_watchdog(std::move(instance));
+
+ TSMutexLock(TSContMutexGet(watchdog->cont));
+ return watchdog;
+}
+
+static void
+arm_watchdog(stats_watchdog *watchdog)
+{
+ TSMutex mutex = TSContMutexGet(watchdog->cont);
+
+ watchdog->action = TSContScheduleOnPool(watchdog->cont,
watchdog->instance->options.wait_timeout_ms, TS_THREAD_POOL_NET);
+ TSMutexUnlock(mutex);
+}
+
+static int render_handler(TSCont contp, TSEvent event, void *edata);
+
+// Each render gets a new continuation, because TSContScheduleOnPool locks the
mutex of the continuation on the calling
+// thread, and render_handler holds that mutex until it returns. A shared
continuation could block the ET_NET thread that
+// schedules the next render.
+static void
+schedule_render(std::shared_ptr<stats_instance> instance)
+{
+ TSCont contp = TSContCreate(render_handler, TSMutexCreate());
+
+ TSContDataSet(contp, new
std::shared_ptr<stats_instance>(std::move(instance)));
+ TSContScheduleOnPool(contp, 0, TS_THREAD_POOL_TASK);
+}
+
+static int
+render_handler(TSCont contp, TSEvent /* event ATS_UNUSED */, void * /* edata
ATS_UNUSED */)
+{
+ auto *ref =
static_cast<std::shared_ptr<stats_instance> *>(TSContDataGet(contp));
+ stats_instance &instance = **ref;
+ bool wanted[FORMAT_COUNT][ENCODING_COUNT] = {};
+ bool waiting = false;
+ snapshot_table fresh;
+
+ {
+ std::lock_guard lock{instance.mutex};
+
+ for (auto const &scrape : instance.waiters) {
+
wanted[static_cast<size_t>(scrape->format)][static_cast<size_t>(scrape->encoding)]
= true;
+ waiting
= true;
+ }
+ }
+ // The watchdog can answer every waiting request before the render starts.
+ if (waiting) {
+ static thread_local int64_t rest = 0;
+ int64_t const start = thread_cpu_ns();
+
+ render_snapshots(instance, wanted, fresh);
+ count(metrics().renders);
+ count_cpu_time(metrics().render_us, start, rest);
+ }
+
+ std::vector<std::shared_ptr<stats_scrape>> ready;
+ stats_watchdog *watchdog = nullptr;
+ stats_watchdog *replacement = nullptr;
+ bool again = false;
+
+ {
+ std::lock_guard lock{instance.mutex};
+
+ for (size_t f = 0; f < FORMAT_COUNT; ++f) {
+ for (size_t e = 0; e < ENCODING_COUNT; ++e) {
+ if (fresh[f][e] != nullptr) {
+ instance.snapshots[f][e] = fresh[f][e];
+ }
+ }
+ }
+
+ // As in scrape_wait, a request takes a render only if the render is
younger than max_age_ms when the request arrives.
+ auto const max_age =
std::chrono::milliseconds{instance.options.max_age_ms};
+ auto rendered = [&fresh, max_age](const
std::shared_ptr<stats_scrape> &scrape) -> snapshot_ptr {
+ auto const &snapshot =
fresh[static_cast<size_t>(scrape->format)][static_cast<size_t>(scrape->encoding)];
+
+ return snapshot != nullptr && scrape->arrived - snapshot->rendered <
max_age ? snapshot : nullptr;
+ };
+ auto answered = std::stable_partition(instance.waiters.begin(),
instance.waiters.end(),
+ [&rendered](const auto &scrape) {
return rendered(scrape) == nullptr; });
+
+ for (auto it = answered; it != instance.waiters.end(); ++it) {
+ (*it)->snapshot = rendered(*it);
+ ready.push_back(std::move(*it));
+ }
+ instance.waiters.erase(answered, instance.waiters.end());
+
+ // Requests for another format or encoding, and requests that arrived too
late for this render, wait for the next one. They
+ // get a new watchdog, because the current one can have started long
before they arrived.
+ again = !instance.waiters.empty();
+ if (again) {
+ replacement = new_watchdog(*ref);
+ } else {
+ instance.rendering = false;
+ }
+ watchdog = std::exchange(instance.watchdog, replacement);
+ }
+
+ if (waiting) {
+ Dbg(dbg_ctl, "Rendered stats instance %p for %zu waiting requests",
&instance, ready.size());
+ }
+ if (watchdog != nullptr) {
+ stop_watchdog(watchdog);
+ }
+ if (replacement != nullptr) {
+ arm_watchdog(replacement);
+ }
+ for (auto const &scrape : ready) {
+ TSHttpTxnReenable(scrape->txnp, TS_EVENT_HTTP_CONTINUE);
+ }
+ if (again) {
+ schedule_render(*ref);
+ }
+ delete ref;
+ TSContDestroy(contp);
+ return 0;
+}
+
+// Answers a request from a snapshot younger than max_age_ms. Otherwise the
request waits for a render, which starts unless
+// one is in flight. Either way, the transaction continues once it has its
answer.
+static void
+scrape_wait(const std::shared_ptr<stats_instance> &instance,
std::shared_ptr<stats_scrape> scrape)
+{
+ auto const now = std::chrono::steady_clock::now();
+ bool waits = false;
+ bool render = false;
+ stats_watchdog *watchdog = nullptr;
+
+ {
+ std::lock_guard lock{instance->mutex};
+ auto const &snapshot =
instance->snapshots[static_cast<size_t>(scrape->format)][static_cast<size_t>(scrape->encoding)];
+
+ if (snapshot != nullptr && now - snapshot->rendered <
std::chrono::milliseconds{instance->options.max_age_ms}) {
+ scrape->snapshot = snapshot;
+ } else {
+ scrape->arrived = now;
+ instance->waiters.push_back(scrape);
+ waits = true;
+ render = !std::exchange(instance->rendering, true);
+ if (instance->watchdog == nullptr) {
+ watchdog = instance->watchdog = new_watchdog(instance);
+ }
+ }
+ }
+
+ if (!waits) {
+ TSHttpTxnReenable(scrape->txnp, TS_EVENT_HTTP_CONTINUE);
+ return;
+ }
+
+ Dbg(dbg_ctl, "Waiting for a render of stats instance %p", instance.get());
+ if (watchdog != nullptr) {
+ arm_watchdog(watchdog);
+ }
+ if (render) {
+ schedule_render(instance);
+ }
+}
+
+// The intercept holds only its scrape, because intercept events can arrive
after the transaction and its remap rule are gone.
+struct scrape_intercept {
+ explicit scrape_intercept(std::shared_ptr<stats_scrape> s) :
scrape(std::move(s)) {}
+ ~scrape_intercept()
+ {
+ if (net_vc != nullptr) {
+ TSVConnClose(net_vc);
+ }
+ if (req_buffer != nullptr) {
+ TSIOBufferDestroy(req_buffer);
+ }
+ if (resp_buffer != nullptr) {
+ TSIOBufferDestroy(resp_buffer);
+ }
+ }
+ scrape_intercept(const scrape_intercept &) = delete;
+ scrape_intercept &operator=(const scrape_intercept &) = delete;
+
+ std::shared_ptr<stats_scrape> scrape;
+ TSVConn net_vc = nullptr;
+ TSVIO write_vio = nullptr;
+ TSIOBuffer req_buffer = nullptr;
+ TSIOBuffer resp_buffer = nullptr;
+ TSIOBufferReader resp_reader = nullptr;
+};
+
+static std::string
+scrape_response_header(const stats_scrape &scrape)
+{
+ if (!scrape.remap) {
+ if (scrape.status != TS_HTTP_STATUS_OK) {
+ return RESP_HEADER_UNAVAILABLE;
+ }
+
+ std::string header{"HTTP/1.0 200 OK\r\nContent-Type: "};
+
+ header.append(format_content_type(scrape.format)).append("\r\n");
+ if (auto const name = encoding_name(scrape.snapshot->encoding);
!name.empty()) {
+ header.append("Content-Encoding: ").append(name).append("\r\n");
+ }
+ header.append("Cache-Control: no-cache\r\n\r\n");
+ return header;
+ }
+
+ std::string header{"HTTP/1.1 "};
+
+ header.append(std::to_string(scrape.status)).append("
").append(TSHttpHdrReasonLookup(scrape.status)).append("\r\n");
+ if (scrape.status == TS_HTTP_STATUS_OK) {
+ header.append("Content-Type:
").append(format_content_type(scrape.format)).append("\r\n");
+ }
+ header.append("Cache-Control: no-store\r\n");
+ header.append(STATS_FORMAT_FIELD).append(":
").append(format_name(scrape.format)).append("\r\n");
+ if (scrape.snapshot != nullptr || scrape.status != TS_HTTP_STATUS_OK) {
+ header.append("Content-Length: ")
+ .append(std::to_string(scrape.snapshot != nullptr ?
scrape.snapshot->bytes : 0))
+ .append("\r\n");
+ }
+ header.append("\r\n");
+ return header;
+}
+
+static void
+scrape_send_response(TSCont contp, scrape_intercept *intercept)
+{
+ const stats_scrape &scrape = *intercept->scrape;
+ std::string const header = scrape_response_header(scrape);
+
+ intercept->resp_buffer = TSIOBufferCreate();
+ intercept->resp_reader = TSIOBufferReaderAlloc(intercept->resp_buffer);
+
+ int64_t bytes = TSIOBufferWrite(intercept->resp_buffer, header.data(),
header.size());
+
+ if (scrape.snapshot != nullptr) {
+ // This shares the blocks of the snapshot rather than copying its data.
+ bytes += TSIOBufferCopy(intercept->resp_buffer, scrape.snapshot->reader,
scrape.snapshot->bytes, 0);
+ }
+ TSVConnShutdown(intercept->net_vc, 1, 0);
+ intercept->write_vio = TSVConnWrite(intercept->net_vc, contp,
intercept->resp_reader, bytes);
+ count(metrics().requests);
+ count(metrics().bytes_out, bytes);
+}
+
+static int
+scrape_intercept_handler(TSCont contp, TSEvent event, void *edata)
+{
+ static thread_local int64_t rest = 0;
+ int64_t const start = thread_cpu_ns();
+ auto *intercept = static_cast<scrape_intercept
*>(TSContDataGet(contp));
+
+ switch (event) {
+ case TS_EVENT_NET_ACCEPT:
+ intercept->net_vc = static_cast<TSVConn>(edata);
+ intercept->req_buffer = TSIOBufferCreate();
+ TSVConnRead(intercept->net_vc, contp, intercept->req_buffer, INT64_MAX);
+ break;
+ case TS_EVENT_VCONN_READ_READY:
+ scrape_send_response(contp, intercept);
+ break;
+ case TS_EVENT_VCONN_WRITE_READY:
+ TSVIOReenable(intercept->write_vio);
+ break;
+ case TS_EVENT_NET_ACCEPT_FAILED:
+ case TS_EVENT_VCONN_EOS:
+ case TS_EVENT_ERROR:
+ case TS_EVENT_VCONN_INACTIVITY_TIMEOUT:
+ case TS_EVENT_VCONN_ACTIVE_TIMEOUT:
+ case TS_EVENT_VCONN_WRITE_COMPLETE:
+ Dbg(dbg_ctl, "Intercept finished on %s", TSHttpEventNameLookup(event));
+ delete intercept;
+ TSContDestroy(contp);
+ break;
+ default:
+ TSReleaseAssert(!"Unexpected Event");
+ }
+ count_cpu_time(metrics().intercept_us, start, rest);
+ return 0;
+}
+
+// The global plugin intercepts the request in its read request header hook,
then waits there for the stats.
+static void
+serve_global_scrape(TSHttpTxn txnp, output_format_t format, encoding_format_t
encoding)
+{
+ auto scrape = std::make_shared<stats_scrape>(txnp, format,
encoding, false);
+ TSCont intercept_cont = TSContCreate(scrape_intercept_handler,
TSMutexCreate());
+
+ TSContDataSet(intercept_cont, new scrape_intercept(scrape));
+ TSHttpTxnIntercept(intercept_cont, txnp);
+ scrape_wait(*global_instance, std::move(scrape));
+}
+
+//
+// Remap plugin.
+//
+
+// The transaction hooks of a GET request to a remap rule.
+struct scrape_txn {
+ std::shared_ptr<stats_instance> instance;
+ std::shared_ptr<stats_scrape> scrape;
+};
+
+static int
+scrape_txn_handler(TSCont contp, TSEvent event, void *edata)
+{
+ auto *txn = static_cast<scrape_txn *>(TSContDataGet(contp));
+
+ if (event == TS_EVENT_HTTP_CACHE_LOOKUP_COMPLETE) {
+ scrape_wait(txn->instance, txn->scrape);
+ return 0;
+ }
+ if (event == TS_EVENT_HTTP_TXN_CLOSE) {
+ delete txn;
+ TSContDestroy(contp);
+ }
+ TSHttpTxnReenable(static_cast<TSHttpTxn>(edata), TS_EVENT_HTTP_CONTINUE);
+ return 0;
+}
+
+static void
+add_allow_field(TSHttpTxn txnp)
+{
+ TSMBuffer bufp;
+ TSMLoc hdr_loc;
+
+ if (TSHttpTxnClientRespGet(txnp, &bufp, &hdr_loc) == TS_SUCCESS) {
+ TSMLoc field_loc;
+
+ // A remap ACL filter that denies the request replaces the 405 with a 403.
+ if (TSHttpHdrStatusGet(bufp, hdr_loc) == TS_HTTP_STATUS_METHOD_NOT_ALLOWED
&&
+ TSMimeHdrFieldCreateNamed(bufp, hdr_loc, TS_MIME_FIELD_ALLOW,
TS_MIME_LEN_ALLOW, &field_loc) == TS_SUCCESS) {
+ TSMimeHdrFieldValueStringSet(bufp, hdr_loc, field_loc, -1,
ALLOWED_METHODS.data(), ALLOWED_METHODS.size());
+ TSMimeHdrFieldAppend(bufp, hdr_loc, field_loc);
+ TSHandleMLocRelease(bufp, hdr_loc, field_loc);
+ }
+ TSHandleMLocRelease(bufp, TS_NULL_MLOC, hdr_loc);
+ }
+}
+
+static int
+allow_handler(TSCont contp, TSEvent event, void *edata)
+{
+ auto txnp = static_cast<TSHttpTxn>(edata);
+
+ if (event == TS_EVENT_HTTP_SEND_RESPONSE_HDR) {
+ add_allow_field(txnp);
+ } else if (event == TS_EVENT_HTTP_TXN_CLOSE) {
+ TSContDestroy(contp);
+ }
+ TSHttpTxnReenable(txnp, TS_EVENT_HTTP_CONTINUE);
+ return 0;
+}
+
+// The settings that a remap rule or its configuration file can set.
+struct stats_settings {
+ std::optional<output_format_t> format;
+ std::optional<int64_t> max_age_ms;
+ std::optional<int64_t> wait_timeout_ms;
+ std::optional<bool> prometheus_help;
+ std::shared_ptr<const PrometheusRules> rules;
+};
+
+static std::string
+yaml_line(const YAML::Node &node)
+{
+ auto const mark = node.Mark();
+
+ return mark.is_null() ? std::string{} : "line " + std::to_string(mark.line +
1) + ": ";
+}
+
+// Reads the configuration file of a remap rule. Returns an empty string, or
a description of the first error.
+static std::string
+load_settings_file(const std::string &path, stats_settings &settings)
+{
+ std::ifstream file{path};
+
+ if (!file) {
+ return std::string{"cannot open the file: "} + strerror(errno);
+ }
+ try {
+ YAML::Node const root = YAML::Load(file);
+
+ if (!root.IsMap()) {
+ return yaml_line(root) + "the file must be a map";
+ }
+ for (auto const &item : root) {
+ std::string const key = item.first.as<std::string>();
+ YAML::Node const value = item.second;
+
+ if (key == "format") {
+ output_format_t format;
+
+ if (!value.IsScalar() || !parse_format(value.Scalar(), format)) {
+ return yaml_line(value) + "format must be json, csv, prometheus or
prometheus_v2";
+ }
+ settings.format = format;
+ } else if (key == "render") {
+ if (!value.IsMap()) {
+ return yaml_line(value) + "render must be a map";
+ }
+ for (auto const &setting : value) {
+ std::string const name = setting.first.as<std::string>();
+ std::string const text = setting.second.IsScalar() ?
setting.second.Scalar() : std::string{};
+ int64_t number = 0;
+ int64_t min = 0;
+ bool valid = false;
+
+ if (name == "max_age_ms") {
+ valid = parse_integer(text, min, MAX_MILLISECONDS,
number);
+ settings.max_age_ms = number;
+ } else if (name == "wait_timeout_ms") {
+ min = 1;
+ valid = parse_integer(text, min,
MAX_MILLISECONDS, number);
+ settings.wait_timeout_ms = number;
+ } else {
+ return yaml_line(setting.first) + "unknown key render." + name;
+ }
+ if (!valid) {
+ return yaml_line(setting.second) + "render." + name + " must be an
integer from " + std::to_string(min) + " to " +
+ std::to_string(MAX_MILLISECONDS) + ", not '" + text + "'";
+ }
+ }
+ } else if (key == "prometheus") {
+ auto rules = std::make_shared<PrometheusRules>();
+
+ if (std::string error = rules->load(value); !error.empty()) {
+ return error;
+ }
+ settings.prometheus_help = rules->help();
+ settings.rules = std::move(rules);
+ } else {
+ return yaml_line(item.first) + "unknown key " + key;
+ }
+ }
+ } catch (const YAML::Exception &e) {
+ return e.what();
+ }
+ return {};
+}
+
+// Registers the configuration file of a remap rule as a child of the remap
configuration. After a change to the file, a
+// configuration reload loads the remap configuration again, and with it the
file.
+static void
+watch_config_file(const std::string &path)
+{
+ TSMgmtString parent = nullptr;
+
+ if (TSMgmtStringGet("proxy.config.url_remap.filename", &parent) ==
TS_SUCCESS) {
+ TSMgmtConfigFileAdd(parent, path.c_str());
+ } else {
+ TSWarning("[%s] Cannot read proxy.config.url_remap.filename, so a
configuration reload does not detect a change to %s",
+ PLUGIN_NAME, path.c_str());
+ }
+ TSfree(parent);
+}
+
+TSReturnCode
+TSRemapInit(TSRemapInterface * /* api_info ATS_UNUSED */, char * /* errbuf
ATS_UNUSED */, int /* errbuf_size ATS_UNUSED */)
+{
+ metrics();
+ return TS_SUCCESS;
+}
+
+TSReturnCode
+TSRemapNewInstance(int argc, char *argv[], void **ih, char *errbuf, int
errbuf_size)
+{
+ static const struct option longopts[] = {
+ {"format", required_argument, nullptr, 'f'},
+ {"integer-counters", no_argument, nullptr, 'i'},
+ {"wrap-counters", no_argument, nullptr, 'w'},
+ {"no-prometheus-help", no_argument, nullptr, 'n'},
+ {"max-age-ms", required_argument, nullptr, 'a'},
+ {"wait-timeout-ms", required_argument, nullptr, 't'},
+ {"config", required_argument, nullptr, 'c'},
+ {"on-config-error", required_argument, nullptr, 'e'},
+ {nullptr, 0, nullptr, 0 }
+ };
+ stats_options options;
+ stats_settings rule;
+ std::string config_path;
+ bool answer_errors = false;
+ int64_t number = 0;
+
+ options.prometheus_epoch = false;
+
+ // argv[0] is the "from" URL. Skip it so that the "to" URL poses as the
program name.
+ --argc;
+ ++argv;
+ for (int opt, option_index = 0; (opt = getopt_long(argc, argv, "", longopts,
&option_index)) != -1;) {
Review Comment:
`getopt_long()` retains process-global parser state, but this entry point
runs once per remap rule. After the first instance (or `TSPluginInit`) leaves
`optind` at the end of its argv, later rules can skip options such as
`--format` and `--config` or parse from the wrong position. Reset getopt before
every instance parse, as the ESI and statichit remap plugins do.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]