Skip to content
Closed
Show file tree
Hide file tree
Changes from 9 commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,8 @@ dist
*.tmp
*.dat
*.db
*.db*
*.wal

# local env files
.env
Expand Down
36 changes: 31 additions & 5 deletions benches.zig
Original file line number Diff line number Diff line change
@@ -1,28 +1,49 @@
const std = @import("std");
const zio = @import("zio");

const Codspeed = @import("codspeed");

const bench_lww = @import("benches/memory/lww.zig");
const bench_skiplist = @import("benches/memory/skiplist.zig");
const bench_table = @import("benches/memory/table.zig");
const bench_runtime = @import("benches/runtime/execution.zig");
const bench_storage = @import("benches/runtime/storage.zig");
const bench_wal = @import("benches/runtime/wal.zig");
pub const LwwRegistry = @import("src/memory/lww.zig").LwwRegistry;
pub const SkipList = @import("src/memory/skiplist.zig").SkipList;
pub const ColumnTable = @import("src/memory/table.zig").ColumnTable;
pub const hlc = @import("src/primitives/hlc.zig");

pub fn main(init: std.process.Init) !void {
const allocator = init.gpa;
const io = init.io;
var runtime = try zio.Runtime.init(allocator, .{});
defer runtime.deinit();
const io = runtime.io();

var args = init.minimal.args.iterate();
var use_codspeed = false;
var storage_only = false;
var runtime_only = false;
var wal_only = false;

while (args.next()) |arg| {
if (std.mem.eql(u8, arg, "--runner")) {
use_codspeed = true;
break;
}
if (std.mem.eql(u8, arg, "--runner")) use_codspeed = true;
if (std.mem.eql(u8, arg, "--storage-only")) storage_only = true;
if (std.mem.eql(u8, arg, "--runtime-only")) runtime_only = true;
if (std.mem.eql(u8, arg, "--wal-only")) wal_only = true;
}

if (storage_only) {
try bench_storage.run(allocator, io);
return;
}
if (runtime_only) {
try bench_runtime.run(allocator, io);
return;
}
if (wal_only) {
try bench_wal.run(allocator, io);
return;
}

if (use_codspeed) {
Expand All @@ -44,10 +65,15 @@ pub fn main(init: std.process.Init) !void {
try codspeed.start("runtime.execution");
try bench_runtime.run(allocator, io);
try codspeed.stop("runtime.execution");

try codspeed.start("runtime.storage");
try bench_storage.run(allocator, io);
try codspeed.stop("runtime.storage");
} else {
try bench_table.run(allocator, io);
try bench_lww.run(allocator, io);
try bench_skiplist.run(allocator, io);
try bench_runtime.run(allocator, io);
try bench_storage.run(allocator, io);
}
}
2 changes: 1 addition & 1 deletion benches/memory/lww.zig
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,7 @@ pub fn run(allocator: Allocator, io: std.Io) !void {
var reg = LwwRegistry.init(allocator);
defer reg.deinit();

var clock = Hlc.init(1, io);
var clock = Hlc.init(io, 1);

var key_bufs: [n_keys][24]u8 = undefined;
var key_slices: [n_keys][]const u8 = undefined;
Expand Down
100 changes: 79 additions & 21 deletions benches/runtime/execution.zig
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@ const loop_mod = @import("../../src/engine/loop.zig");
const InferenceLoop = loop_mod.InferenceLoop;
const RuleDispatcher = loop_mod.RuleDispatcher;
const LwwRegistry = @import("../../src/memory/lww.zig").LwwRegistry;
const Storage = @import("../../src/storage.zig").Storage;
const Arc = @import("../../src/primitives/arc.zig").Arc;
const Hlc = @import("../../src/primitives/hlc.zig").Hlc;
const Mutex = @import("../../src/primitives/mutex.zig").Mutex;
Expand All @@ -23,6 +24,18 @@ const wasm_wire = @import("../../src/wasm/wire.zig");

const n_iters: usize = 100_000;

fn removeBenchmarkFiles(io: std.Io, path: []const u8) void {
std.Io.Dir.cwd().deleteFile(io, path) catch {};
var wal_buf: [160]u8 = undefined;
if (std.fmt.bufPrint(&wal_buf, "{s}.slung.wal", .{path})) |wal_path| {
std.Io.Dir.cwd().deleteFile(io, wal_path) catch {};
} else |_| {}
}

fn elapsedNs(io: std.Io, start: std.Io.Timestamp) u64 {
return @intCast(start.untilNow(io, .awake).toNanoseconds());
}

fn deinitIndices(
allocator: Allocator,
forward: *graph_index.ForwardIndex,
Expand Down Expand Up @@ -152,7 +165,15 @@ pub fn run(allocator: Allocator, io: std.Io) !void {
claim_arc.release();
}

var clock = Hlc.init(1, io);
const storage_path = "runtime-benchmark.db";
removeBenchmarkFiles(io, storage_path);
var storage = try Storage.open(allocator, io, storage_path);
defer {
storage.deinit();
removeBenchmarkFiles(io, storage_path);
}

var clock = Hlc.init(io, 1);
var wasm_module: *zwasm.WasmModule = undefined;
var context = Context.init(
allocator,
Expand All @@ -166,7 +187,9 @@ pub fn run(allocator: Allocator, io: std.Io) !void {
wasm_module,
"bench_ns",
"node-1",
&storage,
);
defer context.deinit();

const env_imports = try wasm_host.createEnvImport(allocator, @intFromPtr(&context));
defer allocator.free(env_imports.source.host_fns);
Expand Down Expand Up @@ -207,13 +230,21 @@ pub fn run(allocator: Allocator, io: std.Io) !void {
.dispatch_fn = WasmRuleDispatcher.dispatch,
};
var loop = InferenceLoop.init(&context, rule_dispatcher, 10, allocator);
defer loop.deinit();

const input = "{\"value\": 42.0}";
const start = std.Io.Clock.awake.now(io);
var total_fired: usize = 0;
var mapping_ns: u64 = 0;
var source_checkpoint_ns: u64 = 0;
var lww_ns: u64 = 0;
var execution_ns: u64 = 0;
var cascade_checkpoint_ns: u64 = 0;

for (0..n_iters) |_| {
var phase_start = std.Io.Clock.awake.now(io);
const mapped = try invokeMapper(wasm_module, allocator, reading_mapper, input);
mapping_ns += elapsedNs(io, phase_start);
defer allocator.free(mapped);

var key_buf: [128]u8 = undefined;
Expand All @@ -222,34 +253,61 @@ pub fn run(allocator: Allocator, io: std.Io) !void {
"{s}:{d}:{d}",
.{ "bench_ns", reading_key.entity, reading_key.component },
);
const ts = clock.send();
const cause = Storage.FactMutation{
.namespace = "bench_ns",
.entity = reading_key.entity,
.component = reading_key.component,
.value = mapped,
.timestamp = ts,
.cause = .{ .cause = reading_key.component, .entity = reading_key.entity, .node = "node-1" },
};

// Source checkpoint: this is the durability boundary for the input.
phase_start = std.Io.Clock.awake.now(io);
_ = try context.persistMutation(cause);
source_checkpoint_ns += elapsedNs(io, phase_start);

// Update LWW cache so rules can read this fact immediately.
phase_start = std.Io.Clock.awake.now(io);
{
var store_guard = lww_arc.getMut().lock();
defer store_guard.deinit();

_ = try store_guard.get().put(
store_key,
clock.send(),
.{ .Bytes = mapped },
.{ .cause = reading_key.component, .entity = reading_key.entity, .node = "node-1" },
);
}

{
var queue_guard = dirty_queue_arc.getMut().lock();
defer queue_guard.deinit();
try queue_guard.get().push(.{ .entity = reading_key.entity, .component = reading_key.component });
_ = try store_guard.get().put(store_key, ts, .{ .Bytes = mapped }, cause.cause);
}
lww_ns += elapsedNs(io, phase_start);

total_fired += try loop.run();
// Rule execution and the final cascade checkpoint are timed separately.
const timing = try loop.runTimed();
total_fired += timing.fired;
execution_ns += timing.execution_ns;
cascade_checkpoint_ns += timing.checkpoint_ns;
}

const elapsed_ns: u64 = @intCast(start.untilNow(io, .awake).toNanoseconds());
const elapsed_ns = elapsedNs(io, start);
const elapsed_s = @as(f64, @floatFromInt(elapsed_ns)) / 1e9;
const ops_per_s = if (elapsed_s > 0) @as(u64, @intFromFloat(@as(f64, n_iters) / elapsed_s)) else 0;

std.debug.print(
" {d} updates in {d:.2}s - {d} updates/s - {d} rules fired\n",
.{ n_iters, elapsed_s, ops_per_s, total_fired },
);
const per_update = struct {
fn value(ns: u64) f64 {
return @as(f64, @floatFromInt(ns)) / @as(f64, @floatFromInt(n_iters)) / 1e6;
}
}.value;

std.debug.print(" {d} updates in {d:.2}s - {d} updates/s - {d} rules fired\n", .{
n_iters, elapsed_s, ops_per_s, total_fired,
});
std.debug.print(" phase avg (ms/update): map {d:.3}, source checkpoint {d:.3}, LWW {d:.3}, execution {d:.3}, cascade checkpoint {d:.3}\n", .{
per_update(mapping_ns),
per_update(source_checkpoint_ns),
per_update(lww_ns),
per_update(execution_ns),
per_update(cascade_checkpoint_ns),
});
std.debug.print(" phase total (ms): map {d:.1}, source checkpoint {d:.1}, LWW {d:.1}, execution {d:.1}, cascade checkpoint {d:.1}\n", .{
per_update(mapping_ns) * @as(f64, @floatFromInt(n_iters)),
per_update(source_checkpoint_ns) * @as(f64, @floatFromInt(n_iters)),
per_update(lww_ns) * @as(f64, @floatFromInt(n_iters)),
per_update(execution_ns) * @as(f64, @floatFromInt(n_iters)),
per_update(cascade_checkpoint_ns) * @as(f64, @floatFromInt(n_iters)),
});
}
84 changes: 84 additions & 0 deletions benches/runtime/storage.zig
Original file line number Diff line number Diff line change
@@ -0,0 +1,84 @@
//! SQLite persistence benchmark without Wasm or inference.

const std = @import("std");
const Allocator = std.mem.Allocator;
const Storage = @import("../../src/storage.zig").Storage;
const types = @import("../../src/types.zig");

const n_iters: usize = 100;
const batch_size: usize = 128;
const value = "{\"value\":42.0}";

fn removeDatabase(io: std.Io, path: []const u8) !void {
std.Io.Dir.cwd().deleteFile(io, path) catch |err| switch (err) {
error.FileNotFound => {},
else => return err,
};
var wal_buf: [128]u8 = undefined;
const wal_path = try std.fmt.bufPrint(&wal_buf, "{s}.slung.wal", .{path});
std.Io.Dir.cwd().deleteFile(io, wal_path) catch |err| switch (err) {
error.FileNotFound => {},
else => return err,
};
var shm_buf: [128]u8 = undefined;
const shm_path = try std.fmt.bufPrint(&shm_buf, "{s}-shm", .{path});
std.Io.Dir.cwd().deleteFile(io, shm_path) catch |err| switch (err) {
error.FileNotFound => {},
else => return err,
};
}

fn mutation(index: usize) Storage.FactMutation {
return .{
.namespace = "storage_bench",
.entity = 1,
.component = 1,
.value = value,
.timestamp = .{ .wall = @intCast(index + 1), .logical = 0, .node_id = 1 },
.cause = .{ .cause = 1, .entity = 1, .node = "bench" },
};
}

fn runSingle(allocator: Allocator, io: std.Io, path: []const u8) !void {
try removeDatabase(io, path);
var storage = try Storage.open(allocator, io, path);
defer storage.deinit();

const start = std.Io.Clock.awake.now(io);
for (0..n_iters) |index| {
_ = try storage.applyMutation(mutation(index));
}
const elapsed_ns: u64 = @intCast(start.untilNow(io, .awake).toNanoseconds());
const elapsed_s = @as(f64, @floatFromInt(elapsed_ns)) / 1e9;
const ops_per_s = if (elapsed_s > 0) @as(f64, @floatFromInt(n_iters)) / elapsed_s else 0;
std.debug.print(" single: {d} updates in {d:.2}s - {d:.0} updates/s\n", .{ n_iters, elapsed_s, ops_per_s });
}

fn runBatch(allocator: Allocator, io: std.Io, path: []const u8) !void {
try removeDatabase(io, path);
var storage = try Storage.open(allocator, io, path);
defer storage.deinit();

var mutations: [batch_size]Storage.FactMutation = undefined;
var accepted: [batch_size]bool = undefined;
const start = std.Io.Clock.awake.now(io);
var completed: usize = 0;
while (completed < n_iters) {
const count = @min(batch_size, n_iters - completed);
for (0..count) |offset| mutations[offset] = mutation(completed + offset);
try storage.applyMutations(mutations[0..count], accepted[0..count]);
completed += count;
}
const elapsed_ns: u64 = @intCast(start.untilNow(io, .awake).toNanoseconds());
const elapsed_s = @as(f64, @floatFromInt(elapsed_ns)) / 1e9;
const ops_per_s = if (elapsed_s > 0) @as(f64, @floatFromInt(n_iters)) / elapsed_s else 0;
std.debug.print(" batch ({d}): {d} updates in {d:.2}s - {d:.0} updates/s\n", .{ batch_size, n_iters, elapsed_s, ops_per_s });
}

pub fn run(allocator: Allocator, io: std.Io) !void {
std.debug.print("\n=== SQLite Storage ===\n", .{});
const path = "storage-benchmark.db";
try runSingle(allocator, io, path);
try runBatch(allocator, io, path);
try removeDatabase(io, path);
}
Loading
Loading