[job] inline 'dispatch' loop
This commit is contained in:
@@ -145,7 +145,6 @@ class Job():
|
|||||||
"""Execute or run the job"""
|
"""Execute or run the job"""
|
||||||
extractor = self.extractor
|
extractor = self.extractor
|
||||||
log = extractor.log
|
log = extractor.log
|
||||||
msg = None
|
|
||||||
|
|
||||||
self._init()
|
self._init()
|
||||||
|
|
||||||
@@ -156,8 +155,7 @@ class Job():
|
|||||||
extractor.sleep(sleep(), "extractor")
|
extractor.sleep(sleep(), "extractor")
|
||||||
|
|
||||||
try:
|
try:
|
||||||
for msg in extractor:
|
msg = self.dispatch(extractor)
|
||||||
self.dispatch(msg)
|
|
||||||
except exception.StopExtraction as exc:
|
except exception.StopExtraction as exc:
|
||||||
if exc.depth > 1 and exc.target != extractor.__class__.subcategory:
|
if exc.depth > 1 and exc.target != extractor.__class__.subcategory:
|
||||||
exc.depth -= 1
|
exc.depth -= 1
|
||||||
@@ -203,31 +201,36 @@ class Job():
|
|||||||
self.status |= s
|
self.status |= s
|
||||||
return self.status
|
return self.status
|
||||||
|
|
||||||
def dispatch(self, msg):
|
def dispatch(self, messages):
|
||||||
"""Call the appropriate message handler"""
|
"""Call the appropriate message handler"""
|
||||||
if msg[0] == Message.Url:
|
msg = None
|
||||||
_, url, kwdict = msg
|
|
||||||
if self.metadata_url:
|
|
||||||
kwdict[self.metadata_url] = url
|
|
||||||
if self.pred_url(url, kwdict):
|
|
||||||
self.update_kwdict(kwdict)
|
|
||||||
self.handle_url(url, kwdict)
|
|
||||||
if FLAGS.FILE is not None:
|
|
||||||
FLAGS.process("FILE")
|
|
||||||
|
|
||||||
elif msg[0] == Message.Directory:
|
for msg in messages:
|
||||||
self.update_kwdict(msg[1])
|
if msg[0] == Message.Url:
|
||||||
self.handle_directory(msg[1])
|
_, url, kwdict = msg
|
||||||
|
if self.metadata_url:
|
||||||
|
kwdict[self.metadata_url] = url
|
||||||
|
if self.pred_url(url, kwdict):
|
||||||
|
self.update_kwdict(kwdict)
|
||||||
|
self.handle_url(url, kwdict)
|
||||||
|
if FLAGS.FILE is not None:
|
||||||
|
FLAGS.process("FILE")
|
||||||
|
|
||||||
elif msg[0] == Message.Queue:
|
elif msg[0] == Message.Directory:
|
||||||
_, url, kwdict = msg
|
self.update_kwdict(msg[1])
|
||||||
if self.metadata_url:
|
self.handle_directory(msg[1])
|
||||||
kwdict[self.metadata_url] = url
|
|
||||||
if self.pred_queue(url, kwdict):
|
elif msg[0] == Message.Queue:
|
||||||
self.update_kwdict(kwdict)
|
_, url, kwdict = msg
|
||||||
self.handle_queue(url, kwdict)
|
if self.metadata_url:
|
||||||
if FLAGS.CHILD is not None:
|
kwdict[self.metadata_url] = url
|
||||||
FLAGS.process("CHILD")
|
if self.pred_queue(url, kwdict):
|
||||||
|
self.update_kwdict(kwdict)
|
||||||
|
self.handle_queue(url, kwdict)
|
||||||
|
if FLAGS.CHILD is not None:
|
||||||
|
FLAGS.process("CHILD")
|
||||||
|
|
||||||
|
return msg
|
||||||
|
|
||||||
def handle_url(self, url, kwdict):
|
def handle_url(self, url, kwdict):
|
||||||
"""Handle Message.Url"""
|
"""Handle Message.Url"""
|
||||||
@@ -948,8 +951,7 @@ class DataJob(Job):
|
|||||||
|
|
||||||
# collect data
|
# collect data
|
||||||
try:
|
try:
|
||||||
for msg in extractor:
|
self.dispatch(extractor)
|
||||||
self.dispatch(msg)
|
|
||||||
except exception.StopExtraction:
|
except exception.StopExtraction:
|
||||||
pass
|
pass
|
||||||
except Exception as exc:
|
except Exception as exc:
|
||||||
|
|||||||
@@ -426,8 +426,7 @@ class ResultJob(job.DownloadJob):
|
|||||||
|
|
||||||
def run(self):
|
def run(self):
|
||||||
self._init()
|
self._init()
|
||||||
for msg in self.extractor:
|
self.dispatch(self.extractor)
|
||||||
self.dispatch(msg)
|
|
||||||
|
|
||||||
def handle_url(self, url, kwdict, fallback=None):
|
def handle_url(self, url, kwdict, fallback=None):
|
||||||
self._update_url(url)
|
self._update_url(url)
|
||||||
|
|||||||
Reference in New Issue
Block a user