diff --git a/tests/test_apk_search.py b/tests/test_apk_search.py index 2629d90..24ab71e 100644 --- a/tests/test_apk_search.py +++ b/tests/test_apk_search.py @@ -84,11 +84,35 @@ class SearchTests(SettingsTest): self.assertLess(time.monotonic() - started, .3) self.assertTrue(result['apps']) self.assertEqual([s['status'] for s in result['sources']], ['ok', 'loading', 'error']) - search.search('other', timeout=.03) + self.assertEqual(search.search('other', timeout=.03)['sources'][1]['status'], 'loading') + search.search('newest', timeout=.03) self.assertEqual(calls, ['']) + release.set() + result = search.search('newest', timeout=2) + self.assertEqual(result['sources'][1]['status'], 'ok') + self.assertEqual(calls, ['', 'newest']) # 'other' was superseded, never run finally: release.set() + def test_set_enabled_does_not_hold_search_lock_in_source(self): + free = [] + def set_enabled(source_id, enabled): + t = threading.Thread(target=lambda: free.append(search._lock.acquire(timeout=1) and not search._lock.release())) + t.start() + t.join() + mod = fake() + mod.set_enabled = set_enabled + with patch.object(search, 'modules', return_value=([mod], [])): + search.set_enabled('one', False) + self.assertEqual(free, [True]) + + def test_limited_source_status(self): + from apk_sources import SourceLimited + mods = [fake('busy', Mock(side_effect=SourceLimited('busy is limiting requests', 60)))] + with patch.object(search, 'modules', return_value=(mods, [])): + status = search.search(timeout=1)['sources'][0] + self.assertEqual((status['status'], status['error']), ('limited', 'busy is limiting requests')) + def test_disable_persists_and_prevents_queries_and_installs(self): mod = fake() with patch.object(search, 'modules', return_value=([mod], [])): diff --git a/ui/apk_sources/__init__.py b/ui/apk_sources/__init__.py index 88df1f4..f7e6b05 100644 --- a/ui/apk_sources/__init__.py +++ b/ui/apk_sources/__init__.py @@ -43,3 +43,11 @@ Rules class SourceError(Exception): """User-readable failure from a source (network, format, verification).""" + + +class SourceLimited(SourceError): + """The source's host asked us to slow down; retry_after is in seconds.""" + + def __init__(self, message, retry_after=None): + super().__init__(message) + self.retry_after = retry_after diff --git a/ui/apk_sources/search.py b/ui/apk_sources/search.py index d0b8962..ed9d622 100644 --- a/ui/apk_sources/search.py +++ b/ui/apk_sources/search.py @@ -15,10 +15,11 @@ import time import unicodedata import frame_host -from apk_sources import SourceError +from apk_sources import SourceError, SourceLimited _lock = threading.RLock() _running = {} +_pending = {} # source id -> newest query waiting for the running one _status = {} TIMEOUT = 12 @@ -84,9 +85,9 @@ def resolve(source_id): def set_enabled(source_id, enabled): module, source = resolve(source_id) + if hasattr(module, 'set_enabled'): + module.set_enabled(source_id, enabled) # never call into a source while holding _lock with _lock: - if hasattr(module, 'set_enabled'): - module.set_enabled(source_id, enabled) values = overrides() values[source_id] = enabled path = settings_path() @@ -208,23 +209,46 @@ def group(entries, query='', vr=None, installable=False): def _launch(module, source, query, limit): + """One search per source at a time; the newest different query runs next.""" key = source['id'] + task = {'event': threading.Event(), 'query': (query, limit), 'started': time.monotonic(), + 'job': (module, source)} with _lock: old = _running.get(key) if old and not old['event'].is_set(): - return old if old['query'] == query else None - task = {'event': threading.Event(), 'query': query, 'started': time.monotonic()} + if old['query'] == task['query']: + return old + queued = _pending.get(key) + if queued and queued['query'] == task['query']: + return queued + _pending[key] = task + return task _running[key] = task + _start(key, task) + return task + + +def _start(key, task): + module, source = task['job'] + query, limit = task['query'] + def run(): + queued = None try: task['entries'] = [dict(e, source=key, source_name=source['name'], trust=source.get('trust')) for e in module.search(source, query, limit=limit) if e.get('free') is True] except Exception as e: task['error'] = str(e) + task['limited'] = isinstance(e, SourceLimited) finally: + with _lock: + queued = _pending.pop(key, None) + if queued: + _running[key] = queued task['event'].set() + if queued: + _start(key, queued) threading.Thread(target=run, daemon=True).start() - return task def search(query='', vr=None, source=None, installable=False, timeout=TIMEOUT, limit=50): @@ -238,11 +262,11 @@ def search(query='', vr=None, source=None, installable=False, timeout=TIMEOUT, l entries, statuses = [], list(errors) for s, task in tasks: status = {'id': s['id'], 'name': s['name']} - if task is None or not task['event'].wait(max(0, task['started'] + timeout - time.monotonic())): + if not task['event'].wait(max(0, task['started'] + timeout - time.monotonic())): # Still working (e.g. first download of a large index); it keeps going and fills the cache. status.update(status='loading') elif 'error' in task: - status.update(status='error', error=task['error']) + status.update(status='limited' if task.get('limited') else 'error', error=task['error']) else: status.update(status='ok') entries.extend(task['entries']) @@ -257,7 +281,7 @@ def warm(): items, _ = registry() for m, s in items: if s['enabled'] and not s.get('page_only'): - _launch(m, s, '', 1) + _launch(m, s, '', 50) # same as the first browse, so that search reuses it def install(source_id, entry_id, version_code=None, progress=None):