mirror of
https://github.com/tahoe-lafs/tahoe-lafs.git
synced 2025-01-27 14:50:03 +00:00
Simplify _scan_remote_* and remove Downloader._download_scan_batch attribute.
Signed-off-by: Daira Hopwood <daira@jacaranda.org>
This commit is contained in:
parent
f14d24b140
commit
5938535b87
@ -515,7 +515,6 @@ class Downloader(QueueMixin, WriteFileMixin):
|
|||||||
self._upload_readonly_dircap = upload_readonly_dircap
|
self._upload_readonly_dircap = upload_readonly_dircap
|
||||||
|
|
||||||
self._turn_delay = self.REMOTE_SCAN_INTERVAL
|
self._turn_delay = self.REMOTE_SCAN_INTERVAL
|
||||||
self._download_scan_batch = {} # path -> [(filenode, metadata)]
|
|
||||||
|
|
||||||
def start_scanning(self):
|
def start_scanning(self):
|
||||||
self._log("start_scanning")
|
self._log("start_scanning")
|
||||||
@ -589,14 +588,8 @@ class Downloader(QueueMixin, WriteFileMixin):
|
|||||||
collective_dirmap_d.addCallback(highest_version)
|
collective_dirmap_d.addCallback(highest_version)
|
||||||
return collective_dirmap_d
|
return collective_dirmap_d
|
||||||
|
|
||||||
def _append_to_batch(self, name, file_node, metadata):
|
def _scan_remote_dmd(self, nickname, dirnode, scan_batch):
|
||||||
if self._download_scan_batch.has_key(name):
|
self._log("_scan_remote_dmd nickname %r" % (nickname,))
|
||||||
self._download_scan_batch[name] += [(file_node, metadata)]
|
|
||||||
else:
|
|
||||||
self._download_scan_batch[name] = [(file_node, metadata)]
|
|
||||||
|
|
||||||
def _scan_remote(self, nickname, dirnode):
|
|
||||||
self._log("_scan_remote nickname %r" % (nickname,))
|
|
||||||
d = dirnode.list()
|
d = dirnode.list()
|
||||||
def scan_listing(listing_map):
|
def scan_listing(listing_map):
|
||||||
for encoded_relpath_u in listing_map.keys():
|
for encoded_relpath_u in listing_map.keys():
|
||||||
@ -607,16 +600,21 @@ class Downloader(QueueMixin, WriteFileMixin):
|
|||||||
local_version = self._get_local_latest(relpath_u)
|
local_version = self._get_local_latest(relpath_u)
|
||||||
remote_version = metadata.get('version', None)
|
remote_version = metadata.get('version', None)
|
||||||
self._log("%r has local version %r, remote version %r" % (relpath_u, local_version, remote_version))
|
self._log("%r has local version %r, remote version %r" % (relpath_u, local_version, remote_version))
|
||||||
|
|
||||||
if local_version is None or remote_version is None or local_version < remote_version:
|
if local_version is None or remote_version is None or local_version < remote_version:
|
||||||
self._log("%r added to download queue" % (relpath_u,))
|
self._log("%r added to download queue" % (relpath_u,))
|
||||||
self._append_to_batch(relpath_u, file_node, metadata)
|
if scan_batch.has_key(relpath_u):
|
||||||
|
scan_batch[relpath_u] += [(file_node, metadata)]
|
||||||
|
else:
|
||||||
|
scan_batch[relpath_u] = [(file_node, metadata)]
|
||||||
|
|
||||||
d.addCallback(scan_listing)
|
d.addCallback(scan_listing)
|
||||||
d.addBoth(self._logcb, "end of _scan_remote")
|
d.addBoth(self._logcb, "end of _scan_remote_dmd")
|
||||||
return d
|
return d
|
||||||
|
|
||||||
def _scan_remote_collective(self):
|
def _scan_remote_collective(self):
|
||||||
self._log("_scan_remote_collective")
|
self._log("_scan_remote_collective")
|
||||||
self._download_scan_batch = {} # XXX
|
scan_batch = {} # path -> [(filenode, metadata)]
|
||||||
|
|
||||||
d = self._collective_dirnode.list()
|
d = self._collective_dirnode.list()
|
||||||
def scan_collective(dirmap):
|
def scan_collective(dirmap):
|
||||||
@ -624,38 +622,31 @@ class Downloader(QueueMixin, WriteFileMixin):
|
|||||||
for dir_name in dirmap:
|
for dir_name in dirmap:
|
||||||
(dirnode, metadata) = dirmap[dir_name]
|
(dirnode, metadata) = dirmap[dir_name]
|
||||||
if dirnode.get_readonly_uri() != self._upload_readonly_dircap:
|
if dirnode.get_readonly_uri() != self._upload_readonly_dircap:
|
||||||
d2.addCallback(lambda ign, dir_name=dir_name: self._scan_remote(dir_name, dirnode))
|
d2.addCallback(lambda ign, dir_name=dir_name:
|
||||||
|
self._scan_remote_dmd(dir_name, dirnode, scan_batch))
|
||||||
def _err(f):
|
def _err(f):
|
||||||
self._log("failed to scan DMD for client %r: %s" % (dir_name, f))
|
self._log("failed to scan DMD for client %r: %s" % (dir_name, f))
|
||||||
# XXX what should we do to make this failure more visible to users?
|
# XXX what should we do to make this failure more visible to users?
|
||||||
d2.addErrback(_err)
|
d2.addErrback(_err)
|
||||||
|
|
||||||
return d2
|
return d2
|
||||||
d.addCallback(scan_collective)
|
d.addCallback(scan_collective)
|
||||||
d.addCallback(self._filter_scan_batch)
|
|
||||||
d.addCallback(self._add_batch_to_download_queue)
|
|
||||||
return d
|
|
||||||
|
|
||||||
def _add_batch_to_download_queue(self, result):
|
def _filter_batch_to_deque(ign):
|
||||||
self._log("result = %r" % (result,))
|
self._log("deque = %r, scan_batch = %r" % (self._deque, scan_batch))
|
||||||
self._log("deque = %r" % (self._deque,))
|
for relpath_u in scan_batch.keys():
|
||||||
self._deque.extend(result)
|
file_node, metadata = max(scan_batch[relpath_u], key=lambda x: x[1]['version'])
|
||||||
self._log("deque after = %r" % (self._deque,))
|
|
||||||
self._count('objects_queued', len(result))
|
|
||||||
|
|
||||||
def _filter_scan_batch(self, result):
|
|
||||||
self._log("_filter_scan_batch")
|
|
||||||
extension = [] # consider whether this should be a dict
|
|
||||||
for relpath_u in self._download_scan_batch.keys():
|
|
||||||
if relpath_u in self._pending:
|
|
||||||
continue
|
|
||||||
file_node, metadata = max(self._download_scan_batch[relpath_u], key=lambda x: x[1]['version'])
|
|
||||||
if self._should_download(relpath_u, metadata['version']):
|
if self._should_download(relpath_u, metadata['version']):
|
||||||
extension += [(relpath_u, file_node, metadata)]
|
self._deque.append( (relpath_u, file_node, metadata) )
|
||||||
else:
|
else:
|
||||||
self._log("Excluding %r" % (relpath_u,))
|
self._log("Excluding %r" % (relpath_u,))
|
||||||
self._count('objects_excluded')
|
self._count('objects_excluded')
|
||||||
self._call_hook(None, 'processed')
|
self._call_hook(None, 'processed')
|
||||||
return extension
|
|
||||||
|
self._log("deque after = %r" % (self._deque,))
|
||||||
|
d.addCallback(_filter_batch_to_deque)
|
||||||
|
return d
|
||||||
|
|
||||||
def _when_queue_is_empty(self):
|
def _when_queue_is_empty(self):
|
||||||
d = task.deferLater(self._clock, self._turn_delay, self._scan_remote_collective)
|
d = task.deferLater(self._clock, self._turn_delay, self._scan_remote_collective)
|
||||||
|
Loading…
x
Reference in New Issue
Block a user