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);
+ }
+}