Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 9 additions & 0 deletions python/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,15 @@ Then, to then get a summary report of all the tests, run the following on anothe
py.test -p ciqueue.pytest_report --queue redis://<host>:6379?build=<build_id>&retry=<n>
```

Workers store rendered pytest reports as compressed JSON in Redis; the reporter never loads Python objects from the queue.
Upgrade workers and the reporter together, using the same `ciqueue` and pytest versions within a build.
Legacy dill records and malformed reports are rejected with a non-zero exit status, not treated as passing tests.
After upgrading, rerun all workers with a fresh build ID; do not reuse an old build's records.

Traceback formatting (for example, `--tb` and `--showlocals`) is determined by the worker's options.
The queue transports standard `TestReport` fields, including captured output, user properties rendered as strings,
and xfail reasons; plugin-specific report attributes are not transported.

## Implementing a new integration

The reference implementation is the minitest one (Ruby).
Expand Down
95 changes: 0 additions & 95 deletions python/ciqueue/_pytest/outcomes.py

This file was deleted.

90 changes: 90 additions & 0 deletions python/ciqueue/_pytest/reports.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,90 @@
"""Transport rendered pytest reports, never exception or traceback objects."""
import json
import math
import zlib

import pytest
from _pytest.reports import TestReport


REPORT_FIELDS = {
'$report_type', 'nodeid', 'location', 'keywords', 'outcome', 'longrepr', 'when',
'sections', 'duration', 'start', 'stop', 'user_properties', 'wasxfail',
}


def dumps(config, reports):
payload = {}
for when, report in reports.items():
data = config.hook.pytest_report_to_serializable(config=config, report=report)
# Plugin-specific attributes are not part of the queue's report protocol.
payload[when] = {key: value for key, value in data.items() if key in REPORT_FIELDS}
# JUnit renders property values as text; arbitrary objects stay on the worker.
payload[when]['user_properties'] = [(name, str(value)) for name, value in report.user_properties]
return zlib.compress(json.dumps(payload, allow_nan=False).encode('utf-8'))


def _location(value, allow_none=False):
return (isinstance(value, list) and len(value) == 3 and
isinstance(value[0], str) and isinstance(value[2], str) and
(type(value[1]) is int or (allow_none and value[1] is None)))


def _pairs(value):
return isinstance(value, list) and all(
isinstance(pair, list) and len(pair) == 2 and
all(isinstance(part, str) for part in pair) for pair in value)


def _validate(data, when, nodeid):
if not isinstance(data, dict) or not data.keys() <= REPORT_FIELDS:
raise ValueError('unsupported report fields')
if (data.get('$report_type') != 'TestReport' or data.get('nodeid') != nodeid or
data.get('when') != when or data.get('outcome') not in ('failed', 'skipped')):
raise ValueError('invalid report identity or outcome')
if not _location(data.get('location'), allow_none=True) or not isinstance(data.get('keywords'), dict):
raise ValueError('invalid report location or keywords')
if not _pairs(data.get('sections', [])) or not _pairs(data.get('user_properties', [])):
raise ValueError('invalid report sections or properties')
for name in ('duration', 'start', 'stop'):
if name in data and (type(data[name]) not in (int, float) or not math.isfinite(data[name])):
raise ValueError('invalid report timing')
if 'wasxfail' in data and (not isinstance(data['wasxfail'], str) or data['outcome'] != 'skipped'):
raise ValueError('xfail metadata requires a skipped outcome')
longrepr = data.get('longrepr')
if isinstance(longrepr, list):
if not _location(longrepr):
raise ValueError('invalid skip representation')
data['longrepr'] = tuple(longrepr)
elif isinstance(longrepr, dict):
if not {'reprcrash', 'reprtraceback', 'sections', 'chain'} <= longrepr.keys():
raise ValueError('invalid traceback representation')
elif not isinstance(longrepr, str):
raise ValueError('missing failure representation')
if data['outcome'] == 'skipped' and 'wasxfail' not in data and not isinstance(data['longrepr'], tuple):
raise ValueError('missing skip location')
data['location'] = tuple(data['location'])


def loads(config, payload, nodeid):
try:
data = json.loads(zlib.decompress(payload).decode('utf-8'))
if not isinstance(data, dict) or not data or not data.keys() <= {'setup', 'call', 'teardown'}:
raise ValueError('expected setup/call/teardown reports')
if 'setup' in data and 'call' in data:
raise ValueError('a non-passing setup cannot have a call report')
reports = {}
for when, report_data in data.items():
_validate(report_data, when, nodeid)
report = config.hook.pytest_report_from_serializable(config=config, data=report_data)
if not isinstance(report, TestReport):
raise ValueError('expected a TestReport')
# Exercise pytest's structured traceback renderer before accepting a record.
# Malformed nested representations must fail here, not during reporting.
_ = report.longreprtext
reports[when] = report
return reports
except (ValueError, TypeError, KeyError, AttributeError, AssertionError, RuntimeError, zlib.error) as error:
raise pytest.UsageError(
'Invalid error report for {}: {}. Expected compressed JSON reports; '
'upgrade workers and reporter together and use a fresh build ID.'.format(nodeid, error)) from error
76 changes: 35 additions & 41 deletions python/ciqueue/pytest.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,10 +5,8 @@
"""
from __future__ import absolute_import
from __future__ import print_function
import zlib
from ciqueue._pytest import test_queue
from ciqueue._pytest import outcomes
import dill
from ciqueue._pytest import reports
import pytest
from _pytest import terminal

Expand Down Expand Up @@ -73,19 +71,21 @@ def _get_progress(self): # pylint: disable=unused-argument

terminal.TerminalReporter._get_progress_information_message = _get_progress # pylint: disable=protected-access

def record(self, item):
# if the test passed, we remove it from the errors queue
# otherwise we add it
if hasattr(item, 'error_reports'):
self.redis.hset(
self.errors_key,
test_queue.key_item(item),
zlib.compress(dill.dumps(item.error_reports)))
def record(self, item, test_failed):
# Serialize before acknowledging so encoding errors cannot lose a failure.
payload = reports.dumps(self.config, item.error_reports) if hasattr(item, 'error_reports') else None
test_name = test_queue.key_item(item)
# A late worker may replace an earlier failure only if it succeeded.
if not self.queue.acknowledge(test_name) and test_failed:
return False
if payload is not None:
self.redis.hset(self.errors_key, test_name, payload)
else:
self.redis.hdel(self.errors_key, test_queue.key_item(item))
self.redis.hdel(self.errors_key, test_name)
return True

def mark_as_skipped(self, call, item, msg):
assert call.when == 'teardown'
def mark_as_skipped(self, report, item, msg):
assert report.when == 'teardown'

stats = self.terminalreporter.stats

Expand All @@ -106,54 +106,48 @@ def clear_out_stats(key):
if self.logxml:
self.logxml.node_reporters_ordered[-1].nodes = []

# the call is converted to a skip
call.excinfo = outcomes.skipped_excinfo(item, msg)
# Render retries locally; no exception or traceback objects go on the wire.
path, lineno, _ = item.location
report.outcome = 'skipped'
report.longrepr = (path, (lineno or 0) + 1, msg)
if hasattr(report, 'wasxfail'):
del report.wasxfail

# clear out the stats like the test never happened
for key in ('passed', 'error', 'failed'):
clear_out_stats(key)

# rollback the testsfailed number like it never happened
item.session.testsfailed -= len([v for k, v in item.error_reports.items()
if not issubclass(v['excinfo'].type, outcomes.Skipped) and k != 'teardown'])
item.session.testsfailed -= sum(
report.failed for when, report in item.error_reports.items() if when != 'teardown')

# and clear out any state on the item like it never happened
if hasattr(item, 'error_reports'):
del item.error_reports

@pytest.hookimpl(tryfirst=True)
@pytest.hookimpl(hookwrapper=True, tryfirst=True)
def pytest_runtest_makereport(self, item, call):
"""This function hooks into pytest's reporting of test results, and pushes a failed test's error report
onto the redis queue. A test can fail in any of the 3 call stages: setup, test, or teardown.
This is captured by pushing a dict of {call_state: error} for each failed test."""
if call.excinfo:
payload = call.__dict__.copy()
payload['excinfo'] = outcomes.swap_in_serializable(payload['excinfo'])

"""Record final reports after pytest has applied skip and xfail outcomes."""
result = yield
report = result.get_result()
if not report.passed:
if not hasattr(item, 'error_reports'):
item.error_reports = {call.when: payload}
else:
item.error_reports[call.when] = payload
item.error_reports = {}
item.error_reports[report.when] = report

if call.when == 'teardown':
if report.when == 'teardown':
test_name = test_queue.key_item(item)
test_failed = outcomes.failed(item)
test_failed = any(report.failed for report in getattr(item, 'error_reports', {}).values())

# Only attempt to requeue if the test failed.
# The method will return `False` if the test couldn't be requeued
if test_failed and self.queue.requeue(test_name):
self.mark_as_skipped(call, item, "WILL_RETRY")
self.mark_as_skipped(report, item, "WILL_RETRY")
self.terminalwriter.write(' WILL_RETRY ', green=True)

# If the test was already acknowledged by another worker (we timed out)
# Then we only record it if it was successful.
elif self.queue.acknowledge(test_name) or not test_failed:
self.record(item)

# The test timed out and failed, mark it as skipped so that it doesn't
# fail the build
else:
self.mark_as_skipped(call, item, "TIMED OUT")
# Ignore a late failure if another worker already acknowledged the test.
elif not self.record(item, test_failed):
self.mark_as_skipped(report, item, "TIMED OUT")
self.terminalwriter.write(' TIMED OUT ', green=True)


Expand Down
32 changes: 12 additions & 20 deletions python/ciqueue/pytest_report.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,12 +6,9 @@

from __future__ import absolute_import
from __future__ import print_function
import zlib
import dill
import pytest
from _pytest import runner
from ciqueue._pytest import test_queue
from ciqueue._pytest import outcomes
from ciqueue._pytest import reports


def pytest_addoption(parser):
Expand Down Expand Up @@ -45,24 +42,19 @@ def pytest_collection_modifyitems(session, config, items): # pylint: disable=un
# store the errors on setup/test/teardown to item.error_reports
key = test_queue.key_item(item)
if key in error_reports:
item.error_reports = dill.loads(zlib.decompress(error_reports[key]))
for _, call_dict in item.error_reports.items():
call_dict['excinfo'] = outcomes.swap_back_original(call_dict['excinfo'])
item.error_reports = reports.loads(config, error_reports[key], key)


@pytest.hookimpl(tryfirst=True)
@pytest.hookimpl(hookwrapper=True, tryfirst=True)
def pytest_runtest_makereport(item, call):
"""This function hooks into pytest's reporting of test results, and replaces the
result of each test's setup/runtest/teardown call with the result from the redis queue"""

# ensure all errors should come off the error-reports queue
"""Replay the worker's final outcome after local skip/xfail hooks have run."""
call.excinfo = None
result = yield
if hasattr(item, 'error_reports') and call.when in item.error_reports:
call.__dict__ = item.error_reports[call.when]

# This is needed to change the location of the failure
# to point to the item definition, otherwise it will display
# the location of where the skip exception was raised within pytest
# https://github.com/pytest-dev/pytest/blob/master/_pytest/skipping.py#L263-L269
if call.excinfo and call.excinfo.type == runner.Skipped:
item._evalskip = True # pylint: disable=protected-access
result.force_result(item.error_reports[call.when])
elif hasattr(item, 'error_reports') and call.when == 'teardown':
# JUnit finalizes metadata on teardown, even when only an earlier phase failed.
previous = item.error_reports.get('call') or item.error_reports['setup']
report = result.get_result()
report.user_properties = previous.user_properties
report.sections = previous.sections
4 changes: 1 addition & 3 deletions python/setup.py
Original file line number Diff line number Diff line change
Expand Up @@ -37,10 +37,8 @@ def get_lua_scripts():
packages=['ciqueue', 'ciqueue._pytest'],
python_requires='>=3.10',
install_requires=[
'dill>=0.2.7',
'pytest>=2.7',
'pytest>=6.2.5',
'redis>=2.10.5',
'tblib>=1.3.2',
'uritools>=2.0.0',
'future>=0.16.0'
],
Expand Down
Loading
Loading