Skip to content

feat(telemetry): add the write-ahead batch cache - #1481

Open
pblazej wants to merge 1 commit into
blaze/telemetry-stack/2-transportfrom
blaze/telemetry-stack/3-cache
Open

pblazej wants to merge 1 commit into
blaze/telemetry-stack/2-transportfrom
blaze/telemetry-stack/3-cache

Conversation

@pblazej

@pblazej pblazej commented Oct 1, 2026 •

Copy link
Copy Markdown
Contributor

Summary

The write-ahead cache between the exporter and the transport.

Changes

  • BatchCache with MemoryCache and FileCache (storage_dir)
  • A batch is committed once its file is fsynced and renamed and the directory fsynced
  • Bounded by size and count, every eviction reported
  • A 413 split is replaced in one transaction
  • SpillCache: keeps a batch the disk refused in memory, and takes no write after the opt-out
  • Test-only fault hooks that fail or crash every write, fsync, rename, delete and publish
Verification

At cf773fde, from a clean checkout (CI's test workflow runs only for PRs into main, so these were run locally; there is no clippy job in CI):

  • cargo fmt -- --check
  • cargo clippy -p livekit-telemetry --all-targets --all-features -- -D warnings
  • cargo check -p livekit-telemetry --all-targets --no-default-features with features [], [net], [uniffi], [net,uniffi]
  • cargo test -p livekit-telemetry: 38 unit; --all-features: 39 unit

cargo doc -D warnings reports links to items #1483 adds (Telemetry::stats, Telemetry::with_cache, Exporter), plus what it reports on 6aba1b68 too (private-item links, ExportError::from_response).

`BatchCache` holds encoded batches between the exporter and the transport,
oldest first: `MemoryCache`, or `FileCache` in `storage_dir`, where a batch
is committed once its file is fsynced, renamed and the directory fsynced.
Both evict the oldest batches past their size and count bounds, report
what they evicted, and replace a batch by its two halves in one
transaction. `SpillCache` keeps a batch the disk refused in memory, counted
as a cache write error, and takes no write once the app opted out.
@pblazej pblazej added the internal to tag changes that don't require changelog documentation label Oct 1, 2026
@pblazej
pblazej force-pushed the blaze/telemetry-stack/3-cache branch from 9ad007b to cf773fd Compare October 1, 2026 13:53
@pblazej
pblazej added this pull request to stack #1486 October 1, 2026 14:23
@pblazej
pblazej marked this pull request as ready for review October 1, 2026 14:32
@pblazej
pblazej requested a review from ladvoc as a code owner October 1, 2026 14:32

@devin-ai-integration devin-ai-integration Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Devin Review found 4 potential issues.

2 flags not posted on this PR by your GitHub settings — view them in Devin Review. (Configure)

Devin Review

Comment on lines +409 to +421
for (id, len) in kept.iter().zip(sizes) {
if total <= self.max_bytes && count <= self.max_batches {
break;
}
// A batch that could not be deleted still takes its room: keep evicting, report only
// what really went.
if self.remove(id).is_ok() {
total -= len;
count -= 1;
removed.push(id.clone());
}
}
Ok(removed)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

🔴 Oversized disk batch disappears before upload

When a batch exceeds max_bytes, prune deletes the batch just committed by push. No batch remains for upload, unlike the single oversized batch retained by MemoryCache.

Learn more

The disk cache uses prune after every successful push and replace. Its loop continues even if the last batch is larger than max_bytes, so it deletes the only copy and returns its ID as an eviction. The memory implementation stops evicting with one batch left. An uploader cannot send an ID that is no longer in pending.

Example: With max_bytes = 25, pushing one 30-byte batch returns Ok([id]), but pending() is empty; the batch never gets an upload attempt.

Recommended fix: Stop FileCache::prune before deleting its last retained batch, matching MemoryCache::evict. Check the behavior for a single oversize push and for replace where both halves together exceed the limit.

Devin Review


Was this helpful? React with 👍 or 👎 to provide feedback.

Comment on lines +382 to +391
// A host-less id ends with its empty host; a split's halves rename the sequence.
let rebinding =
parent.ends_with('-') && half.len() > parent.len() && half.starts_with(parent);
if self.path(parent, EXT).exists() {
if !rebinding {
fs::remove_file(&pend)?;
continue;
}
// Parent first: a crash in between leaves a committed journal, published below.
fs::remove_file(self.path(parent, EXT))?;

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

🔴 Interrupted rebinding loses replacement records

If replace crashes after journaling one of several host-suffixed parts, recover deletes the intact parent. The unwritten parts then have no copy on disk.

Learn more

A replacement journals each member of new separately before deleting the original batch. Recovery normally discards an incomplete replacement when the parent still exists. The rebinding exception instead deletes the parent when any single journaled child extends its host-less ID. That test cannot tell whether all replacement members reached the journal.

Example: Replacing a host-less batch 123-1-2-l- with two parts 123-1-2-l-a and 123-1-2-l-b crashes after the first .pend rename. Reopening the cache publishes a and deletes the original, losing the records intended for b.

Recommended fix: Make forward recovery conditional on an explicit durable marker proving the full replacement was journaled, or restrict rebinding to an atomic one-member operation that cannot represent a partial split. Add a crash test between successive writes of host-suffixed replacement parts.

Devin Review


Was this helpful? React with 👍 or 👎 to provide feedback.

Comment on lines +415 to +419
if self.remove(id).is_ok() {
total -= len;
count -= 1;
removed.push(id.clone());
}

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

🟡 Unreported evictions after directory sync failure

When remove unlinks a batch but directory sync fails, prune treats it as retained. The deleted batch vanishes without an eviction count, and further batches can be evicted.

Learn more

FileCache::remove deletes the file before calling sync_dir and returns the sync error even though the file is already absent. prune only updates its totals and reports an eviction on Ok. On a directory sync failure it therefore leaves the deleted file in its accounting and can continue deleting newer files unnecessarily. The returned eviction list misses the deleted batch.

Example: Two 10-byte batches exceed a 10-byte cap. Removing the older batch succeeds, but its directory sync fails. The loop still counts that batch's 10 bytes and deletes the newer one too; neither deletion is reported if the sync keeps failing.

Recommended fix: Track whether the unlink actually happened separately from whether its directory entry was durably synced. Count and return IDs of files that disappeared even when sync fails, and avoid treating already-unlinked files as occupying the cache.

Devin Review


Was this helpful? React with 👍 or 👎 to provide feedback.

Comment on lines +522 to +525
match fs::remove_file(&path) {
Err(err) if err.kind() != io::ErrorKind::NotFound => Err(err),
_ => self.sync_dir(),
}

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

🟡 Purged cache directory breaks idempotent removal

After the OS purges the cache directory, remove returns an error for an already-gone batch. Its NotFound branch still calls sync_dir, which cannot open the missing directory.

Learn more

The file cache explicitly supports an OS purging its directory while the app runs: write_synced recreates it on the next write. BatchCache::remove promises that an already-gone batch is not an error. If the entire directory was purged, removing its batch returns NotFound, but sync_dir immediately fails to open that same missing directory.

Example: Open a cache, allow the OS to remove its directory, and call remove("123-1-1-l-"). The batch is gone, yet removal returns a NotFound error rather than success.

Recommended fix: Treat NotFound from remove_file as successful removal without syncing a nonexistent directory. Retain the sync for a file that was actually unlinked.

Devin Review


Was this helpful? React with 👍 or 👎 to provide feedback.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

internal to tag changes that don't require changelog documentation

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant