diff --git a/zeppelin-plugins/notebookrepo/filesystem/src/main/java/org/apache/zeppelin/notebook/repo/FileSystemNotebookRepo.java b/zeppelin-plugins/notebookrepo/filesystem/src/main/java/org/apache/zeppelin/notebook/repo/FileSystemNotebookRepo.java index 973c18d412c..9722acfdb23 100644 --- a/zeppelin-plugins/notebookrepo/filesystem/src/main/java/org/apache/zeppelin/notebook/repo/FileSystemNotebookRepo.java +++ b/zeppelin-plugins/notebookrepo/filesystem/src/main/java/org/apache/zeppelin/notebook/repo/FileSystemNotebookRepo.java @@ -55,6 +55,36 @@ public void init(ZeppelinConfiguration zConf, NoteParser noteParser) throws IOEx this.notebookDir = this.fs.makeQualified(new Path(zConf.getNotebookDir())); LOGGER.info("Using folder {} to store notebook", notebookDir); this.fs.tryMkDir(notebookDir); + recoverInterruptedSaves(); + } + + /** + * A save interrupted between its two renames leaves the .bak and .tmp files of a note but + * no .zpln file, so {@link #list} does not pick it up. Restore those notes before they are + * listed. A single leftover .tmp or .bak file is not restored, because it can't be told apart + * from a leftover of a removed or moved note. It is only logged. + */ + private void recoverInterruptedSaves() { + try { + List recovered = + fs.recoverInterruptedWrites(notebookDir, ".zpln", this::isValidNote); + if (!recovered.isEmpty()) { + LOGGER.warn("Recovered {} note(s) whose last save was interrupted: {}", + recovered.size(), recovered); + } + } catch (IOException e) { + LOGGER.warn("Fail to recover notes whose last save was interrupted", e); + } + } + + private boolean isValidNote(Path noteFile, String content) { + try { + noteParser.fromJson(getNoteId(noteFile.getName()), content); + return true; + } catch (IOException | RuntimeException e) { + LOGGER.warn("Content of the temp file of {} is not a valid note", noteFile); + return false; + } } @Override diff --git a/zeppelin-plugins/notebookrepo/filesystem/src/test/java/org/apache/zeppelin/notebook/repo/FileSystemNotebookRepoTest.java b/zeppelin-plugins/notebookrepo/filesystem/src/test/java/org/apache/zeppelin/notebook/repo/FileSystemNotebookRepoTest.java index 1d7f8d19def..bac7b119921 100644 --- a/zeppelin-plugins/notebookrepo/filesystem/src/test/java/org/apache/zeppelin/notebook/repo/FileSystemNotebookRepoTest.java +++ b/zeppelin-plugins/notebookrepo/filesystem/src/test/java/org/apache/zeppelin/notebook/repo/FileSystemNotebookRepoTest.java @@ -35,11 +35,16 @@ import java.io.File; import java.io.IOException; import java.io.OutputStream; +import java.nio.charset.StandardCharsets; import java.nio.file.Files; +import java.nio.file.Paths; import java.util.HashMap; import java.util.Map; +import java.util.stream.Stream; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertTrue; class FileSystemNotebookRepoTest { @@ -151,4 +156,131 @@ void testComplicatedScenarios() throws IOException { hdfsNotebookRepo.save(note, authInfo); assertEquals(1, hdfsNotebookRepo.list(authInfo).size()); } + + @Test + void testSaveLeavesNoTempOrBackupFile() throws IOException { + Note note = createNote("/title_1", "value_1"); + hdfsNotebookRepo.save(note, authInfo); + note.getConfig().put("config_1", "value_2"); + hdfsNotebookRepo.save(note, authInfo); + + assertEquals("value_2", getConfigValue(note)); + try (Stream files = Files.list(Paths.get(notebookDir))) { + assertTrue(files.allMatch(f -> f.getFileName().toString().endsWith(".zpln"))); + } + } + + @Test + void testRecoverFromTempFileWhenSaveStoppedBeforeLastRename() throws IOException { + Note note = createNote("/title_1", "value_1"); + hdfsNotebookRepo.save(note, authInfo); + // The original was moved to .bak and the new content fully written to .tmp + note.getConfig().put("config_1", "value_2"); + Files.move(noteFile(note, ""), noteFile(note, ".bak")); + writeString(noteFile(note, ".tmp"), note.toJson()); + + restartRepo(); + + assertEquals(1, hdfsNotebookRepo.list(authInfo).size()); + assertEquals("value_2", getConfigValue(note)); + assertFalse(Files.exists(noteFile(note, ".bak"))); + } + + @Test + void testRecoverFromBackupFileWhenTempFileIsIncomplete() throws IOException { + Note note = createNote("/title_1", "value_1"); + hdfsNotebookRepo.save(note, authInfo); + Files.move(noteFile(note, ""), noteFile(note, ".bak")); + String json = note.toJson(); + writeString(noteFile(note, ".tmp"), json.substring(0, json.length() / 2)); + + restartRepo(); + + assertEquals(1, hdfsNotebookRepo.list(authInfo).size()); + assertEquals("value_1", getConfigValue(note)); + } + + @Test + void testDoNotRecoverFromTempFileOnly() throws IOException { + // A lone .tmp cannot be told apart from a leftover of a deleted note + Note note = createNote("/title_1", "value_1"); + writeString(noteFile(note, ".tmp"), note.toJson()); + + restartRepo(); + + assertEquals(0, hdfsNotebookRepo.list(authInfo).size()); + assertTrue(Files.exists(noteFile(note, ".tmp"))); + } + + @Test + void testRemovedNoteIsNotRestoredFromLeftoverTmp() throws IOException { + Note note = createNote("/title_1", "value_1"); + hdfsNotebookRepo.save(note, authInfo); + writeString(noteFile(note, ".tmp"), note.toJson()); + hdfsNotebookRepo.remove(note.getId(), note.getPath(), authInfo); + + restartRepo(); + + assertEquals(0, hdfsNotebookRepo.list(authInfo).size()); + } + + @Test + void testMovedNoteIsNotDuplicatedFromLeftoverTmp() throws IOException { + Note note = createNote("/title_1", "value_1"); + hdfsNotebookRepo.save(note, authInfo); + writeString(noteFile(note, ".tmp"), note.toJson()); + hdfsNotebookRepo.move(note.getId(), "/title_1", "/dir/title_2", authInfo); + + restartRepo(); + + long zplnCount; + try (Stream files = Files.walk(Paths.get(notebookDir))) { + zplnCount = files.filter(f -> f.toString().endsWith(".zpln")).count(); + } + assertEquals(1, zplnCount); + } + + @Test + void testKeepExistingNoteWhenTempFileIsLeftOver() throws IOException { + Note note = createNote("/title_1", "value_1"); + hdfsNotebookRepo.save(note, authInfo); + writeString(noteFile(note, ".tmp"), "{"); + + restartRepo(); + + assertEquals(1, hdfsNotebookRepo.list(authInfo).size()); + assertEquals("value_1", getConfigValue(note)); + assertTrue(Files.exists(noteFile(note, ".tmp"))); + assertFalse(Files.exists(noteFile(note, ".bak"))); + } + + private Note createNote(String path, String configValue) { + Note note = new Note(); + note.setZeppelinConfiguration(zConf); + note.setNoteParser(noteParser); + note.setPath(path); + Map config = new HashMap<>(); + config.put("config_1", configValue); + note.setConfig(config); + return note; + } + + private java.nio.file.Path noteFile(Note note, String suffix) throws IOException { + return Paths.get(notebookDir, + hdfsNotebookRepo.buildNoteFileName(note.getId(), note.getPath()) + suffix); + } + + private void writeString(java.nio.file.Path file, String content) throws IOException { + Files.write(file, content.getBytes(StandardCharsets.UTF_8)); + } + + private Object getConfigValue(Note note) throws IOException { + return hdfsNotebookRepo.get(note.getId(), note.getPath(), authInfo) + .getConfig().get("config_1"); + } + + private void restartRepo() throws IOException { + hdfsNotebookRepo = new FileSystemNotebookRepo(); + hdfsNotebookRepo.init(zConf, noteParser); + } } diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/notebook/FileSystemStorage.java b/zeppelin-server/src/main/java/org/apache/zeppelin/notebook/FileSystemStorage.java index 5fc60e74e55..ef170ae7232 100644 --- a/zeppelin-server/src/main/java/org/apache/zeppelin/notebook/FileSystemStorage.java +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/notebook/FileSystemStorage.java @@ -40,8 +40,10 @@ import java.nio.file.attribute.PosixFilePermissions; import java.security.PrivilegedExceptionAction; import java.util.ArrayList; +import java.util.LinkedHashSet; import java.util.List; import java.util.Set; +import java.util.function.BiPredicate; /** @@ -52,6 +54,8 @@ public class FileSystemStorage { private static final Logger LOGGER = LoggerFactory.getLogger(FileSystemStorage.class); private static final String S3A = "s3a"; private static final String FS_DEFAULTFS = "fs.defaultFS"; + static final String TMP_SUFFIX = ".tmp"; + static final String BACKUP_SUFFIX = ".bak"; // only do UserGroupInformation.loginUserFromKeytab one time, otherwise you will still get // your ticket expired. @@ -239,16 +243,136 @@ public void writeFile(final String content, final Path file, boolean writeTempFi public Void call() throws IOException { InputStream in = new ByteArrayInputStream(content.getBytes( zConf.getString(ZeppelinConfiguration.ConfVars.ZEPPELIN_ENCODING))); - Path tmpFile = new Path(file.toString() + ".tmp"); + Path tmpFile = new Path(file.toString() + TMP_SUFFIX); IOUtils.copyBytes(in, fs.create(tmpFile), hadoopConf); fs.setPermission(tmpFile, fsPermission); - fs.delete(file, true); - fs.rename(tmpFile, file); + replaceFile(tmpFile, file); return null; } }); } + /** + * Replaces file with tmpFile without deleting the original first. The original is renamed to + * a backup and deleted only after tmpFile is in place, so an interrupted write always leaves a + * complete copy behind. + */ + private void replaceFile(Path tmpFile, Path file) throws IOException { + Path backupFile = new Path(file.toString() + BACKUP_SUFFIX); + boolean hasOriginal = fs.exists(file); + if (hasOriginal) { + // A backup next to an existing file is left over from an earlier write and is older + // than the file itself. + fs.delete(backupFile, false); + if (!fs.rename(file, backupFile)) { + throw new IOException("Fail to back up " + file + " to " + backupFile); + } + } + if (!fs.rename(tmpFile, file)) { + if (hasOriginal) { + if (fs.rename(backupFile, file)) { + // The original is back in place, so the temp file is no longer needed. Leaving it + // would let it be mistaken for an interrupted write later. + fs.delete(tmpFile, false); + } else { + LOGGER.error("Fail to restore {} from {}, please restore it manually", + file, backupFile); + } + } + throw new IOException("Fail to rename " + tmpFile + " to " + file); + } + if (hasOriginal && !fs.delete(backupFile, false)) { + LOGGER.warn("Fail to delete backup file {}", backupFile); + } + } + + /** + * Restores files under dir (recursively) whose last {@link #writeFile} was interrupted. + * A file is restored only when it is missing and both its temp and backup files exist: + * from the temp file if isComplete accepts its content, otherwise from the backup file. + * A single leftover file is only logged. Leftover files next to an existing file are kept. + * + * @param dir folder to scan recursively + * @param targetSuffix suffix of the files to restore, e.g. ".zpln" + * @param isComplete given the file to restore and the content of its temp file, tells whether + * that content is complete + * @return the restored files + */ + public List recoverInterruptedWrites(final Path dir, final String targetSuffix, + final BiPredicate isComplete) throws IOException { + return callHdfsOperation(new HdfsOperation>() { + @Override + public List call() throws IOException { + List recovered = new ArrayList<>(); + if (!fs.exists(dir)) { + return recovered; + } + Set missingFiles = new LinkedHashSet<>(); + collectMissingFiles(dir, targetSuffix, missingFiles); + for (Path file : missingFiles) { + if (recoverFile(file, isComplete)) { + recovered.add(file); + } + } + return recovered; + } + }); + } + + private void collectMissingFiles(Path folder, String targetSuffix, Set missingFiles) + throws IOException { + for (FileStatus status : fs.listStatus(folder)) { + Path path = status.getPath(); + if (status.isDirectory()) { + collectMissingFiles(path, targetSuffix, missingFiles); + continue; + } + for (String suffix : new String[] {TMP_SUFFIX, BACKUP_SUFFIX}) { + if (path.getName().endsWith(targetSuffix + suffix)) { + String pathString = path.toString(); + Path file = new Path(pathString.substring(0, pathString.length() - suffix.length())); + if (!fs.exists(file)) { + missingFiles.add(file); + } + } + } + } + } + + private boolean recoverFile(Path file, BiPredicate isComplete) + throws IOException { + Path tmpFile = new Path(file.toString() + TMP_SUFFIX); + Path backupFile = new Path(file.toString() + BACKUP_SUFFIX); + // writeFile leaves both files only when it stops between its two renames. A single leftover + // may belong to a file that was deleted or moved on purpose, so it is not restored. + if (!fs.exists(tmpFile) || !fs.exists(backupFile)) { + LOGGER.warn("Found {} without {}, not restoring it automatically. " + + "Rename it manually if it should be restored.", + fs.exists(tmpFile) ? tmpFile : backupFile, file); + return false; + } + // The temp file is newer than the backup, but it may be incomplete if the write stopped + // while it was being written. + if (isComplete.test(file, readContent(tmpFile)) && fs.rename(tmpFile, file)) { + fs.delete(backupFile, false); + LOGGER.warn("Recovered {} from {}", file, tmpFile); + return true; + } + if (fs.rename(backupFile, file)) { + LOGGER.warn("Recovered {} from {}", file, backupFile); + return true; + } + LOGGER.error("Fail to recover {}, please check {} and {} manually", + file, tmpFile, backupFile); + return false; + } + + private String readContent(Path file) throws IOException { + ByteArrayOutputStream bytes = new ByteArrayOutputStream(); + IOUtils.copyBytes(fs.open(file), bytes, hadoopConf); + return bytes.toString(zConf.getString(ZeppelinConfiguration.ConfVars.ZEPPELIN_ENCODING)); + } + public void move(Path src, Path dest) throws IOException { callHdfsOperation(() -> { fs.rename(src, dest); diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/notebook/repo/VFSNotebookRepo.java b/zeppelin-server/src/main/java/org/apache/zeppelin/notebook/repo/VFSNotebookRepo.java index 420a1c148a9..ba739bfaedb 100644 --- a/zeppelin-server/src/main/java/org/apache/zeppelin/notebook/repo/VFSNotebookRepo.java +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/notebook/repo/VFSNotebookRepo.java @@ -24,6 +24,9 @@ import java.net.URISyntaxException; import java.net.URLDecoder; import java.nio.charset.StandardCharsets; +import java.nio.file.AtomicMoveNotSupportedException; +import java.nio.file.Files; +import java.nio.file.StandardCopyOption; import java.util.ArrayList; import java.util.Collections; import java.util.HashMap; @@ -175,10 +178,36 @@ public synchronized void save(Note note, AuthenticationInfo subject) throws IOEx out.close(); } } - noteJson.moveTo(rootNotebookFileObject.resolveFile( + moveReplacing(noteJson, rootNotebookFileObject.resolveFile( buildNoteFileName(note), NameScope.DESCENDENT)); } + /** + * Moves src over dest. FileObject.moveTo deletes dest before renaming src, so a crash in + * between loses the note. For local files, replace dest in a single atomic rename instead. + * Other file systems, and local ones without atomic rename, keep using moveTo. + */ + private void moveReplacing(FileObject src, FileObject dest) throws IOException { + if ("file".equals(src.getName().getScheme())) { + try { + Files.move(src.getPath(), dest.getPath(), + StandardCopyOption.ATOMIC_MOVE, StandardCopyOption.REPLACE_EXISTING); + // The move bypassed VFS, so drop the state it cached for these files. + src.refresh(); + dest.refresh(); + FileObject parent = dest.getParent(); + if (parent != null) { + parent.refresh(); + } + return; + } catch (AtomicMoveNotSupportedException e) { + LOGGER.warn("Atomic move is not supported for {}, falling back to a non-atomic move", + dest.getName(), e); + } + } + src.moveTo(dest); + } + @Override public void move(String noteId, String notePath, diff --git a/zeppelin-server/src/test/java/org/apache/zeppelin/notebook/repo/VFSNotebookRepoTest.java b/zeppelin-server/src/test/java/org/apache/zeppelin/notebook/repo/VFSNotebookRepoTest.java index c7dcc9bdc44..ce279dfc6e6 100644 --- a/zeppelin-server/src/test/java/org/apache/zeppelin/notebook/repo/VFSNotebookRepoTest.java +++ b/zeppelin-server/src/test/java/org/apache/zeppelin/notebook/repo/VFSNotebookRepoTest.java @@ -115,6 +115,27 @@ void testBasics() throws IOException { assertEquals(1, notebookRepo.list(AuthenticationInfo.ANONYMOUS).size()); } + @Test + void testSaveReplacesExistingNote() throws IOException { + Note note = new Note(); + note.setPath("/my_project/my_note1"); + note.setNoteParser(noteParser); + Paragraph p = note.insertNewParagraph(0, AuthenticationInfo.ANONYMOUS); + p.setText("%md hello world"); + notebookRepo.save(note, AuthenticationInfo.ANONYMOUS); + + p.setText("%md hello world2"); + notebookRepo.save(note, AuthenticationInfo.ANONYMOUS); + + assertEquals(1, notebookRepo.list(AuthenticationInfo.ANONYMOUS).size()); + Note savedNote = notebookRepo.get(note.getId(), note.getPath(), AuthenticationInfo.ANONYMOUS); + assertEquals("%md hello world2", savedNote.getParagraphs().get(0).getText()); + File[] files = new File(notebookRepo.rootNotebookFolder, "my_project").listFiles(); + assertEquals(1, files.length); + assertEquals(notebookRepo.buildNoteFileName(note), + "my_project/" + files[0].getName()); + } + @Test void testNoteNameWithColon() throws IOException { assertEquals(0, notebookRepo.list(AuthenticationInfo.ANONYMOUS).size());