aboutsummaryrefslogtreecommitdiffstats
path: root/packages/meshbay-node
diff options
context:
space:
mode:
Diffstat (limited to 'packages/meshbay-node')
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc/disk.py7
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc/upload_handlers.py49
-rw-r--r--packages/meshbay-node/tests/test_upload_disk_full.py134
3 files changed, 187 insertions, 3 deletions
diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/disk.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/disk.py
index 2e484de..e298942 100644
--- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/disk.py
+++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/disk.py
@@ -1,5 +1,6 @@
"""Blocking disk work the session hands to `off_disk`."""
+import shutil
from pathlib import Path
from meshbay_common.protocol import file_chunk_wire
@@ -76,6 +77,12 @@ def _read_and_encrypt(
return file_chunk_wire(gek, plaintext, chunk_index, file_hash, file_id)
+def _free_bytes(directory: Path) -> int:
+ """Space left on the filesystem holding `directory`. Blocking; called
+ through `off_disk`, since a `statvfs` wakes a sleeping disk like any stat."""
+ return shutil.disk_usage(directory).free
+
+
def _append_chunk(tmp_path: Path, chunk_bytes: bytes, first: bool) -> None:
"""Add one chunk to a partial upload. Blocking; called through `off_disk`."""
with open(tmp_path, "wb" if first else "ab") as f:
diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/upload_handlers.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/upload_handlers.py
index 02a50af..b19846c 100644
--- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/upload_handlers.py
+++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/upload_handlers.py
@@ -2,6 +2,7 @@
renamed into place under the name the node chose."""
import asyncio
+import errno
import logging
from pathlib import Path
@@ -16,7 +17,7 @@ from meshbay_node.roots import (
publish_upload,
shell_active,
)
-from meshbay_node.transport.webrtc.disk import _append_chunk
+from meshbay_node.transport.webrtc.disk import _append_chunk, _free_bytes
from meshbay_node.transport.webrtc.limits import LEASE_NONE, LEASE_QUEUED
log = logging.getLogger("meshbay_node.transport.webrtc_server")
@@ -34,6 +35,17 @@ log = logging.getLogger("meshbay_node.transport.webrtc_server")
MAX_UPLOAD_BYTES = 8 * 1024 * 1024 * 1024 # 8 GB per file
GB_BYTES = 1024 * 1024 * 1024
+# What an upload may never take from the operator's disk. A filesystem filled to
+# its last byte is not only a refused upload: the node's own databases, the
+# system's logs and whatever else the operator runs there start failing too.
+# A gigabyte is noise beside any disk a group is hosted on and enough for those
+# to keep writing.
+DISK_RESERVE_BYTES = 1 * GB_BYTES
+
+# Both mean "this write found no room": a full filesystem, or the operator's own
+# disk quota on it. Windows' ERROR_DISK_FULL arrives as ENOSPC too.
+_NO_ROOM = frozenset({errno.ENOSPC, getattr(errno, "EDQUOT", errno.ENOSPC)})
+
class UploadMixin:
def _max_upload_bytes(self) -> int:
@@ -354,6 +366,20 @@ class UploadMixin:
if await off_disk(roots, final_path.exists):
_refuse("File already exists", "already_exists")
return
+ # Asked once, before the first byte, from the size the client
+ # announces: every chunk but the last is the size of this one, so
+ # this is an upper bound within one chunk. Peer-supplied, and that
+ # is fine — overstating it refuses only the sender's own upload.
+ # Without it a file that cannot fit is found out at the chunk that
+ # fills the disk, after the disk is full, and a nightly photo
+ # backup asks again every day. Concurrent uploads can each pass
+ # this; the ENOSPC below is what bounds them.
+ needed = max(0, total_chunks) * len(chunk_bytes)
+ if await off_disk(roots, _free_bytes, target_dir) - needed < DISK_RESERVE_BYTES:
+ log.warning("Upload refused: not enough free space in %s", rel_dir)
+ self._audit("upload_refused", "disk_full")
+ _refuse("The node's disk is full", "disk_full")
+ return
state = uploads.start(user_id, rel_dir, filename, stored_name,
part_path=tmp_path)
elif state is None:
@@ -372,7 +398,21 @@ class UploadMixin:
_refuse("Upload exceeds size limit", "too_large")
return
- await off_disk(roots, _append_chunk, tmp_path, chunk_bytes, chunk_index == 0)
+ try:
+ await off_disk(roots, _append_chunk, tmp_path, chunk_bytes, chunk_index == 0)
+ except OSError as e:
+ if e.errno not in _NO_ROOM:
+ raise
+ # The partial is dropped with its state: it cannot be finished
+ # until the operator makes room, and keeping it holds the very
+ # space that ran out. The client starts over from zero, which is
+ # what it would have to do for a resume against a full disk anyway.
+ uploads.drop(user_id, rel_dir, filename)
+ await off_disk(roots, tmp_path.unlink, True)
+ log.warning("Upload refused: the disk holding %s is full", rel_dir)
+ self._audit("upload_refused", "disk_full")
+ _refuse("The node's disk is full", "disk_full")
+ return
uploads.advance(user_id, rel_dir, filename, chunk_index, len(chunk_bytes))
last = chunk_index + 1 >= total_chunks
@@ -387,7 +427,10 @@ class UploadMixin:
uploads.reserved_names(rel_dir))
except OSError as e:
log.warning("Upload %s could not be published: %s", stored_name, e)
- _refuse("The file could not be stored", "store_failed")
+ if e.errno in _NO_ROOM:
+ _refuse("The node's disk is full", "disk_full")
+ else:
+ _refuse("The file could not be stored", "store_failed")
return
final_path = target_dir / stored_name
diff --git a/packages/meshbay-node/tests/test_upload_disk_full.py b/packages/meshbay-node/tests/test_upload_disk_full.py
new file mode 100644
index 0000000..c977a22
--- /dev/null
+++ b/packages/meshbay-node/tests/test_upload_disk_full.py
@@ -0,0 +1,134 @@
+"""
+A full disk is a stated refusal, not an exception.
+
+Nothing on the upload path used to know about ENOSPC: a write that found no
+room raised out of the handler, the catch-all answered "Request failed", and
+the `.part` stayed behind holding the space that had run out. A client could
+not tell that from any other fault, so the phone's nightly photo backup would
+retry the same file every day. Now the node refuses with `disk_full` — up
+front when the announced size cannot fit beside the reserve, and at the chunk
+where a write actually fails — and keeps nothing it cannot finish.
+"""
+
+import errno
+
+import pytest
+from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey
+from meshbay_common.crypto import generate_gek
+from meshbay_node.indexer.group_index import GroupIndex
+from meshbay_node.transport.webrtc import upload_handlers
+from meshbay_node.transport.webrtc.upload_handlers import DISK_RESERVE_BYTES
+from meshbay_node.transport.webrtc_server import WebRTCPeerSession
+
+from conftest import one_root, sealed_upload
+
+GROUP = "g" * 32
+
+
+def _peer(ctx: dict) -> WebRTCPeerSession:
+ session = WebRTCPeerSession.__new__(WebRTCPeerSession)
+ session._ctx = ctx
+ session._group_id = GROUP
+ session._user_id = "user-1"
+ session._pk_user = ""
+ session.sent = []
+ session._send = session.sent.append
+ session.audited = []
+ session._audit = lambda *a, **k: session.audited.append(a)
+ return session
+
+
+def _group_ctx(tmp_path) -> dict:
+ shared = tmp_path / "shared"
+ shared.mkdir(exist_ok=True)
+ return {"roots": one_root(shared),
+ "index": GroupIndex(group_id=GROUP, sk_node=Ed25519PrivateKey.generate()),
+ "gek": generate_gek()}
+
+
+def _codes(session):
+ return [m.get("code") for m in session.sent if m.get("type") == "error"]
+
+
+def _free(monkeypatch, nbytes: int) -> None:
+ monkeypatch.setattr(upload_handlers, "_free_bytes", lambda _d: nbytes)
+
+
+def _shared(ctx):
+ return ctx["roots"].roots[0].path
+
+
+async def test_a_file_that_cannot_fit_is_refused_before_a_byte_is_written(tmp_path, monkeypatch):
+ ctx = _group_ctx(tmp_path)
+ peer = _peer(ctx)
+ _free(monkeypatch, DISK_RESERVE_BYTES + 10)
+ await peer._do_file_upload(sealed_upload(peer, filename="IMG_0001.jpg",
+ data=b"x" * 8, total_chunks=2))
+ assert _codes(peer) == ["disk_full"]
+ assert list(_shared(ctx).iterdir()) == [], "a refused upload left a file behind"
+ assert ("upload_refused", "disk_full") in peer.audited
+
+
+async def test_the_reserve_is_never_handed_out(tmp_path, monkeypatch):
+ """Exactly enough for the file, nothing for the reserve, is still a refusal:
+ a disk filled to the last byte breaks the node's own databases too."""
+ ctx = _group_ctx(tmp_path)
+ peer = _peer(ctx)
+ _free(monkeypatch, DISK_RESERVE_BYTES + 15)
+ await peer._do_file_upload(sealed_upload(peer, filename="a.jpg",
+ data=b"x" * 8, total_chunks=2))
+ assert _codes(peer) == ["disk_full"]
+
+
+async def test_room_beside_the_reserve_is_accepted(tmp_path, monkeypatch):
+ ctx = _group_ctx(tmp_path)
+ peer = _peer(ctx)
+ _free(monkeypatch, DISK_RESERVE_BYTES + 16)
+ await peer._do_file_upload(sealed_upload(peer, filename="a.jpg",
+ data=b"x" * 8, total_chunks=2))
+ assert _codes(peer) == []
+
+
+@pytest.mark.parametrize("err", [errno.ENOSPC, getattr(errno, "EDQUOT", errno.ENOSPC)])
+async def test_a_write_that_finds_no_room_is_refused_and_cleaned_up(tmp_path, monkeypatch, err):
+ """The check above can be passed by two uploads at once, or by a disk the
+ operator fills meanwhile. The write is the last word."""
+ ctx = _group_ctx(tmp_path)
+ peer = _peer(ctx)
+ _free(monkeypatch, 100 * DISK_RESERVE_BYTES)
+ await peer._do_file_upload(sealed_upload(peer, filename="film.mkv", data=b"first",
+ chunk_index=0, total_chunks=2))
+ assert _codes(peer) == []
+
+ def no_room(*_a):
+ raise OSError(err, "No space left on device")
+ monkeypatch.setattr(upload_handlers, "_append_chunk", no_room)
+ await peer._do_file_upload(sealed_upload(peer, filename="film.mkv", data=b"second",
+ chunk_index=1, total_chunks=2))
+ assert _codes(peer) == ["disk_full"]
+ assert list(_shared(ctx).iterdir()) == [], (
+ "the .part was kept, holding the very space that ran out")
+ assert ctx["partial_uploads"].get("user-1", "shared", "film.mkv") is None
+
+
+async def test_another_write_failure_is_not_called_a_full_disk(tmp_path, monkeypatch):
+ ctx = _group_ctx(tmp_path)
+ peer = _peer(ctx)
+ _free(monkeypatch, 100 * DISK_RESERVE_BYTES)
+
+ def denied(*_a):
+ raise OSError(errno.EACCES, "Permission denied")
+ monkeypatch.setattr(upload_handlers, "_append_chunk", denied)
+ with pytest.raises(OSError):
+ await peer._do_file_upload(sealed_upload(peer, filename="a.jpg", data=b"x"))
+
+
+async def test_the_real_free_space_is_what_is_read(tmp_path):
+ """The un-patched path: a small file into a temp directory goes through."""
+ ctx = _group_ctx(tmp_path)
+ peer = _peer(ctx)
+ if upload_handlers._free_bytes(_shared(ctx)) < DISK_RESERVE_BYTES + 1024:
+ pytest.skip("this machine's temp filesystem is itself nearly full")
+ await peer._do_file_upload(sealed_upload(peer, filename="a.jpg", data=b"hello"))
+ assert _codes(peer) == []
+ assert (_shared(ctx) / "a.jpg").read_bytes() == b"hello"