nio2: NIO.2 file API companion code (streaming 5 GB at -Xmx64m, handle leaks, WatchService, mmap)

Co-Authored-By: Claude Sonnet 5.5 <[email protected]>
Claude-Session: https://claude.ai/code/session_01KqJyCidz3ZgRyHABv2GVJh
This commit is contained in:
2026-09-30 19:14:31 +00:00
co-authored by Claude Sonnet 5.5
parent f4b732799d
commit 1ac2f6077a
30 changed files with 1072 additions and 0 deletions
@@ -0,0 +1,27 @@
package com.ankurm.nio2;
import java.io.IOException;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.ArrayList;
import java.util.List;
import java.util.stream.Stream;
/** Run under a low `ulimit -n` to see what a leak turns into. */
public final class HandleExhaustion {
private HandleExhaustion() {}
public static void main(String[] a) throws IOException {
Path tmp = Files.createTempDirectory("nio2-exhaust");
Path tree = HandleLeakDemo.makeTree(tmp.resolve("tree"));
List<Stream<Path>> leaked = new ArrayList<>();
int i = 0;
try {
for (; ; i++) { Stream<Path> s = Files.walk(tree); s.skip(3).findFirst(); leaked.add(s); }
} catch (Exception e) {
int n = leaked.size();
System.out.printf("failed on call %d after %d leaked walks%n%s: %s%n", i + 1, n,
e.getClass().getName(), e.getMessage() == null ? "" : e.getMessage().replace(tmp.toString(), "<tmp>"));
}
}
}
@@ -0,0 +1,68 @@
package com.ankurm.nio2;
import java.io.IOException;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.ArrayList;
import java.util.List;
import java.util.stream.Stream;
/** Counts this process's open file descriptors (Linux /proc/self/fd) around unclosed Files.lines / Files.walk. */
public final class HandleLeakDemo {
private HandleLeakDemo() {}
public static int openFds() throws IOException {
try (Stream<Path> s = Files.list(Path.of("/proc/self/fd"))) { return (int) s.count(); }
}
static Path makeTree(Path root) throws IOException {
for (int i = 0; i < 5; i++) {
Path d = Files.createDirectories(root.resolve("d" + i).resolve("sub"));
Files.writeString(d.resolve("f.txt"), "x\n");
Files.writeString(root.resolve("d" + i).resolve("g.txt"), "y\n");
}
return root;
}
public static void main(String[] a) throws Exception {
Path tmp = Files.createTempDirectory("nio2-leak");
Path file = Files.writeString(tmp.resolve("one.txt"), "a\nb\nc\n");
Path tree = makeTree(tmp.resolve("tree"));
final int N = 50;
int base = openFds();
System.out.println("open fds at start: " + base);
// 0. Files.walk fully consumed but not closed, run first so later leaks do not hide it.
for (int i = 0; i < N; i++) { Files.walk(tree).count(); }
System.out.printf("after %d fully consumed, unclosed Files.walk: +%d fds%n", N, openFds() - base);
// 1. Files.lines, opened and read partially, never closed. References are kept (as a field or cache would).
List<Stream<String>> keep = new ArrayList<>();
for (int i = 0; i < N; i++) { Stream<String> s = Files.lines(file); s.findFirst(); keep.add(s); }
System.out.printf("after %d unclosed Files.lines + findFirst(): +%d fds%n", N, openFds() - base);
keep.forEach(Stream::close);
System.out.printf("after closing them: +%d fds%n", openFds() - base);
// 2. Files.lines in try-with-resources, same work.
for (int i = 0; i < N; i++) { try (Stream<String> s = Files.lines(file)) { s.findFirst(); } }
System.out.printf("after %d try-with-resources Files.lines: +%d fds%n", N, openFds() - base);
// 3. Files.walk, findFirst(), never closed.
List<Stream<Path>> walks = new ArrayList<>();
for (int i = 0; i < N; i++) { Stream<Path> s = Files.walk(tree); s.skip(3).findFirst(); walks.add(s); }
System.out.printf("after %d unclosed Files.walk + findFirst(): +%d fds%n", N, openFds() - base);
walks.forEach(Stream::close);
System.out.printf("after closing them: +%d fds%n", openFds() - base);
// 3b. The same unclosed Files.lines, but nobody keeps a reference: does the garbage collector clean up?
for (int i = 0; i < N; i++) { Files.lines(file).findFirst(); }
System.out.printf("after %d unreferenced, unclosed Files.lines: +%d fds%n", N, openFds() - base);
System.gc(); Thread.sleep(500);
System.out.printf("after System.gc(): +%d fds%n", openFds() - base);
for (int i = 0; i < N; i++) { Files.walk(tree).skip(3).findFirst(); }
System.out.printf("after %d unreferenced, unclosed Files.walk: +%d fds%n", N, openFds() - base);
System.gc(); Thread.sleep(500);
System.out.printf("after System.gc(): +%d fds%n", openFds() - base);
}
}
@@ -0,0 +1,48 @@
package com.ankurm.nio2;
import java.lang.management.GarbageCollectorMXBean;
import java.lang.management.ManagementFactory;
import java.nio.file.Files;
import java.nio.file.Path;
/** Samples used heap on a daemon thread and reports peak-used, max heap and GC counts. */
public final class HeapProbe implements AutoCloseable {
private volatile boolean run = true;
private volatile long peak;
private final Thread t;
public HeapProbe() {
t = new Thread(() -> {
Runtime rt = Runtime.getRuntime();
while (run) {
long used = rt.totalMemory() - rt.freeMemory();
if (used > peak) peak = used;
try { Thread.sleep(20); } catch (InterruptedException e) { return; }
}
}, "heap-probe");
t.setDaemon(true);
t.start();
}
public long peakUsed() { return peak; }
public String report() {
long gcs = 0, ms = 0;
for (GarbageCollectorMXBean b : ManagementFactory.getGarbageCollectorMXBeans()) {
gcs += b.getCollectionCount(); ms += b.getCollectionTime();
}
return String.format("heap: -Xmx=%d MiB, peak used (sampled every 20 ms)=%d MiB, GC collections=%d, GC time=%d ms",
Runtime.getRuntime().maxMemory() >> 20, peak >> 20, gcs, ms);
}
@Override public void close() { run = false; }
/** Resident set size of this process in MiB (Linux only), or -1. */
public static long rssMiB() {
try {
for (String l : Files.readAllLines(Path.of("/proc/self/status")))
if (l.startsWith("VmRSS:")) return Long.parseLong(l.replaceAll("\\D+", "")) >> 10;
} catch (Exception ignored) { }
return -1;
}
}
@@ -0,0 +1,36 @@
package com.ankurm.nio2;
import java.io.IOException;
import java.io.UncheckedIOException;
import java.nio.charset.StandardCharsets;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.stream.Stream;
/** Files.lines is lazy, so a bad byte fails in the middle of the pipeline, as an unchecked exception. */
public final class LinesCharsetDemo {
private LinesCharsetDemo() {}
public static void main(String[] a) throws IOException {
Path p = Files.createTempFile("nio2-charset", ".txt");
byte[] ok = "good line 1\ngood line 2\n".getBytes(StandardCharsets.UTF_8);
byte[] bad = {(byte) 0xE9, '\n'}; // 0xE9 is 'e-acute' in ISO-8859-1 but not valid UTF-8 here
Files.write(p, ok);
Files.write(p, bad, java.nio.file.StandardOpenOption.APPEND);
int[] seen = {0};
try (Stream<String> s = Files.lines(p)) {
s.forEach(l -> seen[0]++);
} catch (UncheckedIOException e) {
System.out.println("lines processed before failure: " + seen[0]);
System.out.println(e.getClass().getSimpleName() + " caused by " + e.getCause());
}
try (Stream<String> s = Files.lines(p, StandardCharsets.ISO_8859_1)) {
System.out.println("with ISO_8859_1: " + s.count() + " lines, no error");
}
try {
Files.readAllLines(p);
} catch (java.nio.charset.MalformedInputException e) {
System.out.println("readAllLines throws checked " + e.getClass().getSimpleName());
}
}
}
@@ -0,0 +1,46 @@
package com.ankurm.nio2;
import java.io.BufferedWriter;
import java.io.IOException;
import java.nio.file.Files;
import java.nio.file.Path;
/** Generates a deterministic synthetic log file of (at least) a requested size. One line = one record. */
public final class LogFile {
private LogFile() {}
public record Stats(long bytes, long lines, long errors) {}
/** Writes whole lines until at least {@code targetBytes} have been written. */
public static Stats generate(Path target, long targetBytes) throws IOException {
long seed = 42, bytes = 0, lines = 0, errors = 0;
try (BufferedWriter w = Files.newBufferedWriter(target)) {
StringBuilder sb = new StringBuilder(200);
while (bytes < targetBytes) {
seed = seed * 6364136223846793005L + 1442695040888963407L;
int r = (int) (seed >>> 33);
boolean error = r % 50 == 0;
if (error) errors++;
sb.setLength(0);
sb.append("2026-01-01T00:00:").append(String.format("%02d", (lines / 1000) % 60))
.append(' ').append(error ? "ERROR" : "INFO ")
.append(" service-").append(r % 20)
.append(" request handled id=").append(lines)
.append(" user=").append(Integer.toHexString(r))
.append(" payload=").append("x".repeat(40 + r % 40)).append('\n');
w.write(sb.toString());
bytes += sb.length(); // ASCII only, so chars == bytes
lines++;
}
}
return new Stats(bytes, lines, errors);
}
public static void main(String[] a) throws IOException {
Path p = Path.of(a[0]);
long t0 = System.nanoTime();
Stats s = generate(p, Long.parseLong(a[1]));
System.out.printf("generated %s: %,d bytes, %,d lines, %,d ERROR lines in %.1f s%n",
p.getFileName(), s.bytes(), s.lines(), s.errors(), (System.nanoTime() - t0) / 1e9);
}
}
@@ -0,0 +1,67 @@
package com.ankurm.nio2;
import java.io.IOException;
import java.lang.foreign.Arena;
import java.lang.foreign.MemorySegment;
import java.lang.foreign.ValueLayout;
import java.nio.MappedByteBuffer;
import java.nio.channels.FileChannel;
import java.nio.file.Files;
import java.nio.file.Path;
import java.nio.file.StandardOpenOption;
/** Memory-mapped reading: classic MappedByteBuffer (2 GiB cap per mapping) vs FileChannel.map with an Arena. */
public final class MappedDemo {
private MappedDemo() {}
/** Counts '\n' with an Arena-backed mapping of the whole file: no 2 GiB limit, no heap copy. */
public static long countNewlinesArena(Path p) throws IOException {
try (Arena arena = Arena.ofConfined();
FileChannel ch = FileChannel.open(p, StandardOpenOption.READ)) {
MemorySegment seg = ch.map(FileChannel.MapMode.READ_ONLY, 0, ch.size(), arena);
long n = 0;
for (long i = 0, len = seg.byteSize(); i < len; i++) if (seg.get(ValueLayout.JAVA_BYTE, i) == '\n') n++;
return n;
}
}
/** Same with the classic API, in windows of at most 1 GiB. */
public static long countNewlinesWindows(Path p) throws IOException {
long n = 0;
try (FileChannel ch = FileChannel.open(p, StandardOpenOption.READ)) {
long size = ch.size(), win = 1L << 30;
for (long off = 0; off < size; off += win) {
MappedByteBuffer mb = ch.map(FileChannel.MapMode.READ_ONLY, off, Math.min(win, size - off));
for (int i = 0, lim = mb.limit(); i < lim; i++) if (mb.get(i) == '\n') n++;
}
}
return n;
}
public static void main(String[] a) throws IOException {
Path p = Path.of(a[0]);
long size = Files.size(p);
System.out.printf("file %s: %,d bytes, -Xmx=%d MiB, RSS at start=%d MiB%n", p.getFileName(), size,
Runtime.getRuntime().maxMemory() >> 20, HeapProbe.rssMiB());
try (FileChannel ch = FileChannel.open(p, StandardOpenOption.READ)) {
ch.map(FileChannel.MapMode.READ_ONLY, 0, ch.size());
System.out.println("classic map of the whole file: succeeded");
} catch (IllegalArgumentException e) {
System.out.println("classic map of the whole file: IllegalArgumentException: " + e.getMessage());
}
try (HeapProbe probe = new HeapProbe()) {
long t0 = System.nanoTime();
long n2 = countNewlinesArena(p);
System.out.printf("Arena + MemorySegment: %,d newlines in %.1f s; RSS after the Arena was closed=%d MiB%n",
n2, (System.nanoTime() - t0) / 1e9, HeapProbe.rssMiB());
t0 = System.nanoTime();
long n1 = countNewlinesWindows(p);
System.out.printf("classic 1 GiB windows: %,d newlines in %.1f s; RSS right after=%d MiB%n",
n1, (System.nanoTime() - t0) / 1e9, HeapProbe.rssMiB());
System.out.println(probe.report());
System.gc();
try { Thread.sleep(1000); } catch (InterruptedException e) { Thread.currentThread().interrupt(); }
System.out.println("RSS one second after System.gc()=" + HeapProbe.rssMiB() + " MiB");
}
}
}
@@ -0,0 +1,28 @@
package com.ankurm.nio2;
import java.io.IOException;
import java.nio.file.Files;
import java.nio.file.Path;
/** The convenience methods that load the whole file. Run with a small -Xmx to see the failure. */
public final class ReadAllFail {
private ReadAllFail() {}
public static void main(String[] a) throws IOException {
Path p = Path.of(a[1]);
System.out.printf("%s on %s (%,d bytes), -Xmx=%d MiB%n", a[0], p.getFileName(), Files.size(p),
Runtime.getRuntime().maxMemory() >> 20);
try {
Object r = switch (a[0]) {
case "readAllLines" -> Files.readAllLines(p);
case "readString" -> Files.readString(p);
case "readAllBytes" -> Files.readAllBytes(p);
default -> throw new IllegalArgumentException(a[0]);
};
System.out.println("survived: " + (r instanceof java.util.List<?> l ? l.size() + " lines"
: r instanceof String s ? s.length() + " chars" : ((byte[]) r).length + " bytes"));
} catch (OutOfMemoryError e) {
System.out.println("OutOfMemoryError: " + e.getMessage());
}
}
}
@@ -0,0 +1,58 @@
package com.ankurm.nio2;
import java.io.BufferedReader;
import java.io.IOException;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.Map;
import java.util.TreeMap;
import java.util.stream.Stream;
/** Processes a file of any size in constant heap: count lines, ERROR lines, and lines per service. */
public final class StreamLines {
private StreamLines() {}
public record Result(long lines, long errors, Map<String, Long> perService) {}
/** Files.lines: lazy Stream<String>; MUST be closed (it owns a file handle). */
public static Result withFilesLines(Path p) throws IOException {
long[] c = new long[2];
Map<String, Long> per = new TreeMap<>();
try (Stream<String> s = Files.lines(p)) {
s.forEach(l -> tally(l, c, per));
}
return new Result(c[0], c[1], per);
}
/** Files.newBufferedReader: the same thing with a plain loop. */
public static Result withBufferedReader(Path p) throws IOException {
long[] c = new long[2];
Map<String, Long> per = new TreeMap<>();
try (BufferedReader r = Files.newBufferedReader(p)) {
for (String l = r.readLine(); l != null; l = r.readLine()) tally(l, c, per);
}
return new Result(c[0], c[1], per);
}
private static void tally(String line, long[] c, Map<String, Long> per) {
c[0]++;
if (line.startsWith("ERROR", 20)) c[1]++;
int i = line.indexOf("service-");
per.merge(line.substring(i, line.indexOf(' ', i)), 1L, Long::sum);
}
public static void main(String[] a) throws IOException {
Path p = Path.of(a[1]);
long t0 = System.nanoTime();
Result r;
try (HeapProbe probe = new HeapProbe()) {
r = a[0].equals("lines") ? withFilesLines(p) : withBufferedReader(p);
double s = (System.nanoTime() - t0) / 1e9;
System.out.printf("%s over %s (%,d bytes)%n", a[0].equals("lines") ? "Files.lines" : "Files.newBufferedReader",
p.getFileName(), Files.size(p));
System.out.printf("lines=%,d errors=%,d services=%d%n", r.lines(), r.errors(), r.perService().size());
System.out.printf("elapsed %.1f s%n", s);
System.out.println(probe.report());
}
}
}
@@ -0,0 +1,64 @@
package com.ankurm.nio2;
import java.io.IOException;
import java.io.UncheckedIOException;
import java.nio.file.*;
import java.nio.file.attribute.BasicFileAttributes;
import java.util.List;
import java.util.stream.Collectors;
import java.util.stream.Stream;
/** Files.list / walk / find / walkFileTree / DirectoryStream on a small tree. */
public final class WalkFindDemo {
private WalkFindDemo() {}
static List<String> rel(Path root, Stream<Path> s) {
try (s) { return s.map(p -> root.relativize(p).toString()).sorted().collect(Collectors.toList()); }
}
public static void main(String[] a) throws IOException {
Path root = Files.createTempDirectory("nio2-walk");
Files.createDirectories(root.resolve("src/main"));
Files.createDirectories(root.resolve("src/test"));
Files.createDirectories(root.resolve(".git/objects"));
Files.writeString(root.resolve("README.md"), "# demo\n");
Files.writeString(root.resolve("src/main/App.java"), "class App {}\n");
Files.writeString(root.resolve("src/test/AppTest.java"), "class AppTest {}\n" + "x".repeat(2000));
Files.writeString(root.resolve(".git/objects/blob.bin"), "ignored");
System.out.println("Files.list (one level): " + rel(root, Files.list(root)));
System.out.println("Files.walk (everything): " + rel(root, Files.walk(root)));
System.out.println("Files.walk maxDepth=1: " + rel(root, Files.walk(root, 1)));
System.out.println("Files.find *.java: " + rel(root,
Files.find(root, 10, (p, at) -> at.isRegularFile() && p.toString().endsWith(".java"))));
System.out.println("Files.find size > 1000: " + rel(root,
Files.find(root, 10, (p, at) -> at.isRegularFile() && at.size() > 1000)));
List<String> seen = new java.util.ArrayList<>();
Files.walkFileTree(root, new SimpleFileVisitor<>() {
@Override public FileVisitResult preVisitDirectory(Path d, BasicFileAttributes at) {
return d.getFileName().toString().equals(".git") ? FileVisitResult.SKIP_SUBTREE : FileVisitResult.CONTINUE;
}
@Override public FileVisitResult visitFile(Path f, BasicFileAttributes at) {
seen.add(root.relativize(f).toString()); return FileVisitResult.CONTINUE;
}
});
java.util.Collections.sort(seen);
System.out.println("walkFileTree skipping .git: " + seen);
List<String> glob = new java.util.ArrayList<>();
try (DirectoryStream<Path> ds = Files.newDirectoryStream(root.resolve("src/main"), "*.{java,kt}")) {
ds.forEach(p -> glob.add(p.getFileName().toString()));
}
System.out.println("DirectoryStream glob: " + glob);
// A symlink loop: walk() does not follow links by default; with FOLLOW_LINKS it detects the cycle.
Files.createSymbolicLink(root.resolve("src/loop"), root);
System.out.println("walk with a symlink loop, default: " + rel(root, Files.walk(root)).size() + " entries, no error");
try (Stream<Path> s = Files.walk(root, FileVisitOption.FOLLOW_LINKS)) {
s.forEach(p -> { });
} catch (UncheckedIOException e) {
System.out.println("walk with FOLLOW_LINKS: " + e.getCause().getClass().getSimpleName());
}
}
}
@@ -0,0 +1,82 @@
package com.ankurm.nio2;
import java.io.IOException;
import java.nio.file.*;
import java.util.*;
import java.util.concurrent.TimeUnit;
import static java.nio.file.StandardWatchEventKinds.*;
/** What WatchService actually delivers on this platform. Every scenario prints what it saw. */
public final class WatchDemo {
private WatchDemo() {}
/** Drains the key until no new event arrives for quietMs; returns kind -> total repeat-count. */
static Map<String, Integer> drain(WatchService ws, long quietMs, boolean verbose) throws Exception {
Map<String, Integer> tally = new TreeMap<>();
int batches = 0;
WatchKey k;
while ((k = ws.poll(quietMs, TimeUnit.MILLISECONDS)) != null) {
batches++;
List<WatchEvent<?>> evs = k.pollEvents();
for (WatchEvent<?> e : evs) {
tally.merge(e.kind().name(), e.count(), Integer::sum);
if (verbose) System.out.println(" " + e.kind().name() + " count=" + e.count() + " context=" + e.context());
}
k.reset();
}
tally.put("(batches)", batches);
return tally;
}
public static void main(String[] a) throws Exception {
Path dir = Files.createTempDirectory("nio2-watch");
try (WatchService ws = FileSystems.getDefault().newWatchService()) {
dir.register(ws, ENTRY_CREATE, ENTRY_MODIFY, ENTRY_DELETE);
System.out.println("1. create a file, write it twice, delete it (events polled afterwards):");
Path f = dir.resolve("a.txt");
Files.writeString(f, "one");
Files.writeString(f, "two", StandardOpenOption.APPEND);
Files.delete(f);
drain(ws, 300, true);
System.out.println("2. Files.writeString of 100 KiB to a new file (one call):");
Files.writeString(dir.resolve("b.txt"), "y".repeat(100 * 1024));
drain(ws, 300, true);
System.out.println("3. 50 appends to one file with no polling in between, then poll:");
Path g = dir.resolve("g.txt");
Files.writeString(g, "");
drain(ws, 300, false);
for (int i = 0; i < 50; i++) Files.writeString(g, "line\n", StandardOpenOption.APPEND);
System.out.println(" " + drain(ws, 300, false));
System.out.println("4. Files.move (atomic rename) of g.txt to h.txt:");
Files.move(g, dir.resolve("h.txt"), StandardCopyOption.ATOMIC_MOVE);
drain(ws, 300, true);
System.out.println("5. a file created inside a NEW subdirectory (subdirectory not registered):");
Path sub = Files.createDirectory(dir.resolve("sub"));
Files.writeString(sub.resolve("deep.txt"), "hello");
drain(ws, 300, true);
System.out.println("6. burst: create N empty files in a fresh directory before polling once"
+ " (jdk.nio.file.WatchService.maxEventsPerPoll=" + System.getProperty("jdk.nio.file.WatchService.maxEventsPerPoll", "default") + "):");
for (int n : new int[] {100, 512, 513, 5000}) {
Path burst = Files.createDirectory(dir.resolve("burst" + n));
try (WatchService ws2 = FileSystems.getDefault().newWatchService()) {
burst.register(ws2, ENTRY_CREATE);
for (int i = 0; i < n; i++) Files.createFile(burst.resolve("f" + i));
Thread.sleep(500);
WatchKey k = ws2.poll(1, TimeUnit.SECONDS);
List<WatchEvent<?>> evs = k.pollEvents();
long creates = evs.stream().filter(e -> e.kind() == ENTRY_CREATE).count();
long overflow = evs.stream().filter(e -> e.kind() == OVERFLOW).mapToInt(WatchEvent::count).sum();
System.out.printf(" N=%-5d delivered %-4d events: ENTRY_CREATE=%-4d OVERFLOW(count)=%d, files on disk=%d%n",
n, evs.size(), creates, overflow, Files.list(burst).count());
}
}
}
}
}
@@ -0,0 +1,183 @@
package com.ankurm.nio2;
import org.junit.jupiter.api.Assumptions;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
import java.io.IOException;
import java.io.RandomAccessFile;
import java.io.UncheckedIOException;
import java.nio.channels.FileChannel;
import java.nio.charset.MalformedInputException;
import java.nio.file.*;
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.TimeUnit;
import java.util.stream.Stream;
import static java.nio.file.StandardWatchEventKinds.*;
import static org.junit.jupiter.api.Assertions.*;
class Nio2BehaviourTest {
@TempDir Path tmp;
private Path sample() throws IOException {
Path p = tmp.resolve("sample.log");
LogFile.generate(p, 2_000_000);
return p;
}
@Test
void streamingMatchesReadAllLines() throws IOException {
Path p = sample();
List<String> all = Files.readAllLines(p);
assertEquals(all.size(), StreamLines.withFilesLines(p).lines());
assertEquals(all.size(), StreamLines.withBufferedReader(p).lines());
assertEquals(all.stream().filter(l -> l.startsWith("ERROR", 20)).count(), StreamLines.withFilesLines(p).errors());
}
@Test
void generatorStatsMatchTheFile() throws IOException {
Path p = tmp.resolve("g.log");
LogFile.Stats s = LogFile.generate(p, 500_000);
assertEquals(Files.size(p), s.bytes());
try (Stream<String> l = Files.lines(p)) { assertEquals(s.lines(), l.count()); }
}
@Test
void mappedCountersAgreeWithLineCount() throws IOException {
Path p = sample();
long lines = StreamLines.withFilesLines(p).lines();
assertEquals(lines, MappedDemo.countNewlinesArena(p));
assertEquals(lines, MappedDemo.countNewlinesWindows(p));
}
@Test
void classicMapRejectsMoreThan2GiBButArenaMapsIt() throws IOException {
Path sparse = tmp.resolve("sparse.bin");
try (RandomAccessFile f = new RandomAccessFile(sparse.toFile(), "rw")) { f.setLength(3L << 30); } // sparse: no real disk used
try (FileChannel ch = FileChannel.open(sparse, StandardOpenOption.READ)) {
assertThrows(IllegalArgumentException.class, () -> ch.map(FileChannel.MapMode.READ_ONLY, 0, ch.size()));
}
try (var arena = java.lang.foreign.Arena.ofConfined(); FileChannel ch = FileChannel.open(sparse, StandardOpenOption.READ)) {
var seg = ch.map(FileChannel.MapMode.READ_ONLY, 0, ch.size(), arena);
assertEquals(3L << 30, seg.byteSize());
}
}
private static boolean linux() { return Files.isDirectory(Path.of("/proc/self/fd")); }
@Test
void unclosedFilesLinesHoldsAFileDescriptorUntilClosed() throws IOException {
Assumptions.assumeTrue(linux());
Path p = sample();
int base = HandleLeakDemo.openFds();
List<Stream<String>> keep = new ArrayList<>();
for (int i = 0; i < 10; i++) { Stream<String> s = Files.lines(p); s.findFirst(); keep.add(s); }
assertEquals(base + 10, HandleLeakDemo.openFds());
keep.forEach(Stream::close);
assertEquals(base, HandleLeakDemo.openFds());
for (int i = 0; i < 10; i++) { try (Stream<String> s = Files.lines(p)) { s.findFirst(); } }
assertEquals(base, HandleLeakDemo.openFds());
}
@Test
void walkStoppedEarlyHoldsDirectoriesOpenButFullyConsumedWalkDoesNot() throws IOException {
Assumptions.assumeTrue(linux());
Path tree = HandleLeakDemo.makeTree(tmp.resolve("tree"));
int base = HandleLeakDemo.openFds();
Stream<Path> early = Files.walk(tree);
early.skip(3).findFirst();
assertTrue(HandleLeakDemo.openFds() > base, "an unfinished walk should hold directory handles");
early.close();
assertEquals(base, HandleLeakDemo.openFds());
Files.walk(tree).count();
assertEquals(base, HandleLeakDemo.openFds());
}
@Test
void badByteSurfacesLazilyAsUncheckedIOException() throws IOException {
Path p = tmp.resolve("bad.txt");
Files.write(p, new byte[] {'o', 'k', '\n', (byte) 0xE9, '\n'});
try (Stream<String> s = Files.lines(p)) {
UncheckedIOException e = assertThrows(UncheckedIOException.class, () -> s.forEach(l -> { }));
assertInstanceOf(MalformedInputException.class, e.getCause());
}
assertThrows(MalformedInputException.class, () -> Files.readAllLines(p));
}
@Test
void walkFindAndSkipSubtree() throws IOException {
Files.createDirectories(tmp.resolve("a/b"));
Files.writeString(tmp.resolve("a/b/x.java"), "x");
Files.writeString(tmp.resolve("a/y.txt"), "y");
try (Stream<Path> s = Files.find(tmp, 5, (p, at) -> at.isRegularFile() && p.toString().endsWith(".java"))) {
assertEquals(1, s.count());
}
try (Stream<Path> s = Files.walk(tmp, 1)) { assertEquals(2, s.count()); } // tmp itself + "a"
}
@Test
void watchServiceDeliversCreateAndCoalescesRepeatedModify() throws Exception {
try (WatchService ws = FileSystems.getDefault().newWatchService()) {
tmp.register(ws, ENTRY_CREATE, ENTRY_MODIFY);
Path f = tmp.resolve("w.txt");
Files.writeString(f, "x".repeat(100 * 1024));
Thread.sleep(300);
WatchKey k = ws.poll(2, TimeUnit.SECONDS);
assertNotNull(k);
List<WatchEvent<?>> evs = k.pollEvents();
assertEquals(ENTRY_CREATE, evs.get(0).kind());
assertEquals(f.getFileName(), evs.get(0).context());
assertTrue(evs.size() <= 2, "MODIFY events for one file collapse into one event with a count: " + evs);
}
}
@Test
void watchKeyOverflowsPast512PendingEventsAndDropsThem() throws Exception {
Assumptions.assumeTrue(System.getProperty("jdk.nio.file.WatchService.maxEventsPerPoll") == null);
try (WatchService ws = FileSystems.getDefault().newWatchService()) {
tmp.register(ws, ENTRY_CREATE);
for (int i = 0; i < 513; i++) Files.createFile(tmp.resolve("f" + i));
Thread.sleep(500);
WatchKey k = ws.poll(2, TimeUnit.SECONDS);
List<WatchEvent<?>> evs = k.pollEvents();
assertEquals(1, evs.size());
assertEquals(OVERFLOW, evs.get(0).kind());
}
}
@Test
void watchServiceDoesNotWatchSubdirectoriesRecursively() throws Exception {
try (WatchService ws = FileSystems.getDefault().newWatchService()) {
tmp.register(ws, ENTRY_CREATE);
Path sub = Files.createDirectory(tmp.resolve("sub"));
Thread.sleep(200);
WatchKey k = ws.poll(2, TimeUnit.SECONDS);
k.pollEvents(); k.reset();
Files.writeString(sub.resolve("deep.txt"), "x");
assertNull(ws.poll(500, TimeUnit.MILLISECONDS), "no event for a file in an unregistered subdirectory");
}
}
@Test
void readAllLinesFailsWhereStreamingSurvivesUnderSmallHeap() throws Exception {
// Runs the real programs in child JVMs so -Xmx is honoured.
Path p = tmp.resolve("m.log");
LogFile.generate(p, 100_000_000);
String java = Path.of(System.getProperty("java.home"), "bin", "java").toString();
String cp = System.getProperty("java.class.path");
String streamOut = run(java, "-Xmx32m", "-cp", cp, "com.ankurm.nio2.StreamLines", "lines", p.toString());
assertTrue(streamOut.contains("lines="), streamOut);
String failOut = run(java, "-Xmx32m", "-cp", cp, "com.ankurm.nio2.ReadAllFail", "readAllLines", p.toString());
assertTrue(failOut.contains("OutOfMemoryError"), failOut);
}
private static String run(String... cmd) throws Exception {
Process pr = new ProcessBuilder(cmd).redirectErrorStream(true).start();
String s = new String(pr.getInputStream().readAllBytes());
pr.waitFor();
return s;
}
}