Skip to content
Closed
Show file tree
Hide file tree
Changes from 10 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