Skip to content
Open
Show file tree
Hide file tree
Changes from all 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