This fixes a split-RPC bug introduced in 296bacec9 and the resulting shadow 
buffer corruptions. Previously, publishing the successor transaction too early 
allowed a single RPC to be split across multiple transactions. Furthermore, 
overlapping epochs invalidated the t_updates == 0 hydration invariant, as 
threads in the new transaction could edit the live cache while the old one was 
still copying.

The fix reworks transaction handling to make the journal epoch a thread-local 
property for the duration of the RPC (similar to JBD2's current->journal_info).

- Thread-Local Nesting: Outermost starts record the transaction in TLS. Nested 
starts return that same exact transaction and increment a depth counter without 
taking locks.

- Active Threads: t_updates is renamed to t_active_threads to reflect that it 
now counts unique participant threads, not open handles.

- Safe Draining: Successor publishing is pushed until after the previous 
transaction finishes draining. Outermost starts will now safely wait for the 
successor since they hold no node locks yet.

- Pager Carveout: Because a pager thread waiting on the journal would cause a 
deadlock, it receives an explicit exemption. Pager callbacks mark their thread, 
allowing them to bypass the wait and join the draining transaction to safely 
complete their bounded operation.

- Thread-Local Sync: Sync requests are now thread-local. This prevents one 
thread's synchronous commit from forcing all other RPCs sharing the transaction 
to block (which caused severe rm -rf slowdowns).

- Shadow Copy Safety: The needs_copy flag is now cleared under the lock before 
the unlocked copy, preventing the sweep from clobbering a concurrent re-arm.
---
 ext2fs/journal.c | 252 +++++++++++++++++++++++++++++++++++++----------
 ext2fs/journal.h |   9 ++
 ext2fs/pager.c   |  10 ++
 3 files changed, 220 insertions(+), 51 deletions(-)

diff --git a/ext2fs/journal.c b/ext2fs/journal.c
index 44f93514d..cc3aeb35a 100644
--- a/ext2fs/journal.c
+++ b/ext2fs/journal.c
@@ -118,6 +118,9 @@
 #define JOURNAL_WAIT(cond, j) \
     pthread_cond_wait((cond), &(j)->j_state_lock)
 
+#define IS_RUNNING_TID(j, tid) \
+    ((j)->j_running_transaction && (j)->j_running_transaction->t_tid == (tid))
+
 #define JRNL_LIFEBOAT_CAPACITY 512
 
 #define JRNL_LIFEBOAT_ALLOC_MASK_LEN 8
@@ -142,6 +145,30 @@ struct journal_lifeboat
 static struct journal_lifeboat ext2_lifeboat;
 static pthread_t kjournald_tid;
 
+/* The handle this thread currently holds, and how many starts are open
+   on it.  One RPC is one thread, and its metadata updates must all land
+   in one transaction.  The outermost start on a thread is the only one
+   that looks at the journal; every nested start on the same thread
+   returns this transaction.  Only the outermost handle is counted in
+   t_active_threads, so t_active_threadscounts threads, and t_active_threads== 
0 means no
+   thread holds a handle. */
+static __thread diskfs_transaction_t *journal_thread_txn;
+static __thread unsigned int journal_thread_depth;
+
+/* Set while this thread runs a pager callback.  A pager thread must never
+   wait on the journal: an RPC thread inside a transaction may be waiting
+   on this very thread for a page, and the commit waits on that RPC.  An
+   outermost start on a pager thread therefore joins the committing
+   transaction while it drains rather than waiting for the successor. */
+static __thread int journal_thread_is_pager;
+
+/* Set when a nested call on this thread asked for a synchronous commit.
+   The RPC tail reads it through diskfs_journal_needs_sync and decides
+   between commit and stop.  It is per thread so an RPC waits only for
+   its own request; another thread's request on the shared transaction
+   is that thread's business.  Cleared by the outermost start. */
+static __thread int journal_thread_sync;
+
 /**
  * Holds one modified block (4KB) that needs to be written to the journal.
  */
@@ -205,7 +232,7 @@ struct diskfs_transaction
   /* The Log Position */
   uint32_t t_log_start;                /* Where this transaction starts in the 
ring */
 
-  uint32_t t_updates;          /* Refcount: How many threads are in this 
transaction? */
+  uint32_t t_active_threads;           /* Refcount: How many threads are in 
this transaction? */
 
   /* The Map: block_t -> journal_buffer_t (for O(1) lookups).
    * We use a per-transaction map rather than a global, journal-wide map.
@@ -259,6 +286,8 @@ typedef struct journal
   pthread_mutex_t j_state_lock;        /* Protects the pointers below */
   pthread_cond_t j_commit_wait;        /* Cond. var. while waiting for the tx 
to be ready. */
   pthread_cond_t j_flush_wait; /* Cond. var for safely waiting on physical 
flushes */
+  pthread_cond_t j_running_wait;       /* Cond. var for starts waiting while 
the
+                                          running transaction drains for 
commit */
   /* The Transactions */
   diskfs_transaction_t *j_running_transaction; /* Currently filling */
   diskfs_transaction_t *j_committing_transaction;      /* Transaction that is
@@ -897,8 +926,12 @@ journal_record_freed_blocks (block_t start, unsigned long 
count)
   ext->fe_count = count;
 
   JOURNAL_LOCK (ext2_journal);
-  diskfs_transaction_t *txn = ext2_journal->j_running_transaction;
-  if (!txn || txn->t_state != T_RUNNING)
+  /* Record against this thread's transaction when it holds one: the free
+     belongs to the same RPC, and that transaction may be T_LOCKED while
+     j_running_transaction is NULL.  Otherwise use the running one. */
+  diskfs_transaction_t *txn = journal_thread_depth > 0
+    ? journal_thread_txn : ext2_journal->j_running_transaction;
+  if (!txn || (txn->t_state != T_RUNNING && txn->t_state != T_LOCKED))
     {
       JRNL_LOG_DEBUG ("Cannot record freed blocks, no running transaction.");
       /* The transaction was committed by another thread before we locked!
@@ -1008,15 +1041,15 @@ static void
 journal_stop_transaction_locked (journal_t *journal,
                                 diskfs_transaction_t *txn)
 {
-  if (txn->t_updates == 0)
+  if (txn->t_active_threads == 0)
     {
       /* This implies a double-stop or corruption */
       JRNL_LOG_WARN ("Logic Error: Transaction stopped too many times!");
       return;
     }
-  txn->t_updates--;
+  txn->t_active_threads --;
   /* Continue until the map is proven clean while locked */
-  while (txn->t_updates == 0)
+  while (txn->t_active_threads == 0)
     {
       size_t iter = 0;
       journal_buffer_t *jb_exp;
@@ -1031,6 +1064,11 @@ journal_stop_transaction_locked (journal_t *journal,
            {
              /* ALWAYS hydrate from the live VM cache. The lifeboat is for
                 delayed physical I/O, not for sourcing WAL shadow data! */
+             /* Clear the flag here, under the lock, before the copy.  A
+                dirty that lands while we copy unlocked sets it again and
+                the next sweep recopies.  Clearing after the copy could
+                overwrite that re-arm and lose the recopy. */
+             jb_exp->needs_copy = 0;
              jb_exp->jb_next = NULL;
              if (!copy_list_head)
                copy_list_head = jb_exp;
@@ -1047,8 +1085,8 @@ journal_stop_transaction_locked (journal_t *journal,
          break;
        }
 
-      /* Now lockless hydration, we will update t_updates to > 0 so none can 
steal the txn. */
-      txn->t_updates++;
+      /* Now lockless hydration, we will update t_active_threads to > 0 so 
none can steal the txn. */
+      txn->t_active_threads ++;
       JOURNAL_UNLOCK (journal);
 
       journal_buffer_t *curr = copy_list_head;
@@ -1056,10 +1094,10 @@ journal_stop_transaction_locked (journal_t *journal,
        {
          /**
           * Calculate the pointer to the live Mach VM cache for this block.
-          * Because t_updates is 0 AND we hold a lock, we are mathematically
+          * Because t_active_threads is 0 AND we hold a lock, we are 
mathematically
           * guaranteed that no VFS threads are currently mutating this block
           * because if they were mutating it they would have to first obtain
-          * the journal lock AND also increase the t_updates.
+          * the journal lock AND also increase the t_active_threads .
           */
          void *live_cache_ptr = bptr (curr->jb_blocknr);
          /**
@@ -1070,11 +1108,9 @@ journal_stop_transaction_locked (journal_t *journal,
           * WAL so it can overwrite the pager's rushed data during recovery!
           */
          memcpy (curr->jb_shadow_data, live_cache_ptr, block_size);
-         /* We are done with this block even if t_updates reach 0 again before
-          * this transaction is committed. If we get notified that this block
-          * has been modified again journal_dirty_blocks must set needs_copy
-          * back to 1. */
-         curr->needs_copy = 0;
+         /* needs_copy was cleared under the lock when this block was
+          * listed.  If the block is modified again, journal_dirty_block
+          * sets it back to 1 and a later sweep recopies. */
 
          journal_buffer_t *next = curr->jb_next;
          curr->jb_next = NULL;
@@ -1083,10 +1119,10 @@ journal_stop_transaction_locked (journal_t *journal,
 
       JOURNAL_LOCK (journal);
       /* We need to decrement what we incremented. */
-      txn->t_updates--;
+      txn->t_active_threads--;
 
-      /* If t_updates is still 0, the loop repeats to verify no blocks were
-         added. If t_updates > 0, another thread joined while we were unlocked
+      /* If t_active_threads is still 0, the loop repeats to verify no blocks 
were
+         added. If t_active_threads > 0, another thread joined while we were 
unlocked
          and they will handle the final sweep. We safely fall out. */
     }
 }
@@ -1236,7 +1272,7 @@ journal_forget_freed_blocks (journal_t *journal, 
journal_freed_extent_t *ext)
     flush_to_disk ();
 }
 
-/* Install a running transaction with t_updates == 0.
+/* Install a running transaction with t_active_threads == 0.
    The journal lock is held on entry and on return, including failure.
    j_running_transaction must be NULL.  This does not drop the lock, so
    commit can publish the successor with no gap for start to observe. */
@@ -1259,7 +1295,7 @@ journal_create_running_transaction_locked (journal_t 
*journal)
 
   txn->t_tid = journal->j_transaction_sequence++;
   txn->t_state = T_RUNNING;
-  txn->t_updates = 0;
+  txn->t_active_threads = 0;
 
   journal->j_running_transaction = txn;
   JRNL_LOG_DEBUG ("[TRX] Created NEW TID %u", txn->t_tid);
@@ -1599,37 +1635,113 @@ out:
   return err;
 }
 
-/* Join the running transaction.  A successor is installed by commit
-   before the previous one leaves T_RUNNING, so this never allocates.
-   Shutdown is the exception: j_must_exit means commit did not install
-   a successor, and the caller must tolerate NULL. */
+/* Join the running transaction.  Only the outermost start on a thread
+   comes here; nested starts are served from journal_thread_txn.
+
+   While commit drains a transaction (T_LOCKED), j_running_transaction is
+   NULL and the successor does not exist yet.  An RPC thread waits here
+   until the successor is published.  The wait is safe because an RPC
+   takes its outermost start before any node lock, so a waiting starter
+   holds nothing a draining participant needs.  A pager thread joins the
+   draining transaction instead: it must never wait on the journal, and
+   its work is one bounded operation that only extends the drain.
+
+   Nothing else is ever observed here: commit publishes the successor
+   before it drops the lock, so a NULL j_running_transaction with no
+   commit in flight means the successor could not be allocated.  Shutdown
+   is the exception: j_must_exit means commit did not install a
+   successor, and the caller must tolerate NULL. */
 static diskfs_transaction_t *
 journal_join_transaction_locked (journal_t *journal)
 {
   diskfs_transaction_t *txn;
 
+  while (!(txn = journal->j_running_transaction))
+    {
+      diskfs_transaction_t *commit = journal->j_committing_transaction;
+
+      if (journal->j_must_exit)
+       return NULL;
+
+      assert_backtrace (commit);
+
+      if (journal_thread_is_pager && commit->t_state == T_LOCKED)
+       {
+         txn = commit;
+         break;
+       }
+
+      JOURNAL_WAIT (&journal->j_running_wait, journal);
+    }
+
   if (journal->j_must_exit)
     return NULL;
 
-  txn = journal->j_running_transaction;
-  assert_backtrace (txn && txn->t_state == T_RUNNING);
-  txn->t_updates++;
+  assert_backtrace (txn->t_state == T_RUNNING
+                   || (journal_thread_is_pager && txn->t_state == T_LOCKED));
+  txn->t_active_threads ++;
   return txn;
 }
 
+/* Called by the ext2fs pager callbacks on entry (ON = 1) and exit
+   (ON = 0).  See journal_thread_is_pager. */
+void
+journal_thread_set_pager (int on)
+{
+  journal_thread_is_pager = on;
+}
+
+/* Open a handle on the running transaction for this thread.  The outermost
+   start of an RPC may wait while a commit drains the previous transaction;
+   it must therefore be taken before any node lock.  A pager thread never
+   waits here.  Nested starts return the thread's transaction at once. */
 diskfs_transaction_t *
 diskfs_journal_start_transaction (void)
 {
   diskfs_transaction_t *txn;
+
+  if (journal_thread_depth > 0)
+    {
+      /* Nested start.  The outer handle already holds a t_active_threads 
count,
+         so this transaction cannot drain until that handle is released.
+         It may already be T_LOCKED: commit is waiting for us.  Joining it
+         here keeps the whole RPC in one transaction.  No journal lock and
+         no t_active_threads change are needed. */
+      journal_thread_depth++;
+      return journal_thread_txn;
+    }
+
   if (!ext2_journal || ext2_journal->j_must_exit)
     return NULL;
 
   JOURNAL_LOCK (ext2_journal);
   txn = journal_join_transaction_locked (ext2_journal);
   JOURNAL_UNLOCK (ext2_journal);
+
+  if (txn)
+    {
+      journal_thread_txn = txn;
+      journal_thread_depth = 1;
+      journal_thread_sync = 0;
+    }
   return txn;
 }
 
+/* Close one start on this thread.  Returns 1 when TXN was the outermost
+   handle and the caller must release the t_active_threads count it holds.  */
+static int
+journal_thread_release (diskfs_transaction_t *txn)
+{
+  assert_backtrace (journal_thread_depth > 0);
+  assert_backtrace (txn == journal_thread_txn);
+
+  if (--journal_thread_depth > 0)
+    return 0;
+
+  journal_thread_txn = NULL;
+  return 1;
+}
+
 /**
  * Adds a modified filesystem block to the SPECIFIC transaction handle.
  * Defers the actual memory copy until the transaction stops.
@@ -1646,17 +1758,29 @@ journal_dirty_block (diskfs_transaction_t *txn, block_t 
fs_blocknr)
   return err;
 }
 
+/* Two flags.  journal_thread_sync tells this thread's RPC tail to commit
+   rather than stop.  txn->sync_needed tells the outermost stop to wake
+   kjournald when t_active_threads reaches 0, for callers that never commit
+   (the pager reaches diskfs_drop_node with diskfs_synchronous set). */
 void
 diskfs_journal_set_sync (diskfs_transaction_t *txn)
 {
-  if (txn)
-    txn->sync_needed = 1;
+  if (!txn)
+    return;
+
+  assert_backtrace (txn == journal_thread_txn);
+  journal_thread_sync = 1;
+  txn->sync_needed = 1;
 }
 
 int
 diskfs_journal_needs_sync (diskfs_transaction_t *txn)
 {
-  return txn ? txn->sync_needed : 0;
+  if (!txn)
+    return 0;
+
+  assert_backtrace (txn == journal_thread_txn);
+  return journal_thread_sync;
 }
 
 /* Helper to reset the header for a new block */
@@ -1957,28 +2081,37 @@ journal_commit_running_transaction_locked (journal_t 
*journal)
   if (journal->j_free < journal->j_min_free)
     journal_force_checkpoint_locked (journal);
 
+  /* Drain.  No successor exists while txn is T_LOCKED: an outermost start
+     on an RPC thread waits in journal_join_transaction_locked, a pager
+     thread joins txn.  So every thread that can edit a block in this map
+     is counted int_active_threads , and t_active_threads == 0 means the last 
shadow
+     copy taken by diskfs_journal_stop_transaction_locked is final.  The
+     successor is published only after that, so no thread in the next
+     transaction ever edits a block this one has yet to copy.  This
+     thread holds no handle: kjournald never starts one, and
+     diskfs_journal_commit_transaction releases its handle first. */
   journal->j_running_transaction = NULL;
-  if (!journal->j_must_exit)
-    {
-      /* Publish the successor before any unlock.  t_updates is 0: this
-         thread is not a participant. */
-      if (!journal_create_running_transaction_locked (journal))
-       JRNL_LOG_WARN ("Failed to create the successor transaction.");
-    }
-
   txn->t_state = T_LOCKED;
 
-  while (txn->t_updates > 0)
+  while (txn->t_active_threads > 0)
     JOURNAL_WAIT (&journal->j_commit_wait, journal);
 
   txn->t_state = T_FLUSHING;
 
+  /* Publish the successor before any unlock, then release the starters
+     waiting for it.  On shutdown there is no successor; the broadcast
+     lets them see j_must_exit and return NULL. */
+  if (!journal->j_must_exit
+      && !journal_create_running_transaction_locked (journal))
+    JRNL_LOG_WARN ("Failed to create the successor transaction.");
+  pthread_cond_broadcast (&journal->j_running_wait);
+
   journal_ensure_commit_space_locked (journal, txn);
 
   txn->t_log_start = journal_next_block_would_be (journal);
-  /* Unlock for I/O.  The successor is already j_running_transaction, and
-     this transaction is T_LOCKED, so start joins the successor.  Participants
-     of this one have drained.  Its map is stable. */
+  /* Unlock for I/O.  The successor is j_running_transaction, and this
+     transaction is T_FLUSHING with its participants drained.  Its map and
+     shadow data are stable. */
   JOURNAL_UNLOCK (journal);
 
   /* Write Data (I/O) */
@@ -2144,6 +2277,7 @@ journal_create (struct node *journal_inode)
   pthread_cond_init (&j->j_commit_wait, NULL);
   pthread_cond_init (&j->j_flusher_wakeup, NULL);
   pthread_cond_init (&j->j_flush_wait, NULL);
+  pthread_cond_init (&j->j_running_wait, NULL);
 
   j->j_must_exit = 0;
 
@@ -2206,7 +2340,7 @@ journal_quiesce_checkpoints (void)
                                  &ext2_journal->j_state_lock, CLOCK_MONOTONIC, 
&ts);
   if (err)
     {
-      /* If we hit ETIMEDOUT, a VFS thread likely leaked a t_updates refcount
+      /* If we hit ETIMEDOUT, a VFS thread likely leaked a t_active_threads 
refcount
          due to a signal interruption or crash. We MUST bail out without
          clearing the checkpoint list so the WAL replays on next boot! */
       JRNL_LOG_WARN
@@ -2240,23 +2374,30 @@ diskfs_journal_stop_transaction_locked (journal_t 
*journal,
   journal_stop_transaction_locked (journal, txn);
 
   /* Semi auto-commit? */
-  if (txn->t_updates == 0 && (txn->sync_needed ||
+  if (txn->t_active_threads == 0 && (txn->sync_needed ||
       (txn->t_buffer_map.size >= journal->j_max_transaction_buffers)))
     pthread_cond_signal (&journal->j_flusher_wakeup);
 }
 
 /* Ends the caller's participation in the given transaction TXN.
    This informs the journal that the logical operation is complete.
-   If the caller is the final participant (t_updates reaches 0) AND any
-   participant flagged the transaction for a synchronous commit, this function
-   will automatically perform the physical disk flush. Otherwise, this function
-   returns immediately without waiting for the transaction to commit.
+   A nested stop only closes this thread's inner start.  When the outermost
+   handle is released and the caller is the final participant (t_active_threads
+   reaches 0), a transaction flagged for sync, or one that has reached its
+   size limit, wakes kjournald.  This function never waits for the commit:
+   the pager reaches it, and a pager thread that waits deadlocks.
 
    This function consumes TXN. The caller must not use TXN after this call. */
 void
 diskfs_journal_stop_transaction (diskfs_transaction_t *txn)
 {
-  if (!ext2_journal || !txn)
+  if (!txn)
+    return;
+
+  if (!journal_thread_release (txn))
+    return;
+
+  if (!ext2_journal)
     return;
 
   JOURNAL_LOCK (ext2_journal);
@@ -2283,6 +2424,12 @@ diskfs_journal_commit_transaction (diskfs_transaction_t 
*opaque_txn)
   diskfs_transaction_t *txn = (diskfs_transaction_t *) opaque_txn;
   uint32_t tid = txn->t_tid;
 
+  /* Commit belongs at the end of the RPC, on the outermost handle.  A
+     nested commit would wait for t_active_threads to drain while this thread's
+     outer handle holds a count: a self-deadlock.  Catch it here. */
+  int outermost = journal_thread_release (txn);
+  assert_backtrace (outermost);
+
   JRNL_LOG_DEBUG ("Committing tx id: %u.", txn->t_tid);
 
   JOURNAL_LOCK (ext2_journal);
@@ -2295,7 +2442,7 @@ diskfs_journal_commit_transaction (diskfs_transaction_t 
*opaque_txn)
                    tid);
   /* Check if the transaction is currently RUNNING.
      If it is, We "steal" it and become the committer. */
-  if (ext2_journal->j_running_transaction == txn)
+  if (IS_RUNNING_TID (ext2_journal, tid))
     {
       error_t err = journal_commit_running_transaction_locked (ext2_journal);
       if (err)
@@ -2709,7 +2856,10 @@ diskfs_journal_shutdown (void)
   if (!ext2_journal)
     return;
 
-  ext2_journal->j_must_exit = 1;       /* Signal kjournald so it doesn't wake 
up */
+  JOURNAL_LOCK (ext2_journal);
+  ext2_journal->j_must_exit = 1;    /* Signal kjournald so it doesn't wake up 
*/
+  pthread_cond_broadcast (&ext2_journal->j_running_wait); /* Wake waiting 
starters */
+  JOURNAL_UNLOCK (ext2_journal);
 
   /* Sync the disk pager to ensure all shadow data is on disk.  */
   sync_global (1);
diff --git a/ext2fs/journal.h b/ext2fs/journal.h
index 5860b50a1..ab7ab1e6b 100644
--- a/ext2fs/journal.h
+++ b/ext2fs/journal.h
@@ -86,4 +86,13 @@ journal_store_read (block_t start_block, size_t length, void 
**buf,
  */
 void journal_record_freed_blocks (block_t start, unsigned long count);
 
+/**
+ * Marks the calling thread as running a pager callback (ON = 1) or done
+ * with it (ON = 0).  A pager thread must never wait on the journal: an
+ * RPC thread inside a transaction may be waiting on it for a page.  While
+ * marked, diskfs_journal_start_transaction joins the transaction that a
+ * commit is draining instead of waiting for its successor.
+ */
+void journal_thread_set_pager (int on);
+
 #endif //_JOURNAL_H
diff --git a/ext2fs/pager.c b/ext2fs/pager.c
index cba908427..42941150e 100644
--- a/ext2fs/pager.c
+++ b/ext2fs/pager.c
@@ -743,6 +743,10 @@ pager_unlock_page (struct user_pager_info *pager, 
vm_offset_t page)
       struct node *node = pager->node;
       struct disknode *dn = diskfs_node_disknode (node);
 
+      /* Block allocation below starts a journal transaction.  Mark this
+        thread so that start never waits on a draining commit.  */
+      journal_thread_set_pager (1);
+
       pthread_rwlock_wrlock (&dn->alloc_lock);
 
       partial_page = (page + vm_page_size > node->allocsize);
@@ -787,6 +791,8 @@ pager_unlock_page (struct user_pager_info *pager, 
vm_offset_t page)
 
       pthread_rwlock_unlock (&dn->alloc_lock);
 
+      journal_thread_set_pager (0);
+
       if (err == ENOSPC)
        ext2_warning ("This filesystem is out of space.");
       else if (err)
@@ -974,7 +980,11 @@ pager_clear_user_data (struct user_pager_info *upi)
       assert_backtrace (!pager || pager_get_upi (pager) != upi);
       pthread_spin_unlock (&node_to_page_lock);
 
+      /* Dropping the last reference can write the node back, which starts
+        a journal transaction on this pager thread.  */
+      journal_thread_set_pager (1);
       diskfs_nrele_light (upi->node);
+      journal_thread_set_pager (0);
     }
 }
 
-- 
2.55.0


Reply via email to