PY-22549: support multiprocess testing for pytest

In parallel mode tests are still reported sequentially, but out of order and testFinish is never reported for suites, because we do not know when suite is completed, so they must be closed on Java side.

Strategy is also passed to Java side to change duration calculation logic: root duration is wall time, not sum of all children
This commit is contained in:
Ilya.Kazakevich
2019-02-06 01:39:55 +03:00
parent c1e49c2ea3
commit 6d65790513
5 changed files with 279 additions and 148 deletions
@@ -0,0 +1,81 @@
# coding=utf-8
class ParallelTreeManager(object):
"""
Manages output tree by building it from flat test names.
"""
def __init__(self):
super(ParallelTreeManager, self).__init__()
self._max_node_id = 0
self._branches = dict() # key is test name as tuple, value is tuple of test_id, parent_id
def _next_node_id(self):
self._max_node_id += 1
return self._max_node_id
def level_opened(self, test_as_list, func_to_open):
"""
To be called on test start.
:param test_as_list: test name splitted as list
:param func_to_open: func to be called if test can open new level
:return: None if new level opened, or tuple of command client should execute and try opening level again
Command is "open" (open provided level) or "close" (close it). Second item is test name as list
"""
if tuple(test_as_list) in self._branches:
# We have parent, ok
func_to_open()
return None
elif len(test_as_list) == 1:
self._branches[tuple(test_as_list)] = (self._next_node_id(), 0)
func_to_open()
return None
commands = []
parent_id = 0
for i in range(len(test_as_list)):
tmp_parent_as_list = test_as_list[0:i + 1]
try:
parent_id, _ = self._branches[tuple(tmp_parent_as_list)]
except KeyError:
node_id = self._next_node_id()
self._branches[tuple(tmp_parent_as_list)] = (node_id, parent_id)
parent_id = node_id
if tmp_parent_as_list != test_as_list: # Different test opened
commands.append(("open", tmp_parent_as_list))
if commands:
return commands
else:
func_to_open()
# parent_id = "0"
# for branch_n in range(len(test_as_list), 1):
# branch = test_as_list[0:branch_n]
# key = tuple(branch)
# try:
# _, parent_id = self._branches[key]
# except KeyError:
# self._max_node_id += 1
# self._branches[key] = (self._max_node_id, parent_id)
# return "open", branch
# func_to_open()
def level_closed(self, test_as_list, func_to_close):
"""
To be called on test end or failure.
See level_opened doc.
"""
func_to_close()
# Part of contract
def get_node_ids(self, test_name):
"""
:return: (current_node_id, parent_node_id)
"""
return self._branches[tuple(test_name.split("."))]
+19 -5
View File
@@ -1,16 +1,20 @@
# coding=utf-8
import re
import sys
import pytest
from _pytest.config import get_plugin_manager
from _pytest import config
from _jb_runner_tools import jb_start_tests, jb_patch_separator, jb_doc_args, JB_DISABLE_BUFFERING
from teamcity import pytest_plugin
from pkg_resources import iter_entry_points
from _jb_runner_tools import jb_patch_separator, jb_doc_args, JB_DISABLE_BUFFERING, start_protocol, parse_arguments, \
set_parallel_mode
from teamcity import pytest_plugin
if __name__ == '__main__':
path, targets, additional_args = jb_start_tests()
real_prepare_config = config._prepareconfig
path, targets, additional_args = parse_arguments()
sys.argv += additional_args
joined_targets = jb_patch_separator(targets, fs_glue="/", python_glue="::", fs_to_python_glue=".py::")
# When file is launched in pytest it should be file.py: you can't provide it as bare module
@@ -21,11 +25,21 @@ if __name__ == '__main__':
# to prevent "plugin already registered" problem we check it first
plugins_to_load = []
if not get_plugin_manager().hasplugin("pytest-teamcity"):
if "pytest-teamcity" not in map(lambda e:e.name, iter_entry_points(group='pytest11', name=None)):
if "pytest-teamcity" not in map(lambda e: e.name, iter_entry_points(group='pytest11', name=None)):
plugins_to_load.append(pytest_plugin)
args = sys.argv[1:]
if JB_DISABLE_BUFFERING and "-s" not in args:
args += ["-s"]
jb_doc_args("pytest", args)
# We need to preparse numprocesses because user may set it using ini file
config_result = real_prepare_config(args, plugins_to_load)
if getattr(config_result.option, "numprocesses", None):
set_parallel_mode()
config._prepareconfig = lambda _, __: config_result
start_protocol()
pytest.main(args, plugins_to_load)
+66 -142
View File
@@ -2,12 +2,11 @@
"""
Tools to implement runners (https://confluence.jetbrains.com/display/~link/PyCharm+test+runners+protocol)
"""
import atexit
import _jb_utils
import os
import re
import sys
import _jb_utils
from teamcity import teamcity_presence_env_var, messages
# Some runners need it to "detect" TC and start protocol
@@ -21,6 +20,7 @@ JB_DISABLE_BUFFERING = "JB_DISABLE_BUFFERING" in os.environ
# getcwd resolves symlinks, but PWD is not supported by some shells
PROJECT_DIR = os.getenv('PWD', os.getcwd())
def _parse_parametrized(part):
"""
@@ -39,119 +39,38 @@ def _parse_parametrized(part):
return [match.group(1), match.group(2)]
# Monkeypatching TC to pass location hint
class _TreeManager(object):
"""
Manages output tree by building it from flat test names.
"""
class _TreeManagerHolder(object):
def __init__(self):
super(_TreeManager, self).__init__()
# Currently active branch as list. New nodes go to this branch
self.current_branch = []
# node unique name to its nodeId
self._node_ids_dict = {}
# Node id mast be incremented for each new branch
self._max_node_id = 0
def _calculate_relation(self, branch_as_list):
"""
Get relation of branch_as_list to current branch.
:return: tuple. First argument could be: "same", "child", "parent" or "sibling"(need to start new tree)
Second argument is relative path from current branch to child if argument is child
"""
if branch_as_list == self.current_branch:
return "same", None
hierarchy_name_len = len(branch_as_list)
current_branch_len = len(self.current_branch)
if hierarchy_name_len > current_branch_len and branch_as_list[0:current_branch_len] == self.current_branch:
return "child", branch_as_list[current_branch_len:]
if hierarchy_name_len < current_branch_len and self.current_branch[0:hierarchy_name_len] == branch_as_list:
return "parent", None
return "sibling", None
def _add_new_node(self, new_node_name):
"""
Adds new node to branch
"""
self.current_branch.append(new_node_name)
self._max_node_id += 1
self._node_ids_dict[".".join(self.current_branch)] = self._max_node_id
def level_opened(self, test_as_list, func_to_open):
"""
To be called on test start.
:param test_as_list: test name splitted as list
:param func_to_open: func to be called if test can open new level
:return: None if new level opened, or tuple of command client should execute and try opening level again
Command is "open" (open provided level) or "close" (close it). Second item is test name as list
"""
relation, relative_path = self._calculate_relation(test_as_list)
if relation == 'same':
return # Opening same level?
if relation == 'child':
# If one level -- open new level gracefully
if len(relative_path) == 1:
self._add_new_node(relative_path[0])
func_to_open()
return None
else:
# Open previous level
return "open", self.current_branch + relative_path[0:1]
if relation == "sibling":
if self.current_branch:
# Different tree, close whole branch
return "close", self.current_branch
else:
return None
if relation == 'parent':
# Opening parent? Insane
pass
def level_closed(self, test_as_list, func_to_close):
"""
To be called on test end or failure.
See level_opened doc.
"""
relation, relative_path = self._calculate_relation(test_as_list)
if relation == 'same':
# Closing current level
func_to_close()
self.current_branch.pop()
if relation == 'child':
return None
if relation == 'sibling':
pass
if relation == 'parent':
return "close", self.current_branch
self.parallel = "JB_USE_PARALLEL_TREE_MANAGER" in os.environ
self._manager_imp = None
@property
def parent_branch(self):
return self.current_branch[:-1] if self.current_branch else None
def manager(self):
if not self._manager_imp:
self._fill_manager()
return self._manager_imp
def _get_node_id(self, branch):
return self._node_ids_dict[".".join(branch)]
def _fill_manager(self):
if self.parallel:
from _jb_parallel_tree_manager import ParallelTreeManager
self._manager_imp = ParallelTreeManager()
else:
from _jb_serial_tree_manager import SerialTreeManager
self._manager_imp = SerialTreeManager()
@property
def node_ids(self):
"""
:return: (current_node_id, parent_node_id)
"""
current = self._get_node_id(self.current_branch)
parent = self._get_node_id(self.parent_branch) if self.parent_branch else "0"
return str(current), str(parent)
_TREE_MANAGER_HOLDER = _TreeManagerHolder()
TREE_MANAGER = _TreeManager()
def set_parallel_mode():
_TREE_MANAGER_HOLDER.parallel = True
def is_parallel_mode():
return _TREE_MANAGER_HOLDER.parallel
# Monkeypatching TC
_old_service_messages = messages.TeamcityServiceMessages
PARSE_FUNC = None
@@ -159,9 +78,9 @@ PARSE_FUNC = None
class NewTeamcityServiceMessages(_old_service_messages):
_latest_subtest_result = None
def message(self, messageName, **properties):
if messageName in set(["enteredTheMatrix", "testCount"]):
if messageName in {"enteredTheMatrix", "testCount"}:
_old_service_messages.message(self, messageName, **properties)
return
@@ -181,13 +100,13 @@ class NewTeamcityServiceMessages(_old_service_messages):
_old_service_messages.message(self, messageName, **properties)
return
current, parent = _TREE_MANAGER_HOLDER.manager.get_node_ids(properties["name"])
# Shortcut for name
try:
properties["name"] = str(properties["name"]).split(".")[-1]
except IndexError:
pass
current, parent = TREE_MANAGER.node_ids
properties["nodeId"] = str(current)
properties["parentNodeId"] = str(parent)
@@ -230,8 +149,8 @@ class NewTeamcityServiceMessages(_old_service_messages):
return
# closing subtest
test_name = ".".join(TREE_MANAGER.current_branch)
if self._latest_subtest_result in set(["Failure", "Error"]):
test_name = ".".join(_TREE_MANAGER_HOLDER.manager.current_branch)
if self._latest_subtest_result in {"Failure", "Error"}:
self.testFailed(test_name)
if self._latest_subtest_result == "Skip":
self.testIgnored(test_name)
@@ -240,7 +159,7 @@ class NewTeamcityServiceMessages(_old_service_messages):
self._latest_subtest_result = None
def subTestBlockOpened(self, name, subTestResult, flowId=None):
self.testStarted(".".join(TREE_MANAGER.current_branch + [name]))
self.testStarted(".".join(_TREE_MANAGER_HOLDER.manager.current_branch + [name]))
self._latest_subtest_result = subTestResult
def testStarted(self, testName, captureStandardOutput=None, flowId=None, is_suite=False, metainfo=None):
@@ -249,15 +168,17 @@ class NewTeamcityServiceMessages(_old_service_messages):
def _write_start_message():
# testName, captureStandardOutput, flowId
args = {"name": testName, "captureStandardOutput": captureStandardOutput, "metainfo":metainfo}
args = {"name": testName, "captureStandardOutput": captureStandardOutput, "metainfo": metainfo}
if is_suite:
if is_parallel_mode():
args["durationStrategy"] = "explicit_only"
self.message("testSuiteStarted", **args)
else:
self.message("testStarted", **args)
commands = TREE_MANAGER.level_opened(self._test_to_list(testName), _write_start_message)
commands = _TREE_MANAGER_HOLDER.manager.level_opened(self._test_to_list(testName), _write_start_message)
if commands:
self.do_command(commands[0], commands[1])
self.do_commands(commands)
self.testStarted(testName, captureStandardOutput, metainfo=metainfo)
def testFailed(self, testName, message='', details='', flowId=None, comparison_failure=None):
@@ -269,7 +190,7 @@ class NewTeamcityServiceMessages(_old_service_messages):
def _write_finished_message():
# testName, captureStandardOutput, flowId
current, parent = TREE_MANAGER.node_ids
current, parent = _TREE_MANAGER_HOLDER.manager.get_node_ids(testName)
args = {"nodeId": current, "parentNodeId": parent, "name": testName}
# TODO: Doc copy/paste with parent, extract
@@ -280,36 +201,30 @@ class NewTeamcityServiceMessages(_old_service_messages):
args["duration"] = str(duration_ms)
if is_suite:
if is_parallel_mode():
del args["duration"]
args["durationStrategy"] = "explicit_only"
self.message("testSuiteFinished", **args)
else:
self.message("testFinished", **args)
commands = TREE_MANAGER.level_closed(self._test_to_list(testName), _write_finished_message)
commands = _TREE_MANAGER_HOLDER.manager.level_closed(self._test_to_list(testName), _write_finished_message)
if commands:
self.do_command(commands[0], commands[1])
self.do_commands(commands)
self.testFinished(testName, testDuration)
def do_command(self, command, test):
def do_commands(self, commands):
"""
Executes commands, returned by level_closed and level_opened
"""
test_name = ".".join(test)
# By executing commands we open or close suites(branches) since tests(leaves) are always reported by runner
if command == "open":
self.testStarted(test_name, is_suite=True)
else:
self.testFinished(test_name, is_suite=True)
def close_all(self):
"""
Closes all tests
"""
commands = TREE_MANAGER.close_all()
if commands:
self.do_command(commands[0], commands[1])
self.close_all()
for command, test in commands:
test_name = ".".join(test)
# By executing commands we open or close suites(branches) since tests(leaves) are always reported by runner
if command == "open":
self.testStarted(test_name, is_suite=True)
else:
self.testFinished(test_name, is_suite=True)
messages.TeamcityServiceMessages = NewTeamcityServiceMessages
@@ -347,13 +262,25 @@ def jb_patch_separator(targets, fs_glue, python_glue, fs_to_python_glue):
def jb_start_tests():
"""
Parses arguments, starts protocol and returns tuple of arguments
Parses arguments, starts protocol and fixes syspath and returns tuple of arguments
"""
path, targets, additional_args = parse_arguments()
start_protocol()
return path, targets, additional_args
def start_protocol():
properties = {"durationStrategy": "wall_time"} if is_parallel_mode() else dict()
NewTeamcityServiceMessages().message('enteredTheMatrix', **properties)
def parse_arguments():
"""
Parses arguments, fixes syspath and returns tuple of arguments
:return: (string with path or None, list of targets or None, list of additional arguments)
:param func_to_parse function that accepts each part of test name and returns list to be used instead of it.
It may return list with only one element (name itself) if name is the same or split names to several parts
"""
# Handle additional args after --
additional_args = []
try:
@@ -363,12 +290,10 @@ def jb_start_tests():
except ValueError:
pass
utils = _jb_utils.VersionAgnosticUtils()
namespace = utils.get_options(
_jb_utils.OptionDescription('--path', 'Path to file or folder to run'),
_jb_utils.OptionDescription('--target', 'Python target to run', "append"))
del sys.argv[1:] # Remove all args
NewTeamcityServiceMessages().message('enteredTheMatrix')
# PyCharm helpers dir is first dir in sys.path because helper is launched.
# But sys.path should be same as when launched with test runner directly
@@ -380,7 +305,6 @@ def jb_start_tests():
return namespace.path, namespace.target, additional_args
def jb_doc_args(framework_name, args):
"""
Runner encouraged to report its arguments to user with aid of this function
@@ -0,0 +1,112 @@
# coding=utf-8
class SerialTreeManager(object):
"""
Manages output tree by building it from flat test names.
"""
def __init__(self):
super(SerialTreeManager, self).__init__()
# Currently active branch as list. New nodes go to this branch
self.current_branch = []
# node unique name to its nodeId
self._node_ids_dict = {}
# Node id mast be incremented for each new branch
self._max_node_id = 0
def _calculate_relation(self, branch_as_list):
"""
Get relation of branch_as_list to current branch.
:return: tuple. First argument could be: "same", "child", "parent" or "sibling"(need to start new tree)
Second argument is relative path from current branch to child if argument is child
"""
if branch_as_list == self.current_branch:
return "same", None
hierarchy_name_len = len(branch_as_list)
current_branch_len = len(self.current_branch)
if hierarchy_name_len > current_branch_len and branch_as_list[0:current_branch_len] == self.current_branch:
return "child", branch_as_list[current_branch_len:]
if hierarchy_name_len < current_branch_len and self.current_branch[0:hierarchy_name_len] == branch_as_list:
return "parent", None
return "sibling", None
def _add_new_node(self, new_node_name):
"""
Adds new node to branch
"""
self.current_branch.append(new_node_name)
self._max_node_id += 1
self._node_ids_dict[".".join(self.current_branch)] = self._max_node_id
def level_opened(self, test_as_list, func_to_open):
"""
To be called on test start.
:param test_as_list: test name splitted as list
:param func_to_open: func to be called if test can open new level
:return: None if new level opened, or tuple of command client should execute and try opening level again
Command is "open" (open provided level) or "close" (close it). Second item is test name as list
"""
relation, relative_path = self._calculate_relation(test_as_list)
if relation == 'same':
return # Opening same level?
if relation == 'child':
# If one level -- open new level gracefully
if len(relative_path) == 1:
self._add_new_node(relative_path[0])
func_to_open()
return None
else:
# Open previous level
return [("open", self.current_branch + relative_path[0:1])]
if relation == "sibling":
if self.current_branch:
# Different tree, close whole branch
return [("close", self.current_branch)]
else:
return None
if relation == 'parent':
# Opening parent? Insane
pass
def level_closed(self, test_as_list, func_to_close):
"""
To be called on test end or failure.
See level_opened doc.
"""
relation, relative_path = self._calculate_relation(test_as_list)
if relation == 'same':
# Closing current level
func_to_close()
self.current_branch.pop()
if relation == 'child':
return None
if relation == 'sibling':
pass
if relation == 'parent':
return [("close", self.current_branch)]
@property
def parent_branch(self):
return self.current_branch[:-1] if self.current_branch else None
def _get_node_id(self, branch):
return self._node_ids_dict[".".join(branch)]
# Part of contract
# noinspection PyUnusedLocal
def get_node_ids(self, test_name):
"""
:return: (current_node_id, parent_node_id)
"""
current = self._get_node_id(self.current_branch)
parent = self._get_node_id(self.parent_branch) if self.parent_branch else "0"
return str(current), str(parent)
@@ -83,7 +83,7 @@ public class PythonTRunnerConsoleProperties extends SMTRunnerConsoleProperties {
public void onTestFailed(@NotNull final SMTestProxy test) {
super.onTestFailed(test);
SMTestProxy currentTest = test.getParent();
while (currentTest != null) {
while (currentTest != null && currentTest.getParent() != null) {
currentTest.setTestFailed(" ", null, false);
currentTest = currentTest.getParent();
}