Skip to content
Open
Show file tree
Hide file tree
Changes from 2 commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -391,6 +391,7 @@ add_library(moqx_core STATIC
src/admin/ConfigHandler.cpp
src/admin/ConnectionLogsHandler.cpp
src/admin/MetricsHandler.cpp
src/admin/QLogCaptureHandler.cpp
src/admin/TrackMetricsHandler.cpp
src/admin/StateHandler.cpp
src/admin/StateResponse.cpp
Expand All @@ -412,6 +413,7 @@ add_library(moqx_core STATIC
src/relay/PublisherCrossExecFilter.cpp
src/relay/SubscriberCrossExecFilter.cpp
src/logging/LogSetup.cpp
src/logging/QLogCapture.cpp
)

target_include_directories(moqx_core
Expand Down
28 changes: 25 additions & 3 deletions docs/logging.md
Original file line number Diff line number Diff line change
Expand Up @@ -188,13 +188,35 @@ XLOG and qlog have different jobs:
- **XLOG** is the ad-hoc operator/dev surface — code-embedded clues, controlled by `--logging`. Good for "what just went wrong."
- **qlog** is the IETF-standard structured QUIC log — JSON consumed by external tooling like [qvis](https://qvis.quictools.info/) for packet timelines, congestion control, stream events. Good for "trace a single session" investigations.

Per-session qlog enablement and random sampling are planned (the moqx admin surface will let you opt one connection in at a time). For now, enable qlog at the picoquic level for all connections:
Independent channels — enabling one doesn't suppress the other.

qlog covers mvfst listeners only; picoquic listeners write none. Setting a directory enables it:

```yaml
logging:
qlog:
dir: /var/log/moqx/qlog
sample_rate: 0 # fraction of all new connections; 0 = on-demand captures only
```

Each connection's file is `<dir>/<dcid>.qlog`, where `dcid` is the client's original destination connection ID. Every logged event is serialized on the connection's IO thread, so keep `sample_rate` low on a loaded relay and prefer on-demand captures.

### On-demand capture

With a directory set, the admin API qlogs the next N new connections:

```bash
moqx --config c.yaml --qlog-dir /tmp/qlogs --logging=quic.picoquic=INFO
curl -X POST 'localhost:8000/qlog/capture?count=2&seconds=60&mode=cc' # arm
curl localhost:8000/qlog/capture # status + newest files
curl -o c.qlog 'localhost:8000/logs?connection_id=<id>&type=qlog' # fetch one
curl -X DELETE localhost:8000/qlog/capture # disarm
```

Independent channels — enabling one doesn't suppress the other.
- `count` (default 1, max 64) and `seconds` (default 60, max 600) bound the capture; arming again replaces it.
- `mode=cc` (default) keeps congestion control, RTT, loss and pacing events and drops per-packet and per-stream events. `mode=full` keeps everything; use it for short windows.
- Each captured connection logs `qlog capture: logging a new connection` at INFO, so captures line up with the rest of the relay log; `GET /qlog/capture` lists their files by connection ID.

Open the files in [qvis](https://qvis.quictools.info/); it parses them in the browser.

## Troubleshooting

Expand Down
34 changes: 24 additions & 10 deletions src/MoqxRelayServer.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@
*/

#include "MoqxRelayServer.h"
#include "logging/CcQLogger.h"
#include "stats/EventBaseStatsCollector.h"
#include "stats/QuicStatsCollector.h"
#include <moxygen/MoQRelaySession.h>
Expand Down Expand Up @@ -300,6 +301,28 @@ std::shared_ptr<quic::QLogger> MoqxRelayServer::makeQLogger(quic::VantagePoint v
if (qlogDir_.empty()) {
return nullptr;
}
// streaming=true → AsyncFileWriter runs in its own background thread;
// the event-loop thread only pays for JSON serialisation + queue push.
auto fileQLogger = [&] {
return std::make_shared<quic::FileQLogger>(
vantagePoint,
"MOQT",
qlogDir_,
/*prettyJson=*/false,
/*streaming=*/true,
/*compress=*/false
);
};
if (qlogCapture_) {
if (auto mode = qlogCapture_->take()) {
XLOG(INFO) << "qlog capture: logging a new connection (mode="
<< logging::QLogCapture::modeName(*mode) << ")";
if (*mode == logging::QLogCapture::Mode::Cc) {
return std::make_shared<logging::CcQLogger>(vantagePoint, qlogDir_);
}
return fileQLogger();
}
}
if (qlogSampleRate_ <= 0.0f) {
return nullptr;
}
Expand All @@ -311,16 +334,7 @@ std::shared_ptr<quic::QLogger> MoqxRelayServer::makeQLogger(quic::VantagePoint v
return nullptr;
}
}
// streaming=true → AsyncFileWriter runs in its own background thread;
// the event-loop thread only pays for JSON serialisation + queue push.
return std::make_shared<quic::FileQLogger>(
vantagePoint,
"MOQT",
qlogDir_,
/*prettyJson=*/false,
/*streaming=*/true,
/*compress=*/false
);
return fileQLogger();
}

} // namespace openmoq::moqx
6 changes: 6 additions & 0 deletions src/MoqxRelayServer.h
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@

#include "MoqxRelayContext.h"
#include "config/Config.h"
#include "logging/QLogCapture.h"
#include "stats/StatsRegistry.h"
#include <folly/executors/IOThreadPoolExecutor.h>
#include <moxygen/MoQServer.h>
Expand Down Expand Up @@ -41,6 +42,10 @@ class MoqxRelayServer : public moxygen::MoQServer {
qlogSampleRate_ = cfg.sampleRate;
}

void setQLogCapture(std::shared_ptr<logging::QLogCapture> capture) {
qlogCapture_ = std::move(capture);
}

// Preferred entry point: binds the address from the stored ListenerConfig.
void start();

Expand Down Expand Up @@ -72,6 +77,7 @@ class MoqxRelayServer : public moxygen::MoQServer {
bool stopped_{false};
std::string qlogDir_;
float qlogSampleRate_{0.0f};
std::shared_ptr<logging::QLogCapture> qlogCapture_;
};

} // namespace openmoq::moqx
4 changes: 3 additions & 1 deletion src/MoqxServerFactory.h
Original file line number Diff line number Diff line change
Expand Up @@ -30,7 +30,8 @@ inline std::shared_ptr<moxygen::MoQServerBase> makeRelayServer(
folly::IOThreadPoolExecutor* ioExecutor,
std::shared_ptr<stats::StatsRegistry> statsRegistry,
std::shared_ptr<moxygen::MLoggerFactory> mlogFactory = nullptr,
const config::QLogConfig* qlogConfig = nullptr
const config::QLogConfig* qlogConfig = nullptr,
std::shared_ptr<logging::QLogCapture> qlogCapture = nullptr
) {
if (listenerCfg.quicStack == config::QuicStack::Picoquic) {
auto server =
Expand All @@ -54,6 +55,7 @@ inline std::shared_ptr<moxygen::MoQServerBase> makeRelayServer(
}
if (qlogConfig && !qlogConfig->dir.empty()) {
server->setQLogConfig(*qlogConfig);
server->setQLogCapture(std::move(qlogCapture));
}
return server;
}
Expand Down
19 changes: 19 additions & 0 deletions src/admin/AdminResponse.h
Original file line number Diff line number Diff line change
Expand Up @@ -6,9 +6,11 @@

#pragma once

#include <cstdint>
#include <optional>
#include <string>

#include <folly/Conv.h>
#include <folly/io/IOBuf.h>
#include <proxygen/httpserver/ResponseBuilder.h>
#include <proxygen/httpserver/ResponseHandler.h>
Expand Down Expand Up @@ -42,4 +44,21 @@ boolQueryParam(const proxygen::HTTPMessage& req, const std::string& name, bool d
return std::nullopt;
}

// Returns nullopt if the param is present but not an integer within [1, max].
inline std::optional<uint32_t> boundedQueryParam(
const proxygen::HTTPMessage& req,
const std::string& name,
uint32_t defaultValue,
uint32_t max
) {
if (!req.hasQueryParam(name)) {
return defaultValue;
}
auto value = folly::tryTo<uint32_t>(req.getDecodedQueryParam(name));
if (!value || *value < 1 || *value > max) {
return std::nullopt;
}
return *value;
}

} // namespace openmoq::moqx::admin
210 changes: 210 additions & 0 deletions src/admin/QLogCaptureHandler.cpp
Original file line number Diff line number Diff line change
@@ -0,0 +1,210 @@
/*
* Copyright (c) OpenMOQ contributors.
* This source code is licensed under the Apache 2.0 license found in the
* LICENSE file in the root directory of this source tree.
*/

#include "admin/QLogCaptureHandler.h"

#include <algorithm>
#include <chrono>
#include <filesystem>
#include <vector>

#include <sys/stat.h>

#include <folly/Conv.h>
#include <folly/io/Cursor.h>
#include <folly/io/IOBuf.h>
#include <folly/io/IOBufQueue.h>
#include <proxygen/httpserver/ResponseBuilder.h>
#include <proxygen/lib/http/HTTPMessage.h>

#include "admin/AdminResponse.h"
#include "admin/AdminServer.h"
#include "admin/JsonWriter.h"

namespace openmoq::moqx::admin {

namespace {

using logging::QLogCapture;

constexpr size_t kMaxListedFiles = 100;

struct QLogFile {
std::string connectionId;
int64_t bytes;
int64_t modified;
};

// Newest first, at most kMaxListedFiles.
std::vector<QLogFile> listQLogFiles(const std::string& dir) {
std::vector<QLogFile> files;
std::error_code ec;
for (const auto& entry : std::filesystem::directory_iterator(dir, ec)) {
const auto& path = entry.path();
if (path.extension() != ".qlog") {
continue;
}
struct stat st{};
if (::stat(path.c_str(), &st) != 0 || !S_ISREG(st.st_mode)) {
continue;
}
files.push_back({path.stem().string(), static_cast<int64_t>(st.st_size), st.st_mtime});
}
const auto kept = std::min(files.size(), kMaxListedFiles);
std::partial_sort(
files.begin(),
files.begin() + kept,
files.end(),
[](const auto& a, const auto& b) { return a.modified > b.modified; }
);
files.resize(kept);
return files;
}

// Capture status; with a directory, also the newest qlog files in it.
std::unique_ptr<folly::IOBuf>
statusBody(const QLogCapture::Status& status, const std::string* dir = nullptr) {
folly::IOBufQueue queue{folly::IOBufQueue::cacheChainLength()};
folly::io::QueueAppender app{&queue, 1024};
JsonWriter w{app};
w.beginObject();
w.field("armed", status.armed);
w.field("mode", QLogCapture::modeName(status.mode));
w.field("remaining", uint64_t{status.remaining});
w.field("captured", uint64_t{status.captured});
const int64_t expires =
std::chrono::duration_cast<std::chrono::seconds>(status.expiresAt.time_since_epoch()).count();
w.key("expires_at");
if (expires > 0) {
w.intVal(expires);
} else {
w.nullVal();
}
if (dir) {
w.field("dir", *dir);
w.key("files");
w.beginArray();
for (const auto& f : listQLogFiles(*dir)) {
w.beginObject();
w.field("connection_id", f.connectionId);
w.field("bytes", f.bytes);
w.field("modified", f.modified);
w.endObject();
}
w.endArray();
}
w.endObject();
app.write(static_cast<uint8_t>('\n'));
return queue.move();
}

void sendJson(proxygen::ResponseHandler* downstream, std::unique_ptr<folly::IOBuf> body) {
proxygen::ResponseBuilder(downstream)
.status(200, proxygen::HTTPMessage::getDefaultReason(200))
.header("Content-Type", "application/json")
.body(std::move(body))
.sendWithEOM();
}

} // namespace

void registerQLogCaptureRoutes(
AdminServer& adminServer,
std::shared_ptr<QLogCapture> capture,
std::string qlogDir
) {
static const std::string kNotConfigured = "qlog is not configured (logging.qlog.dir)\n";

adminServer.addRoute(
"POST",
"/qlog/capture",
[capture](
std::unique_ptr<proxygen::HTTPMessage> req,
std::unique_ptr<folly::IOBuf> /*body*/,
proxygen::ResponseHandler* downstream,
folly::CancellationToken /*cancelToken*/,
const std::shared_ptr<EgressGate>& /*egress*/
) {
if (!capture) {
sendError(downstream, 503, kNotConfigured);
return;
}
auto count = boundedQueryParam(*req, "count", 1, QLogCapture::kMaxCount);
if (!count) {
sendError(
downstream,
400,
folly::to<std::string>("count must be 1-", QLogCapture::kMaxCount, "\n")
);
return;
}
auto seconds = boundedQueryParam(
*req,
"seconds",
60,
static_cast<uint32_t>(QLogCapture::kMaxDuration.count())
);
if (!seconds) {
sendError(
downstream,
400,
folly::to<std::string>("seconds must be 1-", QLogCapture::kMaxDuration.count(), "\n")
);
return;
}
auto mode = QLogCapture::Mode::Cc;
if (req->hasQueryParam("mode")) {
auto parsed = QLogCapture::parseMode(req->getDecodedQueryParam("mode"));
if (!parsed) {
sendError(downstream, 400, "mode must be cc or full\n");
return;
}
mode = *parsed;
}
capture->arm(*count, std::chrono::seconds(*seconds), mode);
sendJson(downstream, statusBody(capture->status()));
}
);

adminServer.addRoute(
"DELETE",
"/qlog/capture",
[capture](
std::unique_ptr<proxygen::HTTPMessage> /*req*/,
std::unique_ptr<folly::IOBuf> /*body*/,
proxygen::ResponseHandler* downstream,
folly::CancellationToken /*cancelToken*/,
const std::shared_ptr<EgressGate>& /*egress*/
) {
if (!capture) {
sendError(downstream, 503, kNotConfigured);
return;
}
capture->disarm();
sendJson(downstream, statusBody(capture->status()));
}
);

adminServer.addRoute(
"GET",
"/qlog/capture",
[capture, qlogDir = std::move(qlogDir)](
std::unique_ptr<proxygen::HTTPMessage> /*req*/,
std::unique_ptr<folly::IOBuf> /*body*/,
proxygen::ResponseHandler* downstream,
folly::CancellationToken /*cancelToken*/,
const std::shared_ptr<EgressGate>& /*egress*/
) {
if (!capture) {
sendError(downstream, 503, kNotConfigured);
return;
}
sendJson(downstream, statusBody(capture->status(), &qlogDir));
}
);
}

} // namespace openmoq::moqx::admin
Loading
Loading