diff --git a/src/atproto/sync.zig b/src/atproto/sync.zig index f549a37..9a691f6 100644 --- a/src/atproto/sync.zig +++ b/src/atproto/sync.zig @@ -1,4 +1,5 @@ const std = @import("std"); +const PortableAtomic = @import("../core/portable_atomic.zig").PortableAtomic; const config = @import("../core/config.zig"); const log = @import("../core/log.zig"); const http_api = @import("../http/api.zig"); @@ -16,7 +17,7 @@ const notify_threshold_ms = 20 * 60 * 1000; var subscribe_repos_connections: usize = 0; var subscribe_repos_backfills: usize = 0; var get_repo_connections: usize = 0; -var crawler_last_notified_ms: std.atomic.Value(i64) = .init(0); +var crawler_last_notified_ms: PortableAtomic(i64) = .init(0); var crawler_notify_in_flight: std.atomic.Value(bool) = .init(false); pub const SubscribeReposClient = struct { diff --git a/src/core/portable_atomic.zig b/src/core/portable_atomic.zig new file mode 100644 index 0000000..0c3b043 --- /dev/null +++ b/src/core/portable_atomic.zig @@ -0,0 +1,55 @@ +const std = @import("std"); + +/// std.atomic.Value, except types wider than the target's native atomic width +/// (e.g. i64/u64 on 32-bit ARM) fall back to a spinlock-guarded value. +pub fn PortableAtomic(comptime T: type) type { + if (@bitSizeOf(T) <= @bitSizeOf(usize)) return std.atomic.Value(T); + return struct { + lock_state: std.atomic.Value(u32) = .init(0), + raw: T, + + const Self = @This(); + + pub fn init(value: T) Self { + return .{ .raw = value }; + } + + fn acquire(self: *Self) void { + while (self.lock_state.cmpxchgWeak(0, 1, .acquire, .monotonic) != null) { + std.Thread.yield() catch {}; + } + } + + fn release(self: *Self) void { + self.lock_state.store(0, .release); + } + + pub fn load(self: *Self, comptime _: std.builtin.AtomicOrder) T { + self.acquire(); + defer self.release(); + return self.raw; + } + + pub fn store(self: *Self, value: T, comptime _: std.builtin.AtomicOrder) void { + self.acquire(); + defer self.release(); + self.raw = value; + } + + pub fn fetchAdd(self: *Self, operand: T, comptime _: std.builtin.AtomicOrder) T { + self.acquire(); + defer self.release(); + const prev = self.raw; + self.raw +%= operand; + return prev; + } + + pub fn fetchMax(self: *Self, operand: T, comptime _: std.builtin.AtomicOrder) T { + self.acquire(); + defer self.release(); + const prev = self.raw; + self.raw = @max(prev, operand); + return prev; + } + }; +} diff --git a/src/storage/eventlog.zig b/src/storage/eventlog.zig index c22ab9a..598fc2f 100644 --- a/src/storage/eventlog.zig +++ b/src/storage/eventlog.zig @@ -1,4 +1,5 @@ const std = @import("std"); +const PortableAtomic = @import("../core/portable_atomic.zig").PortableAtomic; const Io = std.Io; @@ -7,7 +8,7 @@ var initialized = false; var mutex: Io.Mutex = .init; var changed: Io.Condition = .init; var generation: usize = 0; -var latest_seq: u64 = 0; +var latest_seq: PortableAtomic(u64) = .init(0); pub fn init(io: Io) void { store_io = io; @@ -21,14 +22,14 @@ pub fn publish(seq: u64) void { defer mutex.unlock(store_io); _ = @atomicRmw(usize, &generation, .Add, 1, .release); - _ = @atomicRmw(u64, &latest_seq, .Max, seq, .release); + _ = latest_seq.fetchMax(seq, .release); changed.broadcast(store_io); } pub fn snapshot() Snapshot { return .{ .generation = @atomicLoad(usize, &generation, .acquire), - .latest_seq = @atomicLoad(u64, &latest_seq, .acquire), + .latest_seq = latest_seq.load(.acquire), }; } diff --git a/src/storage/store.zig b/src/storage/store.zig index 4afab4d..06f98e7 100644 --- a/src/storage/store.zig +++ b/src/storage/store.zig @@ -505,11 +505,13 @@ var conn: zqlite.Conn = undefined; var database_path: ?[:0]u8 = null; var initialized = false; var store_io: Io = undefined; +const PortableAtomic = @import("../core/portable_atomic.zig").PortableAtomic; + const DatabaseMutex = struct { inner: Io.Mutex = .init, owner_ptr: std.atomic.Value(usize) = .init(0), owner_len: std.atomic.Value(usize) = .init(0), - acquired_at_ms: std.atomic.Value(i64) = .init(0), + acquired_at_ms: PortableAtomic(i64) = .init(0), fn lockUncancelable(self: *DatabaseMutex, io: Io) void { const owner_ptr = self.owner_ptr.load(.acquire); @@ -542,7 +544,7 @@ const DatabaseMutex = struct { var db_mutex: DatabaseMutex = .{}; var write_lanes: sharded_locks.ShardedLocks(32) = .{}; var blob_lanes: sharded_locks.ShardedLocks(64) = .{}; -var next_seq: std.atomic.Value(u64) = .init(1); +var next_seq: PortableAtomic(u64) = .init(1); var blob_gc_worker_started: std.atomic.Value(bool) = .init(false); const blob_gc_grace_seconds: i64 = 24 * 60 * 60;