From 75839b207e106f4a28bad3b66001c266b3e9837b Mon Sep 17 00:00:00 2001 From: Tam Nguyen Duc <1218621+tamnd@users.noreply.github.com> Date: Thu, 20 Aug 2026 08:43:21 +0700 Subject: [PATCH] Bulk loading and appending The Java client could read anything and write nothing. This adds both of the C ABI's write paths, which between them are the only way values get into a database at all while the engine has no DDL. A loader builds a database out of whole columns. Loader.create refuses a path that already exists, table names the one table and how many rows it has, and columns go in one at a time until finish writes the file. The row count is declared rather than counted, so a column with a value missing is an error and not a shorter table. An edge table name is wanted even for a load that adds no edges, because the engine wants one and guessing a name later is worse than choosing one now. Columns go in as arrays or as java.nio buffers, and which one you pass is the difference between a copy and no copy. A direct buffer reaches the engine through MemorySegment.ofBuffer with nothing crossing the boundary but a pointer. An array is copied into a confined arena first, because native code cannot address a Java array without either a copy or a pause. Linker.Option.critical(true) would allow the array through without either, at the price of blocking the collector for the length of the copy, and on a column of a hundred million values that is not a trade worth making. Measured over a hundred thousand rows, the direct buffer hands a column over in 0.44 ns a row against 1.1 ns for the array, which is 68 microseconds of memcpy for 800 KB and about the bandwidth you would expect. An appender adds rows to a table that already exists, a value at a time in declared column order, ended by endRow. There is no way to append no value because the C ABI has none, and inventing one here would only move the surprise. Closing an appender that was never finished writes what it has: a loop that threw halfway keeps the rows it managed, because throwing away work that succeeded is not a decision a close should make on its own. discard is there for when it should be thrown away. An appended row costs 87 ns against 4.0 ms for the same row as an INSERT statement, which is the whole reason the surface exists. Twenty seven tests over both, covering every column kind, buffer slices rather than whole buffers, edges appended across calls, a refused value leaving no half of a row behind, and every misuse the client catches before the call. Benchmarks in LoadBench and AppendBench, and the README has the numbers and the example the top of it needed. --- README.md | 47 +- .../main/java/dev/zudb/bench/AppendBench.java | 109 +++++ .../main/java/dev/zudb/bench/LoadBench.java | 135 ++++++ .../src/main/java/dev/zudb/bench/Temp.java | 30 ++ zudb-ffm/src/main/java/dev/zudb/ffm/Abi.java | 92 ++++ .../main/java/dev/zudb/ffm/FfmBinding.java | 411 ++++++++++++++++++ .../test/java/dev/zudb/ffm/AppenderTest.java | 287 ++++++++++++ .../test/java/dev/zudb/ffm/LoaderTest.java | 269 ++++++++++++ zudb/src/main/java/dev/zudb/Appender.java | 406 +++++++++++++++++ zudb/src/main/java/dev/zudb/Connection.java | 14 + zudb/src/main/java/dev/zudb/Loader.java | 299 +++++++++++++ .../src/main/java/dev/zudb/spi/ZuBinding.java | 233 ++++++++++ 12 files changed, 2331 insertions(+), 1 deletion(-) create mode 100644 zudb-bench/src/main/java/dev/zudb/bench/AppendBench.java create mode 100644 zudb-bench/src/main/java/dev/zudb/bench/LoadBench.java create mode 100644 zudb-bench/src/main/java/dev/zudb/bench/Temp.java create mode 100644 zudb-ffm/src/test/java/dev/zudb/ffm/AppenderTest.java create mode 100644 zudb-ffm/src/test/java/dev/zudb/ffm/LoaderTest.java create mode 100644 zudb/src/main/java/dev/zudb/Appender.java create mode 100644 zudb/src/main/java/dev/zudb/Loader.java diff --git a/README.md b/README.md index 258127d..3b05be7 100644 --- a/README.md +++ b/README.md @@ -66,6 +66,51 @@ What it is worth, summing one integer column of a hundred thousand rows on an M- A row at a time is a boundary crossing a cell, and a hundred crossings cost about what one borrowed buffer costs. Both surfaces are there because both are the right answer to a different question, but a loop over a million rows should be reading a column. +## Getting rows in + +Two ways, and which one you want follows from whether the database exists yet. + +A loader builds one out of whole columns. It is the fastest way values get in and, while the engine has no DDL, it is the only way a table comes into being at all: + +```java +try (Loader loader = Loader.create(Path.of("social.zu1"))) { + loader.table("Person", "Follows", 3); + loader.column("id", 1L, 2L, 3L); + loader.column("name", "ada", "grace", "alan"); + loader.edges(new int[] {0, 1}, new int[] {1, 2}); + loader.finish(); +} +``` + +Columns go in as arrays or as `java.nio` buffers, and which you pass is the difference between a copy and no copy. A direct buffer is read where it lies, so the engine sees the memory your program already filled and nothing crosses the boundary but a pointer. An array is memory nothing outside the JVM can address, so it is copied off-heap first. `Linker.Option.critical(true)` would let a Java array through without either, at the price of blocking the collector for the length of the copy, and on a column this size that is not a trade worth making. + +An appender adds rows to a table that already exists, a value at a time, with no statement anywhere near it: + +```java +try (Appender rows = conn.appender("Person")) { + rows.append(4L).append("hedy").endRow(); + rows.append(5L).append("katherine").endRow(); + rows.finish(); +} +``` + +Values are written in the order the table declares its columns, which `columnName(int)` will tell you, and a row is a row once `endRow()` has ended it. A value the column will not take ends its row there and rolls back the values already written into it, so a refused append never leaves half a row behind. Closing an appender that was never finished writes what it has anyway, because a loop that threw halfway should keep the rows it managed; `discard()` is there for when it should not. + +What each is worth on an M-series laptop, JDK 25: + +| How | Per row | +|---|---| +| `loader.column(name, direct LongBuffer)` | 0.44 ns | +| `loader.column(name, long[])` | 1.1 ns | +| `loader.column(name, List)` | 108 ns | +| a whole two-column load, write included | 630 ns | +| `appender.append(...).endRow()` | 87 ns | +| the same row as an `INSERT` statement | 4.0 ms | + +The first two lines are the copy: 0.68 ns a row over a hundred thousand rows is 68 microseconds to move 800 KB, which is about what a memcpy costs and about what a direct buffer saves. It is a small share of a load that also writes a file, and it is the whole difference at the boundary itself. + +The last line is the one to read twice. A statement per row parses, plans, runs and commits per row, and none of that work says anything the row before it did not already say. That is what an appender is for. + ## How it binds The Foreign Function and Memory API is the primary path. The downcall handles are written by hand against `zu.h` rather than generated with `jextract`, because the C ABI here is around seventy functions with a stable shape, and a hand-written layer is where the interesting decisions live: which calls are `Linker.Option.critical` because they are short pure accessors, where the out-parameter scratch space comes from so that a query does not allocate, and how a `zu_error` becomes a typed Java exception exactly once. There is no native code in this repository beyond `libzu` itself. @@ -99,7 +144,7 @@ catch (ZuSyntaxException e) { ## What works today -The engine has no DDL yet, so there is no `CREATE NODE TABLE` and nothing in this client writes a schema. What runs against a fresh database is the expression and projection surface: `RETURN`, `UNWIND`, parameters, lists, records, and the temporal types. The example at the top of this file describes the intended shape and needs a graph that some other tool built. +The engine has no DDL yet, so there is no `CREATE NODE TABLE` and no statement in this client writes a schema. A table comes into being through `Loader`, which is why the loader example above builds the graph the example at the top of this file reads. What runs against a fresh database with nothing in it is the expression and projection surface: `RETURN`, `UNWIND`, parameters, lists, records, and the temporal types. ## Building diff --git a/zudb-bench/src/main/java/dev/zudb/bench/AppendBench.java b/zudb-bench/src/main/java/dev/zudb/bench/AppendBench.java new file mode 100644 index 0000000..0f7d69a --- /dev/null +++ b/zudb-bench/src/main/java/dev/zudb/bench/AppendBench.java @@ -0,0 +1,109 @@ +package dev.zudb.bench; + +import dev.zudb.Appender; +import dev.zudb.Connection; +import dev.zudb.Database; +import dev.zudb.Loader; +import java.io.IOException; +import java.nio.file.Files; +import java.nio.file.Path; +import java.util.concurrent.TimeUnit; +import org.openjdk.jmh.annotations.Benchmark; +import org.openjdk.jmh.annotations.BenchmarkMode; +import org.openjdk.jmh.annotations.Fork; +import org.openjdk.jmh.annotations.Level; +import org.openjdk.jmh.annotations.Measurement; +import org.openjdk.jmh.annotations.Mode; +import org.openjdk.jmh.annotations.OutputTimeUnit; +import org.openjdk.jmh.annotations.Scope; +import org.openjdk.jmh.annotations.Setup; +import org.openjdk.jmh.annotations.State; +import org.openjdk.jmh.annotations.TearDown; +import org.openjdk.jmh.annotations.Warmup; + +/** + * What adding a row to a table that already exists costs. + * + *

One benchmark invocation is one row, so the score is the row, which is + * the unit a caller writes. Every iteration starts from a copy of a database + * with one row in it, so a long run measures appending rather than a table + * growing under it. + */ +@State(Scope.Benchmark) +@BenchmarkMode(Mode.AverageTime) +@OutputTimeUnit(TimeUnit.NANOSECONDS) +@Warmup(iterations = 3, time = 2) +@Measurement(iterations = 5, time = 2) +@Fork(value = 1, jvmArgs = {"--enable-native-access=ALL-UNNAMED"}) +public class AppendBench { + + private Path dir; + private Path template; + private long counter; + + private Database db; + private Connection conn; + private Appender appender; + + @Setup + public void build() throws IOException { + dir = Files.createTempDirectory("zu-append-bench"); + // The table an appender appends to has to exist, and a bulk load is the + // only thing that makes one. + template = dir.resolve("template.zu"); + try (Loader loader = Loader.create(template)) { + loader.table("Person", "Knows", 1); + loader.column("id", -1L); + loader.column("name", "seed"); + loader.finish(); + } + } + + @TearDown + public void clean() throws IOException { + Temp.deleteTree(dir); + } + + @Setup(Level.Iteration) + public void open() throws IOException { + Path path = dir.resolve("append-" + counter++ + ".zu"); + Files.copy(template, path); + db = Database.open(path); + conn = db.connect(); + appender = conn.appender("Person"); + } + + @TearDown(Level.Iteration) + public void close() { + if (appender != null) { + appender.close(); + appender = null; + } + if (conn != null) { + conn.close(); + conn = null; + } + if (db != null) { + db.close(); + db = null; + } + } + + /** One row of two columns, written the way a loop that knows its schema writes it. */ + @Benchmark + public void row() { + appender.append(counter++).append("n").endRow(); + } + + /** The same row through the dynamic path, which costs a type test a value. */ + @Benchmark + public void rowOfObjects() { + appender.row(counter++, "n"); + } + + /** What the appender is worth, against the statement it replaces. */ + @Benchmark + public void rowThroughAStatement() { + conn.execute("INSERT (:Person {id: " + counter++ + ", name: 'n'})"); + } +} diff --git a/zudb-bench/src/main/java/dev/zudb/bench/LoadBench.java b/zudb-bench/src/main/java/dev/zudb/bench/LoadBench.java new file mode 100644 index 0000000..369327f --- /dev/null +++ b/zudb-bench/src/main/java/dev/zudb/bench/LoadBench.java @@ -0,0 +1,135 @@ +package dev.zudb.bench; + +import dev.zudb.Loader; +import java.io.IOException; +import java.nio.ByteBuffer; +import java.nio.ByteOrder; +import java.nio.LongBuffer; +import java.nio.file.Files; +import java.nio.file.Path; +import java.util.ArrayList; +import java.util.Comparator; +import java.util.List; +import java.util.concurrent.TimeUnit; +import org.openjdk.jmh.annotations.Benchmark; +import org.openjdk.jmh.annotations.BenchmarkMode; +import org.openjdk.jmh.annotations.Fork; +import org.openjdk.jmh.annotations.Level; +import org.openjdk.jmh.annotations.Measurement; +import org.openjdk.jmh.annotations.Mode; +import org.openjdk.jmh.annotations.OperationsPerInvocation; +import org.openjdk.jmh.annotations.OutputTimeUnit; +import org.openjdk.jmh.annotations.Scope; +import org.openjdk.jmh.annotations.Setup; +import org.openjdk.jmh.annotations.State; +import org.openjdk.jmh.annotations.TearDown; +import org.openjdk.jmh.annotations.Warmup; + +/** + * What building a database out of columns costs, per row of a hundred + * thousand. + * + *

The loader itself is made in an invocation fixture rather than in the + * measured region, because a whole load writes a file and the file would + * drown out everything else. What is measured is handing a column over, which + * is where the difference between an array and a direct buffer lives, and + * {@code wholeLoad} is there so the rest can be read against what a load + * really costs. + */ +@State(Scope.Benchmark) +@BenchmarkMode(Mode.AverageTime) +@OutputTimeUnit(TimeUnit.NANOSECONDS) +@Warmup(iterations = 3, time = 2) +@Measurement(iterations = 5, time = 2) +@Fork(value = 1, jvmArgs = {"--enable-native-access=ALL-UNNAMED"}) +public class LoadBench { + + /** + * How many rows one load carries. A constant rather than a parameter + * because the per-row score is scaled by it, and JMH wants that scale as a + * literal in an annotation. + */ + private static final int ROWS = 100_000; + + private Path dir; + private long counter; + + private long[] ids; + private LongBuffer direct; + private List names; + + private Loader loader; + private Path path; + + @Setup + public void fill() throws IOException { + dir = Files.createTempDirectory("zu-load-bench"); + ids = new long[ROWS]; + for (int i = 0; i < ROWS; i++) { + ids[i] = i; + } + direct = + ByteBuffer.allocateDirect(ROWS * Long.BYTES).order(ByteOrder.nativeOrder()).asLongBuffer(); + direct.put(ids); + direct.flip(); + names = new ArrayList<>(ROWS); + for (int i = 0; i < ROWS; i++) { + names.add("n" + i); + } + } + + @TearDown + public void clean() throws IOException { + Temp.deleteTree(dir); + } + + /** A loader with its table named and no column in it yet. */ + @Setup(Level.Invocation) + public void open() { + path = dir.resolve("load-" + counter++ + ".zu"); + loader = Loader.create(path); + loader.table("Person", "Knows", ROWS); + } + + @TearDown(Level.Invocation) + public void close() throws IOException { + if (loader != null) { + loader.close(); + loader = null; + } + if (path != null) { + Files.deleteIfExists(path); + path = null; + } + } + + /** One integer column as a Java array, which has to be copied off-heap. */ + @Benchmark + @OperationsPerInvocation(ROWS) + public void columnFromArray() { + loader.column("id", ids); + } + + /** The same column as a direct buffer, which is read where it lies. */ + @Benchmark + @OperationsPerInvocation(ROWS) + public void columnFromDirectBuffer() { + loader.column("id", direct.duplicate()); + } + + /** A string column, which has no zero-copy shape and is checked for UTF-8 besides. */ + @Benchmark + @OperationsPerInvocation(ROWS) + public void columnOfStrings() { + loader.column("name", names); + } + + /** Two columns and the write, which is what a load costs a caller. */ + @Benchmark + @OperationsPerInvocation(ROWS) + public void wholeLoad() { + loader.column("id", ids); + loader.column("name", names); + loader.finish(); + } +} diff --git a/zudb-bench/src/main/java/dev/zudb/bench/Temp.java b/zudb-bench/src/main/java/dev/zudb/bench/Temp.java new file mode 100644 index 0000000..fa605bd --- /dev/null +++ b/zudb-bench/src/main/java/dev/zudb/bench/Temp.java @@ -0,0 +1,30 @@ +package dev.zudb.bench; + +import java.io.IOException; +import java.nio.file.Files; +import java.nio.file.Path; +import java.util.Comparator; + +/** The files a benchmark that writes to disk leaves behind. */ +final class Temp { + + private Temp() {} + + /** Removes a directory and everything under it, sidecars included. */ + static void deleteTree(Path root) throws IOException { + if (root == null) { + return; + } + try (var walk = Files.walk(root)) { + walk.sorted(Comparator.reverseOrder()).forEach(Temp::delete); + } + } + + private static void delete(Path path) { + try { + Files.deleteIfExists(path); + } catch (IOException e) { + throw new IllegalStateException(e); + } + } +} diff --git a/zudb-ffm/src/main/java/dev/zudb/ffm/Abi.java b/zudb-ffm/src/main/java/dev/zudb/ffm/Abi.java index 787473d..7717441 100644 --- a/zudb-ffm/src/main/java/dev/zudb/ffm/Abi.java +++ b/zudb-ffm/src/main/java/dev/zudb/ffm/Abi.java @@ -107,6 +107,34 @@ final class Abi { final MethodHandle chunkColNodeOffset; final MethodHandle chunkColValid; + final MethodHandle loaderCreate; + final MethodHandle loaderTable; + final MethodHandle loaderEdges; + final MethodHandle loaderColI64; + final MethodHandle loaderColF64; + final MethodHandle loaderColBool; + final MethodHandle loaderColStr; + final MethodHandle loaderColTemporal; + final MethodHandle loaderFinish; + final MethodHandle loaderFree; + + final MethodHandle appenderOpen; + final MethodHandle appendBool; + final MethodHandle appendI64; + final MethodHandle appendF64; + final MethodHandle appendStr; + final MethodHandle appendBytes; + final MethodHandle appendTemporal; + final MethodHandle appendEndRow; + final MethodHandle appenderFlush; + final MethodHandle appenderBuffered; + final MethodHandle appenderCommitted; + final MethodHandle appenderCols; + final MethodHandle appenderColName; + final MethodHandle appenderDiscard; + final MethodHandle appenderClose; + final MethodHandle appenderFree; + final MethodHandle valueType; final MethodHandle valueBool; final MethodHandle valueI64; @@ -224,6 +252,70 @@ final class Abi { "zu_result_chunk_col_valid", FunctionDescriptor.of(JAVA_INT, ADDRESS, JAVA_LONG, JAVA_INT, ADDRESS)); + loaderCreate = + h("zu_loader_create", FunctionDescriptor.of(JAVA_INT, ADDRESS, SIZE_T, ADDRESS, ADDRESS)); + loaderTable = + h( + "zu_loader_table", + FunctionDescriptor.of( + JAVA_INT, ADDRESS, ADDRESS, SIZE_T, ADDRESS, SIZE_T, JAVA_LONG, ADDRESS)); + loaderEdges = + h( + "zu_loader_edges", + FunctionDescriptor.of(JAVA_INT, ADDRESS, ADDRESS, ADDRESS, JAVA_LONG, ADDRESS)); + loaderColI64 = + h( + "zu_loader_col_i64", + FunctionDescriptor.of(JAVA_INT, ADDRESS, ADDRESS, SIZE_T, ADDRESS, JAVA_LONG, ADDRESS)); + loaderColF64 = + h( + "zu_loader_col_f64", + FunctionDescriptor.of(JAVA_INT, ADDRESS, ADDRESS, SIZE_T, ADDRESS, JAVA_LONG, ADDRESS)); + loaderColBool = + h( + "zu_loader_col_bool", + FunctionDescriptor.of(JAVA_INT, ADDRESS, ADDRESS, SIZE_T, ADDRESS, JAVA_LONG, ADDRESS)); + loaderColStr = + h( + "zu_loader_col_str", + FunctionDescriptor.of( + JAVA_INT, ADDRESS, ADDRESS, SIZE_T, ADDRESS, ADDRESS, JAVA_LONG, ADDRESS)); + loaderColTemporal = + h( + "zu_loader_col_temporal", + FunctionDescriptor.of( + JAVA_INT, ADDRESS, ADDRESS, SIZE_T, JAVA_INT, ADDRESS, JAVA_LONG, ADDRESS)); + loaderFinish = h("zu_loader_finish", FunctionDescriptor.of(JAVA_INT, ADDRESS, ADDRESS)); + loaderFree = h("zu_loader_free", FunctionDescriptor.ofVoid(ADDRESS)); + + appenderOpen = + h( + "zu_appender_open", + FunctionDescriptor.of(JAVA_INT, ADDRESS, ADDRESS, SIZE_T, ADDRESS, ADDRESS)); + appendBool = h("zu_append_bool", FunctionDescriptor.of(JAVA_INT, ADDRESS, JAVA_INT, ADDRESS)); + appendI64 = h("zu_append_i64", FunctionDescriptor.of(JAVA_INT, ADDRESS, JAVA_LONG, ADDRESS)); + appendF64 = h("zu_append_f64", FunctionDescriptor.of(JAVA_INT, ADDRESS, JAVA_DOUBLE, ADDRESS)); + appendStr = + h("zu_append_str", FunctionDescriptor.of(JAVA_INT, ADDRESS, ADDRESS, SIZE_T, ADDRESS)); + appendBytes = + h("zu_append_bytes", FunctionDescriptor.of(JAVA_INT, ADDRESS, ADDRESS, SIZE_T, ADDRESS)); + appendTemporal = + h( + "zu_append_temporal", + FunctionDescriptor.of(JAVA_INT, ADDRESS, JAVA_INT, JAVA_LONG, ADDRESS)); + appendEndRow = h("zu_append_end_row", FunctionDescriptor.of(JAVA_INT, ADDRESS, ADDRESS)); + appenderFlush = h("zu_appender_flush", FunctionDescriptor.of(JAVA_INT, ADDRESS, ADDRESS)); + appenderBuffered = h("zu_appender_buffered", FunctionDescriptor.of(JAVA_INT, ADDRESS, ADDRESS)); + appenderCommitted = + h("zu_appender_committed", FunctionDescriptor.of(JAVA_INT, ADDRESS, ADDRESS)); + appenderCols = h("zu_appender_cols", FunctionDescriptor.of(JAVA_INT, ADDRESS, ADDRESS)); + appenderColName = + h("zu_appender_col_name", FunctionDescriptor.of(ADDRESS, ADDRESS, JAVA_INT, ADDRESS)); + appenderDiscard = h("zu_appender_discard", FunctionDescriptor.of(JAVA_INT, ADDRESS, ADDRESS)); + appenderClose = + h("zu_appender_close", FunctionDescriptor.of(JAVA_INT, ADDRESS, ADDRESS, ADDRESS)); + appenderFree = h("zu_appender_free", FunctionDescriptor.ofVoid(ADDRESS)); + valueType = critical("zu_value_type", FunctionDescriptor.of(JAVA_INT, ADDRESS)); valueBool = h("zu_value_bool", FunctionDescriptor.of(JAVA_INT, ADDRESS, ADDRESS)); valueI64 = h("zu_value_i64", FunctionDescriptor.of(JAVA_INT, ADDRESS, ADDRESS)); diff --git a/zudb-ffm/src/main/java/dev/zudb/ffm/FfmBinding.java b/zudb-ffm/src/main/java/dev/zudb/ffm/FfmBinding.java index 7e5a7a5..8cc5342 100644 --- a/zudb-ffm/src/main/java/dev/zudb/ffm/FfmBinding.java +++ b/zudb-ffm/src/main/java/dev/zudb/ffm/FfmBinding.java @@ -15,12 +15,15 @@ import dev.zudb.Diagnostic; import dev.zudb.Status; import dev.zudb.spi.ZuBinding; +import java.lang.foreign.Arena; import java.lang.foreign.MemorySegment; import java.nio.ByteBuffer; import java.nio.ByteOrder; import java.nio.DoubleBuffer; +import java.nio.IntBuffer; import java.nio.LongBuffer; import java.nio.charset.StandardCharsets; +import java.util.List; /** * The C ABI, called through the Foreign Function and Memory API. @@ -645,6 +648,414 @@ public String valueField(long value, long index) { } } + // ---- bulk load ---- + + @Override + public long loaderCreate(String path) { + Scratch s = Scratch.get(); + MemorySegment sl = s.slots(); + MemorySegment p = s.utf8(path); + clear(sl); + try { + int st = + (int) abi.loaderCreate.invokeExact(p, p.byteSize(), sl.asSlice(OUT, 8), sl.asSlice(ERR, 8)); + check("zu_loader_create", st, sl); + return sl.get(ADDRESS, OUT).address(); + } catch (Throwable t) { + throw fail("zu_loader_create", t); + } + } + + @Override + public void loaderTable(long loader, String nodes, String edges, long rows) { + Scratch s = Scratch.get(); + MemorySegment sl = s.slots(); + MemorySegment n = s.utf8(nodes); + MemorySegment e = edges == null ? MemorySegment.NULL : s.utf8(edges); + long elen = edges == null ? 0 : e.byteSize(); + clear(sl); + try { + int st = + (int) + abi.loaderTable.invokeExact( + ptr(loader), n, n.byteSize(), e, elen, rows, sl.asSlice(ERR, 8)); + check("zu_loader_table", st, sl); + } catch (Throwable t) { + throw fail("zu_loader_table", t); + } + } + + @Override + public void loaderEdges(long loader, IntBuffer from, IntBuffer to) { + int count = from.remaining(); + if (to.remaining() != count) { + throw Diagnostic.misuse( + Status.MISUSE, + "an edge starts somewhere and ends somewhere, and there are " + + count + + " starts against " + + to.remaining() + + " ends") + .toException(); + } + Scratch s = Scratch.get(); + MemorySegment sl = s.slots(); + clear(sl); + try (Arena arena = Arena.ofConfined()) { + MemorySegment f = pass(from, arena); + MemorySegment t = pass(to, arena); + int st = (int) abi.loaderEdges.invokeExact(ptr(loader), f, t, (long) count, sl.asSlice(ERR, 8)); + check("zu_loader_edges", st, sl); + } catch (Throwable t) { + throw fail("zu_loader_edges", t); + } + } + + @Override + public void loaderColumnLongs(long loader, String name, LongBuffer values) { + loaderColumn(abi.loaderColI64, "zu_loader_col_i64", loader, name, values.remaining(), + arena -> pass(values, arena)); + } + + @Override + public void loaderColumnDoubles(long loader, String name, DoubleBuffer values) { + loaderColumn(abi.loaderColF64, "zu_loader_col_f64", loader, name, values.remaining(), + arena -> pass(values, arena)); + } + + @Override + public void loaderColumnBooleans(long loader, String name, IntBuffer values) { + loaderColumn(abi.loaderColBool, "zu_loader_col_bool", loader, name, values.remaining(), + arena -> pass(values, arena)); + } + + @Override + public void loaderColumnStrings(long loader, String name, List values) { + int count = values.size(); + Scratch scratch = Scratch.get(); + MemorySegment sl = scratch.slots(); + clear(sl); + try (Arena arena = Arena.ofConfined()) { + MemorySegment n = utf8(arena, name); + MemorySegment pointers = arena.allocate(ADDRESS, count); + MemorySegment lengths = arena.allocate(Abi.SIZE_T, count); + for (int i = 0; i < count; i++) { + String v = values.get(i); + if (v == null) { + throw Diagnostic.misuse( + Status.MISUSE, + "row " + i + " of column " + name + " is no value at all, and a loaded column" + + " holds a value a row") + .toException(); + } + MemorySegment bytes = utf8(arena, v); + pointers.setAtIndex(ADDRESS, i, bytes); + size(lengths, i, bytes.byteSize()); + } + int st = + (int) + abi.loaderColStr.invokeExact( + ptr(loader), n, n.byteSize(), pointers, lengths, (long) count, sl.asSlice(ERR, 8)); + check("zu_loader_col_str", st, sl); + } catch (Throwable t) { + throw fail("zu_loader_col_str", t); + } + } + + @Override + public void loaderColumnTemporal(long loader, String name, int kind, LongBuffer values) { + int count = values.remaining(); + Scratch scratch = Scratch.get(); + MemorySegment sl = scratch.slots(); + clear(sl); + try (Arena arena = Arena.ofConfined()) { + MemorySegment n = utf8(arena, name); + MemorySegment v = pass(values, arena); + int st = + (int) + abi.loaderColTemporal.invokeExact( + ptr(loader), n, n.byteSize(), kind, v, (long) count, sl.asSlice(ERR, 8)); + check("zu_loader_col_temporal", st, sl); + } catch (Throwable t) { + throw fail("zu_loader_col_temporal", t); + } + } + + @Override + public void loaderFinish(long loader) { + endTransaction(abi.loaderFinish, "zu_loader_finish", loader); + } + + @Override + public void loaderFree(long loader) { + try { + abi.loaderFree.invokeExact(ptr(loader)); + } catch (Throwable t) { + throw fail("zu_loader_free", t); + } + } + + // ---- appending ---- + + @Override + public long appenderOpen(long conn, String table) { + return run(abi.appenderOpen, "zu_appender_open", conn, table); + } + + @Override + public void appendBoolean(long appender, boolean value) { + Scratch s = Scratch.get(); + MemorySegment sl = s.slots(); + clear(sl); + try { + int st = (int) abi.appendBool.invokeExact(ptr(appender), value ? 1 : 0, sl.asSlice(ERR, 8)); + check("zu_append_bool", st, sl); + } catch (Throwable t) { + throw fail("zu_append_bool", t); + } + } + + @Override + public void appendLong(long appender, long value) { + Scratch s = Scratch.get(); + MemorySegment sl = s.slots(); + clear(sl); + try { + int st = (int) abi.appendI64.invokeExact(ptr(appender), value, sl.asSlice(ERR, 8)); + check("zu_append_i64", st, sl); + } catch (Throwable t) { + throw fail("zu_append_i64", t); + } + } + + @Override + public void appendDouble(long appender, double value) { + Scratch s = Scratch.get(); + MemorySegment sl = s.slots(); + clear(sl); + try { + int st = (int) abi.appendF64.invokeExact(ptr(appender), value, sl.asSlice(ERR, 8)); + check("zu_append_f64", st, sl); + } catch (Throwable t) { + throw fail("zu_append_f64", t); + } + } + + @Override + public void appendString(long appender, String value) { + Scratch s = Scratch.get(); + MemorySegment sl = s.slots(); + MemorySegment v = s.utf8(value); + clear(sl); + try { + int st = + (int) abi.appendStr.invokeExact(ptr(appender), v, v.byteSize(), sl.asSlice(ERR, 8)); + check("zu_append_str", st, sl); + } catch (Throwable t) { + throw fail("zu_append_str", t); + } + } + + @Override + public void appendBytes(long appender, ByteBuffer value) { + Scratch s = Scratch.get(); + MemorySegment sl = s.slots(); + clear(sl); + try (Arena arena = Arena.ofConfined()) { + MemorySegment v = pass(value, arena); + int st = (int) abi.appendBytes.invokeExact(ptr(appender), v, v.byteSize(), sl.asSlice(ERR, 8)); + check("zu_append_bytes", st, sl); + } catch (Throwable t) { + throw fail("zu_append_bytes", t); + } + } + + @Override + public void appendTemporal(long appender, int kind, long count) { + Scratch s = Scratch.get(); + MemorySegment sl = s.slots(); + clear(sl); + try { + int st = (int) abi.appendTemporal.invokeExact(ptr(appender), kind, count, sl.asSlice(ERR, 8)); + check("zu_append_temporal", st, sl); + } catch (Throwable t) { + throw fail("zu_append_temporal", t); + } + } + + @Override + public void appendEndRow(long appender) { + endTransaction(abi.appendEndRow, "zu_append_end_row", appender); + } + + @Override + public void appenderFlush(long appender) { + endTransaction(abi.appenderFlush, "zu_appender_flush", appender); + } + + @Override + public long appenderBuffered(long appender) { + return counter(abi.appenderBuffered, "zu_appender_buffered", appender); + } + + @Override + public long appenderCommitted(long appender) { + return counter(abi.appenderCommitted, "zu_appender_committed", appender); + } + + @Override + public int appenderColumns(long appender) { + Scratch s = Scratch.get(); + MemorySegment sl = s.slots(); + try { + int st = (int) abi.appenderCols.invokeExact(ptr(appender), sl.asSlice(OUT, 8)); + check("zu_appender_cols", st, null); + return sl.get(JAVA_INT, OUT); + } catch (Throwable t) { + throw fail("zu_appender_cols", t); + } + } + + @Override + public String appenderColumnName(long appender, int col) { + Scratch s = Scratch.get(); + MemorySegment sl = s.slots(); + try { + MemorySegment out = + (MemorySegment) abi.appenderColName.invokeExact(ptr(appender), col, sl.asSlice(LEN, 8)); + return utf8(out.address(), sl.get(JAVA_LONG, LEN)); + } catch (Throwable t) { + throw fail("zu_appender_col_name", t); + } + } + + @Override + public long appenderDiscard(long appender) { + return counter(abi.appenderDiscard, "zu_appender_discard", appender); + } + + @Override + public long appenderClose(long appender) { + Scratch s = Scratch.get(); + MemorySegment sl = s.slots(); + clear(sl); + try { + int st = + (int) abi.appenderClose.invokeExact(ptr(appender), sl.asSlice(OUT, 8), sl.asSlice(ERR, 8)); + check("zu_appender_close", st, sl); + return sl.get(JAVA_LONG, OUT); + } catch (Throwable t) { + throw fail("zu_appender_close", t); + } + } + + @Override + public void appenderFree(long appender) { + try { + abi.appenderFree.invokeExact(ptr(appender)); + } catch (Throwable t) { + throw fail("zu_appender_free", t); + } + } + + /** One of the {@code (handle, uint64_t *out)} calls that cannot fail with an error. */ + private long counter(java.lang.invoke.MethodHandle mh, String what, long handle) { + Scratch s = Scratch.get(); + MemorySegment sl = s.slots(); + try { + int st = (int) mh.invokeExact(ptr(handle), sl.asSlice(OUT, 8)); + check(what, st, null); + return sl.get(JAVA_LONG, OUT); + } catch (Throwable t) { + throw fail(what, t); + } + } + + /** The shape every one-array loader column call has. */ + private void loaderColumn( + java.lang.invoke.MethodHandle mh, + String what, + long loader, + String name, + int count, + java.util.function.Function values) { + Scratch scratch = Scratch.get(); + MemorySegment sl = scratch.slots(); + clear(sl); + try (Arena arena = Arena.ofConfined()) { + MemorySegment n = utf8(arena, name); + MemorySegment v = values.apply(arena); + int st = + (int) mh.invokeExact(ptr(loader), n, n.byteSize(), v, (long) count, sl.asSlice(ERR, 8)); + check(what, st, sl); + } catch (Throwable t) { + throw fail(what, t); + } + } + + /** + * A buffer where a native function can read it. + * + *

A direct buffer already is that, and {@link MemorySegment#ofBuffer} + * addresses exactly the region between its position and its limit, so nothing + * is copied and a load of a hundred million values costs the call. A heap + * buffer is memory nothing outside this JVM can address and has to be copied + * off-heap first. That is the difference between passing a {@code long[]} and + * passing a direct {@code LongBuffer}, and it is the reason both are offered. + */ + private static MemorySegment pass(LongBuffer values, Arena arena) { + if (values.isDirect()) { + return MemorySegment.ofBuffer(values); + } + long[] copy = new long[values.remaining()]; + values.duplicate().get(copy); + return arena.allocateFrom(JAVA_LONG, copy); + } + + private static MemorySegment pass(DoubleBuffer values, Arena arena) { + if (values.isDirect()) { + return MemorySegment.ofBuffer(values); + } + double[] copy = new double[values.remaining()]; + values.duplicate().get(copy); + return arena.allocateFrom(JAVA_DOUBLE, copy); + } + + private static MemorySegment pass(IntBuffer values, Arena arena) { + if (values.isDirect()) { + return MemorySegment.ofBuffer(values); + } + int[] copy = new int[values.remaining()]; + values.duplicate().get(copy); + return arena.allocateFrom(JAVA_INT, copy); + } + + private static MemorySegment pass(ByteBuffer values, Arena arena) { + if (values.isDirect()) { + return MemorySegment.ofBuffer(values); + } + byte[] copy = new byte[values.remaining()]; + values.duplicate().get(copy); + return arena.allocateFrom(JAVA_BYTE, copy); + } + + /** A string as UTF-8 in an arena, without a terminator, since the length goes beside it. */ + private static MemorySegment utf8(Arena arena, String s) { + byte[] bytes = s.getBytes(StandardCharsets.UTF_8); + MemorySegment out = arena.allocate(bytes.length); + MemorySegment.copy(bytes, 0, out, JAVA_BYTE, 0, bytes.length); + return out; + } + + /** Writes one {@code size_t}, which is not the same width everywhere. */ + private static void size(MemorySegment array, long index, long value) { + if (Abi.SIZE_T.byteSize() == 8) { + array.setAtIndex(JAVA_LONG, index, value); + } else { + array.setAtIndex(JAVA_INT, index, (int) value); + } + } + private long openOrCreate( java.lang.invoke.MethodHandle mh, String what, diff --git a/zudb-ffm/src/test/java/dev/zudb/ffm/AppenderTest.java b/zudb-ffm/src/test/java/dev/zudb/ffm/AppenderTest.java new file mode 100644 index 0000000..9232e18 --- /dev/null +++ b/zudb-ffm/src/test/java/dev/zudb/ffm/AppenderTest.java @@ -0,0 +1,287 @@ +package dev.zudb.ffm; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import dev.zudb.Appender; +import dev.zudb.Connection; +import dev.zudb.Database; +import dev.zudb.Loader; +import dev.zudb.Result; +import dev.zudb.Value; +import dev.zudb.ZuException; +import java.nio.file.Path; +import java.util.ArrayList; +import java.util.List; +import java.util.stream.Collectors; +import org.junit.jupiter.api.BeforeAll; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.io.TempDir; + +/** + * Adding rows to a table that already exists. + * + *

Every test here builds its own database first, because the engine has no + * DDL and a bulk load is the only thing that makes a table for an appender to + * append to. + */ +class AppenderTest { + + @TempDir Path dir; + + @BeforeAll + static void engine() { + Libzu.require(); + } + + /** A two column Person table with three rows in it, at a path of its own. */ + private Path people(String name) { + Path path = dir.resolve(name + ".zu"); + try (Loader loader = Loader.create(path)) { + loader.table("Person", "Knows", 3); + loader.column("id", 1L, 2L, 3L); + loader.column("name", "ada", "grace", "alan"); + loader.finish(); + } + return path; + } + + @Test + void rowsGoInAndComeBack() { + Path path = people("rows"); + try (Database db = Database.open(path); + Connection conn = db.connect()) { + try (Appender rows = conn.appender("Person")) { + rows.append(4L).append("hedy").endRow(); + rows.append(5L).append("katherine").endRow(); + assertEquals(2L, rows.finish()); + assertTrue(rows.isFinished()); + } + + try (Result r = conn.query("MATCH (p:Person) RETURN p.name ORDER BY p.id")) { + assertEquals( + List.of("ada", "grace", "alan", "hedy", "katherine"), + r.stream().map(row -> row.getString(0)).collect(Collectors.toList())); + } + } + } + + @Test + void aRowIsARowOnceEndRowHasEndedIt() { + Path path = people("ended"); + try (Database db = Database.open(path); + Connection conn = db.connect(); + Appender rows = conn.appender("Person")) { + assertEquals(0L, rows.buffered()); + rows.append(4L).append("hedy"); + assertEquals(0L, rows.buffered()); + rows.endRow(); + assertEquals(1L, rows.buffered()); + assertEquals(0L, rows.committed()); + rows.flush(); + assertEquals(0L, rows.buffered()); + assertEquals(1L, rows.committed()); + } + } + + @Test + void theAppenderSaysWhatItIsWriting() { + Path path = people("columns"); + try (Database db = Database.open(path); + Connection conn = db.connect(); + Appender rows = conn.appender("Person")) { + assertEquals(2, rows.columns()); + assertEquals("id", rows.columnName(0)); + assertEquals("name", rows.columnName(1)); + ZuException e = assertThrows(ZuException.class, () -> rows.columnName(2)); + assertTrue(e.getMessage().contains("only 2"), e.getMessage()); + } + } + + @Test + void whatIsDiscardedNeverArrives() { + Path path = people("discard"); + try (Database db = Database.open(path); + Connection conn = db.connect()) { + try (Appender rows = conn.appender("Person")) { + rows.append(4L).append("hedy").endRow(); + rows.flush(); + rows.append(5L).append("katherine").endRow(); + assertEquals(1L, rows.discard()); + assertEquals(0L, rows.buffered()); + // What an earlier flush wrote is written, and a discard does not + // reach back to it. + assertEquals(1L, rows.committed()); + } + + try (Result r = conn.query("MATCH (p:Person) RETURN count(*)")) { + assertEquals(4L, r.row(0).getLong(0)); + } + } + } + + @Test + void closingWithoutFinishingKeepsWhatWasWritten() { + // A loop that threw halfway keeps the rows it managed. Throwing away work + // that succeeded is not a decision a close gets to make on its own. + Path path = people("kept"); + try (Database db = Database.open(path); + Connection conn = db.connect()) { + try (Appender rows = conn.appender("Person")) { + rows.append(4L).append("hedy").endRow(); + assertFalse(rows.isFinished()); + } + + try (Result r = conn.query("MATCH (p:Person) RETURN count(*)")) { + assertEquals(4L, r.row(0).getLong(0)); + } + } + } + + @Test + void aValueTheColumnWillNotTakeEndsItsRowAndLeavesNoHalfOfIt() { + Path path = people("refused"); + try (Database db = Database.open(path); + Connection conn = db.connect()) { + try (Appender rows = conn.appender("Person")) { + rows.append(4L); + // The second column holds strings, and a double is not one. + assertThrows(ZuException.class, () -> rows.append(1.5)); + assertEquals(0L, rows.buffered()); + rows.append(4L).append("hedy").endRow(); + assertEquals(1L, rows.finish()); + } + + try (Result r = conn.query("MATCH (p:Person) RETURN count(*)")) { + assertEquals(4L, r.row(0).getLong(0)); + } + } + } + + @Test + void aRowOfObjectsIsARowOfValues() { + Path path = people("objects"); + try (Database db = Database.open(path); + Connection conn = db.connect()) { + try (Appender rows = conn.appender("Person")) { + rows.row(4, "hedy"); + rows.row(5L, "katherine"); + assertEquals(2L, rows.finish()); + } + + try (Result r = conn.query("MATCH (p:Person) WHERE p.id > 3 RETURN count(*)")) { + assertEquals(2L, r.row(0).getLong(0)); + } + } + } + + @Test + void aRowOfSomethingNoColumnHoldsSaysWhichClass() { + Path path = people("wrongclass"); + try (Database db = Database.open(path); + Connection conn = db.connect(); + Appender rows = conn.appender("Person")) { + ZuException e = assertThrows(ZuException.class, () -> rows.row(4L, List.of("hedy"))); + assertTrue(e.getMessage().contains("no column holds one of those"), e.getMessage()); + } + } + + @Test + void thereIsNoNullToAppend() { + Path path = people("nulls"); + try (Database db = Database.open(path); + Connection conn = db.connect(); + Appender rows = conn.appender("Person")) { + assertThrows(ZuException.class, () -> rows.append((String) null)); + ZuException e = assertThrows(ZuException.class, () -> rows.row(4L, null)); + assertTrue(e.getMessage().contains("no null to append"), e.getMessage()); + } + } + + @Test + void aTableThatIsNotThereIsRefusedAtTheOpen() { + Path path = people("missing"); + try (Database db = Database.open(path); + Connection conn = db.connect()) { + assertThrows(ZuException.class, () -> conn.appender("Nobody")); + } + } + + @Test + void aClosedAppenderSaysSoRatherThanCrashing() { + Path path = people("closed"); + try (Database db = Database.open(path); + Connection conn = db.connect()) { + Appender rows = conn.appender("Person"); + rows.close(); + rows.close(); + assertThrows(ZuException.class, () -> rows.append(4L)); + } + } + + @Test + void aTemporalGoesInAsTheCountItsKindImplies() { + Path path = dir.resolve("dates.zu"); + try (Loader loader = Loader.create(path)) { + loader.table("Event", "Before", 1); + loader.column("id", 1L); + loader.temporalColumn("on", Value.Temporal.Kind.DATE, java.time.LocalDate.EPOCH.toEpochDay()); + loader.finish(); + } + + try (Database db = Database.open(path); + Connection conn = db.connect()) { + try (Appender rows = conn.appender("Event")) { + rows.append(2L).append(java.time.LocalDate.of(1843, 8, 10)).endRow(); + assertEquals(1L, rows.finish()); + } + + try (Result r = conn.query("MATCH (e:Event) RETURN e.on ORDER BY e.id")) { + assertEquals(2, r.rows()); + assertEquals(java.time.LocalDate.EPOCH, r.row(0).getTemporal(0).toLocalDate()); + assertEquals(java.time.LocalDate.of(1843, 8, 10), r.row(1).getTemporal(0).toLocalDate()); + } + } + } + + @Test + void manyRowsAcrossManyFlushes() { + Path path = people("many"); + int count = 5_000; + try (Database db = Database.open(path); + Connection conn = db.connect()) { + try (Appender rows = conn.appender("Person")) { + for (int i = 0; i < count; i++) { + rows.append(100L + i).append("n" + i).endRow(); + } + assertEquals(count, rows.finish()); + } + + try (Result r = conn.query("MATCH (p:Person) WHERE p.id >= 100 RETURN count(*)")) { + assertEquals(count, r.row(0).getLong(0)); + } + } + } + + @Test + void anAppenderAndAQueryOnOneConnectionDoNotTreadOnEachOther() { + Path path = people("interleaved"); + List counts = new ArrayList<>(); + try (Database db = Database.open(path); + Connection conn = db.connect()) { + try (Appender rows = conn.appender("Person")) { + for (int i = 0; i < 3; i++) { + rows.append(10L + i).append("n" + i).endRow(); + rows.flush(); + try (Result r = conn.query("MATCH (p:Person) RETURN count(*)")) { + counts.add(r.row(0).getLong(0)); + } + } + rows.finish(); + } + } + assertEquals(List.of(4L, 5L, 6L), counts); + } +} diff --git a/zudb-ffm/src/test/java/dev/zudb/ffm/LoaderTest.java b/zudb-ffm/src/test/java/dev/zudb/ffm/LoaderTest.java new file mode 100644 index 0000000..671fb9f --- /dev/null +++ b/zudb-ffm/src/test/java/dev/zudb/ffm/LoaderTest.java @@ -0,0 +1,269 @@ +package dev.zudb.ffm; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import dev.zudb.Connection; +import dev.zudb.Database; +import dev.zudb.Loader; +import dev.zudb.Result; +import dev.zudb.Value; +import dev.zudb.ZuException; +import java.nio.ByteBuffer; +import java.nio.ByteOrder; +import java.nio.LongBuffer; +import java.nio.file.Files; +import java.nio.file.Path; +import java.time.LocalDate; +import java.util.ArrayList; +import java.util.HashSet; +import java.util.List; +import java.util.stream.Collectors; +import org.junit.jupiter.api.BeforeAll; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.io.TempDir; + +/** Building a database out of columns, which is the only way a table comes into being. */ +class LoaderTest { + + @TempDir Path dir; + + @BeforeAll + static void engine() { + Libzu.require(); + } + + @Test + void aTableOfTwoColumnsReadsBack() { + Path path = dir.resolve("people.zu"); + try (Loader loader = Loader.create(path)) { + loader.table("Person", "Knows", 3); + loader.column("id", 1L, 2L, 3L); + loader.column("name", "ada", "grace", "alan"); + loader.finish(); + assertTrue(loader.isFinished()); + } + + try (Database db = Database.open(path); + Connection conn = db.connect(); + Result r = conn.query("MATCH (p:Person) RETURN p.id, p.name ORDER BY p.id")) { + List ids = new ArrayList<>(); + List names = new ArrayList<>(); + r.forEach( + row -> { + ids.add(row.getLong(0)); + names.add(row.getString(1)); + }); + assertEquals(List.of(1L, 2L, 3L), ids); + assertEquals(List.of("ada", "grace", "alan"), names); + } + } + + @Test + void everyColumnKindGoesInAndComesBack() { + Path path = dir.resolve("kinds.zu"); + try (Loader loader = Loader.create(path)) { + loader.table("Thing", "Near", 2); + loader.column("n", 7L, 8L); + loader.column("d", 1.5, 2.5); + loader.column("ok", true, false); + loader.column("s", "one", "two"); + loader.temporalColumn( + "born", + Value.Temporal.Kind.DATE, + LocalDate.of(1815, 12, 10).toEpochDay(), + LocalDate.of(1906, 12, 9).toEpochDay()); + loader.finish(); + } + + try (Database db = Database.open(path); + Connection conn = db.connect(); + Result r = + conn.query("MATCH (t:Thing) RETURN t.n, t.d, t.ok, t.s, t.born ORDER BY t.n")) { + assertEquals(2, r.rows()); + assertEquals(7L, r.row(0).getLong(0)); + assertEquals(1.5, r.row(0).getDouble(1)); + assertTrue(r.row(0).getBoolean(2)); + assertEquals("one", r.row(0).getString(3)); + assertEquals(LocalDate.of(1815, 12, 10), r.row(0).getTemporal(4).toLocalDate()); + assertEquals(8L, r.row(1).getLong(0)); + assertEquals(2.5, r.row(1).getDouble(1)); + assertFalse(r.row(1).getBoolean(2)); + assertEquals("two", r.row(1).getString(3)); + assertEquals(LocalDate.of(1906, 12, 9), r.row(1).getTemporal(4).toLocalDate()); + } + } + + @Test + void aDirectBufferIsReadWhereItLies() { + // The same table filled from both shapes, once through an array that has + // to be copied off-heap and once through a direct buffer that does not, + // so the path this whole surface exists for is the one under test. + Path path = dir.resolve("direct.zu"); + LongBuffer values = + ByteBuffer.allocateDirect(4 * Long.BYTES).order(ByteOrder.nativeOrder()).asLongBuffer(); + values.put(new long[] {10L, 20L, 30L, 40L}); + values.flip(); + + try (Loader loader = Loader.create(path)) { + loader.table("Point", "Near", 4); + loader.column("x", values); + loader.column("y", 1L, 2L, 3L, 4L); + loader.finish(); + } + + try (Database db = Database.open(path); + Connection conn = db.connect(); + Result r = conn.query("MATCH (p:Point) RETURN sum(p.x), sum(p.y)")) { + assertEquals(100L, r.row(0).getLong(0)); + assertEquals(10L, r.row(0).getLong(1)); + } + } + + @Test + void aSliceOfABufferIsTheSliceAndNotTheWholeThing() { + Path path = dir.resolve("slice.zu"); + LongBuffer all = LongBuffer.wrap(new long[] {99L, 1L, 2L, 3L, 99L}); + all.position(1); + all.limit(4); + + try (Loader loader = Loader.create(path)) { + loader.table("Slice", "Near", 3); + loader.column("v", all); + loader.finish(); + } + + try (Database db = Database.open(path); + Connection conn = db.connect(); + Result r = conn.query("MATCH (s:Slice) RETURN sum(s.v)")) { + assertEquals(6L, r.row(0).getLong(0)); + } + } + + @Test + void edgesJoinRowsToRows() { + Path path = dir.resolve("follows.zu"); + try (Loader loader = Loader.create(path)) { + loader.table("User", "Follows", 3); + loader.column("id", 1L, 2L, 3L); + loader.edges(new int[] {0, 1}, new int[] {1, 2}); + loader.finish(); + } + + try (Database db = Database.open(path); + Connection conn = db.connect(); + Result r = + conn.query("MATCH (a:User)-[:Follows]->(b:User) RETURN a.id, b.id ORDER BY a.id")) { + assertEquals(2, r.rows()); + assertEquals(1L, r.row(0).getLong(0)); + assertEquals(2L, r.row(0).getLong(1)); + assertEquals(2L, r.row(1).getLong(0)); + assertEquals(3L, r.row(1).getLong(1)); + } + } + + @Test + void edgesAppendAcrossCalls() { + Path path = dir.resolve("appended.zu"); + try (Loader loader = Loader.create(path)) { + loader.table("Node", "Link", 4); + loader.column("id", 0L, 1L, 2L, 3L); + loader.edges(new int[] {0}, new int[] {1}); + loader.edges(new int[] {1, 2}, new int[] {2, 3}); + loader.finish(); + } + + try (Database db = Database.open(path); + Connection conn = db.connect(); + Result r = conn.query("MATCH (:Node)-[:Link]->(:Node) RETURN count(*)")) { + assertEquals(3L, r.row(0).getLong(0)); + } + } + + @Test + void aColumnWithAValueMissingIsRefused() { + Path path = dir.resolve("short.zu"); + try (Loader loader = Loader.create(path)) { + loader.table("Person", "Knows", 3); + assertThrows(ZuException.class, () -> loader.column("id", 1L, 2L)); + } + } + + @Test + void anEdgeThatStartsSomewhereHasToEndSomewhere() { + Path path = dir.resolve("lopsided.zu"); + try (Loader loader = Loader.create(path)) { + loader.table("Node", "Link", 2); + ZuException e = + assertThrows(ZuException.class, () -> loader.edges(new int[] {0, 1}, new int[] {1})); + assertTrue(e.getMessage().contains("2 starts against 1 ends"), e.getMessage()); + } + } + + @Test + void aPathThatExistsIsRefused() throws Exception { + Path path = dir.resolve("taken.zu"); + Files.writeString(path, "not a database"); + assertThrows(ZuException.class, () -> Loader.create(path)); + } + + @Test + void aLoaderClosedWithoutFinishingWroteNothing() { + Path path = dir.resolve("abandoned.zu"); + try (Loader loader = Loader.create(path)) { + loader.table("Person", "Knows", 1); + loader.column("id", 1L); + assertFalse(loader.isFinished()); + } + // What is left is the empty file the loader created and nothing else. No + // half of a table got in, and there is no Person to find. + assertTrue(Files.exists(path)); + try (Database db = Database.open(path); + Connection conn = db.connect(); + Result r = conn.query("MATCH (p:Person) RETURN p.id")) { + assertEquals(0, r.rows()); + } + } + + @Test + void aClosedLoaderSaysSoRatherThanCrashing() { + Loader loader = Loader.create(dir.resolve("closed.zu")); + loader.close(); + loader.close(); + assertThrows(ZuException.class, () -> loader.table("Person", "Knows", 1)); + } + + @Test + void aStringColumnTakesWhatUtf8Takes() { + Path path = dir.resolve("utf8.zu"); + List names = List.of("", "ada", "éàü", "🐍", "a".repeat(300)); + try (Loader loader = Loader.create(path)) { + loader.table("Name", "Alias", names.size()); + loader.column("s", names); + loader.finish(); + } + + try (Database db = Database.open(path); + Connection conn = db.connect(); + Result r = conn.query("MATCH (n:Name) RETURN n.s")) { + assertEquals( + new HashSet<>(names), + r.stream().map(row -> row.getString(0)).collect(Collectors.toSet())); + } + } + + @Test + void aStringColumnWithAHoleInItIsRefusedBeforeTheCall() { + Path path = dir.resolve("hole.zu"); + try (Loader loader = Loader.create(path)) { + loader.table("Name", "Alias", 2); + List withHole = new ArrayList<>(); + withHole.add("ada"); + withHole.add(null); + ZuException e = assertThrows(ZuException.class, () -> loader.column("s", withHole)); + assertTrue(e.getMessage().contains("row 1"), e.getMessage()); + } + } +} diff --git a/zudb/src/main/java/dev/zudb/Appender.java b/zudb/src/main/java/dev/zudb/Appender.java new file mode 100644 index 0000000..5bdcf4a --- /dev/null +++ b/zudb/src/main/java/dev/zudb/Appender.java @@ -0,0 +1,406 @@ +package dev.zudb; + +import dev.zudb.spi.ZuBinding; +import java.nio.ByteBuffer; +import java.time.Duration; +import java.time.LocalDate; +import java.time.LocalDateTime; +import java.time.LocalTime; +import java.time.Period; +import java.time.ZoneOffset; +import java.util.concurrent.atomic.AtomicLong; + +/** + * Adds rows to a table that already exists, a value at a time, without a + * statement anywhere near it. + * + *

A statement is the wrong shape for bulk writes. Every row would be + * parsed, bound, planned and run, and none of that work says anything the row + * before it did not already say. An appender skips all of it: values go + * straight into the column they belong to, buffered until there are enough to + * write, and the cost of a row is the cost of the values in it. + * + *

{@code
+ * try (Appender rows = conn.appender("Person")) {
+ *   rows.append(4L).append("hedy").endRow();
+ *   rows.append(5L).append("katherine").endRow();
+ * }
+ * }
+ * + *

Values are written in the order the table declares its columns, which + * {@link #columnName(int)} will tell you, and a row is a row once + * {@link #endRow()} has ended it. A value the column will not take ends its + * row there and rolls back the values already written into it, so a refused + * append never leaves half a row behind. + * + *

There is no way to append no value. The C ABI has no null to append and + * this offers none rather than inventing one. + * + *

Rows are written in batches, so a row that has been ended is not yet a + * row anybody else can see. {@link #flush()} writes what is buffered, + * {@link #finish()} writes the rest and says how many rows went in altogether, + * and {@link #close()} on an appender that was never finished writes what it + * has anyway. That last one is deliberate: a loop that threw halfway keeps the + * rows it managed, because throwing away work that succeeded is not a decision + * a close should make on its own. Call {@link #discard()} to actually throw + * them away. + */ +public final class Appender implements AutoCloseable { + + private static final long NANOS = 1_000_000_000L; + + private final ZuBinding zu; + private final AtomicLong handle; + private long finished = -1; + + Appender(ZuBinding zu, long handle) { + this.zu = zu; + this.handle = new AtomicLong(handle); + } + + /** + * Appends a boolean. + * + * @param value what to write + * @return this appender + */ + public Appender append(boolean value) { + zu.appendBoolean(open(), value); + return this; + } + + /** + * Appends an integer. + * + * @param value what to write + * @return this appender + */ + public Appender append(long value) { + zu.appendLong(open(), value); + return this; + } + + /** + * Appends a double. + * + * @param value what to write + * @return this appender + */ + public Appender append(double value) { + zu.appendDouble(open(), value); + return this; + } + + /** + * Appends a string. + * + * @param value what to write, which is not allowed to be null because there + * is nothing this could write for it + * @return this appender + */ + public Appender append(String value) { + if (value == null) { + throw new ZuProgrammingException( + Diagnostic.misuse(Status.MISUSE, "append(null): there is no null to append")); + } + zu.appendString(open(), value); + return this; + } + + /** + * Appends bytes. + * + * @param value what to write + * @return this appender + */ + public Appender append(byte[] value) { + return append(ByteBuffer.wrap(value)); + } + + /** + * Appends bytes, read where they lie when the buffer is direct. + * + * @param value what to write, between the buffer's position and its limit + * @return this appender + */ + public Appender append(ByteBuffer value) { + zu.appendBytes(open(), value); + return this; + } + + /** + * Appends a date. + * + * @param value what to write + * @return this appender + */ + public Appender append(LocalDate value) { + return append(Value.Temporal.Kind.DATE, value.toEpochDay()); + } + + /** + * Appends a time of day. + * + * @param value what to write + * @return this appender + */ + public Appender append(LocalTime value) { + return append(Value.Temporal.Kind.LOCAL_TIME, value.toNanoOfDay()); + } + + /** + * Appends a datetime. + * + * @param value what to write + * @return this appender + */ + public Appender append(LocalDateTime value) { + return append( + Value.Temporal.Kind.LOCAL_DATETIME, + nanos(value.toEpochSecond(ZoneOffset.UTC), value.getNano())); + } + + /** + * Appends a span of months. + * + * @param value what to write, whose days are refused rather than turned into + * a length of time they do not have + * @return this appender + */ + public Appender append(Period value) { + if (value.getDays() != 0) { + throw new ZuProgrammingException( + Diagnostic.misuse( + Status.MISUSE, + "append(" + + value + + "): a year-month duration holds months, and days are a duration of their own")); + } + return append(Value.Temporal.Kind.DURATION_YEAR_MONTH, value.toTotalMonths()); + } + + /** + * Appends a span of time. + * + * @param value what to write + * @return this appender + */ + public Appender append(Duration value) { + return append(Value.Temporal.Kind.DURATION_DAY_TIME, value.toNanos()); + } + + /** + * Appends a temporal read out of a result, unchanged. + * + * @param value what to write + * @return this appender + */ + public Appender append(Value.Temporal value) { + return append(value.kind(), value.count()); + } + + /** + * Appends a temporal as a kind and the count in the unit that kind implies. + * + *

{@link Value.Temporal.Kind#ZONED_TIME} and + * {@link Value.Temporal.Kind#ZONED_DATETIME} are refused, because a stored + * column has nowhere to keep the offset that makes those two what they are. + * + * @param kind which of the seven + * @param count days for a date, months for a year-month duration, + * nanoseconds for the other five + * @return this appender + */ + public Appender append(Value.Temporal.Kind kind, long count) { + zu.appendTemporal(open(), kind.value(), count); + return this; + } + + /** + * Ends the row being written, which is what makes it a row. + * + * @return this appender + */ + public Appender endRow() { + zu.appendEndRow(open()); + return this; + } + + /** + * One whole row, for the caller who has it as objects already. + * + *

Each value is dispatched on the class it turns out to be, which costs a + * type test a value and is the shape a program reading from somewhere + * dynamic wants. A loop that knows what it is writing calls the + * {@code append} overloads and pays nothing. + * + * @param values one a column, in the order the table declares them + * @return this appender + * @throws ZuProgrammingException if one of them is a class no column holds + */ + public Appender row(Object... values) { + for (Object value : values) { + if (value instanceof Long v) { + append(v.longValue()); + } else if (value instanceof Integer v) { + append(v.longValue()); + } else if (value instanceof Short v) { + append(v.longValue()); + } else if (value instanceof Byte v) { + append(v.longValue()); + } else if (value instanceof Double v) { + append(v.doubleValue()); + } else if (value instanceof Float v) { + append(v.doubleValue()); + } else if (value instanceof Boolean v) { + append(v.booleanValue()); + } else if (value instanceof String v) { + append(v); + } else if (value instanceof byte[] v) { + append(v); + } else if (value instanceof ByteBuffer v) { + append(v); + } else if (value instanceof LocalDate v) { + append(v); + } else if (value instanceof LocalTime v) { + append(v); + } else if (value instanceof LocalDateTime v) { + append(v); + } else if (value instanceof Period v) { + append(v); + } else if (value instanceof Duration v) { + append(v); + } else if (value instanceof Value.Temporal v) { + append(v); + } else if (value == null) { + throw new ZuProgrammingException( + Diagnostic.misuse(Status.MISUSE, "row(..., null, ...): there is no null to append")); + } else { + throw new ZuProgrammingException( + Diagnostic.misuse( + Status.MISUSE, + "row(..., " + value.getClass().getName() + ", ...): no column holds one of those")); + } + } + return endRow(); + } + + /** + * Writes the rows that have been ended and not yet written. + * + * @return this appender + */ + public Appender flush() { + zu.appenderFlush(open()); + return this; + } + + /** + * Rows that have been ended and not yet written. + * + * @return the count, which a flush takes back to nought + */ + public long buffered() { + return zu.appenderBuffered(open()); + } + + /** + * Rows written across every flush so far. + * + * @return the count + */ + public long committed() { + return zu.appenderCommitted(open()); + } + + /** + * How many columns a row has. + * + * @return the count + */ + public int columns() { + return zu.appenderColumns(open()); + } + + /** + * What a column is called, so a program can check it is writing what it + * thinks it is. + * + * @param column the index, from nought + * @return the name + * @throws ZuProgrammingException if there is no such column + */ + public String columnName(int column) { + String name = zu.appenderColumnName(open(), column); + if (name == null) { + throw new ZuProgrammingException( + Diagnostic.misuse( + Status.MISUSE, + "there is no column " + column + " here, only " + columns() + " of them")); + } + return name; + } + + /** + * Throws away the rows that have been ended and not yet written. + * + *

Rows an earlier flush wrote are written and this does not reach them. + * + * @return how many rows were thrown away + */ + public long discard() { + return zu.appenderDiscard(open()); + } + + /** + * Writes what is left and spends the appender. + * + *

The only thing left to do with it afterwards is close it, and closing + * an appender that was finished frees it and writes nothing more. + * + * @return how many rows this appender wrote in all + */ + public long finish() { + finished = zu.appenderClose(open()); + return finished; + } + + /** + * Whether {@link #finish()} has run. + * + * @return true once the appender has been spent + */ + public boolean isFinished() { + return finished >= 0; + } + + /** + * Frees the appender, writing what is still buffered if + * {@link #finish()} never ran. + * + *

What it cannot do is tell you whether that last write worked, because + * a close has nowhere to report to. A program that needs to know calls + * {@link #finish()} and closes afterwards, which is what the count it hands + * back is for. Closing twice does nothing the second time. + */ + @Override + public void close() { + long h = handle.getAndSet(0); + if (h != 0) { + zu.appenderFree(h); + } + } + + private static long nanos(long seconds, int nano) { + return Math.addExact(Math.multiplyExact(seconds, NANOS), nano); + } + + private long open() { + long h = handle.get(); + if (h == 0) { + throw new ZuClosedException( + Diagnostic.misuse(Status.MISUSE_CLOSED, "this appender is closed")); + } + return h; + } +} diff --git a/zudb/src/main/java/dev/zudb/Connection.java b/zudb/src/main/java/dev/zudb/Connection.java index 81f15ef..5bc9463 100644 --- a/zudb/src/main/java/dev/zudb/Connection.java +++ b/zudb/src/main/java/dev/zudb/Connection.java @@ -70,6 +70,20 @@ public Statement prepare(String statement) { return new Statement(zu, zu.prepare(open(), statement)); } + /** + * Opens an appender on a table, which is how rows get in without a + * statement anywhere near them. + * + *

The appender writes through this connection for as long as it is open, + * so close it before the connection goes back to a pool. + * + * @param table the table to write to, which has to exist already + * @return the appender, which the caller closes + */ + public Appender appender(String table) { + return new Appender(zu, zu.appenderOpen(open(), table)); + } + /** * A second connection on the database this one is already on, made without * a path. diff --git a/zudb/src/main/java/dev/zudb/Loader.java b/zudb/src/main/java/dev/zudb/Loader.java new file mode 100644 index 0000000..aa48d07 --- /dev/null +++ b/zudb/src/main/java/dev/zudb/Loader.java @@ -0,0 +1,299 @@ +package dev.zudb; + +import dev.zudb.spi.ZuBinding; +import java.nio.DoubleBuffer; +import java.nio.IntBuffer; +import java.nio.LongBuffer; +import java.nio.file.Path; +import java.util.Arrays; +import java.util.List; +import java.util.concurrent.atomic.AtomicLong; + +/** + * Builds a database out of columns, which is the fastest way values get in and + * the only way a table comes into being. + * + *

A loader writes a new file. It refuses a path that already exists, + * because a bulk load builds a database rather than adding to one: what adds + * to one is {@link Appender}. The order is fixed and short. Say what the table + * is called and how many rows it has, hand over one column at a time and the + * edges if there are any, then {@link #finish()}. + * + *

{@code
+ * try (Loader loader = Loader.create(Path.of("people.zu"))) {
+ *   loader.table("Person", "Knows", 3);
+ *   loader.column("id", 1L, 2L, 3L);
+ *   loader.column("name", "ada", "grace", "alan");
+ *   loader.finish();
+ * }
+ * }
+ * + *

Every column is passed as an array or as a {@link java.nio.Buffer}, and + * which one you pass is the difference between a copy and no copy. A direct + * buffer is read where it lies: the engine sees the memory your program + * already filled and nothing crosses the boundary but a pointer. An array, or + * a heap buffer, is memory nothing outside the JVM can address, so it is + * copied off-heap first. For a few thousand rows that costs nothing worth + * measuring. For a hundred million it is the whole cost, and a program at that + * size should fill a {@link java.nio.ByteBuffer#allocateDirect} and view it. + * + *

The count each column carries has to match the row count the table was + * declared with. A column with a value missing is an error rather than a + * shorter table, which is the mistake this refuses to make quietly. + * + *

A loader that is closed without finishing wrote nothing, and leaves the + * empty file it created for the caller to remove. + */ +public final class Loader implements AutoCloseable { + + private final ZuBinding zu; + private final AtomicLong handle; + private boolean finished; + + private Loader(ZuBinding zu, long handle) { + this.zu = zu; + this.handle = new AtomicLong(handle); + } + + /** + * Starts a load into a database that does not exist yet. + * + * @param path where to write it, which must not exist + * @return the loader, which the caller closes + * @throws ZuException if the path exists or cannot be written + */ + public static Loader create(Path path) { + ZuBinding zu = Zu.binding(); + return new Loader(zu, zu.loaderCreate(path.toString())); + } + + /** + * Names the one table this loader builds and says how many rows it has. + * + *

Both names are wanted, even for a load that adds no edges at all. A + * node table comes with the edge table between its rows whether or not + * anything is in it, and naming it here rather than guessing a name later + * is the difference between a schema you chose and one that happened. + * + * @param nodes the node table + * @param edges the edge table between its rows + * @param rows how many rows every column of the node table will carry, + * which is given rather than counted so that a column with a value + * missing is an error and not a shorter table + * @return this loader + */ + public Loader table(String nodes, String edges, long rows) { + zu.loaderTable(open(), nodes, edges, rows); + return this; + } + + /** + * Adds edges as the row each one starts at and the row it ends at. + * + *

This appends, so call it as often as you like. The loader sorts and + * deduplicates at {@link #finish()}, so neither order nor a repeat matters. + * + * @param from the row each edge starts at + * @param to the row each edge ends at + * @return this loader + */ + public Loader edges(int[] from, int[] to) { + return edges(IntBuffer.wrap(from), IntBuffer.wrap(to)); + } + + /** + * Adds edges from two buffers, which are read where they lie when they are + * direct. + * + * @param from the row each edge starts at + * @param to the row each edge ends at + * @return this loader + */ + public Loader edges(IntBuffer from, IntBuffer to) { + zu.loaderEdges(open(), from, to); + return this; + } + + /** + * A column of integers. + * + * @param name the column + * @param values one a row + * @return this loader + */ + public Loader column(String name, long... values) { + return column(name, LongBuffer.wrap(values)); + } + + /** + * A column of integers, read where they lie when the buffer is direct. + * + * @param name the column + * @param values one a row, between the buffer's position and its limit + * @return this loader + */ + public Loader column(String name, LongBuffer values) { + zu.loaderColumnLongs(open(), name, values); + return this; + } + + /** + * A column of doubles. + * + * @param name the column + * @param values one a row + * @return this loader + */ + public Loader column(String name, double... values) { + return column(name, DoubleBuffer.wrap(values)); + } + + /** + * A column of doubles, read where they lie when the buffer is direct. + * + * @param name the column + * @param values one a row, between the buffer's position and its limit + * @return this loader + */ + public Loader column(String name, DoubleBuffer values) { + zu.loaderColumnDoubles(open(), name, values); + return this; + } + + /** + * A column of booleans. + * + * @param name the column + * @param values one a row + * @return this loader + */ + public Loader column(String name, boolean... values) { + int[] ints = new int[values.length]; + for (int i = 0; i < values.length; i++) { + ints[i] = values[i] ? 1 : 0; + } + return booleanColumn(name, IntBuffer.wrap(ints)); + } + + /** + * A column of booleans as the ints the C ABI carries them as, where anything + * that is not nought is true. + * + *

Separately named because a buffer of ints could as easily be meant as a + * column of integers, and guessing which is not something a bulk load should + * do. + * + * @param name the column + * @param values one a row, between the buffer's position and its limit + * @return this loader + */ + public Loader booleanColumn(String name, IntBuffer values) { + zu.loaderColumnBooleans(open(), name, values); + return this; + } + + /** + * A column of strings. + * + * @param name the column + * @param values one a row, none of them null + * @return this loader + */ + public Loader column(String name, String... values) { + return column(name, Arrays.asList(values)); + } + + /** + * A column of strings. + * + *

There is no zero-copy shape for this one. Every string is encoded and + * checked for UTF-8 on the way in, which is the price of never reading back + * a value no query could have returned. + * + * @param name the column + * @param values one a row, none of them null + * @return this loader + */ + public Loader column(String name, List values) { + zu.loaderColumnStrings(open(), name, values); + return this; + } + + /** + * A column of temporals, all of one kind, each row the count in the unit + * that kind implies. + * + *

This is {@link Row#getTemporal(int)} read backwards: a value that came + * out as 19782 days goes back in as 19782 days. + * {@link Value.Temporal.Kind#ZONED_TIME} and + * {@link Value.Temporal.Kind#ZONED_DATETIME} are refused, because a stored + * column has nowhere to keep the offset that makes those two what they are. + * + * @param name the column + * @param kind which of the seven every row of it is + * @param counts one a row + * @return this loader + */ + public Loader temporalColumn(String name, Value.Temporal.Kind kind, long... counts) { + return temporalColumn(name, kind, LongBuffer.wrap(counts)); + } + + /** + * A column of temporals, read where they lie when the buffer is direct. + * + * @param name the column + * @param kind which of the seven every row of it is + * @param counts one a row, between the buffer's position and its limit + * @return this loader + */ + public Loader temporalColumn(String name, Value.Temporal.Kind kind, LongBuffer counts) { + zu.loaderColumnTemporal(open(), name, kind.value(), counts); + return this; + } + + /** + * Writes it all. + * + *

The database is on disk when this returns, and {@link Database#open} + * on the same path reads it. The loader is spent afterwards and the only + * thing left to do with it is close it. + */ + public void finish() { + zu.loaderFinish(open()); + finished = true; + } + + /** + * Whether {@link #finish()} has run. + * + * @return true once the database is on disk + */ + public boolean isFinished() { + return finished; + } + + /** + * Frees the loader, which writes nothing that {@link #finish()} did not + * already write. + * + *

So a try-with-resources whose body threw leaves the empty file the + * loader created and no half-built database, which is the outcome that + * cannot be misread. Closing twice does nothing the second time. + */ + @Override + public void close() { + long h = handle.getAndSet(0); + if (h != 0) { + zu.loaderFree(h); + } + } + + private long open() { + long h = handle.get(); + if (h == 0) { + throw new ZuClosedException( + Diagnostic.misuse(Status.MISUSE_CLOSED, "this loader is closed")); + } + return h; + } +} diff --git a/zudb/src/main/java/dev/zudb/spi/ZuBinding.java b/zudb/src/main/java/dev/zudb/spi/ZuBinding.java index 32b5e8d..ba8429e 100644 --- a/zudb/src/main/java/dev/zudb/spi/ZuBinding.java +++ b/zudb/src/main/java/dev/zudb/spi/ZuBinding.java @@ -3,7 +3,9 @@ import dev.zudb.Diagnostic; import java.nio.ByteBuffer; import java.nio.DoubleBuffer; +import java.nio.IntBuffer; import java.nio.LongBuffer; +import java.util.List; /** * The C ABI, as Java. One method per call in {@code zu.h}, named for it, and @@ -578,4 +580,235 @@ public interface ZuBinding { * @return the name, never null */ String valueField(long value, long index); + + // ---- bulk load ---- + + /** + * Starts a load, which builds a database that does not exist yet. + * + * @param path the file to make, which must not exist + * @return the loader handle + */ + long loaderCreate(String path); + + /** + * Names the one table this load builds and how many rows it holds. + * + * @param loader the loader + * @param nodes what the node table is called + * @param edges what the relationship table is called, which the engine wants + * even for a load that adds no edges at all + * @param rows how many rows every column will carry + */ + void loaderTable(long loader, String nodes, String edges, long rows); + + /** + * Adds edges, as the row each starts at and the row it ends at. Appends, so + * it may be called as often as the caller likes. + * + * @param loader the loader + * @param from the starting row of each edge + * @param to the ending row of each edge + */ + void loaderEdges(long loader, IntBuffer from, IntBuffer to); + + /** + * Adds a column of integers. + * + *

A direct buffer is read where it lies and nothing is copied on this side + * of the boundary. A heap buffer is copied off-heap first, because a native + * function cannot be handed a Java array without either a copy or a pause + * long enough to matter on a column this size. + * + * @param loader the loader + * @param name what the column is called + * @param values the values, of which {@code remaining()} are read + */ + void loaderColumnLongs(long loader, String name, LongBuffer values); + + /** + * Adds a column of doubles. + * + * @param loader the loader + * @param name what the column is called + * @param values the values, of which {@code remaining()} are read + */ + void loaderColumnDoubles(long loader, String name, DoubleBuffer values); + + /** + * Adds a column of booleans, one {@code int} a row, where anything not zero + * is true. + * + * @param loader the loader + * @param name what the column is called + * @param values the values, of which {@code remaining()} are read + */ + void loaderColumnBooleans(long loader, String name, IntBuffer values); + + /** + * Adds a column of strings. Every one is checked for UTF-8 by the engine as + * it arrives rather than read back later as something no query could return. + * + * @param loader the loader + * @param name what the column is called + * @param values the values, none of which may be null + */ + void loaderColumnStrings(long loader, String name, List values); + + /** + * Adds a column of dates, times, datetimes or durations, as one kind and the + * count each row holds in the unit that kind implies. + * + * @param loader the loader + * @param name what the column is called + * @param kind which of the {@code ZU_TEMPORAL_} kinds every row is + * @param values the counts, of which {@code remaining()} are read + */ + void loaderColumnTemporal(long loader, String name, int kind, LongBuffer values); + + /** + * Writes it all. The database is on disk when this returns. + * + * @param loader the loader + */ + void loaderFinish(long loader); + + /** + * Releases a loader. A loader freed before it finished wrote nothing. + * + * @param loader the loader + */ + void loaderFree(long loader); + + // ---- appending ---- + + /** + * Opens an appender on a table that already exists. + * + * @param conn the connection to write through + * @param table what the table is called + * @return the appender handle + */ + long appenderOpen(long conn, String table); + + /** + * Appends one boolean to the row being written. + * + * @param appender the appender + * @param value the value + */ + void appendBoolean(long appender, boolean value); + + /** + * Appends one integer to the row being written. + * + * @param appender the appender + * @param value the value + */ + void appendLong(long appender, long value); + + /** + * Appends one double to the row being written. + * + * @param appender the appender + * @param value the value + */ + void appendDouble(long appender, double value); + + /** + * Appends one string to the row being written. + * + * @param appender the appender + * @param value the value + */ + void appendString(long appender, String value); + + /** + * Appends one run of bytes to the row being written. + * + * @param appender the appender + * @param value the bytes, of which {@code remaining()} are read + */ + void appendBytes(long appender, ByteBuffer value); + + /** + * Appends one date, time, datetime or duration to the row being written. + * + * @param appender the appender + * @param kind which of the {@code ZU_TEMPORAL_} kinds it is + * @param count how many of the unit that kind implies + */ + void appendTemporal(long appender, int kind, long count); + + /** + * Ends the row being written, which is what makes it a row. + * + * @param appender the appender + */ + void appendEndRow(long appender); + + /** + * Writes what is buffered. + * + * @param appender the appender + */ + void appenderFlush(long appender); + + /** + * Rows ended and not yet written. + * + * @param appender the appender + * @return the count + */ + long appenderBuffered(long appender); + + /** + * Rows written across every flush. + * + * @param appender the appender + * @return the count + */ + long appenderCommitted(long appender); + + /** + * How many values a row carries. + * + * @param appender the appender + * @return the count + */ + int appenderColumns(long appender); + + /** + * What one of those values is called. + * + * @param appender the appender + * @param col the column, counting from zero + * @return the name, or null out of range + */ + String appenderColumnName(long appender, int col); + + /** + * Throws away what is buffered. Rows an earlier flush wrote are written and + * this does not reach them. + * + * @param appender the appender + * @return how many rows were thrown away + */ + long appenderDiscard(long appender); + + /** + * Flushes what is left and spends the appender. + * + * @param appender the appender + * @return how many rows it wrote in all + */ + long appenderClose(long appender); + + /** + * Releases an appender. Writes what is still buffered, and cannot say + * whether that worked, which is what {@link #appenderClose(long)} is for. + * + * @param appender the appender + */ + void appenderFree(long appender); }