mirror of
https://github.com/saphid/frame-control.git
synced 2026-10-06 04:04:21 +02:00
Store: newest query runs after a source's current search; no source calls under the search lock
A search for a different query while a source is busy now queues (newest wins) instead of being dropped, and warm() uses the browse limit so the first browse reuses it. set_enabled calls the source before taking search._lock. A SourceLimited error reports the source as 'limited'. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
This commit is contained in:
1 parent
ea73145fbc
commit
3149a6e269
3 files changed
+66
-10
No files matched your search
@@ -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], [])):
|
||||
|
||||
@@ -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
|
||||
@@ -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):
|
||||
|
||||
Reference in new issue
Block a user