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]

Reply via email to