This is an automated email from the ASF dual-hosted git repository.

tbonelee pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/zeppelin.git


The following commit(s) were added to refs/heads/master by this push:
     new 2ef78bcc2a [ZEPPELIN-6594] Serialize NoteManager.addNote and 
removeFolder with the other tree mutations
2ef78bcc2a is described below

commit 2ef78bcc2a59944cd829debd2ad2e66f07aff08a
Author: HwangRock <[email protected]>
AuthorDate: Mon Oct 5 22:31:02 2026 +0900

    [ZEPPELIN-6594] Serialize NoteManager.addNote and removeFolder with the 
other tree mutations
    
    ### What is this PR for?
    
    `NoteManager` serializes its tree mutations — `saveNote`, `removeNote`, 
`moveNote`, `moveFolder` — on a single `synchronized (this)` monitor so they 
never interleave on the folder tree, the noteId→path mapping, and the notebook 
repo. `addNote` and `removeFolder` were left out of this serialization.
    
    **As-Is:** Because `addNote` and `removeFolder` do not take the monitor, 
they can run concurrently with a `moveNote` that holds it, corrupting shared 
state:
    
    - A note being moved to another folder while its original folder is removed 
becomes unreachable — the note is effectively lost.
    - Adding a note at a path while another note is concurrently moved to that 
same path lets both pass the "path is available" check and commit, leaving two 
different notes at one path.
    - Removing the same folder concurrently throws a `NullPointerException`.
    
    **To-Be:** `addNote` and `removeFolder` take the same `synchronized (this)` 
monitor as the other mutations, so the whole class shares one serialization 
point and these interleavings can no longer happen. Reads are unaffected — 
`processNote` resolves against the volatile `noteTree` and never takes the 
monitor.
    
    ### What type of PR is it?
    
    Bug Fix
    
    ### What is the Jira issue?
    
    [ZEPPELIN-6594](https://issues.apache.org/jira/browse/ZEPPELIN-6594)
    
    ### How should this be tested?
    
    Added a deterministic gate-based test 
(`NoteManagerMutationSerializationTest`) that pins the monitor inside a parked 
`moveNote` and asserts the other mutation blocks on it, plus two concurrency 
stress tests in `NoteManagerTest`.
    
    Running the same tests before and after the fix:
    
    ```
    # Before the fix (same tests run against master's pre-fix behavior)
    Tests run: 11, Failures: 4, Errors: 0, Skipped: 0
    
      testAddNoteWaitsForInFlightMoveNote
        The mutation completed while moveNote held the NoteManager monitor; it 
is not serialized with moveNote
      testRemoveFolderWaitsForInFlightMoveNote
        The mutation completed while moveNote held the NoteManager monitor; it 
is not serialized with moveNote
      testConcurrentAddNoteAndMoveNoteOnSamePath
        exactly one operation may claim /dst_1/note ==> expected: <1> but was: 
<2>
      testConcurrentRemoveFolderAndMoveNote
        java.lang.NullPointerException: Cannot invoke
        "NoteManager$Folder.getNoteInfoRecursively()" because "folder" is null
    
    # After the fix (this branch)
    Tests run: 11, Failures: 0, Errors: 0, Skipped: 0
    ```
    
    ```
    ./mvnw test -pl zeppelin-server --am 
-Dtest=NoteManagerTest,NoteManagerMoveResaveRaceTest,NoteManagerMutationSerializationTest
    ```
    
    ### Questions:
    
    - Does the license files need to be updated? No
    - Is there breaking changes for older versions? No
    - Does this need documentation? No
    
    
    Closes #5530 from HwangRock/ZEPPELIN-6594-serialize-notemanager-mutations.
    
    Signed-off-by: ChanHo Lee <[email protected]>
---
 .../org/apache/zeppelin/notebook/NoteManager.java  |  32 ++--
 .../NoteManagerMutationSerializationTest.java      | 184 +++++++++++++++++++++
 .../apache/zeppelin/notebook/NoteManagerTest.java  | 108 ++++++++++++
 .../notebook/repo/VFSNotebookRepoWithMoveGate.java |  78 +++++++++
 4 files changed, 388 insertions(+), 14 deletions(-)

diff --git 
a/zeppelin-server/src/main/java/org/apache/zeppelin/notebook/NoteManager.java 
b/zeppelin-server/src/main/java/org/apache/zeppelin/notebook/NoteManager.java
index c31cad72b8..dfdae2fac3 100644
--- 
a/zeppelin-server/src/main/java/org/apache/zeppelin/notebook/NoteManager.java
+++ 
b/zeppelin-server/src/main/java/org/apache/zeppelin/notebook/NoteManager.java
@@ -196,8 +196,10 @@ public class NoteManager {
   }
 
   public void addNote(Note note, AuthenticationInfo subject) throws 
IOException {
-    addOrUpdateNoteNode(this.noteTree, new NoteInfo(note), true);
-    noteCache.putNote(note);
+    synchronized (this) {
+      addOrUpdateNoteNode(this.noteTree, new NoteInfo(note), true);
+      noteCache.putNote(note);
+    }
   }
 
   /**
@@ -334,21 +336,23 @@ public class NoteManager {
    */
   public List<NoteInfo> removeFolder(String folderPath, AuthenticationInfo 
subject) throws IOException {
 
-    // update notebookrepo
-    this.notebookRepo.remove(folderPath, subject);
+    synchronized (this) {
+      // update notebookrepo
+      this.notebookRepo.remove(folderPath, subject);
 
-    // update filesystem tree
-    NoteTree tree = this.noteTree;
-    Folder folder = getFolder(tree, folderPath);
-    List<NoteInfo> noteInfos = 
folder.getParent().removeFolder(folder.getName(), subject);
+      // update filesystem tree
+      NoteTree tree = this.noteTree;
+      Folder folder = getFolder(tree, folderPath);
+      List<NoteInfo> noteInfos = 
folder.getParent().removeFolder(folder.getName(), subject);
 
-    // update notesInfo and evict the deleted notes from the cache, mirroring 
removeNote
-    for (NoteInfo noteInfo : noteInfos) {
-      tree.notesInfo.remove(noteInfo.getId());
-      this.noteCache.removeNote(noteInfo.getId());
-    }
+      // update notesInfo and evict the deleted notes from the cache, 
mirroring removeNote
+      for (NoteInfo noteInfo : noteInfos) {
+        tree.notesInfo.remove(noteInfo.getId());
+        this.noteCache.removeNote(noteInfo.getId());
+      }
 
-    return noteInfos;
+      return noteInfos;
+    }
   }
 
   /**
diff --git 
a/zeppelin-server/src/test/java/org/apache/zeppelin/notebook/NoteManagerMutationSerializationTest.java
 
b/zeppelin-server/src/test/java/org/apache/zeppelin/notebook/NoteManagerMutationSerializationTest.java
new file mode 100644
index 0000000000..daf7473905
--- /dev/null
+++ 
b/zeppelin-server/src/test/java/org/apache/zeppelin/notebook/NoteManagerMutationSerializationTest.java
@@ -0,0 +1,184 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *    http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.zeppelin.notebook;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.junit.jupiter.api.Assertions.fail;
+
+import java.io.File;
+import java.io.IOException;
+import java.nio.file.Files;
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.List;
+import java.util.concurrent.TimeUnit;
+
+import org.apache.commons.io.FileUtils;
+import org.apache.zeppelin.conf.ZeppelinConfiguration;
+import org.apache.zeppelin.notebook.repo.VFSNotebookRepoWithMoveGate;
+import org.apache.zeppelin.user.AuthenticationInfo;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+/**
+ * Reproduction test for ZEPPELIN-6594. {@link NoteManager#moveNote} calls
+ * {@code notebookRepo.move} while holding the NoteManager monitor. {@code 
addNote} and
+ * {@code removeFolder} must take the same monitor, otherwise they run 
concurrently with
+ * {@code moveNote} and corrupt the folder tree and the noteId to path mapping.
+ *
+ * <p>The monitor is pinned deterministically with {@link 
VFSNotebookRepoWithMoveGate}, which
+ * parks a {@code moveNote} call inside {@code notebookRepo.move}, and the 
test then checks that
+ * the other mutation is blocked on the monitor until the gate is released.
+ */
+class NoteManagerMutationSerializationTest {
+
+  private static final long JOIN_TIMEOUT_MILLIS = 30_000L;
+  private static final long GATE_ARRIVAL_TIMEOUT_SECONDS = 30L;
+  private static final long BLOCKED_TIMEOUT_MILLIS = 5_000L;
+
+  private File notebookDir;
+  private ZeppelinConfiguration zConf;
+  private NoteParser noteParser;
+  private VFSNotebookRepoWithMoveGate notebookRepo;
+  private NoteManager noteManager;
+
+  @BeforeEach
+  void setUp() throws Exception {
+    notebookDir = 
Files.createTempDirectory("notebookDir").toAbsolutePath().toFile();
+    zConf = ZeppelinConfiguration.load();
+    
zConf.setProperty(ZeppelinConfiguration.ConfVars.ZEPPELIN_NOTEBOOK_DIR.getVarName(),
+        notebookDir.getAbsolutePath());
+    noteParser = new GsonNoteParser(zConf);
+    notebookRepo = new VFSNotebookRepoWithMoveGate();
+    notebookRepo.init(zConf, noteParser);
+    noteManager = new NoteManager(notebookRepo, zConf);
+  }
+
+  @AfterEach
+  void tearDown() throws IOException {
+    FileUtils.deleteDirectory(notebookDir);
+  }
+
+  /**
+   * Given a moveNote parked inside notebookRepo.move (monitor held), when 
another thread calls
+   * removeFolder, then it must wait for the monitor and only finish after the 
move resumes.
+   */
+  @Test
+  void testRemoveFolderWaitsForInFlightMoveNote() throws Exception {
+    Note moving = createAndSave("/src/moving");
+    Note victim = createAndSave("/victim/note");
+
+    List<Throwable> errors = runBlockedMutationAfterParkedMove(
+        moving.getId(), "/dst/moving",
+        () -> noteManager.removeFolder("/victim", 
AuthenticationInfo.ANONYMOUS));
+
+    assertTrue(errors.isEmpty(), () -> "Mutations threw: " + errors);
+    assertFalse(noteManager.getNotesInfo().containsKey(victim.getId()));
+    assertEquals("/dst/moving", 
noteManager.getNotesInfo().get(moving.getId()));
+  }
+
+  /**
+   * Given a moveNote parked inside notebookRepo.move (monitor held), when 
another thread calls
+   * addNote, then it must wait for the monitor and only finish after the move 
resumes.
+   */
+  @Test
+  void testAddNoteWaitsForInFlightMoveNote() throws Exception {
+    Note moving = createAndSave("/src/moving");
+    Note added = newNote("/other/added");
+
+    List<Throwable> errors = runBlockedMutationAfterParkedMove(
+        moving.getId(), "/dst/moving",
+        () -> noteManager.addNote(added, AuthenticationInfo.ANONYMOUS));
+
+    assertTrue(errors.isEmpty(), () -> "Mutations threw: " + errors);
+    assertEquals("/other/added", 
noteManager.getNotesInfo().get(added.getId()));
+    assertEquals("/dst/moving", 
noteManager.getNotesInfo().get(moving.getId()));
+  }
+
+  private interface Mutation {
+    void run() throws IOException;
+  }
+
+  private List<Throwable> runBlockedMutationAfterParkedMove(
+      String movingNoteId, String newPath, Mutation mutation) throws Exception 
{
+    List<Throwable> errors = Collections.synchronizedList(new ArrayList<>());
+    notebookRepo.armGate();
+
+    Thread mover = new Thread(() -> {
+      try {
+        noteManager.moveNote(movingNoteId, newPath, 
AuthenticationInfo.ANONYMOUS);
+      } catch (Throwable t) {
+        errors.add(t);
+      }
+    }, "mutation-serialization-mover");
+    mover.start();
+
+    assertTrue(
+        notebookRepo.awaitArrival(GATE_ARRIVAL_TIMEOUT_SECONDS, 
TimeUnit.SECONDS),
+        "moveNote never reached the gated notebookRepo.move()");
+
+    Thread other = new Thread(() -> {
+      try {
+        mutation.run();
+      } catch (Throwable t) {
+        errors.add(t);
+      }
+    }, "mutation-serialization-other");
+    other.start();
+
+    try {
+      awaitBlockedOnMonitor(other);
+    } finally {
+      notebookRepo.release();
+    }
+    mover.join(JOIN_TIMEOUT_MILLIS);
+    other.join(JOIN_TIMEOUT_MILLIS);
+    assertFalse(mover.isAlive(), "moveNote did not finish within the timeout");
+    assertFalse(other.isAlive(), "The other mutation did not finish within the 
timeout");
+    return errors;
+  }
+
+  private static void awaitBlockedOnMonitor(Thread thread) throws 
InterruptedException {
+    long deadline = System.nanoTime() + 
TimeUnit.MILLISECONDS.toNanos(BLOCKED_TIMEOUT_MILLIS);
+    while (System.nanoTime() < deadline) {
+      Thread.State state = thread.getState();
+      if (state == Thread.State.BLOCKED) {
+        return;
+      }
+      if (state == Thread.State.TERMINATED) {
+        fail("The mutation completed while moveNote held the NoteManager 
monitor; "
+            + "it is not serialized with moveNote");
+      }
+      Thread.sleep(10);
+    }
+    fail("The mutation neither blocked on the NoteManager monitor nor 
finished");
+  }
+
+  private Note newNote(String notePath) {
+    return new Note(notePath, "test", null, null, null, null, null, zConf, 
noteParser);
+  }
+
+  private Note createAndSave(String notePath) throws IOException {
+    Note note = newNote(notePath);
+    noteManager.saveNote(note, AuthenticationInfo.ANONYMOUS);
+    return note;
+  }
+}
diff --git 
a/zeppelin-server/src/test/java/org/apache/zeppelin/notebook/NoteManagerTest.java
 
b/zeppelin-server/src/test/java/org/apache/zeppelin/notebook/NoteManagerTest.java
index eaed222f9e..9f6f97f8cf 100644
--- 
a/zeppelin-server/src/test/java/org/apache/zeppelin/notebook/NoteManagerTest.java
+++ 
b/zeppelin-server/src/test/java/org/apache/zeppelin/notebook/NoteManagerTest.java
@@ -31,10 +31,13 @@ import java.util.List;
 import java.util.Map;
 import java.util.concurrent.ConcurrentHashMap;
 import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.CyclicBarrier;
 import java.util.concurrent.ExecutorService;
 import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
 import java.util.concurrent.TimeUnit;
 import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.concurrent.atomic.AtomicInteger;
 
 import static org.junit.jupiter.api.Assertions.assertEquals;
 import static org.junit.jupiter.api.Assertions.assertFalse;
@@ -267,6 +270,111 @@ class NoteManagerTest {
     }
   }
 
+  @Test
+  void testConcurrentAddNoteAndMoveNoteOnSamePath() throws Exception {
+    int rounds = 100, adders = 4;
+    for (int r = 0; r < rounds; r++) {
+      String targetPath = "/dst_" + r + "/note";
+      Note moving = createNote("/src_" + r + "/note");
+      noteManager.saveNote(moving);
+
+      List<Throwable> failures = Collections.synchronizedList(new 
ArrayList<>());
+      AtomicInteger winners = new AtomicInteger();
+      List<Runnable> tasks = new ArrayList<>();
+      for (int i = 0; i < adders; i++) {
+        Note added = createNote(targetPath);
+        tasks.add(() -> {
+          try {
+            noteManager.addNote(added, AuthenticationInfo.ANONYMOUS);
+            winners.incrementAndGet();
+          } catch (NotePathAlreadyExistsException e) {
+            // expected for every loser
+          } catch (Throwable t) {
+            failures.add(t);
+          }
+        });
+      }
+      tasks.add(() -> {
+        try {
+          noteManager.moveNote(moving.getId(), targetPath, 
AuthenticationInfo.ANONYMOUS);
+          winners.incrementAndGet();
+        } catch (NotePathAlreadyExistsException e) {
+          // expected when an addNote won
+        } catch (Throwable t) {
+          failures.add(t);
+        }
+      });
+      runConcurrently(tasks);
+
+      assertTrue(failures.isEmpty(), () -> "Unexpected failures: " + failures);
+      assertEquals(1, winners.get(), "exactly one operation may claim " + 
targetPath);
+      assertEquals(1, noteManager.getNotesInfo().values().stream()
+          .filter(targetPath::equals).count(), "notesInfo must map one note to 
" + targetPath);
+    }
+  }
+
+  @Test
+  void testConcurrentRemoveFolderAndMoveNote() throws Exception {
+    int rounds = 100;
+    for (int r = 0; r < rounds; r++) {
+      String folder = "/folder_" + r;
+      Note inFolder = createNote(folder + "/in_folder");
+      Note moving = createNote("/src_" + r + "/moving");
+      noteManager.saveNote(inFolder);
+      noteManager.saveNote(moving);
+
+      List<Throwable> failures = Collections.synchronizedList(new 
ArrayList<>());
+      List<Runnable> tasks = new ArrayList<>();
+      for (int i = 0; i < 2; i++) {
+        tasks.add(() -> {
+          try {
+            noteManager.removeFolder(folder, AuthenticationInfo.ANONYMOUS);
+          } catch (IOException e) {
+            // the folder was already removed by the other removeFolder call
+          } catch (Throwable t) {
+            failures.add(t);
+          }
+        });
+      }
+      tasks.add(() -> {
+        try {
+          noteManager.moveNote(moving.getId(), folder + "/moving", 
AuthenticationInfo.ANONYMOUS);
+        } catch (Throwable t) {
+          failures.add(t);
+        }
+      });
+      runConcurrently(tasks);
+
+      assertTrue(failures.isEmpty(), () -> "Unexpected failures: " + failures);
+      assertFalse(noteManager.getNotesInfo().containsKey(inFolder.getId()));
+      // every remaining mapping entry must still resolve to a note in the tree
+      for (String noteId : noteManager.getNotesInfo().keySet()) {
+        assertNotNull(noteManager.processNote(noteId, note -> note),
+            "notesInfo entry " + noteId + " no longer resolves to a note");
+      }
+    }
+  }
+
+  private void runConcurrently(List<Runnable> tasks) throws Exception {
+    CyclicBarrier start = new CyclicBarrier(tasks.size());
+    ExecutorService pool = Executors.newFixedThreadPool(tasks.size());
+    try {
+      List<Future<?>> futures = new ArrayList<>();
+      for (Runnable task : tasks) {
+        futures.add(pool.submit(() -> {
+          start.await();
+          task.run();
+          return null;
+        }));
+      }
+      for (Future<?> future : futures) {
+        future.get(30, TimeUnit.SECONDS);
+      }
+    } finally {
+      pool.shutdownNow();
+    }
+  }
+
   abstract class ConcurrentTask {
     private ExecutorService threadPool;
     private int noteNum;
diff --git 
a/zeppelin-server/src/test/java/org/apache/zeppelin/notebook/repo/VFSNotebookRepoWithMoveGate.java
 
b/zeppelin-server/src/test/java/org/apache/zeppelin/notebook/repo/VFSNotebookRepoWithMoveGate.java
new file mode 100644
index 0000000000..c750dc303d
--- /dev/null
+++ 
b/zeppelin-server/src/test/java/org/apache/zeppelin/notebook/repo/VFSNotebookRepoWithMoveGate.java
@@ -0,0 +1,78 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *    http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.zeppelin.notebook.repo;
+
+import java.io.IOException;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
+import org.apache.zeppelin.user.AuthenticationInfo;
+
+/**
+ * Test-only subclass of {@link VFSNotebookRepo} that parks the first note
+ * {@code move(noteId, ...)} call after arming. {@code NoteManager#moveNote} 
calls this method
+ * while holding the NoteManager monitor, so a parked call pins the monitor in 
the held state
+ * and lets a test check whether other mutations wait for it.
+ */
+public class VFSNotebookRepoWithMoveGate extends VFSNotebookRepo {
+
+  private static final long GATE_SELF_TIMEOUT_SECONDS = 30;
+
+  private final AtomicBoolean armed = new AtomicBoolean(false);
+  private volatile CountDownLatch arrivedLatch;
+  private volatile CountDownLatch releaseLatch;
+
+  /**
+   * Arm the gate. Only the next note {@code move()} call parks; later calls 
pass through.
+   */
+  public void armGate() {
+    arrivedLatch = new CountDownLatch(1);
+    releaseLatch = new CountDownLatch(1);
+    armed.set(true);
+  }
+
+  /**
+   * Wait for the gated {@code move()} call to arrive and park. Returns false 
if it does not
+   * arrive within the timeout so the caller can fail with a clear message 
instead of hanging.
+   */
+  public boolean awaitArrival(long timeout, TimeUnit unit) throws 
InterruptedException {
+    return arrivedLatch.await(timeout, unit);
+  }
+
+  /**
+   * Let the parked {@code move()} call resume.
+   */
+  public void release() {
+    releaseLatch.countDown();
+  }
+
+  @Override
+  public void move(String noteId, String notePath, String newNotePath,
+                   AuthenticationInfo subject) throws IOException {
+    if (armed.compareAndSet(true, false)) {
+      arrivedLatch.countDown();
+      try {
+        // Self-timeout so a test that forgets release() fails fast instead of 
hanging.
+        releaseLatch.await(GATE_SELF_TIMEOUT_SECONDS, TimeUnit.SECONDS);
+      } catch (InterruptedException e) {
+        Thread.currentThread().interrupt();
+      }
+    }
+    super.move(noteId, notePath, newNotePath, subject);
+  }
+}

Reply via email to