Skip to content
Merged
Show file tree
Hide file tree
Changes from all 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
8 changes: 4 additions & 4 deletions .github/workflows/python-app-arm64.yml
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
# This workflow will run on an external ubicloud arm64 runner on latest ubuntu (ubicloud standard)
# This workflow will run on a GitHub-hosted arm64 runner on Ubuntu 24.04

name: Python application (ubuntu, arm64)

Expand All @@ -13,15 +13,15 @@ on:
permissions:
contents: read

# One run per branch: pushing again supersedes the run already going, which
# matters most on the paid arm64 runner. A run for `main` is never cancelled.
# One run per branch: pushing again supersedes the run already going.
# A run for `main` is never cancelled.
concurrency:
group: ${{ github.workflow }}-${{ github.ref }}
cancel-in-progress: ${{ github.ref != 'refs/heads/main' }}

jobs:
build:
runs-on: ubicloud-standard-2-arm
runs-on: ubuntu-24.04-arm

steps:
- uses: actions/checkout@v5
Expand Down
5 changes: 5 additions & 0 deletions AGENTS.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
# Agent Guidelines for Caterva2

- **Trailing Newline**: Always ensure every created or modified file ends with a single trailing newline (`\n`) to satisfy pre-commit's `end-of-file-fixer`.
- **Pre-commit Compliance**: Ensure all code changes adhere to repository pre-commit hooks (formatting, linting, and whitespace).
- **Documentation**: Maintain documentation integrity, preserving existing comments and docstrings unless explicitly directed otherwise.
185 changes: 106 additions & 79 deletions caterva2/services/server.py
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@
import shutil
import string
import tarfile
import threading
import time
import traceback
import types
Expand Down Expand Up @@ -100,6 +101,29 @@ def dataset_lock(abspath) -> asyncio.Lock:
return lock


_thread_locks: weakref.WeakValueDictionary[str, threading.Lock] = weakref.WeakValueDictionary()
_thread_locks_guard = threading.Lock()


def dataset_thread_lock(abspath) -> threading.Lock:
"""Serialize threadpool writes on one array within this process.

Python-Blosc2's file locking (`locking=True`) executes in C without releasing
the GIL during open/lock. If two threads in the same process attempt to
open/lock the same file concurrently, one can block on the OS file lock
while holding Python's GIL, deadlocking the process. Guarding operations on
the same array with a threading.Lock ensures the waiting thread yields the
GIL cleanly.
"""
key = str(abspath)
with _thread_locks_guard:
lock = _thread_locks.get(key)
if lock is None:
lock = threading.Lock()
_thread_locks[key] = lock
return lock


mimetypes.add_type("text/markdown", ".md") # Because in macOS this is not by default
mimetypes.add_type("application/x-ipynb+json", ".ipynb")

Expand Down Expand Up @@ -1550,10 +1574,12 @@ def publish_dataset(abspath: pathlib.Path, path: pathlib.Path) -> str:
with contextlib.suppress(Exception):
fs.rm(staging)
raise
array = blosc2.open(abspath, mode="a", locking=True)
with array.schunk.holding_lock():
array.schunk.vlmeta[PUBLISHED_URL] = destination
array.schunk.vlmeta[FILL_STATE] = PUBLISHED
with dataset_thread_lock(abspath):
array = blosc2.open(abspath, mode="a", locking=True)
with array.schunk.holding_lock():
array.schunk.vlmeta[PUBLISHED_URL] = destination
array.schunk.vlmeta[FILL_STATE] = PUBLISHED
del array
return destination


Expand All @@ -1568,82 +1594,83 @@ def store_chunk(abspath: pathlib.Path, nchunk: int, chunk: bytes) -> dict:
both find the slot free would otherwise both write it, and the second would
move every chunk that came after the first.
"""
try:
array = blosc2.open(abspath, mode="a", locking=True)
except Exception as exc:
srv_utils.raise_bad_request(f"{abspath.name} cannot be opened for writing: {exc}")
if not isinstance(array, blosc2.NDArray):
srv_utils.raise_bad_request(
f"{abspath.name} is not an NDArray, so it has no chunks of a shape to write into"
)
schunk = array.schunk
if not 0 <= nchunk < schunk.nchunks:
srv_utils.raise_not_found(f"{abspath.name} has no chunk {nchunk}")
try:
nbytes, _, blocksize = blosc2.get_cbuffer_sizes(chunk)
typesize = chunk_typesize(chunk)
except Exception:
srv_utils.raise_bad_request("the body is not a Blosc2 chunk")
# A chunk of another geometry would be stored and then read as nonsense, so
# it is refused here rather than left for whoever reads the array next
if nbytes != schunk.chunksize:
srv_utils.raise_bad_request(
f"the chunk holds {nbytes} bytes where this array's chunks hold {schunk.chunksize}"
)
if blocksize != schunk.blocksize:
srv_utils.raise_bad_request(
f"the chunk is split into blocks of {blocksize} bytes where this array's are "
f"{schunk.blocksize}; compress it against the array's blocks"
)
# The one part of the geometry the sizes do not carry, and the one whose
# mismatch is silent: the shuffle filters read and write on a stride of it,
# so a chunk compressed against another typesize decompresses to the right
# number of bytes with every one of them in the wrong place -- no error
# anywhere, just an array of scrambled values
if typesize != filter_typesize(schunk.typesize):
srv_utils.raise_bad_request(
f"the chunk was compressed with a typesize of {typesize} where this array's is "
f"{schunk.typesize}; compress it against the array's dtype"
)
complete = False
with schunk.holding_lock():
if not chunk_is_unwritten(schunk, nchunk):
raise fastapi.HTTPException(
status_code=409, detail=f"chunk {nchunk} of {abspath.name} was already written"
with dataset_thread_lock(abspath):
try:
array = blosc2.open(abspath, mode="a", locking=True)
except Exception as exc:
srv_utils.raise_bad_request(f"{abspath.name} cannot be opened for writing: {exc}")
if not isinstance(array, blosc2.NDArray):
srv_utils.raise_bad_request(
f"{abspath.name} is not an NDArray, so it has no chunks of a shape to write into"
)
schunk.update_chunk(nchunk, chunk)
if FILL_NONCE not in schunk.vlmeta:
# What names *this* array, as against another one that came to sit at
# the same path with the same size. A client caching the array reads
# it from api/info and can tell the two apart, which a size and an
# mtime cannot always do. Written once, by whichever writer arrived
# first, and never again
schunk.vlmeta[FILL_NONCE] = uuid.uuid4().hex
# Said out loud rather than left to be inferred from the absence of
# it, and free here: the same locked region, the same trailer
schunk.vlmeta[FILL_STATE] = FILLING
written, nchunks = count_written(abspath)
state = schunk.vlmeta.get(FILL_STATE, FILLING)
if written == nchunks and state == FILLING:
# Exactly once, whichever writer got here: the lock is held, so of two
# writers that both see the array complete only one makes this move,
# and that one owns the publishing. Recorded even where there is
# nowhere to publish to, because "every slot is claimed" is worth
# saying on its own: it is what tells a reader the array can no
# longer change under a cache of it
complete = bool(settings.publish_root)
state = PUBLISHING if complete else COMPLETE
schunk.vlmeta[FILL_STATE] = state
# Drop the handle before anything reads the file again: a handle left open
# over a frame another one writes is the stale-handle hazard, and it is silent
del array, schunk
return {
"nchunk": nchunk,
"written": written,
"nchunks": nchunks,
"state": state,
"publish": complete,
}
schunk = array.schunk
if not 0 <= nchunk < schunk.nchunks:
srv_utils.raise_not_found(f"{abspath.name} has no chunk {nchunk}")
try:
nbytes, _, blocksize = blosc2.get_cbuffer_sizes(chunk)
typesize = chunk_typesize(chunk)
except Exception:
srv_utils.raise_bad_request("the body is not a Blosc2 chunk")
# A chunk of another geometry would be stored and then read as nonsense, so
# it is refused here rather than left for whoever reads the array next
if nbytes != schunk.chunksize:
srv_utils.raise_bad_request(
f"the chunk holds {nbytes} bytes where this array's chunks hold {schunk.chunksize}"
)
if blocksize != schunk.blocksize:
srv_utils.raise_bad_request(
f"the chunk is split into blocks of {blocksize} bytes where this array's are "
f"{schunk.blocksize}; compress it against the array's blocks"
)
# The one part of the geometry the sizes do not carry, and the one whose
# mismatch is silent: the shuffle filters read and write on a stride of it,
# so a chunk compressed against another typesize decompresses to the right
# number of bytes with every one of them in the wrong place -- no error
# anywhere, just an array of scrambled values
if typesize != filter_typesize(schunk.typesize):
srv_utils.raise_bad_request(
f"the chunk was compressed with a typesize of {typesize} where this array's is "
f"{schunk.typesize}; compress it against the array's dtype"
)
complete = False
with schunk.holding_lock():
if not chunk_is_unwritten(schunk, nchunk):
raise fastapi.HTTPException(
status_code=409, detail=f"chunk {nchunk} of {abspath.name} was already written"
)
schunk.update_chunk(nchunk, chunk)
if FILL_NONCE not in schunk.vlmeta:
# What names *this* array, as against another one that came to sit at
# the same path with the same size. A client caching the array reads
# it from api/info and can tell the two apart, which a size and an
# mtime cannot always do. Written once, by whichever writer arrived
# first, and never again
schunk.vlmeta[FILL_NONCE] = uuid.uuid4().hex
# Said out loud rather than left to be inferred from the absence of
# it, and free here: the same locked region, the same trailer
schunk.vlmeta[FILL_STATE] = FILLING
written, nchunks = count_written(abspath)
state = schunk.vlmeta.get(FILL_STATE, FILLING)
if written == nchunks and state == FILLING:
# Exactly once, whichever writer got here: the lock is held, so of two
# writers that both see the array complete only one makes this move,
# and that one owns the publishing. Recorded even where there is
# nowhere to publish to, because "every slot is claimed" is worth
# saying on its own: it is what tells a reader the array can no
# longer change under a cache of it
complete = bool(settings.publish_root)
state = PUBLISHING if complete else COMPLETE
schunk.vlmeta[FILL_STATE] = state
# Drop the handle before anything reads the file again: a handle left open
# over a frame another one writes is the stale-handle hazard, and it is silent
del array, schunk
return {
"nchunk": nchunk,
"written": written,
"nchunks": nchunks,
"state": state,
"publish": complete,
}


@app.post("/api/chunk/{path:path}")
Expand Down
4 changes: 3 additions & 1 deletion caterva2/tests/test_hdf5_tree.py
Original file line number Diff line number Diff line change
Expand Up @@ -44,7 +44,9 @@ def fill_h5_public(client):
dest_dir = pathlib.Path(TEST_STATE_DIR) / "server/public"
dest_dir.mkdir(parents=True, exist_ok=True)
fname = "test_tree.h5"
_make_h5(dest_dir / fname)
target = dest_dir / fname
if not target.exists():
_make_h5(target)
return fname, client.get(TEST_CATERVA2_ROOT)


Expand Down
54 changes: 54 additions & 0 deletions caterva2/tests/test_peers.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,9 +9,11 @@
import asyncio
import json
import os
import pathlib
import shutil
import signal
import socket
import ssl
import subprocess
import sys
import time
Expand Down Expand Up @@ -1020,6 +1022,58 @@ async def afetch(self, slice_, **kwargs):
assert len(proxy.calls) == 2


def test_caddy_fixture_negotiates_http2(tmp_path):
"""The optional local fixture is real TLS+h2, with direct Uvicorn as h1."""
caddy = shutil.which("caddy")
if caddy is None:
pytest.skip("Caddy is not installed")

upstream_port, h2_port = _unused_tcp_ports(2)
server_dir = tmp_path / "server"
(server_dir / "public").mkdir(parents=True)
server = _start(server_dir, upstream_port)
env = dict(
os.environ,
CATERVA2_UPSTREAM=f"127.0.0.1:{upstream_port}",
CATERVA2_H2_ADDRESS=f"localhost:{h2_port}",
XDG_DATA_HOME=str(tmp_path / "caddy-data"),
XDG_CONFIG_HOME=str(tmp_path / "caddy-config"),
)
caddyfile = pathlib.Path(__file__).parents[2] / "examples" / "benchmarks" / "http2" / "Caddyfile"
proxy = subprocess.Popen(
[caddy, "run", "--config", str(caddyfile)],
env=env,
stdout=subprocess.DEVNULL,
stderr=subprocess.DEVNULL,
)
try:
root_ca = tmp_path / "caddy-data" / "caddy" / "pki" / "authorities" / "local" / "root.crt"
for _ in range(50):
if proxy.poll() is not None:
raise RuntimeError(f"Caddy exited during startup with status {proxy.returncode}")
if root_ca.exists():
try:
ssl_context = ssl.create_default_context(cafile=str(root_ca))
with httpx.Client(http2=True, verify=ssl_context, timeout=1) as client:
response = client.get(f"https://localhost:{h2_port}/api/roots")
if response.is_success:
break
except httpx.TransportError:
pass
time.sleep(0.1)
else:
raise RuntimeError("Caddy HTTP/2 fixture did not start")

assert response.http_version == "HTTP/2"
direct = httpx.get(f"http://127.0.0.1:{upstream_port}/api/roots", timeout=1)
assert direct.http_version == "HTTP/1.1"
finally:
proxy.send_signal(signal.SIGTERM)
proxy.wait(timeout=10)
server.send_signal(signal.SIGTERM)
server.wait(timeout=10)


def test_concurrent_fetches_of_different_datasets_dont_serialize(two_dataset_peers):
"""Interleaved concurrent fetches of two different datasets under a tiny
shared quota: correctness under load is the regression net (per
Expand Down
9 changes: 9 additions & 0 deletions examples/benchmarks/http2/Caddyfile
Original file line number Diff line number Diff line change
@@ -0,0 +1,9 @@
{
# The local CA is deliberate: this fixture tests real TLS ALPN, not h2c.
local_certs
}

{$CATERVA2_H2_ADDRESS:localhost:8443} {
tls internal
reverse_proxy {$CATERVA2_UPSTREAM:127.0.0.1:8000}
}
Loading
Loading