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
1 change: 1 addition & 0 deletions ChangeLog.md
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

### Added
- `--write-cache=on|off` for `rawstor-vhost` and `rawstor-vhost-qemu` (default `off`, write-through): advertises `VIRTIO_BLK_F_CONFIG_WCE` and honors the guest live-toggling it via `SET_CONFIG`. With write-cache off, every write is made durable (`sync=true`) since the guest treats a completed write as already durable and won't issue a `FLUSH`.
- `librawstor`'s telemetry report now also covers `blk`-backed sessions (e.g. `file://`), with their own submission/callback-latency and top-N slowest-requests section, not just `ost://` sessions.

### Changed
- `rawstor_object_pwrite()`/`rawstor_object_pwritev()` gained a `sync` parameter — when true, the write is durable on stable storage by the time the callback reports success. Breaking C API change; existing callers need to pass a `sync` argument (`false` preserves the old behavior).
Expand Down
132 changes: 127 additions & 5 deletions src/blk_session.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@

#include "object.hpp"
#include "target.hpp"
#include "telemetry.hpp"

#include <rawio/awaitable.hpp>

Expand Down Expand Up @@ -70,7 +71,36 @@ rawstd::Task<size_t> Session::pread(void* buf, size_t size, off_t offset) {
(intmax_t)offset
);

co_return co_await _queue.pread(fd(), buf, size, offset);
// telemetry: a blk session has no round-trip to a remote peer, so
// there is no rtt to report -- just slat (submitting to the queue)
// and clat (the usually negligible gap between completion and the
// caller resuming). See telemetry.hpp's blk namespace doc comment.
rawstor::telemetry::TimePoint t_created = rawstor::telemetry::now();
rawstor::telemetry::blk::op_started();

rawio::Awaitable<size_t> awaitable = _queue.pread(fd(), buf, size, offset);
rawstor::telemetry::TimePoint t_submitted = rawstor::telemetry::now();
rawstor::telemetry::blk::record_slat(t_submitted - t_created);
Comment on lines +78 to +83

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[medium] blk::op_started() runs before the queue submit call, so a synchronous submission throw leaks the in-flight counter.

Why

uring_queue.cpp's pread/preadv/pwrite/pwritev/fsync throw ENOBUFS synchronously (ring full) before returning the Awaitable. op_started() already ran; the try/catch only wraps co_await awaitable, not its construction, so op_finished() never runs.

Suggested change
rawstor::telemetry::TimePoint t_created = rawstor::telemetry::now();
rawstor::telemetry::blk::op_started();
rawio::Awaitable<size_t> awaitable = _queue.pread(fd(), buf, size, offset);
rawstor::telemetry::TimePoint t_submitted = rawstor::telemetry::now();
rawstor::telemetry::blk::record_slat(t_submitted - t_created);
rawstor::telemetry::TimePoint t_created = rawstor::telemetry::now();
rawio::Awaitable<size_t> awaitable = _queue.pread(fd(), buf, size, offset);
rawstor::telemetry::blk::op_started();
rawstor::telemetry::TimePoint t_submitted = rawstor::telemetry::now();
rawstor::telemetry::blk::record_slat(t_submitted - t_created);


size_t result;
try {
result = co_await awaitable;
} catch (...) {
// clat/lat only mean something for an op that actually completed
// -- a failed submission has nothing useful to measure past slat.
rawstor::telemetry::blk::op_finished();
throw;
}

rawstor::telemetry::TimePoint t_now = rawstor::telemetry::now();
rawstor::telemetry::TimePoint clat = t_now - t_submitted;
rawstor::telemetry::blk::record_clat(clat);
rawstor::telemetry::blk::record_op(
t_now - t_created, t_submitted - t_created, clat, "pread", size, offset
);
rawstor::telemetry::blk::op_finished();

co_return result;
}

rawstd::Task<size_t>
Expand All @@ -80,7 +110,30 @@ Session::preadv(iovec* iov, unsigned int niov, size_t size, off_t offset) {
(intmax_t)offset
);

co_return co_await _queue.preadv(fd(), iov, niov, offset);
rawstor::telemetry::TimePoint t_created = rawstor::telemetry::now();
rawstor::telemetry::blk::op_started();

rawio::Awaitable<size_t> awaitable = _queue.preadv(fd(), iov, niov, offset);
rawstor::telemetry::TimePoint t_submitted = rawstor::telemetry::now();
rawstor::telemetry::blk::record_slat(t_submitted - t_created);

size_t result;
try {
result = co_await awaitable;
} catch (...) {
rawstor::telemetry::blk::op_finished();
throw;
}

rawstor::telemetry::TimePoint t_now = rawstor::telemetry::now();
rawstor::telemetry::TimePoint clat = t_now - t_submitted;
rawstor::telemetry::blk::record_clat(clat);
rawstor::telemetry::blk::record_op(
t_now - t_created, t_submitted - t_created, clat, "preadv", size, offset
);
rawstor::telemetry::blk::op_finished();

co_return result;
}

rawstd::Task<size_t>
Expand All @@ -90,7 +143,31 @@ Session::pwrite(const void* buf, size_t size, off_t offset, bool sync) {
fd(), size, (intmax_t)offset, sync
);

co_return co_await _queue.pwrite(fd(), buf, size, offset, sync);
rawstor::telemetry::TimePoint t_created = rawstor::telemetry::now();
rawstor::telemetry::blk::op_started();

rawio::Awaitable<size_t> awaitable =
_queue.pwrite(fd(), buf, size, offset, sync);
rawstor::telemetry::TimePoint t_submitted = rawstor::telemetry::now();
rawstor::telemetry::blk::record_slat(t_submitted - t_created);

size_t result;
try {
result = co_await awaitable;
} catch (...) {
rawstor::telemetry::blk::op_finished();
throw;
}

rawstor::telemetry::TimePoint t_now = rawstor::telemetry::now();
rawstor::telemetry::TimePoint clat = t_now - t_submitted;
rawstor::telemetry::blk::record_clat(clat);
rawstor::telemetry::blk::record_op(
t_now - t_created, t_submitted - t_created, clat, "pwrite", size, offset
);
rawstor::telemetry::blk::op_finished();

co_return result;
}

rawstd::Task<size_t> Session::pwritev(
Expand All @@ -101,13 +178,58 @@ rawstd::Task<size_t> Session::pwritev(
fd(), size, (intmax_t)offset, sync
);

co_return co_await _queue.pwritev(fd(), iov, niov, offset, sync);
rawstor::telemetry::TimePoint t_created = rawstor::telemetry::now();
rawstor::telemetry::blk::op_started();

rawio::Awaitable<size_t> awaitable =
_queue.pwritev(fd(), iov, niov, offset, sync);
rawstor::telemetry::TimePoint t_submitted = rawstor::telemetry::now();
rawstor::telemetry::blk::record_slat(t_submitted - t_created);

size_t result;
try {
result = co_await awaitable;
} catch (...) {
rawstor::telemetry::blk::op_finished();
throw;
}

rawstor::telemetry::TimePoint t_now = rawstor::telemetry::now();
rawstor::telemetry::TimePoint clat = t_now - t_submitted;
rawstor::telemetry::blk::record_clat(clat);
rawstor::telemetry::blk::record_op(
t_now - t_created, t_submitted - t_created, clat, "pwritev", size,
offset
);
rawstor::telemetry::blk::op_finished();

co_return result;
}

rawstd::Task<void> Session::flush() {
rawstd_debug("%s(): fd = %d\n", __FUNCTION__, fd());

co_await _queue.fsync(fd(), /*datasync=*/true);
rawstor::telemetry::TimePoint t_created = rawstor::telemetry::now();
rawstor::telemetry::blk::op_started();

rawio::Awaitable<int> awaitable = _queue.fsync(fd(), /*datasync=*/true);
rawstor::telemetry::TimePoint t_submitted = rawstor::telemetry::now();
rawstor::telemetry::blk::record_slat(t_submitted - t_created);

try {
co_await awaitable;
} catch (...) {
rawstor::telemetry::blk::op_finished();
throw;
}

rawstor::telemetry::TimePoint t_now = rawstor::telemetry::now();
rawstor::telemetry::TimePoint clat = t_now - t_submitted;
rawstor::telemetry::blk::record_clat(clat);
rawstor::telemetry::blk::record_op(
t_now - t_created, t_submitted - t_created, clat, "flush", 0, 0
);
rawstor::telemetry::blk::op_finished();
}

} // namespace blk
Expand Down
10 changes: 6 additions & 4 deletions src/connection.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -41,10 +41,12 @@ class Connection final {

// Every data-path/metadata method's terminal path -- success or final
// failure -- runs through here exactly once; records the cross-retry
// call-to-completion latency. Per-attempt telemetry, including the
// top-N slowest-requests sample, lives in ost::SessionOp::_dispatch()
// instead -- Connection is transport-agnostic and has nothing else to
// report here.
// call-to-completion latency, shared across every backend. Per-attempt
// telemetry, including the top-N slowest-requests sample, is split by
// backend and lives in ost::SessionOp::_dispatch() / blk::Session's
// own data-path methods (telemetry::ost / telemetry::blk) instead --
// Connection is transport-agnostic and has nothing else to report
// here.
void _finish(rawstor::telemetry::TimePoint t_call);

// Shared retry-loop body for every data-path/metadata method: tries
Expand Down
14 changes: 7 additions & 7 deletions src/ost_session.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -149,7 +149,7 @@ class SessionOp {
// A string literal (e.g. "pread"/"pwrite"/"flush"), not owned; size
// and offset are 0 for ops without either (flush). Set once at
// construction by each subclass, purely for _dispatch()'s
// telemetry::record_op() call.
// telemetry::ost::record_op() call.
const char* _op_name;
size_t _op_size;
off_t _op_offset;
Expand Down Expand Up @@ -190,7 +190,7 @@ class SessionOp {
t_response_ready = rawstor::telemetry::now();
slat = _t_send_done - _t_created;
rtt = t_response_ready - _t_send_done;
rawstor::telemetry::record_rtt(rtt);
rawstor::telemetry::ost::record_rtt(rtt);
}

_result = result;
Expand All @@ -202,8 +202,8 @@ class SessionOp {
if (timed) {
rawstor::telemetry::TimePoint t_now = rawstor::telemetry::now();
rawstor::telemetry::TimePoint clat = t_now - t_response_ready;
rawstor::telemetry::record_clat(clat);
rawstor::telemetry::record_op(
rawstor::telemetry::ost::record_clat(clat);
rawstor::telemetry::ost::record_op(
t_now - _t_created, slat, rtt, clat, _op_name, _op_size,
_op_offset
);
Expand Down Expand Up @@ -249,7 +249,7 @@ class SessionOp {

if (!error) {
_t_send_done = rawstor::telemetry::now();
rawstor::telemetry::record_slat(_t_send_done - _t_created);
rawstor::telemetry::ost::record_slat(_t_send_done - _t_created);
} else {
_dispatch(0, error);
}
Expand Down Expand Up @@ -773,12 +773,12 @@ void Session::_add_op(const std::shared_ptr<SessionOp>& op) {
return;
}
_ops[op->cid()] = op;
rawstor::telemetry::op_started();
rawstor::telemetry::ost::op_started();
}

void Session::_remove_op(uint16_t cid) {
_ops.erase(cid);
rawstor::telemetry::op_finished();
rawstor::telemetry::ost::op_finished();
}

Session::Session(Private p, rawio::Queue& queue, const rawstd::URI& location) :
Expand Down
Loading
Loading