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]

Reply via email to