dhruvarya-db opened a new pull request, #3275:
URL: https://github.com/apache/iceberg-rust/pull/3275
## Which issue does this PR close?
- Closes #3274.
This case was raised during review of #3268.
## What changes are included in this PR?
An equality-delete file that **fails to load** (missing, unreadable, or
unparseable) can leave an entry permanently stuck in
the `Loading` state inside `DeleteFilter`. Any task that later waits on that
entry blocks
**forever** on a notification that will never fire. This PR makes a failed
load clean up
after itself and wake its waiters, so the failure surfaces as an error
instead of a
hang.
### Background: how equality-delete loading works
Equality-delete predicates are produced by one task and consumed by others
through a small
state machine keyed by delete-file path:
- `EqDelState::Loading(Notify)` — a task is parsing this file; waiters hold
the `Notify`.
- `EqDelState::Loaded(Predicate)` — the predicate is ready.
The producer side (`insert_equality_delete`) marks the entry `Loading`, then
spawns a task
that waits on a `oneshot` channel for the parsed predicate, stores it as
`Loaded`, and calls
`notify_waiters()`. The parsed predicate is sent from the delete-file loader
(`caching_delete_file_loader.rs`, the `FreshEqDel` branch).
The consumer side (`get_equality_delete_predicate_for_delete_file_path`)
reads the entry; if
it is `Loading`, it waits on the `Notify` until the producer signals
completion.
### The bug
The spawned producer task did:
```rust
let eq_del = eq_del.await.unwrap(); // <-- panics if the sender was dropped
// ... insert Loaded ...
notify.notify_waiters();
```
The predicate is sent here (`caching_delete_file_loader.rs`):
```rust
let predicate =
Self::parse_equality_deletes_record_batch_stream(batch_stream,
equality_ids)
.await?; // <-- early return on parse error DROPS
`sender`
sender.send(predicate) ...
```
If parsing the delete file fails, the `?` returns early **before**
`sender.send(...)`,
dropping the sender. On the producer task, `eq_del.await` then resolves to
`Err(RecvError)`,
`.unwrap()` **panics**, and the two lines that follow — the `Loaded` insert
and
`notify_waiters()` — **never run**. The entry is stranded in `Loading`, and
its `Notify` is
never signalled.
#### Exact sequence of events
```mermaid
sequenceDiagram
participant P as Producer task (insert_equality_delete)
participant L as Loader (FreshEqDel branch)
participant S as DeleteFilter state
participant W as Waiter (get_equality_delete_predicate...)
L->>S: insert(path, Loading(notify))
L->>P: spawn task, awaiting oneshot rx
W->>S: read(path) -> Loading(notify)
W->>W: notify.notified().await (parked)
Note over L: parse_equality_deletes_...().await? ❌ parse error
L--xP: sender dropped (early return, never sent)
P->>P: eq_del.await -> Err(RecvError)
Note over P: .unwrap() PANICS
Note over S: entry stays Loading forever,<br/>notify_waiters() never
called
Note over W: parked forever — HANG 🔒
```
State of the map entry:
```
success path failure path (this bug)
Loading ──────────────────► Loaded Loading ──► (parse error) ──►
Loading 🔒
send + notify sender dropped,
forever
task panics
```
#### When it actually hangs
Within a single `load_deletes` call the parse error also propagates up the
results stream
via `?`, which tears down that call and cancels its sibling tasks — so the
stuck entry is
usually not observed there. The hang bites when a **`DeleteFilter` is shared
across
concurrent or sequential loads** (an `ArrowReader` is `Clone` and its clones
share one
loader/filter, and many data files typically reference the same
equality-delete file):
1. Loader task A claims delete file `E`, marks it `Loading`, and fails to
parse it.
2. A panics; `E` is never signalled and never leaves `Loading`.
3. A later task B references `E`, sees it already claimed, and waits on its
`Notify`.
4. B blocks forever.
Secondary effects: the runtime task panics on an ordinary I/O/parse error,
and a reused
`ArrowReader` stays poisoned for that file across subsequent scans.
### The fix
Handle the dropped sender instead of unwrapping it:
- On `Ok(predicate)`: unchanged — store `Loaded` and `notify_waiters()`.
- On `Err(_)`: remove the stale `Loading` marker, then `notify_waiters()`.
Waiters that wake to a **missing** entry now return `None` (the post-`await`
match no longer
treats that as `unreachable!`). `build_equality_delete_predicate` turns that
`None` into a
clear `"Missing predicate for equality delete file '…'"` error, so the
failure is reported
rather than deadlocked. Removing the entry (rather than caching a permanent
failure) also
lets a later, independent load of the same file try again.
```mermaid
sequenceDiagram
participant P as Producer task
participant S as DeleteFilter state
participant W as Waiter
Note over P: eq_del.await -> Err(RecvError)
P->>S: remove(path) (drop stale Loading)
P->>W: notify.notify_waiters()
W->>S: re-read(path) -> None
W-->>W: return None ➜ surfaced as "missing predicate" error ✅
```
## Are these changes tested?
Each test is bounded by `tokio::time::timeout`. With the fix reverted, all
three fail after 5s with a "hung" assertion instead of hanging the test suite.
- `test_get_equality_delete_predicate_resolves_when_load_fails`: a unit test
on `DeleteFilter`. It drops the predicate sender without sending, which is what
a failed parse does, and checks that the waiter gets `None`.
- `test_cloned_readers_error_on_failed_equality_delete_load`: builds one
`ArrowReader`, clones it, and reads a task whose equality-delete file can't be
read. Both the original and the clone must return an error. The clones share
one `DeleteFilter`, so this tests the shared-reader case from the description
directly.
- `test_equality_delete_load_retries_after_failure`: after a failed read, it
restores the delete file and reads again with the same reader. The deletes must
be applied (`[1, 2, 3, 4, 5]` minus `3`), which shows that removing the stale
entry lets the reader recover.
## AI Disclosure
The code was written using Claude but was manually reviewed by me.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]