From 689df816d693ed2f103bca479e82587f12cdb9e8 Mon Sep 17 00:00:00 2001 From: goodboy Date: Mon, 27 Jul 2026 16:58:00 -0400 Subject: [PATCH] Drop redundant `NativeDB` writer locks `datad`'s persistent feed task already owns each series writer. The actor-local and file locks duplicated this SC contract and left `.parquet.lock` sidecars that upstream indexed as data. - preserve validated temp writes, `fsync()`, and atomic replacement - keep exact-suffix indexing for legacy lock and temp sidecars - prove full and incremental writes create no lock artifacts Prompt-IO: ai/prompt-io/opencode/20260727T204711Z_84a6d47b_prompt_io.md (this patch was generated in some part by `opencode` using `gpt-5.6-sol` (`openai`)) --- .../20260727T204711Z_84a6d47b_prompt_io.md | 40 ++++++ ...20260727T204711Z_84a6d47b_prompt_io.raw.md | 30 ++++ piker/storage/nativedb.py | 129 ++++-------------- tests/test_storage_nativedb.py | 54 +++----- 4 files changed, 121 insertions(+), 132 deletions(-) create mode 100644 ai/prompt-io/opencode/20260727T204711Z_84a6d47b_prompt_io.md create mode 100644 ai/prompt-io/opencode/20260727T204711Z_84a6d47b_prompt_io.raw.md diff --git a/ai/prompt-io/opencode/20260727T204711Z_84a6d47b_prompt_io.md b/ai/prompt-io/opencode/20260727T204711Z_84a6d47b_prompt_io.md new file mode 100644 index 00000000..bc058266 --- /dev/null +++ b/ai/prompt-io/opencode/20260727T204711Z_84a6d47b_prompt_io.md @@ -0,0 +1,40 @@ +--- +model: gpt-5.6-sol +provider: openai +service: opencode +session: unavailable +timestamp: 2026-07-27T20:47:11Z +git_ref: 84a6d47b +scope: code +substantive: true +raw_file: 20260727T204711Z_84a6d47b_prompt_io.raw.md +--- + +## Prompt + +The user identified that `backfiller_deep_fixes` introduced persistent +`.parquet.lock` files, questioned whether Parquet or Arrow already handled +writer coordination, and clarified that piker's `tractor` structured +concurrency design already provides single-writer ownership through the +persistent `datad` feed task. The user then asked to open the backfiller +worktree and resolve the redundant locking. + +## Response summary + +Removed redundant actor-local and filesystem locking while preserving +atomic crash-safe Parquet replacement. Replaced the synthetic contention +test with regressions for the actual cross-branch failure: writes create +no lock sidecars, and legacy lock artifacts are ignored during indexing. + +## Files changed + +- `piker/storage/nativedb.py` - rely on `datad` writer ownership +- `tests/test_storage_nativedb.py` - lock-sidecar regressions +- `ai/prompt-io/opencode/20260727T204711Z_84a6d47b_prompt_io.raw.md` + - unedited response record +- `ai/prompt-io/opencode/20260727T204711Z_84a6d47b_prompt_io.md` + - provenance metadata and response summary + +## Human edits + +None - generated changes have not been edited by the human. diff --git a/ai/prompt-io/opencode/20260727T204711Z_84a6d47b_prompt_io.raw.md b/ai/prompt-io/opencode/20260727T204711Z_84a6d47b_prompt_io.raw.md new file mode 100644 index 00000000..a158c740 --- /dev/null +++ b/ai/prompt-io/opencode/20260727T204711Z_84a6d47b_prompt_io.raw.md @@ -0,0 +1,30 @@ +--- +model: gpt-5.6-sol +provider: openai +service: opencode +timestamp: 2026-07-27T20:47:11Z +git_ref: 84a6d47b +diff_cmd: git diff HEAD~1..HEAD +--- + +> `git diff HEAD~1..HEAD -- piker/storage/nativedb.py` + +Removed actor-local and filesystem writer locks from NativeDB. The +persistent `datad` feed task already owns writes for each series, while +its timeframe children write distinct files. NativeDB retains validated +temporary writes, file and directory synchronization, and atomic target +replacement for crash durability. The exact `.parquet` suffix filter is +retained so legacy lock sidecars are harmless. + +> `git diff HEAD~1..HEAD -- tests/test_storage_nativedb.py` + +Removed the synthetic cross-client lock-contention test. Added regression +coverage proving replacement and incremental writes create no lock +sidecars. Extended indexing coverage with the legacy +`.parquet.lock` artifact that caused upstream `flake_update` to crash. + +Verification generated with the patch: + +- NativeDB and history regression set: 20 passed +- Ruff: passed +- `git diff --check`: passed diff --git a/piker/storage/nativedb.py b/piker/storage/nativedb.py index 73fd5850..54e89116 100644 --- a/piker/storage/nativedb.py +++ b/piker/storage/nativedb.py @@ -51,15 +51,8 @@ YET! # - https://github.com/spslater/borgapi # - https://nixos.wiki/wiki/ZFS from __future__ import annotations -from collections.abc import AsyncIterator from contextlib import asynccontextmanager as acm from datetime import datetime -from fcntl import ( - flock, - LOCK_EX, - LOCK_NB, - LOCK_UN, -) import os from pathlib import Path from tempfile import NamedTemporaryFile @@ -70,7 +63,6 @@ import polars as pl from pendulum import ( from_timestamp, ) -import trio from piker import config from piker import tsp @@ -151,10 +143,6 @@ class NativeStorageClient: # series' cache from tsdb reads self._dfs: dict[str, dict[str, pl.DataFrame]] = {} - self._write_locks: dict[ - tuple[str, int], - trio.Lock, - ] = {} @property def address(self) -> str: @@ -273,50 +261,6 @@ class NativeStorageClient: {}, )[fqme] = df - def _get_write_lock( - self, - fqme: str, - timeframe: int, - - ) -> trio.Lock: - ''' - Return the actor-local lock for one durable series. - - ''' - key: tuple[str, int] = (fqme, timeframe) - return self._write_locks.setdefault( - key, - trio.Lock(), - ) - - @acm - async def _open_file_lock( - self, - fqme: str, - timeframe: int, - - ) -> AsyncIterator[None]: - ''' - Serialize a series read-merge-write across client processes. - - ''' - path: Path = self.mk_path(fqme, timeframe) - lock_path: Path = path.with_name(f'.{path.name}.lock') - with lock_path.open('a+b') as lock_file: - while True: - try: - flock( - lock_file.fileno(), - LOCK_EX | LOCK_NB, - ) - break - except BlockingIOError: - await trio.sleep(0.01) - try: - yield - finally: - flock(lock_file.fileno(), LOCK_UN) - async def read_ohlcv( self, fqme: str, @@ -499,46 +443,39 @@ class NativeStorageClient: Merge an incremental frame into durable OHLCV history. Incoming rows replace stored rows at matching timestamps. - Writes for each series are serialized so concurrent backfills - can not overwrite each other's read-merge-write cycle. + Writer ownership belongs to the persistent ``datad`` feed task. ''' - lock: trio.Lock = self._get_write_lock( + if isinstance(ohlcv, np.ndarray): + incoming: pl.DataFrame = tsp.np2pl(ohlcv) + else: + incoming = ohlcv + self._validate_ohlcv(incoming) + + path: Path = self.mk_path(fqme, timeframe) + if path.exists(): + stored: pl.DataFrame = pl.read_parquet(path) + merged: pl.DataFrame = pl.concat( + [stored, incoming], + how='diagonal_relaxed', + ) + merged = ( + merged + .unique( + subset='time', + keep='last', + ) + .sort('time') + ) + else: + merged = incoming + + self._validate_ohlcv(merged) + return self._write_ohlcv( fqme, + merged, timeframe, ) - async with lock: - async with self._open_file_lock(fqme, timeframe): - if isinstance(ohlcv, np.ndarray): - incoming: pl.DataFrame = tsp.np2pl(ohlcv) - else: - incoming = ohlcv - self._validate_ohlcv(incoming) - - path: Path = self.mk_path(fqme, timeframe) - if path.exists(): - stored: pl.DataFrame = pl.read_parquet(path) - merged: pl.DataFrame = pl.concat( - [stored, incoming], - how='diagonal_relaxed', - ) - merged = ( - merged - .unique( - subset='time', - keep='last', - ) - .sort('time') - ) - else: - merged = incoming - - self._validate_ohlcv(merged) - return self._write_ohlcv( - fqme, - merged, - timeframe, - ) async def write_ohlcv( self, @@ -552,17 +489,11 @@ class NativeStorageClient: to (local) disk. ''' - lock: trio.Lock = self._get_write_lock( + return self._write_ohlcv( fqme, + ohlcv, timeframe, ) - async with lock: - async with self._open_file_lock(fqme, timeframe): - return self._write_ohlcv( - fqme, - ohlcv, - timeframe, - ) async def delete_ts( self, diff --git a/tests/test_storage_nativedb.py b/tests/test_storage_nativedb.py index a1f4b674..46ec0c20 100644 --- a/tests/test_storage_nativedb.py +++ b/tests/test_storage_nativedb.py @@ -2,11 +2,6 @@ NativeDB durability and history-preservation regressions. ''' -from fcntl import ( - flock, - LOCK_EX, - LOCK_UN, -) import os from pathlib import Path @@ -309,15 +304,22 @@ def test_series_do_not_interfere( assert y_60['time'].to_list() == [60, 120] -def test_index_files_ignores_lock_and_stale_temp_files( +def test_index_files_ignores_sidecar_files( tmp_path: Path, ) -> None: ''' - Crash leftovers and writer locks are not durable series entries. + Legacy writer locks and crash leftovers are not series entries. + + Earlier deep-fix revisions created persistent + ``.parquet.lock`` files beside each series. The upstream indexer + mistakes those sidecars for Parquet data and crashes while parsing + their period. Arrange both a legacy lock and stale temporary file, + then prove exact-suffix indexing exposes only the durable series. ''' client = NativeStorageClient(tmp_path) run(client.write_ohlcv('x.test', mk_ohlcv((60,)), 60)) + (tmp_path / '.x.test.ohlcv60s.parquet.lock').touch() (tmp_path / 'x.test.ohlcv60s.parquet.crash.tmp').touch() index = client.index_files() @@ -326,39 +328,25 @@ def test_index_files_ignores_lock_and_stale_temp_files( assert index['x.test']['period'] == 60 -def test_contended_file_lock_yields_to_trio( +def test_writes_create_no_lock_sidecars( tmp_path: Path, ) -> None: ''' - Cross-client lock contention does not block the actor loop. + Actor-owned NativeDB writes must not create lock sidecars. + + ``datad`` already gives each persistent feed one parent history + writer, with its child tasks writing distinct timeframe files. A + redundant filesystem lock previously leaked ``.parquet.lock`` + files into NativeDB and made upstream ``flake_update`` crash during + startup. Exercise replacement and incremental writes, then prove + no lock artifact exists and the merged history remains intact. ''' client = NativeStorageClient(tmp_path) run(client.write_ohlcv('x.test', mk_ohlcv((60,)), 60)) + run(client.update_ohlcv('x.test', mk_ohlcv((120,)), 60)) path = client.mk_path('x.test', 60) - lock_path = path.with_name(f'.{path.name}.lock') - - async def main() -> None: - done = trio.Event() - - async def update() -> None: - await client.update_ohlcv( - 'x.test', - mk_ohlcv((120,)), - 60, - ) - done.set() - - with lock_path.open('a+b') as lock_file: - flock(lock_file.fileno(), LOCK_EX) - async with trio.open_nursery() as nursery: - nursery.start_soon(update) - await trio.sleep(0.03) - assert not done.is_set() - flock(lock_file.fileno(), LOCK_UN) - with trio.fail_after(0.5): - await done.wait() - - trio.run(main) stored = pl.read_parquet(path) + + assert not list(tmp_path.glob('*.lock')) assert stored['time'].to_list() == [60, 120]