diff --git a/examples/simple_crawler/crawler.py b/examples/simple_crawler/crawler.py index 6c53dc90..d60d0574 100644 --- a/examples/simple_crawler/crawler.py +++ b/examples/simple_crawler/crawler.py @@ -86,16 +86,16 @@ def run(self, params): collection = connections.mongodb_jobs.simple_crawler_urls - print - print "Crawl stats" - print "===========" - print "URLs queued: %s" % collection.find().count() - print "URLs successfully crawled: %s" % collection.find({"fetched_date": {"$exists": True}}).count() - print "URLs redirected: %s" % collection.find({"redirected_to": {"$exists": True}}).count() - print "Bytes fetched: %s" % (list(collection.aggregate( + print() + print( "Crawl stats") + print( "===========") + print( "URLs queued: %s" % collection.find().count()) + print( "URLs successfully crawled: %s" % collection.find({"fetched_date": {"$exists": True}}).count()) + print( "URLs redirected: %s" % collection.find({"redirected_to": {"$exists": True}}).count()) + print( "Bytes fetched: %s" % (list(collection.aggregate( {"$group": {"_id": None, "sum": {"$sum": "$html_length"}}} - )) or [{}])[0].get("sum", 0) - print + )) or [{}])[0].get("sum", 0)) + print() class Reset(Task): diff --git a/mrq/basetasks/cleaning.py b/mrq/basetasks/cleaning.py index 76919eaa..9615f612 100644 --- a/mrq/basetasks/cleaning.py +++ b/mrq/basetasks/cleaning.py @@ -1,3 +1,4 @@ +from builtins import str from mrq.queue import Queue from mrq.task import Task from mrq.job import Job diff --git a/mrq/basetasks/utils.py b/mrq/basetasks/utils.py index 5a70bcc4..9c9bbb81 100644 --- a/mrq/basetasks/utils.py +++ b/mrq/basetasks/utils.py @@ -1,3 +1,6 @@ +from __future__ import print_function +from future.utils import itervalues +from builtins import str from mrq.task import Task from mrq.queue import Queue from bson import ObjectId @@ -50,7 +53,7 @@ def build_query(self): if self.params.get("params"): params_dict = json.loads(self.params.get("params")) # pylint: disable=no-member - for key in params_dict.keys(): + for key in params_dict: query["params.%s" % key] = params_dict[key] return query @@ -76,7 +79,7 @@ def perform_action(self, action, query, destination_queue): else: tasks_defs = get_current_config().get("tasks", {}) - tasks_ttls = [cfg.get("result_ttl", 0) for cfg in tasks_defs.values()] + tasks_ttls = [cfg.get("result_ttl", 0) for cfg in itervalues(tasks_defs)] result_ttl = max([default_job_timeout] + tasks_ttls) @@ -92,7 +95,7 @@ def perform_action(self, action, query, destination_queue): # In this case we could also loose some jobs that were queued after # the MongoDB update. They will be "lost" and requeued later like the other case # after the Redis BLPOP - if query.keys() == ["queue"]: + if list(query.keys()) == ["queue"]: Queue(query["queue"]).empty() elif action in ("requeue", "requeue_retry"): @@ -135,6 +138,6 @@ def perform_action(self, action, query, destination_queue): Queue(destination_queue or queue).enqueue_job_ids( [str(x) for x in jobs_by_queue[queue]]) - print stats + print(stats) return stats diff --git a/mrq/bin/mrq_run.py b/mrq/bin/mrq_run.py index 1cda61ad..eedf6e4c 100755 --- a/mrq/bin/mrq_run.py +++ b/mrq/bin/mrq_run.py @@ -1,4 +1,6 @@ #!/usr/bin/env python +from __future__ import print_function + import os # Needed to make getaddrinfo() work in pymongo on Mac OS X @@ -40,13 +42,13 @@ def main(): # mrq-run taskpath a 1 b 2 => {"a": "1", "b": "2"} for group in utils.group_iter(cfg["taskargs"], n=2): if len(group) != 2: - print "Number of arguments wasn't even" + print("Number of arguments wasn't even") sys.exit(1) params[group[0]] = group[1] if cfg["queue"]: ret = queue_job(cfg["taskpath"], params, queue=cfg["queue"]) - print ret + print(ret) else: worker_class = load_class_by_path(cfg["worker_class"]) job = worker_class.job_class(None) @@ -58,7 +60,7 @@ def main(): job.datestarted = datetime.datetime.utcnow() set_current_job(job) ret = job.perform() - print json_stdlib.dumps(ret, cls=MongoJSONEncoder) # pylint: disable=no-member + print(json_stdlib.dumps(ret, cls=MongoJSONEncoder)) # pylint: disable=no-member # This shouldn't be needed as the process will exit and close any remaining sockets # connections.redis.connection_pool.disconnect() diff --git a/mrq/bin/mrq_worker.py b/mrq/bin/mrq_worker.py index 85ce8e29..ec58575c 100755 --- a/mrq/bin/mrq_worker.py +++ b/mrq/bin/mrq_worker.py @@ -1,5 +1,6 @@ #!/usr/bin/env python import os +from builtins import str # Needed to make getaddrinfo() work in pymongo on Mac OS X # Docs mention it's a better choice for Linux as well. @@ -8,15 +9,19 @@ os.environ["GEVENT_RESOLVER"] = "ares" from gevent import monkey -monkey.patch_all() +monkey.patch_all(subprocess=False) import sys import tempfile import signal -import subprocess32 as subprocess import psutil import argparse +try: + import subprocess32 as subprocess +except: + import subprocess + sys.path.insert(0, os.getcwd()) from mrq import config diff --git a/mrq/config.py b/mrq/config.py index bacf7dd9..45e11a01 100644 --- a/mrq/config.py +++ b/mrq/config.py @@ -1,3 +1,5 @@ +from __future__ import print_function +from builtins import str import argparse import os import sys @@ -325,7 +327,7 @@ def add_parser_args(parser, config_type): '--report_file', default="", action='store', - type=unicode, + type=str, help='Filepath of a json dump of the worker status. Disabled if none') parser.add_argument( @@ -414,7 +416,7 @@ def get_config( # line from_args = {} if "args" in sources: - for k, v in parser.parse_args().__dict__.iteritems(): + for k, v in parser.parse_args().__dict__.items(): if default_config[k] != v: from_args[k] = v @@ -434,7 +436,7 @@ def get_config( sys.path.insert(0, os.path.dirname(config_file)) config_module = __import__(os.path.basename(config_file.replace(".py", ""))) sys.path.pop(0) - for k, v in config_module.__dict__.iteritems(): + for k, v in config_module.__dict__.items(): # We only keep variables starting with an uppercase character. if k[0].isupper(): @@ -443,7 +445,7 @@ def get_config( # Merge the config in the order given by the user merged_config = default_config - config_keys = set(default_config.keys() + from_file.keys()) + config_keys = set(list(default_config.keys()) + list(from_file.keys())) for part in sources: for name in config_keys: @@ -475,11 +477,11 @@ def print_profiling(): atexit.register(print_profiling) if merged_config["version"]: - print "MRQ version: %s" % VERSION - print "Python version: %s" % sys.version + print("MRQ version: %s" % VERSION) + print("Python version: %s" % sys.version) sys.exit(1) if "no_import_patch" in from_args: - print "WARNING: --no_import_patch will be deprecated in MRQ 1.0!" + print("WARNING: --no_import_patch will be deprecated in MRQ 1.0!") return merged_config diff --git a/mrq/context.py b/mrq/context.py index 18b02f49..6474224c 100644 --- a/mrq/context.py +++ b/mrq/context.py @@ -1,7 +1,12 @@ +from future import standard_library +standard_library.install_aliases() +from builtins import next +from builtins import map +from past.builtins import basestring from .logger import Logger import gevent import gevent.pool -import urlparse +import urllib.parse import time import pymongo import traceback @@ -108,12 +113,12 @@ def versiontuple(v): return tuple(map(int, (v.split(".")))) if attr.startswith("redis"): - if type(config_obj) in [str, unicode]: + if isinstance(config_obj, basestring): import redis as pyredis - urlparse.uses_netloc.append('redis') - redis_url = urlparse.urlparse(config_obj) + urllib.parse.uses_netloc.append('redis') + redis_url = urllib.parse.urlparse(config_obj) log.info("%s: Connecting to Redis at %s..." % (attr, redis_url.hostname)) @@ -124,7 +129,8 @@ def versiontuple(v): db=int((redis_url.path or "").replace("/", "") or "0"), password=redis_url.password, max_connections=int(config.get("redis_max_connections")), - timeout=int(config.get("redis_timeout")) + timeout=int(config.get("redis_timeout")), + decode_responses=True ) return pyredis.StrictRedis(connection_pool=redis_pool) @@ -134,7 +140,7 @@ def versiontuple(v): elif attr.startswith("mongodb"): - if type(config_obj) in [str, unicode]: + if isinstance(config_obj, basestring): if attr == "mongodb_logs" and config_obj == "1": return connections.mongodb_jobs @@ -218,7 +224,7 @@ def inner_func(*args): try: ret = func(*args) - except Exception, exc: + except Exception as exc: trace = traceback.format_exc() log.error("Error in subpool: %s \n%s" % (exc, trace)) raise @@ -271,7 +277,7 @@ def inner_func(*args): try: ret = func(*args) - except Exception, exc: + except Exception as exc: trace = traceback.format_exc() log.error("Error in subpool: %s \n%s" % (exc, trace)) raise diff --git a/mrq/dashboard/app.py b/mrq/dashboard/app.py index 8e7c880d..ef8fecaf 100644 --- a/mrq/dashboard/app.py +++ b/mrq/dashboard/app.py @@ -1,3 +1,7 @@ +from __future__ import print_function +from future import standard_library +standard_library.install_aliases() +from future.utils import iteritems from gevent import monkey monkey.patch_all() @@ -41,7 +45,7 @@ @requires_auth def root(): return render_template("index.html", MRQ_CONFIG={ - k: v for k, v in cfg.items() if k in WHITELISTED_MRQ_CONFIG_KEYS + k: v for k, v in iteritems(cfg) if k in WHITELISTED_MRQ_CONFIG_KEYS }) @@ -149,10 +153,10 @@ def build_api_datatables_query(req): try: params_dict = json.loads(req.args.get("params")) - for key in params_dict.keys(): + for key in params_dict: query["params.%s" % key] = params_dict[key] except Exception as e: # pylint: disable=broad-except - print "Error will converting form JSON: %s" % e + print("Error will converting form JSON: %s" % e) return query @@ -301,7 +305,7 @@ def api_job_traceback(job_id): @app.route('/api/jobaction', methods=["POST"]) @requires_auth def api_job_action(): - params = {k: v for k, v in request.form.iteritems()} + params = {k: v for k, v in iteritems(request.form)} if params.get("status") and "-" in params.get("status"): params["status"] = params.get("status").split("-") return jsonify({"job_id": queue_job("mrq.basetasks.utils.JobAction", diff --git a/mrq/job.py b/mrq/job.py index aa1c6c1e..0de601b3 100644 --- a/mrq/job.py +++ b/mrq/job.py @@ -1,3 +1,7 @@ +from future import standard_library +standard_library.install_aliases() +from builtins import str +from builtins import object import datetime from bson import ObjectId import time @@ -10,12 +14,12 @@ from collections import defaultdict import traceback import sys -import urlparse +import urllib.parse import re import linecache import fnmatch import encodings -import copy_reg +import copyreg from . import context @@ -54,7 +58,10 @@ def __init__(self, job_id, queue=None, start=False, fetch=False): if job_id is None: self.id = None else: - self.id = ObjectId(job_id) + if isinstance(job_id, bytes): + self.id = ObjectId(job_id.decode('utf-8')) + else: + self.id = ObjectId(job_id) self.data = None self.saved = True @@ -463,10 +470,10 @@ def set_current_io(self, io_data): def trace_memory_clean_caches(self): """ Avoid polluting results with some builtin python caches """ - urlparse.clear_cache() + urllib.parse.clear_cache() re.purge() linecache.clearcache() - copy_reg.clear_extension_cache() + copyreg.clear_extension_cache() if hasattr(fnmatch, "purge"): fnmatch.purge() # pylint: disable=no-member diff --git a/mrq/logger.py b/mrq/logger.py index 7a019a85..d99511e0 100644 --- a/mrq/logger.py +++ b/mrq/logger.py @@ -1,9 +1,17 @@ +from __future__ import print_function +from builtins import object +from future.utils import iteritems from collections import defaultdict import datetime - +import sys +PY3 = sys.version_info > (3,) def _encode_if_unicode(string): + + if PY3: + return string + if isinstance(string, unicode): return string.encode("utf-8", "replace") else: @@ -11,6 +19,10 @@ def _encode_if_unicode(string): def _decode_if_str(string): + + if PY3: + return str(string) + if isinstance(string, str): return string.decode("utf-8", "replace") else: @@ -62,9 +74,9 @@ def log(self, level, *args, **kwargs): if not self.quiet: try: - print _encode_if_unicode(formatted) + print(_encode_if_unicode(formatted)) except UnicodeDecodeError: - print formatted + print(formatted) if self.collection is False: return @@ -88,10 +100,10 @@ def flush(self, w=0): inserts = [{ "worker": k, "logs": "\n".join(v) + "\n" - } for k, v in self.buffer["workers"].iteritems()] + [{ + } for k, v in iteritems(self.buffer["workers"])] + [{ "job": k, "logs": "\n".join(v) + "\n" - } for k, v in self.buffer["jobs"].iteritems()] + } for k, v in iteritems(self.buffer["jobs"])] if len(inserts) == 0: return diff --git a/mrq/monkey.py b/mrq/monkey.py index bddd8392..76b3c494 100644 --- a/mrq/monkey.py +++ b/mrq/monkey.py @@ -1,3 +1,6 @@ +from future import standard_library +standard_library.install_aliases() +from past.builtins import basestring from .context import get_current_job, get_current_worker import time import random @@ -341,7 +344,7 @@ def connect(old_method, self, *args, **kwargs): return ret - from httplib import HTTPConnection, HTTPSConnection + from http.client import HTTPConnection, HTTPSConnection patch_method(HTTPConnection, "request", request) patch_method(HTTPConnection, "connect", connect) @@ -394,7 +397,7 @@ def _Cursor__send_message(self, *args, **kwargs): collection = self._Cursor__collection.name # pylint: disable=no-member if collection == "$cmd": - items = self._Cursor__spec.items() # pylint: disable=no-member + items = list(self._Cursor__spec.items()) # pylint: disable=no-member if len(items) > 0: subtype, collection = items[0] diff --git a/mrq/queue.py b/mrq/queue.py index a0fa2608..665956cd 100644 --- a/mrq/queue.py +++ b/mrq/queue.py @@ -1,9 +1,21 @@ +from __future__ import division + +from builtins import range +from builtins import object +from past.utils import old_div from .redishelpers import redis_zaddbyscore, redis_zpopbyscore, redis_lpopsafe from .redishelpers import redis_group_command import time from bson import ObjectId from . import context from . import job as jobmodule +import binascii + +import sys +PY3 = sys.version_info > (3,) +from builtins import bytes +from future import standard_library +standard_library.install_aliases() class Queue(object): @@ -126,14 +138,14 @@ def serialize_job_ids(self, job_ids): elif isinstance(job_ids[0], ObjectId): return [x.binary for x in job_ids] else: - return [x.decode('hex') for x in job_ids] + return [bytes.fromhex(str(x)) for x in job_ids] def unserialize_job_ids(self, job_ids): """ Unserialize job_ids stored in Redis """ if len(job_ids) == 0 or self.use_large_ids: return job_ids else: - return [x.encode('hex') for x in job_ids] + return [binascii.hexlify(x).decode('ascii') for x in job_ids] def size(self): """ Returns the total number of jobs on the queue """ @@ -207,7 +219,7 @@ def get_sorted_graph( raise Exception("Not a sorted queue") with context.connections.redis.pipeline(transaction=exact) as pipe: - interval = float(stop - start) / slices + interval = old_div(float(stop - start), slices) for i in range(0, slices): pipe.zcount(self.redis_key, (start + i * interval), @@ -228,7 +240,7 @@ def all_active(cls): prefix = context.get_current_config()["redis_prefix"] queues = [] - for key in context.connections.redis.keys(): + for key in context.connections.redis: if key.startswith(prefix): queues.append(Queue(key[len(prefix) + 3:])) @@ -239,7 +251,7 @@ def all_known(cls, ): """ List all previously known queues """ # raw queues we know exist from the config + known queues in redis - return set(context.get_current_config().get("raw_queues", {}).keys() + cls.redis_known_queues().keys()) + return set(list(context.get_current_config().get("raw_queues", {}).keys()) + list(cls.redis_known_queues().keys())) @classmethod def all(cls): @@ -274,8 +286,8 @@ def enqueue_job_ids(self, job_ids): job_ids = {x: now for x in self.serialize_job_ids(job_ids)} else: - serialized_job_ids = self.serialize_job_ids(job_ids.keys()) - values = job_ids.values() + serialized_job_ids = self.serialize_job_ids(list(job_ids.keys())) + values = list(job_ids.values()) job_ids = {k: values[i] for i, k in enumerate(serialized_job_ids)} context.connections.redis.zadd(self.redis_key, **job_ids) diff --git a/mrq/redishelpers.py b/mrq/redishelpers.py index 0eea95ad..c1b4f398 100644 --- a/mrq/redishelpers.py +++ b/mrq/redishelpers.py @@ -1,3 +1,4 @@ +from builtins import range from .utils import memoize from . import context diff --git a/mrq/scheduler.py b/mrq/scheduler.py index d6982b03..9fdadd7b 100644 --- a/mrq/scheduler.py +++ b/mrq/scheduler.py @@ -1,3 +1,6 @@ +from builtins import str +from builtins import object +from future.utils import iteritems from .context import log, queue_job import datetime import ujson as json @@ -8,7 +11,7 @@ def _hash_task(task): params = task.get("params") if params: - params = json.dumps(sorted(task["params"].items(), key=lambda x: x[0])) # pylint: disable=no-member + params = json.dumps(sorted(list(task["params"].items()), key=lambda x: x[0])) # pylint: disable=no-member full = [str(task.get(x)) for x in ["path", "interval", "dailytime", "queue"]] @@ -39,7 +42,7 @@ def sync_tasks(self, tasks): self.collection.remove({"_id": task["_id"]}) log.debug("Scheduler: deleted %s" % task["hash"]) - for h, task in tasks_by_hash.iteritems(): + for h, task in iteritems(tasks_by_hash): task["hash"] = h task["datelastqueued"] = datetime.datetime.fromtimestamp(0) if task.get("dailytime"): diff --git a/mrq/task.py b/mrq/task.py index 22dd0202..b3dd7cb5 100644 --- a/mrq/task.py +++ b/mrq/task.py @@ -1,3 +1,4 @@ +from builtins import object class Task(object): diff --git a/mrq/utils.py b/mrq/utils.py index 8af6a455..032346fc 100644 --- a/mrq/utils.py +++ b/mrq/utils.py @@ -1,3 +1,8 @@ +from __future__ import division +from builtins import str +from builtins import range +from builtins import object +from past.utils import old_div import re import importlib import time @@ -32,7 +37,7 @@ def group_iter(iterator, n=2): if isinstance(iterator, list): length = len(iterator) - for i in range(int(math.ceil(float(length) / n))): + for i in range(int(math.ceil(old_div(float(length), n)))): yield iterator[i * n: (i + 1) * n] else: @@ -135,7 +140,7 @@ def wait_for_net_service(server, port, timeout=None, poll_interval=0.1): except Exception as err: # catch timeout exception from underlying network library # this one is different from socket.timeout - if not isinstance(err.args, tuple) or err[0] != errno.ETIMEDOUT: + if not isinstance(err.args, tuple) or err.args[0] != errno.ETIMEDOUT: pass # raise else: s.close() @@ -177,7 +182,9 @@ def default(self, obj): # pylint: disable=E0202 if isinstance(obj, (datetime.datetime, datetime.date)): return obj.isoformat() elif isinstance(obj, ObjectId): - return unicode(obj) + return str(obj) + elif isinstance(obj, bytes): + return obj.decode('utf-8') return json.JSONEncoder.default(self, obj) diff --git a/mrq/worker.py b/mrq/worker.py index 7fb9378d..cbf242c0 100644 --- a/mrq/worker.py +++ b/mrq/worker.py @@ -1,3 +1,8 @@ +from future import standard_library +standard_library.install_aliases() +from builtins import str +from builtins import bytes +from future.utils import iteritems import gevent import gevent.pool import os @@ -10,7 +15,7 @@ import sys import json as json_stdlib import ujson as json -import BaseHTTPServer +import http.server from bson import ObjectId from collections import defaultdict @@ -241,7 +246,6 @@ def get_worker_report(self, with_memory=False): its jobs. """ greenlets = [] - for greenlet in self.gevent_pool: g = {} short_stack = [] @@ -300,17 +304,17 @@ def get_worker_report(self, with_memory=False): io = None if self._traced_io: io = {} - for k, v in self._traced_io.items(): + for k, v in iteritems(self._traced_io): if k == "total": io[k] = v else: - io[k] = sorted(v.items(), reverse=True, key=lambda x: x[1]) + io[k] = sorted(list(v.items()), reverse=True, key=lambda x: x[1]) used_pool_slots = self.pool_size - self.gevent_pool.free_count() return { "status": self.status, - "config": {k: v for k, v in self.config.iteritems() if k in whitelisted_config}, + "config": {k: v for k, v in iteritems(self.config) if k in whitelisted_config}, "done_jobs": self.done_jobs, "pool_usage_average": self.pool_usage_average.next(used_pool_slots), "datestarted": self.datestarted, @@ -343,7 +347,7 @@ def report_worker(self, w=0): if self.config["report_file"]: with open(self.config["report_file"], "wb") as f: - f.write(json.dumps(report, ensure_ascii=False)) # pylint: disable=no-member + f.write(bytes(json.dumps(report, ensure_ascii=False), 'utf-8')) # pylint: disable=no-member if "_id" in report: del report["_id"] @@ -378,7 +382,7 @@ def admin_routes(env, start_response): res = "" if path in ["/", "/report", "/report_mem"]: report = self.get_worker_report(with_memory=(path == "/report_mem")) - res = json_stdlib.dumps(report, cls=MongoJSONEncoder) + res = bytes(json_stdlib.dumps(report, cls=MongoJSONEncoder), 'utf-8') elif path == "/wait_for_idle": self.idle_wait_count = 0 self.idle_event.clear() diff --git a/requirements-base.txt b/requirements-base.txt index b7b679b0..efd4fd76 100644 --- a/requirements-base.txt +++ b/requirements-base.txt @@ -1,14 +1,14 @@ argparse>=1.1 -redis>=2.10.3 +redis>=2.10.5 pymongo>=3.0.1 -gevent>=1.1rc3 +gevent==1.1.1 ujson==1.33 hiredis>=0.1.5 psutil==1.2.1 objgraph==1.8.1 termcolor==1.1.0 - subprocess32==3.2.7; python_version < '3.2' - supervisor==3.0; python_version < '3.0' -git+git://github.com/Supervisor/supervisor.git@c18aecf1641d8953767e7010be8bae1924a133bf#egg=Supervisor; python_version >= '3.0' \ No newline at end of file +git+git://github.com/Supervisor/supervisor.git@c18aecf1641d8953767e7010be8bae1924a133bf#egg=Supervisor; python_version >= '3.0' +future==0.15.2 +importlib==1.0.3; python_version < '2.7' diff --git a/requirements-dev.txt b/requirements-dev.txt index cbc8fc85..9a4791c2 100644 --- a/requirements-dev.txt +++ b/requirements-dev.txt @@ -1,13 +1,16 @@ -pytest==2.7.0 +pytest==2.9.1 pylint==1.5.2 pytest-cov==1.8.1 -pytest-httpbin==0.0.6 pytest-instafail==0.3.0 +pytest-html==1.8.0 +pytest-httpbin==0.2.0 +pytest-circleci==0.0.2 +subprocess32==3.2.5; python_version < '3.2' # git+https://github.com/srlindsay/gevent-profiler@master#egg=gevent-profiler==0.2 # Used to test IO tracing -requests==2.4.3 +requests==2.9.1 mkdocs==0.15.3 -mistune==0.7.3 \ No newline at end of file +mistune==0.7.3 diff --git a/requirements-heroku.txt b/requirements-heroku.txt index 0c1b578c..04e6699c 100644 --- a/requirements-heroku.txt +++ b/requirements-heroku.txt @@ -1 +1 @@ -uwsgi==2.0.2 \ No newline at end of file +uwsgi==2.0.2 diff --git a/tests/conftest.py b/tests/conftest.py index 04b0ff0d..5cfacb09 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -1,3 +1,10 @@ +from __future__ import print_function +from future import standard_library +standard_library.install_aliases() +from builtins import range +from builtins import str +from builtins import object +from past.builtins import basestring import pytest import os try: @@ -9,7 +16,7 @@ import time import re import json -import urllib2 +import urllib.request, urllib.error, urllib.parse sys.path.append(os.getcwd()) @@ -59,7 +66,7 @@ def start(self, cmdline=None, env=None, expected_children=0): self.cmdline = cmdline # print cmdline - self.process = subprocess.Popen(re.split(r"\s+", cmdline) if type(cmdline) in [str, unicode] else cmdline, + self.process = subprocess.Popen(re.split(r"\s+", cmdline) if isinstance(cmdline, basestring) else cmdline, shell=False, close_fds=True, env=env, cwd=os.getcwd(), stdout=stdout) if self.quiet: @@ -151,7 +158,7 @@ def start(self, flush=True, deps=True, trace=True, **kwargs): if processes > 0: processes += 1 - print cmdline + print(cmdline) ProcessFixture.start(self, cmdline=cmdline, env=kwargs.get("env"), expected_children=processes) def start_deps(self, flush=True): @@ -234,13 +241,13 @@ def send_task_cli(self, path, params, queue=None, **kwargs): out = subprocess.check_output(cli).strip() if not queue: - return json.loads(out) + return json.loads(out.decode('utf-8')) return out def get_report(self, with_memory=False): wait_for_net_service("127.0.0.1", 20020, poll_interval=0.01) - f = urllib2.urlopen("http://127.0.0.1:20020/report%s" % ("_mem" if with_memory else "")) - data = json.load(f) + f = urllib.request.urlopen("http://127.0.0.1:20020/report%s" % ("_mem" if with_memory else "")) + data = json.loads(f.read().decode('utf-8')) f.close() return data diff --git a/tests/fixtures/standalone_script1.py b/tests/fixtures/standalone_script1.py index e7b0bf5e..0b867bd0 100644 --- a/tests/fixtures/standalone_script1.py +++ b/tests/fixtures/standalone_script1.py @@ -1,8 +1,9 @@ +from __future__ import print_function from mrq.context import setup_context, run_task, get_current_config # Autoconfigure MRQ's environment setup_context() -print run_task("tests.tasks.general.Add", {"a": 41, "b": 1}) +print(run_task("tests.tasks.general.Add", {"a": 41, "b": 1})) -print get_current_config()["name"] \ No newline at end of file +print(get_current_config()["name"]) \ No newline at end of file diff --git a/tests/tasks/general.py b/tests/tasks/general.py index 37b132a4..e2bdd0a3 100644 --- a/tests/tasks/general.py +++ b/tests/tasks/general.py @@ -1,8 +1,11 @@ +from future import standard_library +standard_library.install_aliases() +from builtins import range from mrq.task import Task from mrq.context import (log, retry_current_job, connections, get_current_config, get_current_job, subpool_map, abort_current_job, set_current_job_progress) from mrq.job import queue_job -import urllib2 +import urllib.request, urllib.error, urllib.parse import json import time import copy @@ -38,7 +41,7 @@ class Fetch(Task): def run(self, params): - f = urllib2.urlopen(params.get("url")) + f = urllib.request.urlopen(params.get("url")) t = f.read() f.close() diff --git a/tests/tasks/io.py b/tests/tasks/io.py index acb9ae8d..a891e3a1 100644 --- a/tests/tasks/io.py +++ b/tests/tasks/io.py @@ -1,7 +1,9 @@ +from future import standard_library +standard_library.install_aliases() from mrq.task import Task from mrq.context import connections, log -import urllib2 - +import urllib.request, urllib.error, urllib.parse +from future.moves.urllib.request import urlopen class TestIo(Task): @@ -44,17 +46,14 @@ def _run(self, params): elif params["test"] == "urllib2-get": - fp = urllib2.urlopen(params["params"]["url"]) + fp = urlopen(params["params"]["url"]) return fp.read() elif params["test"] == "urllib2-post": - return urllib2.urlopen(params["params"]["url"], data="x=x").read() + return urlopen(params["params"]["url"], data="x=x").read() elif params["test"] == "requests-get": import requests return requests.get(params["params"]["url"], verify=False).text - - - diff --git a/tests/tasks/logger.py b/tests/tasks/logger.py index 2bf8f2eb..f09ec539 100644 --- a/tests/tasks/logger.py +++ b/tests/tasks/logger.py @@ -1,18 +1,20 @@ from mrq.task import Task from mrq.context import log +import sys +PY3 = sys.version_info > (3,) class Simple(Task): def run(self, params): + # Some systems may be configured like this. - if params.get("utf8_sys_stdout"): + if not PY3 and params.get("utf8_sys_stdout"): import codecs import sys UTF8Writer = codecs.getwriter('utf8') sys.stdout = UTF8Writer(sys.stdout) - if params["class_name"] == "unicode": log.info(u"caf\xe9") elif params["class_name"] == "string": diff --git a/tests/tasks/redis.py b/tests/tasks/redis.py index ede9f2de..2c229f8e 100644 --- a/tests/tasks/redis.py +++ b/tests/tasks/redis.py @@ -1,3 +1,5 @@ +from __future__ import print_function +from builtins import range from mrq.task import Task from mrq.context import connections, subpool_map import gevent @@ -23,10 +25,10 @@ def run(self, params): get_clients = lambda: [c for c in connections.redis.client_list() if c.get("cmd") != "client"] def inner(i): - print "Greenlet #%s, %s clients so far" % (id(gevent.getcurrent()), len(get_clients())) + print("Greenlet #%s, %s clients so far" % (id(gevent.getcurrent()), len(get_clients()))) return connections.redis.get("test") if params["subpool_size"]: - subpool_map(params["subpool_size"], inner, range(0, params["subpool_size"] * 5)) + subpool_map(params["subpool_size"], inner, list(range(0, params["subpool_size"] * 5))) else: inner(0) diff --git a/tests/test_context.py b/tests/test_context.py index bafc2808..9c288d04 100644 --- a/tests/test_context.py +++ b/tests/test_context.py @@ -94,5 +94,4 @@ def test_context_setup(): out, err = process.communicate() - assert out.endswith("42\ntestname1\n") - + assert out.endswith(b"42\ntestname1\n") diff --git a/tests/test_disconnects.py b/tests/test_disconnects.py index 9ebba560..76a63777 100644 --- a/tests/test_disconnects.py +++ b/tests/test_disconnects.py @@ -1,3 +1,4 @@ +from builtins import range import time from mrq.job import Job import pytest diff --git a/tests/test_general.py b/tests/test_general.py index 70e86020..89b51db2 100644 --- a/tests/test_general.py +++ b/tests/test_general.py @@ -1,5 +1,8 @@ +from future import standard_library +standard_library.install_aliases() +from builtins import str from bson import ObjectId -import urllib2 +import urllib.request, urllib.error, urllib.parse import json import time @@ -20,7 +23,7 @@ def test_general_simple_task_one(worker): assert worker_report["done_jobs"] == 1 # Test the HTTP admin API - admin_worker = json.load(urllib2.urlopen("http://localhost:20020")) + admin_worker = json.loads(urllib.request.urlopen("http://localhost:20020").read().decode('utf-8')) assert admin_worker["_id"] == str(db_workers[0]["_id"]) assert admin_worker["status"] == "wait" @@ -167,4 +170,3 @@ def test_general_exception_status(worker): assert job1["status"] == "failed" assert "raise" in job1["traceback"] assert "xyz" in job1["traceback"] - diff --git a/tests/test_interrupts.py b/tests/test_interrupts.py index 5404658a..5a3fc23a 100644 --- a/tests/test_interrupts.py +++ b/tests/test_interrupts.py @@ -1,5 +1,6 @@ import time import datetime +from builtins import str from mrq.job import Job from mrq.queue import Queue from bson import ObjectId diff --git a/tests/test_io_hooks.py b/tests/test_io_hooks.py index 943e0443..876908b6 100644 --- a/tests/test_io_hooks.py +++ b/tests/test_io_hooks.py @@ -1,3 +1,4 @@ +from __future__ import print_function import json @@ -14,7 +15,7 @@ def test_io_hooks_nothing(worker): job_events = [x for x in events if x.get("job")] for evt in job_events: - print evt + print(evt) # Only update should be the result in mongodb. assert len(job_events) == 1 * 2 @@ -44,7 +45,7 @@ def test_io_hooks_redis(worker): job_events = [x for x in events if x.get("job")] for evt in job_events: - print evt + print(evt) assert len(job_events) == 2 * 2 @@ -84,7 +85,7 @@ def test_io_hooks_mongodb(worker): job_events = [x for x in events if x.get("job")] for evt in job_events: - print evt + print(evt) assert len(job_events) == 4 * 2 diff --git a/tests/test_jobinspect.py b/tests/test_jobinspect.py index 7c696972..9327f371 100644 --- a/tests/test_jobinspect.py +++ b/tests/test_jobinspect.py @@ -1,6 +1,11 @@ +from __future__ import print_function +from future import standard_library +standard_library.install_aliases() +from builtins import str +from builtins import range import time import ujson as json -import urllib2 +import urllib.request, urllib.error, urllib.parse import os import pytest @@ -15,7 +20,7 @@ def test_current_job_inspect(worker): time.sleep(1) # Test the HTTP admin API - admin_worker = json.load(urllib2.urlopen("http://localhost:20020")) + admin_worker = json.loads(urllib.request.urlopen("http://localhost:20020").read().decode('utf-8')) assert admin_worker["status"] == "full" assert len(admin_worker["jobs"]) == 1 @@ -30,7 +35,7 @@ def test_current_job_inspect(worker): time.sleep(3) - admin_worker = json.load(urllib2.urlopen("http://localhost:20020")) + admin_worker = json.loads(urllib.request.urlopen("http://localhost:20020").read().decode('utf-8')) assert admin_worker["status"] == "wait" assert len(admin_worker["jobs"]) == 0 @@ -91,13 +96,12 @@ def test_current_job_trace_io(worker, p_testtype, p_testparams, p_type, p_data, if os.path.isfile(report_file): with open(report_file, "rb") as f: try: - read = f.read() + read = f.read().decode('utf-8') admin_worker = json.loads(read) except: admin_worker = {} if len(admin_worker.get("jobs", [])) > 0: io = admin_worker["jobs"][0].get("io") - print io # Don't take MRQ's IOs as regular IO if io: if io["type"] == "mongodb" and io["data"]["collection"] in ["mrq.mrq_jobs", "mrq.mrq_logs"]: @@ -107,7 +111,7 @@ def test_current_job_trace_io(worker, p_testtype, p_testparams, p_type, p_data, time.sleep(0.05) - print io + print(io) assert io assert io["type"] == p_type assert io["data"] == p_data diff --git a/tests/test_memoryleaks.py b/tests/test_memoryleaks.py index 5746c9f5..cccfffd3 100644 --- a/tests/test_memoryleaks.py +++ b/tests/test_memoryleaks.py @@ -1,3 +1,5 @@ +from __future__ import print_function +from builtins import range import time @@ -41,7 +43,7 @@ def get_diff_after_jobs(worker, n_tasks, leak, sleep=0): diff = mem_stop - mem_start - print "Memory diff for %s tasks was %s" % (n_tasks, diff) + print("Memory diff for %s tasks was %s" % (n_tasks, diff)) return diff diff --git a/tests/test_parallel.py b/tests/test_parallel.py index 4aacc0d6..dc7c6864 100644 --- a/tests/test_parallel.py +++ b/tests/test_parallel.py @@ -1,3 +1,4 @@ +from builtins import range import time import pytest from mrq.context import connections @@ -23,7 +24,7 @@ def test_parallel_100sleeps(worker, p_flags): assert total_time < 15 # ... and return correct results - assert result == range(100) + assert result == list(range(100)) @pytest.mark.parametrize(["p_greenlets"], [ diff --git a/tests/test_performance.py b/tests/test_performance.py index e422b253..3c81d8e4 100644 --- a/tests/test_performance.py +++ b/tests/test_performance.py @@ -1,9 +1,17 @@ +from __future__ import division +from __future__ import print_function +from builtins import str +from builtins import range +from past.utils import old_div import time from mrq.queue import Queue import pytest -import subprocess import os +try: + import subprocess32 as subprocess +except: + import subprocess @pytest.mark.parametrize(["p_max_latency", "p_min_observed_latency", "p_max_observed_latency"], [ [1, 0.021, 1], @@ -23,10 +31,10 @@ def get_latency(): # This is the latency induced by our test system & general task work # We're on the same machine so even in different processes time.time() should be pretty reliable base_latency = get_latency() - print "Base latency: %ss" % base_latency + print("Base latency: %ss" % base_latency) min_latency = min([get_latency() for _ in range(0, 20)]) - print "FYI, min latency = %ss" % min_latency + print("FYI, min latency = %ss" % min_latency) # Sleep a while with an idle worker to make the poll interval go up latencies = [] @@ -35,12 +43,12 @@ def get_latency(): latency = get_latency() - min_latency - print "Observed latency (corrected): %ss" % latency + print("Observed latency (corrected): %ss" % latency) latencies.append(latency) - avg_latency = float(sum(latencies)) / len(latencies) - print "Average observed latency: %ss" % avg_latency + avg_latency = old_div(float(sum(latencies)), len(latencies)) + print("Average observed latency: %ss" % avg_latency) assert p_min_observed_latency <= avg_latency < p_max_observed_latency @@ -75,10 +83,10 @@ def benchmark_task(worker, taskpath, taskparams, tasks=1000, greenlets=50, proce ), queues=queues, trace=False) # Warm up the workers with one simple task. - print "Warming up workers..." + print("Warming up workers...") worker.send_tasks("tests.tasks.general.Add", [{"a": i, "b": 0, "sleep": 0} for i in range(greenlets * min(1, processes))]) - print "Starting benchmark..." + print("Starting benchmark...") start_time = time.time() # result = worker.send_tasks("tests.tasks.general.Add", @@ -91,7 +99,7 @@ def benchmark_task(worker, taskpath, taskparams, tasks=1000, greenlets=50, proce total_time = time.time() - start_time - print "%s tasks done with %s greenlets and %s processes in %0.3f seconds : %0.2f jobs/second!" % (tasks, greenlets, processes, total_time, tasks / total_time) + print("%s tasks done with %s greenlets and %s processes in %0.3f seconds : %0.2f jobs/second!" % (tasks, greenlets, processes, total_time, old_div(tasks, total_time))) assert total_time < max_seconds @@ -120,7 +128,7 @@ def test_performance_simpleadds_regular(worker, p_processes): max_seconds=max_seconds) # ... and return correct results - assert result == range(n_tasks) + assert result == list(range(n_tasks)) @pytest.mark.parametrize(["p_queue", "p_greenlets"], [x1 + x2 for x1 in [ @@ -195,7 +203,7 @@ def test_performance_writeconcern(worker_mongodb_with_journal): max_seconds=max_seconds ) - print total_time_acknowledged + print(total_time_acknowledged) result, total_time_unacknowledged = benchmark_task( worker, @@ -212,8 +220,8 @@ def test_performance_writeconcern(worker_mongodb_with_journal): max_seconds=max_seconds ) - print "total_time_acknowledged: ", total_time_acknowledged - print "total_time_unacknowledged: ", total_time_unacknowledged + print("total_time_acknowledged: ", total_time_acknowledged) + print("total_time_unacknowledged: ", total_time_unacknowledged) # Make sure it's faster. assert total_time_unacknowledged < total_time_acknowledged * 0.9 @@ -254,7 +262,7 @@ def test_performance_queue_cancel_requeue(worker): queue_time = time.time() - start_time - print "Queued %s tasks in %s seconds (%s/s)" % (n_tasks, queue_time, float(n_tasks) / queue_time) + print("Queued %s tasks in %s seconds (%s/s)" % (n_tasks, queue_time, old_div(float(n_tasks), queue_time))) assert queue_time < 2 assert Queue("noexec").size() == n_tasks @@ -272,7 +280,7 @@ def test_performance_queue_cancel_requeue(worker): ) assert res["cancelled"] == n_tasks queue_time = time.time() - start_time - print "Cancelled %s tasks in %s seconds (%s/s)" % (n_tasks, queue_time, float(n_tasks) / queue_time) + print("Cancelled %s tasks in %s seconds (%s/s)" % (n_tasks, queue_time, old_div(float(n_tasks), queue_time))) assert queue_time < 5 assert worker.mongodb_jobs.mrq_jobs.find( {"status": "cancel"}).count() == n_tasks @@ -291,7 +299,7 @@ def test_performance_queue_cancel_requeue(worker): ) queue_time = time.time() - start_time - print "Requeued %s tasks in %s seconds (%s/s)" % (n_tasks, queue_time, float(n_tasks) / queue_time) + print("Requeued %s tasks in %s seconds (%s/s)" % (n_tasks, queue_time, old_div(float(n_tasks), queue_time))) assert queue_time < 2 assert worker.mongodb_jobs.mrq_jobs.find( {"status": "queued"}).count() == n_tasks diff --git a/tests/test_ratelimit.py b/tests/test_ratelimit.py index 1b052348..da8b8e41 100644 --- a/tests/test_ratelimit.py +++ b/tests/test_ratelimit.py @@ -1,3 +1,4 @@ +from builtins import range from mrq.helpers import ratelimit import time diff --git a/tests/test_retry.py b/tests/test_retry.py index 3a70df15..266bf90b 100644 --- a/tests/test_retry.py +++ b/tests/test_retry.py @@ -1,3 +1,4 @@ +from builtins import str from mrq.job import Job import datetime from mrq.queue import Queue diff --git a/tests/test_scheduler.py b/tests/test_scheduler.py index 285bfa48..0fdc209b 100644 --- a/tests/test_scheduler.py +++ b/tests/test_scheduler.py @@ -1,3 +1,5 @@ +from __future__ import print_function +from builtins import str import time import pytest import datetime @@ -54,7 +56,7 @@ def test_scheduler_simple(worker, p_flags): # Only 3 should have been replaced and ran immediately again because they # have different config. inserts = list(collection.find()) - print inserts + print(inserts) assert len(inserts) == 3, inserts @@ -77,7 +79,7 @@ def test_scheduler_dailytime(worker, p_flags): time.sleep(3) inserts = list(collection.find()) assert len(inserts) == 2 - print inserts + print(inserts) assert collection.find({"params.b": "test"}).count() == 1 # Then a second time once the dailytime passes diff --git a/tests/test_subpool.py b/tests/test_subpool.py index f78b2df4..337bca6b 100644 --- a/tests/test_subpool.py +++ b/tests/test_subpool.py @@ -1,5 +1,9 @@ +from __future__ import print_function +from future import standard_library +standard_library.install_aliases() +from builtins import range from bson import ObjectId -import urllib2 +import urllib.request, urllib.error, urllib.parse import json import time import os @@ -81,18 +85,18 @@ def iterator(n): def inner_func(i): time.sleep(1) - print "inner_func: %s" % i + print("inner_func: %s" % i) if i == 4: raise Exception("Inner exception!") return i * 2 with pytest.raises(Exception): for res in subpool_imap(10, inner_func, iterator(10)): - print "Got %s" % res + print("Got %s" % res) for res in subpool_imap(2, inner_func, iterator(1)): - print "Got %s" % res + print("Got %s" % res) with pytest.raises(Exception): for res in subpool_imap(2, inner_func, iterator(5)): - print "Got %s" % res + print("Got %s" % res)