mirror of
https://github.com/saphid/frame-control.git
synced 2026-10-06 03:00:18 +02:00
Store: close second-round races in F-Droid loads, APK reuse and pruning
- The locked publication step also refuses a v1 index once v2 was accepted, so an overlapping v1 fallback can't replace a v2 cache at an equal timestamp. - A cached APK is touched before hashing; if it vanishes, it's downloaded again. - Only the app prunes (at start and after store downloads), since claims are in-process; the CLIs never prune. - The CLI joins background refreshes on error exits too. - The Windows lock loop retries only contention errors. The concurrent-publication test now uses real flock contention. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
This commit is contained in:
1 parent
8315c7c3aa
commit
a97e8a6183
5 files changed
+141
-21
No files matched your search
@@ -71,6 +71,8 @@ Authenticated reduced indexes and APKs live under
|
|||||||
`frame_host.cache_dir('apk-sources')`; indexes refresh after 24 hours. An
|
`frame_host.cache_dir('apk-sources')`; indexes refresh after 24 hours. An
|
||||||
expired index is still served (marked stale in the store) while it refreshes in
|
expired index is still served (marked stale in the store) while it refreshes in
|
||||||
the background; a failed refresh is retried after 10 minutes.
|
the background; a failed refresh is retried after 10 minutes.
|
||||||
|
Only the running Frame Control app prunes cached APKs (at start and after store
|
||||||
|
downloads); the command-line tools never do.
|
||||||
The existing catalogue's unverified index cache is never treated as authenticated.
|
The existing catalogue's unverified index cache is never treated as authenticated.
|
||||||
|
|
||||||
Rollback protection: each repository's newest accepted signed index timestamp
|
Rollback protection: each repository's newest accepted signed index timestamp
|
||||||
|
|||||||
+112
-13
@@ -318,28 +318,109 @@ class Repositories(unittest.TestCase):
|
|||||||
self.assertEqual(fdroid.add_repo(URL)['fingerprint'], pin)
|
self.assertEqual(fdroid.add_repo(URL)['fingerprint'], pin)
|
||||||
|
|
||||||
def test_concurrent_processes_cannot_publish_an_older_index_last(self):
|
def test_concurrent_processes_cannot_publish_an_older_index_last(self):
|
||||||
|
# Threads with their own in-memory locks, as separate processes would have; each
|
||||||
|
# _state_file_lock() opens its own file description, so flock contends for real.
|
||||||
import threading
|
import threading
|
||||||
self.files['entry.jar'], _ = entry_jar(50)
|
self.files['entry.jar'], _ = entry_jar(50)
|
||||||
source = fdroid.add_repo(URL)
|
source = fdroid.add_repo(URL)
|
||||||
self.files['entry.jar'], _ = entry_jar(100)
|
jars = {'older': entry_jar(100)[0], 'newer': entry_jar(200)[0]}
|
||||||
newer, _ = entry_jar(200)
|
|
||||||
inner = []
|
|
||||||
def fetch(url, path, maximum):
|
def fetch(url, path, maximum):
|
||||||
if url.endswith('index-v2.json') and not inner:
|
name = threading.current_thread().name
|
||||||
# Another process (its own in-memory locks) loads a newer index after our early check.
|
if name in jars and url.endswith('entry.jar'):
|
||||||
inner.append(1)
|
Path(path).write_bytes(jars[name])
|
||||||
self.files['entry.jar'] = newer
|
else:
|
||||||
self.assertEqual(len(fdroid._load(source, force=True)[0]), 1)
|
|
||||||
self.fetch(url, path, maximum)
|
self.fetch(url, path, maximum)
|
||||||
self.fetch_mock.side_effect = fetch
|
self.fetch_mock.side_effect = fetch
|
||||||
writes = []
|
inside, go, order = threading.Event(), threading.Event(), []
|
||||||
real_write = fdroid._write
|
real_write = fdroid._write
|
||||||
with patch.object(fdroid, '_source_lock', lambda source_id: threading.Lock()), \
|
def write(path, value):
|
||||||
patch.object(fdroid, '_write', lambda path, value: (writes.append(path.name), real_write(path, value))):
|
if path.name.endswith('.json') and 'apps' in value and threading.current_thread().name == 'older':
|
||||||
with self.assertRaisesRegex(SourceError, 'older'):
|
inside.set() # the older load has passed its locked recheck; hold it there
|
||||||
|
go.wait(5)
|
||||||
|
order.append(threading.current_thread().name)
|
||||||
|
real_write(path, value)
|
||||||
|
errors = {}
|
||||||
|
def load():
|
||||||
|
try:
|
||||||
fdroid._load(source, force=True)
|
fdroid._load(source, force=True)
|
||||||
|
except SourceError as e:
|
||||||
|
errors[threading.current_thread().name] = str(e)
|
||||||
|
with patch.object(fdroid, '_source_lock', lambda source_id: threading.Lock()), \
|
||||||
|
patch.object(fdroid, '_write', write):
|
||||||
|
older = threading.Thread(target=load, name='older')
|
||||||
|
older.start()
|
||||||
|
self.assertTrue(inside.wait(5))
|
||||||
|
newer = threading.Thread(target=load, name='newer')
|
||||||
|
newer.start()
|
||||||
|
newer.join(.5)
|
||||||
|
self.assertTrue(newer.is_alive()) # blocked on the file lock, not publishing
|
||||||
|
go.set()
|
||||||
|
older.join(5)
|
||||||
|
newer.join(5)
|
||||||
|
self.assertEqual(errors, {})
|
||||||
|
self.assertEqual(order, ['older', 'older', 'newer', 'newer']) # cache+state, one load at a time
|
||||||
self.assertEqual(fdroid._state(source)['timestamp'], 200)
|
self.assertEqual(fdroid._state(source)['timestamp'], 200)
|
||||||
self.assertEqual(writes, [source['id'] + '.json', 'apk-repo-state.json']) # only the newer index published
|
|
||||||
|
def test_overlapping_v1_load_cannot_replace_accepted_v2(self):
|
||||||
|
import threading
|
||||||
|
v1 = json.loads(zipfile.ZipFile(FIXTURES / 'index-v1.jar').read('index-v1.json'))
|
||||||
|
v1['repo'] = {'timestamp': 100}
|
||||||
|
self.files['index-v1.jar'], pin = signed_jar('index-v1.json', json.dumps(v1).encode())
|
||||||
|
self.files['entry.jar'], _ = entry_jar(100) # the same timestamp as the v1 index
|
||||||
|
source = dict(id='overlap', name='Overlap', url=URL, fingerprint=None)
|
||||||
|
self.v1 = True
|
||||||
|
inner = []
|
||||||
|
def fetch(url, path, maximum):
|
||||||
|
if url.endswith('index-v1.jar') and not inner:
|
||||||
|
inner.append(1) # the v1 load passed its fallback check; a v2 load finishes now
|
||||||
|
self.v1 = False
|
||||||
|
t = threading.Thread(target=lambda: inner.append(fdroid._load(source, force=True)))
|
||||||
|
with patch.object(fdroid, '_source_lock', lambda source_id: threading.Lock()):
|
||||||
|
t.start()
|
||||||
|
t.join(5)
|
||||||
|
self.v1 = True
|
||||||
|
self.fetch(url, path, maximum)
|
||||||
|
self.fetch_mock.side_effect = fetch
|
||||||
|
with self.assertRaisesRegex(SourceError, 'v2'):
|
||||||
|
fdroid._load(source, force=True)
|
||||||
|
self.assertEqual(inner[1][1], pin)
|
||||||
|
self.assertTrue(fdroid._state(source)['v2'])
|
||||||
|
self.v1 = False
|
||||||
|
count = self.fetch_mock.call_count
|
||||||
|
apps, _ = fdroid._load(dict(source, fingerprint=pin))
|
||||||
|
self.assertEqual((apps['org.example.app']['version_code'], self.fetch_mock.call_count), (2, count)) # v2 cache stayed
|
||||||
|
|
||||||
|
def test_apk_removed_during_cache_check_is_downloaded_again(self):
|
||||||
|
source = self.add()
|
||||||
|
first = fdroid.download(source, 'org.example.app', 1)
|
||||||
|
count = self.fetch_mock.call_count
|
||||||
|
real = fdroid._sha256
|
||||||
|
def removed_first(path):
|
||||||
|
if str(path) == first['apk'] and not hashed:
|
||||||
|
hashed.append(1)
|
||||||
|
os.remove(first['apk']) # deleted between the existence check and the open
|
||||||
|
return real(path)
|
||||||
|
hashed = []
|
||||||
|
with patch.object(fdroid, '_sha256', removed_first):
|
||||||
|
again = fdroid.download(source, 'org.example.app', 1)
|
||||||
|
self.assertEqual(self.fetch_mock.call_count, count + 1)
|
||||||
|
self.assertEqual(Path(again['apk']).read_bytes(), (FIXTURES / 'example.apk').read_bytes())
|
||||||
|
|
||||||
|
def test_cached_apk_is_touched_before_hashing(self):
|
||||||
|
source = self.add()
|
||||||
|
first = fdroid.download(source, 'org.example.app', 1)
|
||||||
|
os.utime(first['apk'], (1, 1))
|
||||||
|
count = self.fetch_mock.call_count
|
||||||
|
real = fdroid._sha256
|
||||||
|
def prune_first(path):
|
||||||
|
if str(path) == first['apk']:
|
||||||
|
with patch.object(_web, 'APK_CAP', 0):
|
||||||
|
_web.prune() # a pruner running now sees a just-used APK
|
||||||
|
return real(path)
|
||||||
|
with patch.object(fdroid, '_sha256', prune_first):
|
||||||
|
fdroid.download(source, 'org.example.app', 1)
|
||||||
|
self.assertEqual(self.fetch_mock.call_count, count)
|
||||||
|
self.assertTrue(Path(first['apk']).exists())
|
||||||
|
|
||||||
def test_cli_search_waits_for_background_refresh(self):
|
def test_cli_search_waits_for_background_refresh(self):
|
||||||
source = self.add()
|
source = self.add()
|
||||||
@@ -354,6 +435,24 @@ class Repositories(unittest.TestCase):
|
|||||||
self.assertFalse(fdroid.stale(source))
|
self.assertFalse(fdroid.stale(source))
|
||||||
self.assertNotIn(source['id'], fdroid._refreshing)
|
self.assertNotIn(source['id'], fdroid._refreshing)
|
||||||
|
|
||||||
|
def test_cli_error_still_waits_for_background_refresh(self):
|
||||||
|
import threading
|
||||||
|
source = self.add()
|
||||||
|
self.expire(source)
|
||||||
|
release = threading.Event()
|
||||||
|
def slow(url, path, maximum):
|
||||||
|
release.wait(5)
|
||||||
|
self.fetch(url, path, maximum)
|
||||||
|
self.fetch_mock.side_effect = slow
|
||||||
|
threading.Timer(.3, release.set).start()
|
||||||
|
with patch.object(sys, 'argv', ['fdroid.py', 'download', source['id'], 'org.missing']), \
|
||||||
|
patch('sys.stderr', io.StringIO()) as err, self.assertRaises(SystemExit):
|
||||||
|
fdroid.main()
|
||||||
|
self.assertIn('no Lepton-compatible version', err.getvalue())
|
||||||
|
self.assertTrue(release.is_set())
|
||||||
|
self.assertEqual(fdroid._refreshing, {}) # joined before exiting
|
||||||
|
self.assertFalse(fdroid.stale(source))
|
||||||
|
|
||||||
def test_no_v1_fallback_once_v2_accepted(self):
|
def test_no_v1_fallback_once_v2_accepted(self):
|
||||||
source = self.add()
|
source = self.add()
|
||||||
self.v1 = True
|
self.v1 = True
|
||||||
|
|||||||
@@ -85,7 +85,11 @@ def cache():
|
|||||||
|
|
||||||
|
|
||||||
def prune():
|
def prune():
|
||||||
"""Trim the download caches: APKs to APK_CAP by mtime, orphaned .part/temp files, old listings."""
|
"""Trim the download caches: APKs to APK_CAP by mtime, orphaned .part/temp files, old listings.
|
||||||
|
|
||||||
|
Only the long-running app prunes (at start and after store downloads): claim() is
|
||||||
|
in-process, so the CLIs never prune and so can't delete an APK the app is installing.
|
||||||
|
"""
|
||||||
now, apks = time.time(), []
|
now, apks = time.time(), []
|
||||||
for folder in (str(frame_host.cache_dir('apk-sources')), cache()):
|
for folder in (str(frame_host.cache_dir('apk-sources')), cache()):
|
||||||
try:
|
try:
|
||||||
@@ -220,7 +224,6 @@ def apk(url, hosts, digest=None, name=None):
|
|||||||
raise SourceError('Download is not an APK')
|
raise SourceError('Download is not an APK')
|
||||||
path = os.path.join(cache(), actual + '.apk')
|
path = os.path.join(cache(), actual + '.apk')
|
||||||
os.replace(tmp, path)
|
os.replace(tmp, path)
|
||||||
prune()
|
|
||||||
return {'apk': path, 'obb': [], 'sha256': actual, 'verified': bool(digest)}
|
return {'apk': path, 'obb': [], 'sha256': actual, 'verified': bool(digest)}
|
||||||
except urllib.error.HTTPError as e:
|
except urllib.error.HTTPError as e:
|
||||||
if e.code in (403, 429):
|
if e.code in (403, 429):
|
||||||
|
|||||||
@@ -2,6 +2,7 @@
|
|||||||
import argparse
|
import argparse
|
||||||
import base64
|
import base64
|
||||||
import contextlib
|
import contextlib
|
||||||
|
import errno
|
||||||
import hashlib
|
import hashlib
|
||||||
from html.parser import HTMLParser
|
from html.parser import HTMLParser
|
||||||
import json
|
import json
|
||||||
@@ -295,8 +296,9 @@ def _state_file_lock():
|
|||||||
f.seek(0)
|
f.seek(0)
|
||||||
msvcrt.locking(f.fileno(), msvcrt.LK_LOCK, 1)
|
msvcrt.locking(f.fileno(), msvcrt.LK_LOCK, 1)
|
||||||
break
|
break
|
||||||
except OSError:
|
except OSError as e: # LK_LOCK gives up after ~10 s of contention; keep waiting
|
||||||
pass # LK_LOCK gives up after ~10 s; keep waiting
|
if e.errno not in (errno.EACCES, errno.EDEADLK):
|
||||||
|
raise
|
||||||
try:
|
try:
|
||||||
yield
|
yield
|
||||||
finally:
|
finally:
|
||||||
@@ -555,6 +557,8 @@ def _load(source, force=False):
|
|||||||
apps = _reduce(raw, source)
|
apps = _reduce(raw, source)
|
||||||
with _state_file_lock(): # recheck: another process may have accepted a newer index meanwhile
|
with _state_file_lock(): # recheck: another process may have accepted a newer index meanwhile
|
||||||
_check_timestamp(source, timestamp)
|
_check_timestamp(source, timestamp)
|
||||||
|
if not v2 and _state(source).get('v2'): # a v2 index was accepted while we fetched v1
|
||||||
|
raise SourceError('repository now has a signed v2 index; refusing the older v1 index')
|
||||||
_write(cache, {'version': CACHE_VERSION, 'url': source['url'], 'fingerprint': pin, 'apps': apps})
|
_write(cache, {'version': CACHE_VERSION, 'url': source['url'], 'fingerprint': pin, 'apps': apps})
|
||||||
_accept(source, timestamp, v2)
|
_accept(source, timestamp, v2)
|
||||||
_stale.discard(source['id'])
|
_stale.discard(source['id'])
|
||||||
@@ -647,7 +651,13 @@ def download(source, entry_id, version_code=None):
|
|||||||
sha = version['sha256']
|
sha = version['sha256']
|
||||||
path = frame_host.cache_dir('apk-sources', sha + '.apk')
|
path = frame_host.cache_dir('apk-sources', sha + '.apk')
|
||||||
try:
|
try:
|
||||||
if not (path.exists() and _sha256(path) == sha and _web.touch(path)): # touch: recently used, not pruned
|
reuse = False
|
||||||
|
if _web.touch(path): # before hashing: a just-used APK is never pruned
|
||||||
|
try:
|
||||||
|
reuse = _sha256(path) == sha
|
||||||
|
except FileNotFoundError: # removed anyway (e.g. by hand): download again
|
||||||
|
pass
|
||||||
|
if not reuse:
|
||||||
path.parent.mkdir(parents=True, exist_ok=True)
|
path.parent.mkdir(parents=True, exist_ok=True)
|
||||||
fd, tmp = tempfile.mkstemp(dir=str(path.parent), suffix='.part')
|
fd, tmp = tempfile.mkstemp(dir=str(path.parent), suffix='.part')
|
||||||
os.close(fd)
|
os.close(fd)
|
||||||
@@ -659,7 +669,6 @@ def download(source, entry_id, version_code=None):
|
|||||||
finally:
|
finally:
|
||||||
if os.path.exists(tmp):
|
if os.path.exists(tmp):
|
||||||
os.unlink(tmp)
|
os.unlink(tmp)
|
||||||
_web.prune()
|
|
||||||
return {'apk': str(path), 'obb': [], 'sha256': sha, 'verified': True}
|
return {'apk': str(path), 'obb': [], 'sha256': sha, 'verified': True}
|
||||||
except SourceLimited as e:
|
except SourceLimited as e:
|
||||||
raise _limited(source, e) from e
|
raise _limited(source, e) from e
|
||||||
@@ -682,6 +691,14 @@ def main():
|
|||||||
p.add_argument('source')
|
p.add_argument('source')
|
||||||
p.add_argument('query' if command == 'search' else 'package')
|
p.add_argument('query' if command == 'search' else 'package')
|
||||||
args = parser.parse_args()
|
args = parser.parse_args()
|
||||||
|
try:
|
||||||
|
_cli(parser, args)
|
||||||
|
finally:
|
||||||
|
for thread in list(_refreshing.values()): # daemon refreshes would die with the CLI, even on errors
|
||||||
|
thread.join()
|
||||||
|
|
||||||
|
|
||||||
|
def _cli(parser, args):
|
||||||
try:
|
try:
|
||||||
if args.command == 'add':
|
if args.command == 'add':
|
||||||
result = add_repo(args.url, args.fingerprint, args.name)
|
result = add_repo(args.url, args.fingerprint, args.name)
|
||||||
@@ -694,8 +711,6 @@ def main():
|
|||||||
if not source:
|
if not source:
|
||||||
raise SourceError('unknown repository id; use list')
|
raise SourceError('unknown repository id; use list')
|
||||||
result = search(source, args.query) if args.command == 'search' else download(source, args.package)
|
result = search(source, args.query) if args.command == 'search' else download(source, args.package)
|
||||||
for thread in list(_refreshing.values()): # daemon refreshes would die with the CLI
|
|
||||||
thread.join()
|
|
||||||
print(json.dumps(result, indent=2))
|
print(json.dumps(result, indent=2))
|
||||||
except SourceError as e:
|
except SourceError as e:
|
||||||
parser.exit(1, 'error: ' + str(e) + '\n')
|
parser.exit(1, 'error: ' + str(e) + '\n')
|
||||||
|
|||||||
@@ -324,6 +324,7 @@ def install(source_id, entry_id, version_code=None, progress=None):
|
|||||||
except OSError as e:
|
except OSError as e:
|
||||||
raise SourceError('The downloaded APK disappeared before installing; try again') from e
|
raise SourceError('The downloaded APK disappeared before installing; try again') from e
|
||||||
try:
|
try:
|
||||||
|
_web.prune() # the app is the only pruner (see _web.prune)
|
||||||
result = frame_android.install(downloaded['apk'], **kwargs)
|
result = frame_android.install(downloaded['apk'], **kwargs)
|
||||||
finally:
|
finally:
|
||||||
_web.release(downloaded['apk'])
|
_web.release(downloaded['apk'])
|
||||||
|
|||||||
Reference in new issue
Block a user