Add serialized download queue with transfer speed and ETA
This commit is contained in:
+33
-17
@@ -70,12 +70,12 @@ def capabilities(model):
|
||||
class Catalog:
|
||||
def __init__(self,root):
|
||||
self.root=Path(root);self.lock=threading.RLock();self.job=None;self.cancel=threading.Event()
|
||||
self.cache={};self.cache_lock=threading.Lock();self.history=[]
|
||||
self.cache={};self.cache_lock=threading.Lock();self.history=[];self.pending=[];self.closing=False;self.samples=[]
|
||||
path=self.root/'downloads.json'
|
||||
if path.exists():
|
||||
self.history=json.loads(path.read_text())
|
||||
for job in self.history:
|
||||
if job['state']=='downloading':job.update(state='interrupted',error='Deck wurde neu gestartet. Datei erneut auswählen und herunterladen.')
|
||||
if job['state'] in ('downloading','queued'):job.update(state='interrupted',error='Deck wurde neu gestartet. Datei erneut auswählen und herunterladen.')
|
||||
else:
|
||||
for entry in self.root.glob('*/entry.json'):
|
||||
try:
|
||||
@@ -91,6 +91,9 @@ class Catalog:
|
||||
job=next((x for x in self.history if x['id']==job_id),None)
|
||||
if not job:raise ValueError('Download nicht gefunden.')
|
||||
if job['state']=='downloading':raise ValueError('Laufenden Download zuerst abbrechen.')
|
||||
if job['state']=='queued':
|
||||
self.pending=[x for x in self.pending if x[0]['id']!=job_id]
|
||||
job['state']='cancelled'
|
||||
job['dismissed']=True
|
||||
if self.job and self.job.get('id')==job_id:self.job=None
|
||||
self._save_history()
|
||||
@@ -161,7 +164,8 @@ class Catalog:
|
||||
def start(self,repo,filename,revision,kind):
|
||||
if kind not in KINDS:raise ValueError('Ungültiger Bereich.')
|
||||
with self.lock:
|
||||
if self.job and self.job['state']=='downloading':raise ValueError('Ein Download läuft bereits.')
|
||||
if self.closing:raise ValueError('Deck wird beendet.')
|
||||
if len(self.pending)>=20:raise ValueError('Warteschlange voll (20 Dateien).')
|
||||
with self.cache_lock:self.cache.pop(repo,None)
|
||||
data=self.files(repo)
|
||||
if data['gated']:raise ValueError('Zugangsbeschränkte Modelle werden noch nicht unterstützt.')
|
||||
@@ -169,18 +173,25 @@ class Catalog:
|
||||
item=next((x for x in data['files'] if x['name']==filename),None)
|
||||
if not item:raise ValueError('Datei nicht verfügbar.')
|
||||
self.root.mkdir(parents=True,exist_ok=True,mode=0o700)
|
||||
if item['size']>shutil.disk_usage(self.root).free-10*1024**3:raise ValueError('Nicht genug Platz mit 10 GiB freier Reserve.')
|
||||
reserved=sum(x['total']-x['bytes'] for x in self.history if x['state'] in ('downloading','queued'))
|
||||
if item['size']+reserved>shutil.disk_usage(self.root).free-10*1024**3:raise ValueError('Nicht genug Platz mit 10 GiB freier Reserve.')
|
||||
ident=hashlib.sha256((repo+revision+filename).encode()).hexdigest()
|
||||
if any(x.get('entry_id')==ident and x['state'] in ('downloading','queued') for x in self.history):raise ValueError('Datei läuft bereits oder steht in der Warteschlange.')
|
||||
target=self.root/ident
|
||||
if (target/'entry.json').exists():raise ValueError('Datei bereits in der Bibliothek.')
|
||||
target.mkdir(exist_ok=True,mode=0o700)
|
||||
self.cancel.clear()
|
||||
self.job=dict(id=uuid.uuid4().hex,entry_id=ident,state='downloading',repo=repo,file=filename,kind=kind,revision=revision,created_at=time.time(),bytes=0,total=item['size'],error=None)
|
||||
self.history.append(self.job);self._save_history()
|
||||
threading.Thread(target=self._download,args=(data,item,target,kind),daemon=True).start()
|
||||
return dict(self.job)
|
||||
def stop(self):
|
||||
self.cancel.set()
|
||||
job=dict(id=uuid.uuid4().hex,entry_id=ident,state='queued',repo=repo,file=filename,kind=kind,revision=revision,created_at=time.time(),bytes=0,total=item['size'],error=None,bytes_per_second=None,eta_seconds=None)
|
||||
self.history.append(job);self.pending.append((job,data,item,target,kind));self._next();self._save_history()
|
||||
return dict(job)
|
||||
def _next(self):
|
||||
if self.closing or (self.job and self.job['state']=='downloading') or not self.pending:return
|
||||
self.job,data,item,target,kind=self.pending.pop(0)
|
||||
self.job['state']='downloading';self.cancel.clear();self.samples=[(time.monotonic(),0)]
|
||||
threading.Thread(target=self._download,args=(data,item,target,kind),daemon=True).start()
|
||||
def stop(self,shutdown=False):
|
||||
with self.lock:
|
||||
if shutdown:self.closing=True
|
||||
self.cancel.set()
|
||||
return {'cancellation_requested':True}
|
||||
def _download(self,data,item,target,kind):
|
||||
partial=target/'download.part';dest=target/('model'+PurePosixPath(item['name']).suffix)
|
||||
@@ -195,17 +206,22 @@ class Catalog:
|
||||
received+=len(chunk)
|
||||
if received>item['size'] or shutil.disk_usage(target).free<10*1024**3:raise ValueError('Größe oder Speicherreserve überschritten.')
|
||||
out.write(chunk);digest.update(chunk)
|
||||
with self.lock:self.job['bytes']=received
|
||||
with self.lock:
|
||||
now=time.monotonic();self.samples.append((now,received))
|
||||
while len(self.samples)>2 and self.samples[1][0]<now-10:self.samples.pop(0)
|
||||
elapsed=now-self.samples[0][0];speed=(received-self.samples[0][1])/elapsed if elapsed>=.25 else None
|
||||
self.job.update(bytes=received,bytes_per_second=speed,eta_seconds=(item['size']-received)/speed if speed else None)
|
||||
if received!=item['size']:raise ValueError('Unvollständiger Download.')
|
||||
if item['sha256'] and digest.hexdigest()!=item['sha256']:raise ValueError('SHA-256-Prüfung fehlgeschlagen.')
|
||||
partial.replace(dest)
|
||||
entry=dict(repo=data['repo'],revision=data['revision'],file=item['name'],size=received,sha256=digest.hexdigest(),upstream_hash_verified=bool(item['sha256']),kind=kind,downloaded_at=time.time(),state='downloaded',runtime_ready=False)
|
||||
temp=target/'entry.tmp';temp.write_text(json.dumps(entry));temp.replace(target/'entry.json')
|
||||
with self.lock:self.job['state']='complete'
|
||||
terminal='complete';error=None
|
||||
except Exception as exc:
|
||||
with self.lock:
|
||||
self.job['state']='cancelled' if isinstance(exc,InterruptedError) else 'failed'
|
||||
self.job['error']=str(exc) if isinstance(exc,ValueError) else ('Abgebrochen.' if isinstance(exc,InterruptedError) else 'Download fehlgeschlagen; Verbindung oder Anbieter prüfen.')
|
||||
terminal='cancelled' if isinstance(exc,InterruptedError) else 'failed'
|
||||
error=str(exc) if isinstance(exc,ValueError) else ('Abgebrochen.' if isinstance(exc,InterruptedError) else 'Download fehlgeschlagen; Verbindung oder Anbieter prüfen.')
|
||||
partial.unlink(missing_ok=True)
|
||||
finally:
|
||||
with self.lock:self._save_history()
|
||||
with self.lock:
|
||||
self.job.update(state=terminal,error=error,bytes_per_second=None,eta_seconds=None)
|
||||
self._save_history();self._next()
|
||||
|
||||
Reference in New Issue
Block a user