Files
parsedmarc/tests/test_parallel.py
T
07bca1ad28 Make the output and mailbox integrations optional extras (#888)
* Make the output and mailbox integrations optional extras (#883)

Breaking change for the next major release: pip install parsedmarc now
installs the parsing core plus a working core CLI (file, IMAP, Maildir,
and mbox input; CSV/JSON, Splunk HEC, webhook, and syslog output).
Everything else moves behind an extra: elastic, opensearch, kafka, s3,
gelf, loganalytics, msgraph, and gmail, joining the existing postgresql
extra, with an umbrella [all] that deliberately excludes postgresql
(psycopg's binary wheels do not exist on every platform, so
parsedmarc[all] must never fail to install there).

cli.py imports the six SDK-dependent output modules behind the #884
TYPE_CHECKING/try-except guard; a configured section whose extra is
missing fails fast with a ConfigurationError naming the section and the
exact pip install command — including the msgraph and gmail_api mailbox
sections (detected via parsedmarc.mail's placeholder classes) and
postgresql (checked before the constructor so the startup retry loop
does not retry a missing dependency for a minute). The Azure/kiota Graph
error types fall back to never-raised sentinel classes.

The Docker image installs [all,postgresql], so container users see no
change. CI lint installs [build,all,postgresql]; the unit-test job
installs [build,all], deliberately without postgresql so
test_postgres.py's absent-psycopg arm stays exercised. The
never-imported dateparser dependency is dropped in favor of declaring
python-dateutil, which utils.py actually imports; pytz moves to the
build extra for the one test that uses it.

Verified live: a no-extras wheel install imports, parses samples, and
reports the install hint for each gated section; a [all] install
restores every integration; the Docker image builds with every SDK
importable.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

* Patch psycopg presence in the PostgreSQL CLI wiring tests

CI's unit-test job deliberately installs [build,all] without the
postgresql extra, so parsedmarc.cli.postgres.psycopg is None there and
the new missing-extra presence check correctly made _main exit 1 before
the wiring under test ran. The tests simulate the SDK being available
(PostgreSQLClient is mocked at the SDK boundary), so the module-level
psycopg handle is now patched present in setUp. Verified against a
simulated psycopg-absent environment as well as the local full install.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

* Address Copilot review: narrow guards to ModuleNotFoundError, fix docs

- The optional-integration and Graph error-type import guards now catch
  ModuleNotFoundError instead of ImportError, so only a genuinely absent
  package reads as a missing extra; a broken-but-present SDK fails
  loudly with its real error instead of masquerading as one. The test
  blocker raises ModuleNotFoundError accordingly — the exact exception a
  missing package produces.
- _missing_extra_hint docstring no longer calls every gated integration
  an output module (it also serves the msgraph/gmail_api mailbox
  sections).
- Fix the pre-existing passsword typo in usage.md's kafka section; the
  INI key the code reads is password (cli.py _parse_config).

Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com>
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

* Quote extras specs in copy-paste install commands

From Copilot's second review round: zsh treats an unquoted .[build,all]
as a glob and fails with 'no matches found', so the commands shown in
AGENTS.md, CONTRIBUTING.md, dashboards/README.md, and the bootstrap
script's comment are now quoted. The CI workflows keep the unquoted
form: they run under bash, which passes unmatched globs through
literally. The suggestion to change the 'Choosing what to install'
heading level was rejected — it is a subsection of 'Installing
parsedmarc', matching the file's existing hierarchy.

Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com>
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

* Fix upgrade command in the changelog

* Documentation review: accuracy, spelling, grammar, and clarity pass

A full prose review of docs/source, README, CONTRIBUTING, and the
dashboards README, with every accuracy claim verified against the code
before changing it. Highlights:

- usage.md: documented six missing [general] options (the CSV/JSON
  filename options, prettify_json, normalize_timespan_threshold_hours),
  the required kafka smtp_tls_topic, [imap] timeout/max_retries, and
  the postgresql env-var prefix; corrected the maildir_path default
  (None, not INBOX — cli.py Namespace defaults), the mailbox
  check_timeout option name, the systemd restart interval (RestartSec
  is 5m), and merged the duplicate silent entry; quoted every
  copy-paste extras spec for zsh safety.
- elasticsearch.md: fixed an invalid openssl command (rsa:4096 -nodes),
  the dashboards filename (opensearch_dashboards.ndjson, matching the
  file the link serves), and assorted grammar.
- davmail.md: the service-enable command now enables davmail.service
  (was parsedmarc.service — a copy-paste error that left DavMail
  unenabled), plus a view typo and DavMail capitalization.
- output.md: the example schema reference is RFC 7489 Appendix C
  (7480 is RDAP). kibana.md: SPF relies on the SMTP envelope, not
  session headers (RFC 7208). dmarc.md: DKM -> DKIM.
- README: the intro now also names the OpenSearch/Grafana stack,
  matching the feature list. CONTRIBUTING: pre-PR checks now include
  ruff format --check and pyright, matching CI's lint job.
- dashboards/README: the service table and seed description now include
  the PostgreSQL backend the compose stack runs.

Sample data blocks, the CLI-help mirror block, and released CHANGELOG
entries were deliberately left untouched.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

* Docstring review: accuracy, spelling, grammar, and clarity pass

Every docstring in parsedmarc/, parsedmarc/mail/, the maps maintainer
scripts, and the test suite reviewed with each claim verified against
the code it documents. Text-only — no behavior changes. Highlights:

- Copy-paste errors corrected: parsed_smtp_tls_reports_to_csv and
  splunk/loganalytics save functions described aggregate or failure
  reports they do not handle; LogAnalyticsException claimed to be an
  Elasticsearch error.
- Docstring/behavior mismatches: parse_report_email's report_type
  enumeration omitted smtp_tls; parse_failure_report typed msg_date as
  str (it is datetime); strip_attachment_payloads claimed payloads are
  replaced with None (the key is deleted); kafkaclient's failure and
  SMTP TLS savers claimed per-record slicing while sending the whole
  list in one message (docstrings now describe reality — whether
  slicing was intended is flagged for follow-up); the postgres savers
  claimed to take parse_report_file's return value but receive the
  inner report dict; elastic/opensearch save functions' Raises listed
  only AlreadySaved.
- None-as-semantic-state documented where missing (get_base_domain,
  get_ip_address_country), enumeration completeness fixed
  (get_ip_address_info's 9 result keys, maps script outputs, TSV
  columns), and the stale 44-industry-types count corrected to the
  46 the authoritative README list defines.
- Test docstrings aligned with what the tests actually assert,
  including two that overstated coverage of the elastic/opensearch
  address-list tests.
- Two argparse help strings fixed: file_path now names SMTP TLS report
  files alongside aggregate and failure, mirrored into usage.md's
  CLI-help block; --offline's doubled spaces removed (rendered help
  unchanged).
- elasticsearch.md's security claim corrected against Elastic's docs:
  security is enabled and auto-configured on first startup since 8.0
  (not "8.7 secure mode"), so the settings are verified, not
  hand-written.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

---------

Co-authored-by: Claude Fable 5 <noreply@anthropic.com>
Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com>
2026-08-28 17:33:18 -04:00

300 lines
11 KiB
Python

"""Tests for parsedmarc.parallel"""
import functools
import logging
import os
import tempfile
import unittest
from unittest.mock import patch
import parsedmarc
from parsedmarc.parallel import (
_init_worker_logging,
_parse_report_email_job,
_parse_report_file_job,
parallel_map,
)
# Stable sample files reused from tests/test_cli.py's TestDirectoryFilePaths,
# plus a third plain aggregate sample, so parity/order tests exercise more
# than one worker submission.
SAMPLE_PATHS = [
"samples/aggregate/!example.com!1538204542!1538463818.xml",
"samples/aggregate/!large-example.com!1711897200!1711983600.xml",
"samples/aggregate/example.net!example.com!1529366400!1529452799.xml",
]
def _echo_job(x):
"""Trivial module-level (spawn-picklable) worker used to exercise
parallel_map's scheduling behavior without involving real parsing."""
return x
class _CountingIterable:
"""Wraps a range so tests can observe how many items parallel_map has
pulled from a lazily-iterated jobs source at any point during
iteration, without materializing the whole sequence up front."""
def __init__(self, n):
self.n = n
self.pulled = 0
def __iter__(self):
for i in range(self.n):
self.pulled += 1
yield i
class _ParallelTestCase(unittest.TestCase):
"""Common env setup shared by parallel.py tests: offline mode, no DNS,
and a cold IP address cache."""
def setUp(self):
self._env_patcher = patch.dict(
os.environ, {"GITHUB_ACTIONS": "true"}, clear=False
)
self._env_patcher.start()
self.addCleanup(self._env_patcher.stop)
# Earlier test modules (e.g. tests/test_init.py) may run with DNS
# enabled locally and warm the shared module-level
# parsedmarc.IP_ADDRESS_CACHE. get_ip_address_info consults the
# cache before the offline check, so a warm cache would leak
# DNS-enriched entries into offline parses run in this process
# while spawned workers start cold, breaking parity.
parsedmarc.IP_ADDRESS_CACHE.clear()
class TestParallelMapParseReportFile(_ParallelTestCase):
"""parallel_map + _parse_report_file_job over real sample files must
behave identically to calling parse_report_file sequentially, and
must preserve submission order in its results."""
def test_results_match_sequential_parsing_in_order(self):
expected = [
(path, parsedmarc.parse_report_file(path, offline=True))
for path in SAMPLE_PATHS
]
job = functools.partial(
_parse_report_file_job, config=parsedmarc.ParserConfig(offline=True)
)
results = list(parallel_map(job, SAMPLE_PATHS, n_procs=2))
self.assertEqual(len(results), len(SAMPLE_PATHS))
# Order must match submission order (SAMPLE_PATHS), not completion
# order.
self.assertEqual([path for path, _ in results], SAMPLE_PATHS)
for (path, report), (expected_path, expected_report) in zip(results, expected):
self.assertEqual(path, expected_path)
self.assertNotIsInstance(report, Exception)
self.assertEqual(report, expected_report)
class TestParallelMapJunkFile(_ParallelTestCase):
"""A worker crash on an unparseable file must surface as an Exception
*value* in the result tuple, not hang the whole run. This is a
regression guard for the old Pipe/Process CLI worker, which left the
parent blocked forever on conn.recv() when a child died from a
non-ParserError exception."""
def test_junk_file_yields_exception_value_and_completes(self):
with tempfile.NamedTemporaryFile(suffix=".xml", delete=False, mode="wb") as tf:
tf.write(b"not a report")
junk_path = tf.name
self.addCleanup(os.remove, junk_path)
job = functools.partial(
_parse_report_file_job, config=parsedmarc.ParserConfig(offline=True)
)
results = list(parallel_map(job, [junk_path, junk_path], n_procs=2))
self.assertEqual(len(results), 2)
for path, result in results:
self.assertEqual(path, junk_path)
self.assertIsInstance(result, Exception)
class TestParallelMapBoundedLaziness(unittest.TestCase):
"""The jobs iterable must never be materialized up front. At any point
during iteration, the number of items pulled from the source should
stay within window_factor * n_procs of the number of results already
yielded, bounding memory use for very large inputs (e.g. a 20,000
message mbox)."""
def test_consumption_stays_bounded_and_results_are_complete_and_ordered(self):
n = 20
window_factor = 2
n_procs = 2
window_size = window_factor * n_procs
jobs = _CountingIterable(n)
results = []
gen = parallel_map(
_echo_job, jobs, n_procs=n_procs, window_factor=window_factor
)
for result in gen:
results.append(result)
# +1 buffer: the generator may pull one extra job before it
# can submit-then-harvest on a given step.
self.assertLessEqual(jobs.pulled, window_size + len(results) + 1)
self.assertEqual(results, list(range(n)))
class TestParallelMapShouldStop(unittest.TestCase):
"""should_stop lets a caller (e.g. a CLI handling SIGTERM) end a run
early without raising, while still yielding results already
completed in the submission window."""
def test_should_stop_ends_iteration_early(self):
n = 20
jobs = _CountingIterable(n)
results = list(
parallel_map(_echo_job, jobs, n_procs=2, should_stop=lambda: True)
)
self.assertLess(len(results), n)
self.assertLess(jobs.pulled, n)
# Results still seen so far must be a prefix of submission order.
self.assertEqual(results, list(range(len(results))))
class TestParallelMapValidation(unittest.TestCase):
"""parallel_map is a reusable helper, so it validates n_procs itself
with a clear message instead of surfacing ProcessPoolExecutor's
max_workers error later - and it must do so eagerly at the call, not
on first iteration of the returned iterator (callers that pass the
iterator elsewhere before consuming it would otherwise see the error
far from the bad argument)."""
def test_n_procs_below_one_raises_value_error_eagerly(self):
with self.assertRaises(ValueError):
parallel_map(_echo_job, [1, 2], n_procs=0)
class TestParallelMapEmptyJobs(unittest.TestCase):
"""An empty jobs iterable must return immediately without spawning a
process pool."""
def test_empty_jobs_yields_nothing_and_spawns_no_pool(self):
with patch("parsedmarc.parallel.ProcessPoolExecutor") as mock_executor:
results = list(parallel_map(_echo_job, [], n_procs=2))
self.assertEqual(results, [])
mock_executor.assert_not_called()
class TestWorkerLogging(_ParallelTestCase):
"""Worker processes must reconstruct the parent parsedmarc logger's
level and FileHandler(s) so records emitted during parsing (e.g.
parse_report_file's "Parsing <path>" debug line) aren't silently
dropped just because they happened in a child process."""
def setUp(self):
super().setUp()
from parsedmarc.log import logger as plog
self._saved_handlers = list(plog.handlers)
self._saved_level = plog.level
def tearDown(self):
from parsedmarc.log import logger as plog
for handler in list(plog.handlers):
if handler not in self._saved_handlers:
plog.removeHandler(handler)
if isinstance(handler, logging.FileHandler):
handler.close()
plog.handlers[:] = self._saved_handlers
plog.setLevel(self._saved_level)
super().tearDown()
def test_worker_debug_log_reaches_parent_file_handler(self):
from parsedmarc.log import configure_logging
with tempfile.NamedTemporaryFile(suffix=".log", delete=False) as tf:
log_path = tf.name
self.addCleanup(lambda: os.path.exists(log_path) and os.remove(log_path))
configure_logging(logging.DEBUG, log_path)
sample = SAMPLE_PATHS[0]
job = functools.partial(
_parse_report_file_job, config=parsedmarc.ParserConfig(offline=True)
)
results = list(parallel_map(job, [sample], n_procs=2))
self.assertEqual(len(results), 1)
path, report = results[0]
self.assertEqual(path, sample)
self.assertNotIsInstance(report, Exception)
with open(log_path) as f:
contents = f.read()
# parse_report_file logs `Parsing {file_path}` at DEBUG
# (parsedmarc/__init__.py) -- this line only appears if the
# worker process's reconstructed logger actually wrote to the
# parent's log file.
self.assertIn(f"Parsing {sample}", contents)
class TestInitWorkerLogging(unittest.TestCase):
"""_init_worker_logging must be usable directly as a pool initializer:
with no log files it sets the level (configure_logging still ensures
a console handler), and with log files it attaches a FileHandler per
path."""
def setUp(self):
from parsedmarc.log import logger as plog
self._saved_handlers = list(plog.handlers)
self._saved_level = plog.level
def tearDown(self):
from parsedmarc.log import logger as plog
for handler in list(plog.handlers):
if handler not in self._saved_handlers:
plog.removeHandler(handler)
if isinstance(handler, logging.FileHandler):
handler.close()
plog.handlers[:] = self._saved_handlers
plog.setLevel(self._saved_level)
def test_no_log_files_sets_level_only(self):
from parsedmarc.log import logger as plog
_init_worker_logging(logging.WARNING, [])
self.assertEqual(plog.level, logging.WARNING)
def test_log_files_attach_file_handlers(self):
from parsedmarc.log import logger as plog
with tempfile.NamedTemporaryFile(suffix=".log", delete=False) as tf:
log_path = tf.name
self.addCleanup(lambda: os.path.exists(log_path) and os.remove(log_path))
_init_worker_logging(logging.DEBUG, [log_path])
self.assertEqual(plog.level, logging.DEBUG)
file_handlers = [h for h in plog.handlers if isinstance(h, logging.FileHandler)]
self.assertTrue(any(h.baseFilename == log_path for h in file_handlers))
for h in file_handlers:
h.close()
class TestParseReportEmailJob(_ParallelTestCase):
"""_parse_report_email_job must return a ParserError as a value (never
raise it) on an invalid message."""
def test_invalid_email_returns_parser_error_value(self):
result = _parse_report_email_job(
b"not a valid email", config=parsedmarc.ParserConfig(offline=True)
)
self.assertIsInstance(result, parsedmarc.ParserError)
if __name__ == "__main__":
unittest.main(verbosity=2)