Skip to content
Merged
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
123 changes: 110 additions & 13 deletions src/drivers/adt/adt_driver.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -59,17 +59,19 @@ DbfFieldType classify_adt_field(std::uint16_t raw_type) {
}
}

// Retry-with-backoff lock helper, same pattern as CdxDriver.
// Retry-with-backoff lock helper, same pattern as CdxDriver. `kind` lets a
// reader take the header region shared, so open/refresh wait out an appender
// instead of failing with ERROR_LOCK_VIOLATION (mapped to 5000).
util::Result<platform::ByteLock>
acquire_with_retry_(platform::File& f,
std::uint64_t offset,
std::uint64_t length,
int max_retries = 200)
acquire_with_retry_(platform::File& f,
std::uint64_t offset,
std::uint64_t length,
platform::LockKind kind = platform::LockKind::Exclusive,
int max_retries = 200)
{
util::Error last_err{};
for (int i = 0; i < max_retries; ++i) {
auto lk = platform::ByteLock::try_acquire(f, offset, length,
platform::LockKind::Exclusive);
auto lk = platform::ByteLock::try_acquire(f, offset, length, kind);
if (lk) return std::move(lk).value();
last_err = lk.error();
std::this_thread::sleep_for(
Expand All @@ -91,6 +93,13 @@ AdtDriver::open(const std::string& path, DriverOpenMode mode) {
if (!fres) return fres.error();
file_ = std::move(fres).value();

// Coordinate the header read with concurrent appenders (who hold 0..399
// exclusive). Shared lock + retry so open does not fail with
// ERROR_LOCK_VIOLATION while another connection appends.
auto hdr_lock = acquire_with_retry_(file_, 0, 400,
platform::LockKind::Shared);
if (!hdr_lock) return hdr_lock.error();

// Read the 400-byte ADT file header.
std::uint8_t hdr[400]{};
auto got = file_.read_at(0, hdr, sizeof(hdr));
Expand All @@ -112,6 +121,12 @@ AdtDriver::open(const std::string& path, DriverOpenMode mode) {
return util::Error{5103, 0, "ADT record_length is zero", path};
}

// Defensive cap: a .DAT (ADT-format) file can carry a header rec_count_
// larger than the data bytes actually present (truncation, crash, partial
// copy). Cap it so every read stays inside the file instead of raising a
// "short read" / "out of range" 5000.
cap_record_count_from_size_();

// Read all field descriptors (200 bytes each after the 400-byte header).
std::uint32_t num_fields = (hdr_len_ - 400) / 200;
std::vector<std::uint8_t> fd_buf(num_fields * 200, 0);
Expand Down Expand Up @@ -232,18 +247,71 @@ void AdtDriver::denormalize_deletion_flag_(std::uint8_t* buf) noexcept {

util::Result<std::vector<std::uint8_t>>
AdtDriver::read_record_raw(std::uint32_t recno) {
if (recno == 0 || recno > rec_count_) {
if (recno == 0) {
return util::Error{5000, 0, "record number out of range", ""};
}
if (recno > rec_count_) {
// A peer may have appended; refresh (shared) + re-cap before failing.
if (auto rh = refresh_record_count_shared_(); !rh) return rh.error();
if (recno > rec_count_) {
return util::Error{5000, 0, "record number out of range", ""};
}
}
std::vector<std::uint8_t> buf(rec_len_, 0);

// Fast path: the record is already in the read-ahead block -> serve it
// with a memcpy, no syscall. The cache holds raw on-disk bytes; normalise
// the deletion flag on the returned copy (mirrors the non-cached tail).
if (read_cache_first_ != 0 &&
recno >= read_cache_first_ &&
recno < read_cache_first_ + read_cache_recs_) {
std::size_t pos = static_cast<std::size_t>(recno - read_cache_first_) *
rec_len_;
std::memcpy(buf.data(), read_cache_.data() + pos, rec_len_);
normalize_deletion_flag_(buf.data());
return buf;
}

// Miss: fetch the ALIGNED block that contains recno. Aligning (rather
// than starting at recno) keeps backward scans and local random reads
// hitting the cache too, and bounds a record to exactly one block.
std::uint32_t blk_recs = rec_len_ != 0
? static_cast<std::uint32_t>(kReadAheadBytes / rec_len_)
: 1u;
if (blk_recs == 0) blk_recs = 1;
std::uint32_t first = ((recno - 1) / blk_recs) * blk_recs + 1;
std::uint32_t last = first + blk_recs - 1;
if (last > rec_count_) last = rec_count_;
std::uint32_t nrecs = last - first + 1;

std::uint64_t offset = static_cast<std::uint64_t>(hdr_len_) +
static_cast<std::uint64_t>(recno - 1) *
static_cast<std::uint64_t>(first - 1) *
static_cast<std::uint64_t>(rec_len_);
auto got = file_.read_at(offset, buf.data(), buf.size());
if (!got) return got.error();
if (got.value() < buf.size()) {
std::size_t block_bytes = static_cast<std::size_t>(nrecs) * rec_len_;
read_cache_.assign(block_bytes, 0);
auto got = file_.read_at(offset, read_cache_.data(), block_bytes);
if (!got) { invalidate_read_cache_(); return got.error(); }

// A short read still yields whole records up to what landed; the target
// recno sits at or after `first`, so only a read that stops before it is
// a real failure.
std::uint32_t got_recs = rec_len_ != 0
? static_cast<std::uint32_t>(got.value() / rec_len_)
: 0u;
if (got_recs == 0 || recno >= first + got_recs) {
// Re-check size in case a truncate raced; re-cap and re-test.
invalidate_read_cache_();
cap_record_count_from_size_();
if (recno > rec_count_) {
return util::Error{5000, 0, "record number out of range", ""};
}
return util::Error{5000, 0, "short read on ADT record body", ""};
}
read_cache_first_ = first;
read_cache_recs_ = got_recs;

std::size_t pos = static_cast<std::size_t>(recno - first) * rec_len_;
std::memcpy(buf.data(), read_cache_.data() + pos, rec_len_);
normalize_deletion_flag_(buf.data());
return buf;
}
Expand All @@ -254,9 +322,16 @@ AdtDriver::write_record_raw(std::uint32_t recno,
if (mode_ == DriverOpenMode::ReadOnly) {
return util::Error{5000, 0, "table opened read-only", ""};
}
if (recno == 0 || recno > rec_count_) {
invalidate_read_cache_(); // record body about to change on disk
if (recno == 0) {
return util::Error{5000, 0, "record number out of range", ""};
}
if (recno > rec_count_) {
if (auto rh = refresh_record_count_shared_(); !rh) return rh.error();
if (recno > rec_count_) {
return util::Error{5000, 0, "record number out of range", ""};
}
}
if (n != rec_len_) {
return util::Error{5000, 0, "record buffer length mismatch", ""};
}
Expand Down Expand Up @@ -284,6 +359,7 @@ AdtDriver::append_record_raw(const std::uint8_t* buf, std::size_t n) {
if (mode_ == DriverOpenMode::ReadOnly) {
return util::Error{5000, 0, "table opened read-only", ""};
}
invalidate_read_cache_(); // rec_count_ / trailing block change
if (n != rec_len_) {
return util::Error{5000, 0, "record buffer length mismatch", ""};
}
Expand Down Expand Up @@ -327,9 +403,29 @@ util::Result<void> AdtDriver::refresh_record_count_() {
"ADT header truncated during refresh", ""};
}
rec_count_ = read_u32_le(buf);
cap_record_count_from_size_();
return {};
}

util::Result<void> AdtDriver::refresh_record_count_shared_() {
// Wait out any exclusive header lock held by an appender, then refresh
// (and cap). Mirrors the pattern used by CdxDriver.
auto lk = acquire_with_retry_(file_, 0, 400, platform::LockKind::Shared);
(void)lk;
return refresh_record_count_();
}

void AdtDriver::cap_record_count_from_size_() {
auto szr = file_.size();
if (!szr) return;
auto sz = szr.value();
std::uint64_t data_sz = (sz > hdr_len_) ? (sz - hdr_len_) : 0ULL;
std::uint32_t phys = (rec_len_ > 0)
? static_cast<std::uint32_t>(data_sz / rec_len_)
: 0u;
if (rec_count_ > phys) rec_count_ = phys;
}

util::Result<void> AdtDriver::rewrite_header_() {
// Only the rec_count field needs updating on ordinary append/zap.
std::uint8_t buf[4]{};
Expand Down Expand Up @@ -357,6 +453,7 @@ util::Result<void> AdtDriver::zap() {
if (mode_ == DriverOpenMode::ReadOnly) {
return util::Error{5000, 0, "table opened read-only", ""};
}
invalidate_read_cache_(); // file truncated below
rec_count_ = 0;
if (auto r = rewrite_header_(); !r) return r.error();
// Truncate to the header length so stale record bytes don't linger.
Expand Down
28 changes: 28 additions & 0 deletions src/drivers/adt/adt_driver.h
Original file line number Diff line number Diff line change
Expand Up @@ -62,10 +62,32 @@ class AdtDriver final : public IDriver {
(void)refresh_record_count_();
}

// Drop the sequential read-ahead block (the engine calls this when a peer
// may have changed the file under us, mirroring CdxDriver).
void invalidate_read_cache() noexcept override { invalidate_read_cache_(); }

private:
util::Result<void> refresh_record_count_();
// Refresh the header count while an appender may hold 0..399 exclusive:
// take the region shared (with retry) first, so the read waits instead of
// failing.
util::Result<void> refresh_record_count_shared_();
// Clamp rec_count_ to the records that physically fit in the file.
void cap_record_count_from_size_();
util::Result<void> rewrite_header_();

// Sequential read-ahead block cache (raw on-disk bytes). The reindex /
// PACK scans read every record 1..N once per tag; without this each read
// was a positioned syscall + heap alloc, making an ADT reindex of a large
// table (ESTAELEC ~441k) take minutes vs ~seconds for DBF/CDX. Mirrors
// CdxDriver: a miss fetches one ALIGNED block of up to kReadAheadBytes,
// hits are served by memcpy. Invalidated on any local write/append/zap.
static constexpr std::size_t kReadAheadBytes = 64u * 1024u;
void invalidate_read_cache_() noexcept {
read_cache_first_ = 0;
read_cache_recs_ = 0;
}

// Translate a record buffer between ADT on-disk format and the
// DBF-convention format exposed to the engine layer. The only
// byte that differs is byte 0 (deletion flag).
Expand All @@ -78,6 +100,12 @@ class AdtDriver final : public IDriver {
std::uint32_t rec_count_ = 0;
std::uint32_t rec_len_ = 0;
std::uint32_t hdr_len_ = 0;

// Read-ahead block cache state (raw file bytes). read_cache_first_ == 0
// means empty/invalid. Records [first, first+recs) live in read_cache_.
std::vector<std::uint8_t> read_cache_;
std::uint32_t read_cache_first_ = 0;
std::uint32_t read_cache_recs_ = 0;
};

} // namespace openads::drivers::adt
Loading