1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
|
"""
An indexer has one `progress`, and a whole-root walk owns it while it runs.
Found while planning a progress band for the operator (2026-09-14), by reading
who writes `progress`. Two writers had no business doing it during a scan:
- A watchdog burst — a file dropped into a folder while a 900 GB root was being
hashed — added its size to the scan's total, and cleared `scanning` when its
own hash finished. Anything showing progress would have gone blank with hours
of hashing left.
- A plug rescan took no lock, so it walked its root beside an added root's
scan: both reset the same counters, and both read the drive in turn.
"""
import asyncio
import os
from pathlib import Path
import pytest
from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey
from meshbay_node.indexer.indexer import DirectoryIndexer
from meshbay_node.roots import RootSet
pytestmark = pytest.mark.asyncio
GROUP = "g" * 32
def _set(*specs) -> RootSet:
return RootSet.build([s if isinstance(s, dict) else {"path": str(s), "name": s.name}
for s in specs])
def _names(idx) -> list[str]:
return sorted(e.name for e in idx.index.entries)
async def _until(predicate, timeout: float = 3.0) -> bool:
deadline = asyncio.get_running_loop().time() + timeout
while not predicate():
if asyncio.get_running_loop().time() > deadline:
return False
await asyncio.sleep(0.02)
return True
class _Held(DirectoryIndexer):
"""Stops hashing under `hold_under` once `hold_after` of its files are done."""
hold_under: Path | None = None
hold_after = 0
def __init__(self, *args, **kwargs):
super().__init__(*args, **kwargs)
self.gate = asyncio.Event()
self.at_gate = asyncio.Event()
self.rescanned: list[str] = []
self._seen = 0
async def _hash_or_cached(self, root, file_path):
if self.hold_under is not None and file_path.is_relative_to(self.hold_under):
if self._seen == self.hold_after:
self.at_gate.set()
await self.gate.wait()
self._seen += 1
return await super()._hash_or_cached(root, file_path)
async def _rescan_root(self, root):
self.rescanned.append(root.name)
return await super()._rescan_root(root)
def _tree(tmp_path: Path) -> tuple[Path, Path, Path]:
one, two, big = tmp_path / "one", tmp_path / "two", tmp_path / "big"
for d in (one, two, big):
d.mkdir()
(one / "a.txt").write_bytes(b"first root")
(two / "b.txt").write_bytes(b"removable root")
(big / "x.bin").write_bytes(os.urandom(3000))
(big / "y.bin").write_bytes(os.urandom(5000))
return one, two, big
async def test_a_burst_during_a_scan_leaves_the_scan_on_screen(tmp_path):
one, _, big = _tree(tmp_path)
idx = _Held(roots=_set(one), group_id=GROUP, sk_node=Ed25519PrivateKey.generate(),
gek=None, debounce_secs=0.01)
await idx.initial_scan()
idx._loop = asyncio.get_running_loop()
idx.hold_under, idx.hold_after = big, 1
try:
await idx.retarget(_set(one, big), wait=False)
await asyncio.wait_for(idx.at_gate.wait(), 5)
total, scanned = idx.progress.total_bytes, idx.progress.scanned_bytes
assert total == 8000 and scanned in (3000, 5000)
dropped = one / "dropped.txt"
dropped.write_bytes(os.urandom(700))
idx._schedule_update(dropped)
assert await _until(lambda: "dropped.txt" in _names(idx))
assert await _until(lambda: idx._burst_inflight == 0)
assert idx.progress.scanning is True, (
"the dropped file's hash finishing turned the scan's progress off")
assert (idx.progress.total_bytes, idx.progress.scanned_bytes) == (total, scanned), (
"the dropped file was counted into the scan under way")
idx.gate.set()
await asyncio.wait_for(asyncio.gather(*list(idx._scan_tasks)), 5)
assert await _until(lambda: not idx.progress.scanning)
assert _names(idx) == ["a.txt", "dropped.txt", "x.bin", "y.bin"]
finally:
idx.gate.set()
await idx.stop()
def _with_two_ejected(one: Path, two: Path, *more: Path) -> RootSet:
return _set(one, {"path": str(two), "name": "two", "ejected": True}, *more)
async def test_a_plug_rescan_waits_for_the_scan_under_way(tmp_path):
one, two, big = _tree(tmp_path)
idx = _Held(roots=_set(one, two), group_id=GROUP,
sk_node=Ed25519PrivateKey.generate(), gek=None)
await idx.initial_scan()
idx.eject_root("two")
idx.hold_under = big
try:
await idx.retarget(_with_two_ejected(one, two, big), wait=False)
await asyncio.wait_for(idx.at_gate.wait(), 5)
plug = asyncio.create_task(idx.plug_root("two"))
await asyncio.sleep(0.1)
assert idx.rescanned == [], "the plug walked its root beside the scan under way"
idx.gate.set()
await asyncio.wait_for(plug, 5)
assert idx.rescanned == ["two"]
assert _names(idx) == ["a.txt", "b.txt", "x.bin", "y.bin"]
finally:
idx.gate.set()
await idx.stop()
async def test_a_root_ejected_again_while_its_plug_waited_is_not_emptied(tmp_path):
one, two, big = _tree(tmp_path)
idx = _Held(roots=_set(one, two), group_id=GROUP,
sk_node=Ed25519PrivateKey.generate(), gek=None)
await idx.initial_scan()
idx.eject_root("two")
idx.hold_under = big
try:
await idx.retarget(_with_two_ejected(one, two, big), wait=False)
await asyncio.wait_for(idx.at_gate.wait(), 5)
plug = asyncio.create_task(idx.plug_root("two"))
await asyncio.sleep(0.1)
idx.eject_root("two")
idx.gate.set()
await asyncio.wait_for(plug, 5)
assert idx.rescanned == []
assert "b.txt" in _names(idx), "a root ejected while it waited was emptied"
finally:
idx.gate.set()
await idx.stop()
|