jetstream v2 in zig stream.waow.tech
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188"""nbdkit block backend that can discard every write after the last flush.
This is not a Stream storage substitute. The production Linux binary stilluses RocksDB, std.Io, and ext4 through /dev/nbdN. This backend models a drivewith volatile write cache: writes update ``live``; NBD FLUSH/FUA copies dirtyextents to ``durable``. Killing nbdkit and starting it again reconstructs``live`` exclusively from ``durable``, which is the power-loss operation.
Environment: STRICT_NBD_DURABLE durable image path STRICT_NBD_LIVE volatile image path STRICT_NBD_SIZE bytes (used when creating a new image) STRICT_NBD_RECEIPT append-only flush receipt path"""
import osimport shutilimport threading
import nbdkit
API_VERSION = 2
durable_path = os.environ["STRICT_NBD_DURABLE"]live_path = os.environ["STRICT_NBD_LIVE"]receipt_path = os.environ["STRICT_NBD_RECEIPT"]size = int(os.environ["STRICT_NBD_SIZE"])live_fd = -1durable_fd = -1dirty = []flush_count = 0lock = threading.RLock()
def _write_all(fd, data, offset): view = memoryview(data) written = 0 while written < len(view): n = os.pwrite(fd, view[written:], offset + written) if n <= 0: raise OSError("short pwrite") written += n
def _mark_dirty(offset, count): dirty.append((offset, offset + count))
def _coalesced_dirty(): if not dirty: return [] ordered = sorted(dirty) out = [list(ordered[0])] for start, end in ordered[1:]: prior = out[-1] if start <= prior[1]: prior[1] = max(prior[1], end) else: out.append([start, end]) return out
def _persist_dirty(): global flush_count for start, end in _coalesced_dirty(): offset = start while offset < end: chunk = os.pread(live_fd, min(1024 * 1024, end - offset), offset) if not chunk: raise OSError("short pread while flushing") _write_all(durable_fd, chunk, offset) offset += len(chunk) os.fsync(durable_fd) dirty.clear() flush_count += 1 receipt_fd = os.open(receipt_path, os.O_WRONLY | os.O_CREAT | os.O_APPEND, 0o600) try: _write_all(receipt_fd, f"flush {flush_count}\n".encode("ascii"), os.lseek(receipt_fd, 0, os.SEEK_END)) os.fsync(receipt_fd) finally: os.close(receipt_fd)
def _initialize(): global live_fd, durable_fd if not os.path.exists(durable_path): fd = os.open(durable_path, os.O_RDWR | os.O_CREAT | os.O_EXCL, 0o600) try: os.posix_fallocate(fd, 0, size) os.fsync(fd) finally: os.close(fd) elif os.path.getsize(durable_path) != size: raise RuntimeError("durable image size does not match STRICT_NBD_SIZE") shutil.copyfile(durable_path, live_path) live_fd = os.open(live_path, os.O_RDWR) durable_fd = os.open(durable_path, os.O_RDWR)
def cleanup(): # Deliberately never flush here. SIGKILL is the normal power-cut path, but # even orderly plugin teardown must not accidentally make dirty data safe. if live_fd >= 0: os.close(live_fd) if durable_fd >= 0: os.close(durable_fd)
def open(readonly): if readonly: raise RuntimeError("strict NBD oracle requires a writable export") return 1
def get_size(_handle): return size
def can_flush(_handle): return True
def can_fua(_handle): return nbdkit.FUA_NATIVE
def can_trim(_handle): return True
def can_zero(_handle): return True
def pread(_handle, buf, offset, flags): assert flags == 0 with lock: data = os.pread(live_fd, len(buf), offset) if len(data) != len(buf): raise OSError("short pread") buf[:] = data
def pwrite(_handle, buf, offset, flags): with lock: _write_all(live_fd, buf, offset) _mark_dirty(offset, len(buf)) if flags & nbdkit.FLAG_FUA: _persist_dirty()
def flush(_handle, flags): assert flags == 0 with lock: _persist_dirty()
def _zero_range(count, offset): zeros = bytes(min(1024 * 1024, count)) remaining = count while remaining: chunk = zeros[: min(len(zeros), remaining)] _write_all(live_fd, chunk, offset) offset += len(chunk) remaining -= len(chunk)
def trim(_handle, count, offset, flags): # Logical zeroing is a conservative implementation of discard semantics. with lock: _zero_range(count, offset) _mark_dirty(offset, count) if flags & nbdkit.FLAG_FUA: _persist_dirty()
def zero(_handle, count, offset, flags): with lock: _zero_range(count, offset) _mark_dirty(offset, count) if flags & nbdkit.FLAG_FUA: _persist_dirty()
# The nbdkit Python API intentionally has no load callback. Top-level code is# the documented initialization mechanism; cleanup above is invoked at exit._initialize()