Backfill combined DKIM/SPF fields automatically at startup

migrate_indexes() now backfills dkim_results_combined and
spf_results_combined on aggregate documents saved by older versions,
so ES/OS users get historical data in the reworked alignment tables
without running the documented _update_by_query by hand. The backfill
is submitted as a non-blocking background task
(wait_for_completion=false, conflicts=proceed) guarded by a cheap
count query, making repeated startups a fast no-op once an index is
backfilled; any cluster error is logged as a warning and retried at
the next startup rather than raised. The manual command remains
documented for users who upgrade dashboards without pointing the new
parsedmarc at the cluster or who want to control write-load timing.

The legacy published_policy.fo long-to-text reindex migration in the
OpenSearch module is kept ahead of the new backfill, for clusters
upgraded from very old data.

Verified end-to-end against the live dev environment: a real CLI
startup backfilled 9 stripped OpenSearch documents (logged with task
ID) while the already-backfilled Elasticsearch side stayed silent,
and a second startup was silent on both engines.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
Sean Whalen
2026-07-22 09:14:59 -04:00
co-authored by Claude Fable 5
parent e892de794b
commit bbee148d2a
6 changed files with 424 additions and 32 deletions
+1 -1
View File
@@ -12,7 +12,7 @@
- **The Elasticsearch/OpenSearch aggregate dashboards' over-time charts (and the Grafana ES dashboard's summary pies and time series) bucketed on the multi-valued `date_range` field**; a date histogram counts a report once per value, double-counting any report whose begin and end dates fall in different buckets. All date histograms and time-range filters now use the single-valued `date_begin`, matching the report-begin semantics of the PostgreSQL (`begin_date`) and Splunk (`_time` = interval begin) dashboards.
- **Aggregate-report policy and authentication result words are now normalized to lowercase** ([#288](https://github.com/domainaware/parsedmarc/issues/288)): reporters that emit mixed-case values such as `Pass` no longer create duplicate result categories in outputs and dashboards.
- **The results email (SMTP and Microsoft Graph) is no longer sent when no reports were parsed** ([#200](https://github.com/domainaware/parsedmarc/issues/200)): previously an empty run — an empty inbox, or one where every message was invalid — still emailed a zip of headers-only CSVs. The email step is now skipped with an INFO log when the run produced no aggregate, failure, or SMTP TLS reports.
- **The DKIM/SPF alignment-detail tables on the Kibana/OpenSearch Dashboards, Grafana (Elasticsearch), and Splunk aggregate dashboards no longer show a selector × domain × result cross-product** ([#169](https://github.com/domainaware/parsedmarc/issues/169)). Elasticsearch/OpenSearch flatten the `dkim_results`/`spf_results` object arrays (they are dynamic-mapped as `object`, not `nested`), so stacking terms aggregations on their subfields produced every combination of selector/domain/result across a report's signatures, each repeating the full message count. Aggregate documents now also carry `dkim_results_combined` and `spf_results_combined` — one `"selector / domain / result"` (`"scope / domain / result"`) string per auth result — and the dashboards aggregate those instead; the Splunk detail panels now pair the values with `mvzip`/`mvexpand`. Documents saved by older versions can be backfilled with a documented `_update_by_query` command (see the Elasticsearch docs page). The Grafana "DKIM Alignment Details" panel's dmarcian.com DKIM-checker data link was removed because it required the separate domain/selector columns. The SMTP TLS visualizations have the same class of defect and are tracked separately.
- **The DKIM/SPF alignment-detail tables on the Kibana/OpenSearch Dashboards, Grafana (Elasticsearch), and Splunk aggregate dashboards no longer show a selector × domain × result cross-product** ([#169](https://github.com/domainaware/parsedmarc/issues/169)). Elasticsearch/OpenSearch flatten the `dkim_results`/`spf_results` object arrays (they are dynamic-mapped as `object`, not `nested`), so stacking terms aggregations on their subfields produced every combination of selector/domain/result across a report's signatures, each repeating the full message count. Aggregate documents now also carry `dkim_results_combined` and `spf_results_combined` — one `"selector / domain / result"` (`"scope / domain / result"`) string per auth result — and the dashboards aggregate those instead; the Splunk detail panels now pair the values with `mvzip`/`mvexpand`. Documents saved by older versions are now backfilled automatically at startup (a non-blocking, idempotent background task); the documented `_update_by_query` command (see the Elasticsearch docs page) remains available for running the backfill manually. The Grafana "DKIM Alignment Details" panel's dmarcian.com DKIM-checker data link was removed because it required the separate domain/selector columns. The SMTP TLS visualizations have the same class of defect and are tracked separately.
- **Corrected the dead `_SPFResult.results` (plural) field declaration to `result`**, matching what was always written to it.
## 10.2.4
+22 -8
View File
@@ -235,14 +235,28 @@ result paired, which the dashboards' alignment-detail tables aggregate on.
Reports saved by older versions lack these fields and will not appear in
those tables.
Running the following once per cluster backfills the fields on existing
documents. It is idempotent (documents that already have the fields are
skipped), so it is safe to re-run. It works identically on OpenSearch;
just adjust the URL and credentials. The query matches only documents
that have at least one DKIM or SPF auth result and lack the corresponding
combined field; documents with no auth results are skipped, because an
`exists` query cannot see an empty array, and for search purposes an
empty `dkim_results_combined` is identical to an absent one.
parsedmarc now backfills this automatically. On startup, it runs a cheap
count query against each configured aggregate index pattern to check for
documents that have DKIM or SPF results but are missing the corresponding
combined field. If any are found, it submits the backfill as a background
`_update_by_query` task (`wait_for_completion=false`), so startup is never
blocked on it; progress is logged, including the task ID. The check itself
is idempotent — once an index is fully backfilled, later startups see a
count of 0 and log nothing further — and it works the same way on
OpenSearch. Any error talking to the cluster (for example, no indexes yet
on a fresh install) is logged as a warning and retried on the next startup,
rather than aborting parsedmarc.
If you upgrade the dashboards without pointing the new parsedmarc version
at the cluster, or you'd rather control when the write load happens, you
can still run the backfill manually. It is idempotent (documents that
already have the fields are skipped), so it is safe to re-run. It works
identically on OpenSearch; just adjust the URL and credentials. The query
matches only documents that have at least one DKIM or SPF auth result and
lack the corresponding combined field; documents with no auth results are
skipped, because an `exists` query cannot see an empty array, and for
search purposes an empty `dkim_results_combined` is identical to an
absent one.
```bash
curl -X POST "http://localhost:9200/dmarc_aggregate*/_update_by_query?conflicts=proceed&wait_for_completion=false" \
+111 -10
View File
@@ -41,6 +41,64 @@ _SERVERLESS = False
# settings (e.g. ``refresh_interval``) are accepted and pass through.
_SERVERLESS_REJECTED_SETTINGS = frozenset({"number_of_shards", "number_of_replicas"})
# Guard query for the dkim_results_combined/spf_results_combined backfill
# (see ``migrate_indexes``). Matches only documents that have at least one
# DKIM or SPF auth result and are missing the corresponding combined field.
# Empty arrays are invisible to ``exists``, so documents with zero DKIM/SPF
# results are correctly skipped, making this idempotent: once a document is
# backfilled it no longer matches, and re-running against an already
# up-to-date index counts 0.
_COMBINED_BACKFILL_QUERY: dict[str, Any] = {
"bool": {
"minimum_should_match": 1,
"should": [
{
"bool": {
"must": [{"exists": {"field": "dkim_results.domain"}}],
"must_not": [{"exists": {"field": "dkim_results_combined"}}],
}
},
{
"bool": {
"must": [{"exists": {"field": "spf_results.domain"}}],
"must_not": [{"exists": {"field": "spf_results_combined"}}],
}
},
],
}
}
# Painless script that (re)derives dkim_results_combined/spf_results_combined
# from dkim_results/spf_results, matching the format written by
# save_aggregate_report_to_elasticsearch(): "{selector} / {domain} / {result}"
# per DKIM result and "{scope} / {domain} / {result}" per SPF result.
_COMBINED_BACKFILL_SCRIPT = (
"List dk = new ArrayList(); "
"def dr = ctx._source.dkim_results; "
"if (dr != null) { "
"if (!(dr instanceof List)) { dr = [dr]; } "
"for (e in dr) { "
"if (e == null) { continue; } "
'def sel = e.selector != null ? e.selector : "none"; '
'def dom = e.domain != null ? e.domain : "none"; '
'def res = e.result != null ? e.result : "none"; '
'dk.add(sel + " / " + dom + " / " + res); '
"} } "
"ctx._source.dkim_results_combined = dk; "
"List sp = new ArrayList(); "
"def sr = ctx._source.spf_results; "
"if (sr != null) { "
"if (!(sr instanceof List)) { sr = [sr]; } "
"for (e in sr) { "
"if (e == null) { continue; } "
'def sc = e.scope != null ? e.scope : "mfrom"; '
'def dom = e.domain != null ? e.domain : "none"; '
'def res = e.result != null ? e.result : (e.results != null ? e.results : "none"); '
'sp.add(sc + " / " + dom + " / " + res); '
"} } "
"ctx._source.spf_results_combined = sp;"
)
class _PolicyOverride(InnerDoc):
# The elasticsearch.dsl 8.x type stubs use dataclass_transform and only
@@ -487,23 +545,66 @@ def migrate_indexes(
failure_indexes: list[str] | None = None,
):
"""
Updates index mappings
Backfills the ``dkim_results_combined``/``spf_results_combined`` fields
(added for issue #169) on aggregate report documents that were saved
before those fields existed.
This is a no-op kept for API compatibility (``cli.py`` calls it on
startup). The only migration this function ever performed was
re-typing ``published_policy.fo`` from ``long`` to ``text``, which
applied exclusively to indices still carrying the legacy
Elasticsearch 6-era ``"doc"`` mapping type. The 8.x client can only
reach servers (Elasticsearch 8.x/9.x) whose indices were created on
Elasticsearch 7.x or later and are therefore typeless, so that
migration path is unreachable and has been removed.
For each name in ``aggregate_indexes``, this submits an
``update_by_query`` against the ``f"{name}*"`` index pattern (the real
indexes are date-suffixed) as a non-blocking background task
(``wait_for_completion=False``), so it never delays parsedmarc startup.
Submission is guarded by a cheap ``count`` query that only matches
documents with DKIM/SPF results but no combined field, so once an index
is fully backfilled, later calls are a fast no-op. Any error talking to
the cluster (e.g. no indexes yet on a fresh install, or a transient
connection issue) is caught and logged as a warning rather than raised;
the backfill is simply retried on the next startup, and the manual
``_update_by_query`` command documented in
``docs/source/elasticsearch.md`` remains available in the meantime.
Args:
aggregate_indexes (list): A list of aggregate index names
(accepted for API compatibility; unused)
failure_indexes (list): A list of failure index names
(accepted for API compatibility; unused)
"""
if not aggregate_indexes:
return
client = connections.get_connection()
for name in aggregate_indexes:
pattern = f"{name}*"
try:
count_response = client.count(
index=pattern,
query=_COMBINED_BACKFILL_QUERY,
ignore_unavailable=True,
allow_no_indices=True,
)
count = count_response["count"]
if not count:
continue
update_response = client.update_by_query(
index=pattern,
query=_COMBINED_BACKFILL_QUERY,
script={"source": _COMBINED_BACKFILL_SCRIPT, "lang": "painless"},
conflicts="proceed",
wait_for_completion=False,
ignore_unavailable=True,
allow_no_indices=True,
)
task_id = update_response.get("task")
logger.info(
"Backfilling dkim_results_combined/spf_results_combined on "
f"{count} existing documents in {pattern} (task {task_id})"
)
except Exception as e:
logger.warning(
"Failed to check/submit the dkim_results_combined/"
f"spf_results_combined backfill for {pattern}: {e}. This "
"will be retried at the next startup; the manual "
"_update_by_query command in the documentation remains "
"available in the meantime."
)
def save_aggregate_report_to_elasticsearch(
+124 -3
View File
@@ -34,6 +34,65 @@ class OpenSearchError(Exception):
"""Raised when an OpenSearch error occurs"""
# Guard query for the dkim_results_combined/spf_results_combined backfill
# (see ``migrate_indexes``). Matches only documents that have at least one
# DKIM or SPF auth result and are missing the corresponding combined field.
# Empty arrays are invisible to ``exists``, so documents with zero DKIM/SPF
# results are correctly skipped, making this idempotent: once a document is
# backfilled it no longer matches, and re-running against an already
# up-to-date index counts 0.
_COMBINED_BACKFILL_QUERY: dict[str, Any] = {
"bool": {
"minimum_should_match": 1,
"should": [
{
"bool": {
"must": [{"exists": {"field": "dkim_results.domain"}}],
"must_not": [{"exists": {"field": "dkim_results_combined"}}],
}
},
{
"bool": {
"must": [{"exists": {"field": "spf_results.domain"}}],
"must_not": [{"exists": {"field": "spf_results_combined"}}],
}
},
],
}
}
# Painless script that (re)derives dkim_results_combined/spf_results_combined
# from dkim_results/spf_results, matching the format written by
# save_aggregate_report_to_opensearch(): "{selector} / {domain} / {result}"
# per DKIM result and "{scope} / {domain} / {result}" per SPF result.
_COMBINED_BACKFILL_SCRIPT = (
"List dk = new ArrayList(); "
"def dr = ctx._source.dkim_results; "
"if (dr != null) { "
"if (!(dr instanceof List)) { dr = [dr]; } "
"for (e in dr) { "
"if (e == null) { continue; } "
'def sel = e.selector != null ? e.selector : "none"; '
'def dom = e.domain != null ? e.domain : "none"; '
'def res = e.result != null ? e.result : "none"; '
'dk.add(sel + " / " + dom + " / " + res); '
"} } "
"ctx._source.dkim_results_combined = dk; "
"List sp = new ArrayList(); "
"def sr = ctx._source.spf_results; "
"if (sr != null) { "
"if (!(sr instanceof List)) { sr = [sr]; } "
"for (e in sr) { "
"if (e == null) { continue; } "
'def sc = e.scope != null ? e.scope : "mfrom"; '
'def dom = e.domain != null ? e.domain : "none"; '
'def res = e.result != null ? e.result : (e.results != null ? e.results : "none"); '
'sp.add(sc + " / " + dom + " / " + res); '
"} } "
"ctx._source.spf_results_combined = sp;"
)
class _PolicyOverride(InnerDoc):
type = Text()
comment = Text()
@@ -406,7 +465,27 @@ def migrate_indexes(
failure_indexes: list[str] | None = None,
):
"""
Updates index mappings
Runs index migrations and backfills.
First, the legacy ``published_policy.fo`` migration: indexes where that
field was mapped as ``long`` (data indexed by very old parsedmarc
releases under the Elasticsearch 6-era ``doc`` mapping type) are rebuilt
as a ``-v2`` index with the text/keyword shape.
Second, the ``dkim_results_combined``/``spf_results_combined`` backfill
(added for issue #169) for aggregate report documents that were saved
before those fields existed. For each name in ``aggregate_indexes``,
this submits an ``update_by_query`` against the ``f"{name}*"`` index
pattern (the real indexes are date-suffixed) as a non-blocking
background task (``wait_for_completion=False``), so it never delays
parsedmarc startup. Submission is guarded by a cheap ``count`` query
that only matches documents with DKIM/SPF results but no combined
field, so once an index is fully backfilled, later calls are a fast
no-op. Any error talking to the cluster (e.g. no indexes yet on a
fresh install, or a transient connection issue) is caught and logged
as a warning rather than raised; the backfill is simply retried on the
next startup, and the manual ``_update_by_query`` command documented
in ``docs/source/elasticsearch.md`` remains available in the meantime.
Args:
aggregate_indexes (list): A list of aggregate index names
@@ -414,9 +493,10 @@ def migrate_indexes(
(accepted for API compatibility; no migrations are
currently needed for failure indexes)
"""
if not aggregate_indexes:
return
version = 2
if aggregate_indexes is None:
aggregate_indexes = []
for aggregate_index_name in aggregate_indexes:
if not Index(aggregate_index_name).exists():
continue
@@ -446,6 +526,47 @@ def migrate_indexes(
reindex(connections.get_connection(), aggregate_index_name, new_index_name)
Index(aggregate_index_name).delete()
client = connections.get_connection()
for name in aggregate_indexes:
pattern = f"{name}*"
try:
count_response = client.count(
index=pattern,
body={"query": _COMBINED_BACKFILL_QUERY},
ignore_unavailable=True,
allow_no_indices=True,
)
count = count_response["count"]
if not count:
continue
update_response = client.update_by_query(
index=pattern,
body={
"query": _COMBINED_BACKFILL_QUERY,
"script": {
"source": _COMBINED_BACKFILL_SCRIPT,
"lang": "painless",
},
},
conflicts="proceed",
wait_for_completion=False,
ignore_unavailable=True,
allow_no_indices=True,
)
task_id = update_response.get("task")
logger.info(
"Backfilling dkim_results_combined/spf_results_combined on "
f"{count} existing documents in {pattern} (task {task_id})"
)
except Exception as e:
logger.warning(
"Failed to check/submit the dkim_results_combined/"
f"spf_results_combined backfill for {pattern}: {e}. This "
"will be retried at the next startup; the manual "
"_update_by_query command in the documentation remains "
"available in the meantime."
)
def save_aggregate_report_to_opensearch(
aggregate_report: dict[str, Any],
+64
View File
@@ -16,6 +16,7 @@ from parsedmarc.elastic import (
AlreadySaved,
ElasticsearchError,
create_indexes,
migrate_indexes,
save_aggregate_report_to_elasticsearch,
save_failure_report_to_elasticsearch,
save_smtp_tls_report_to_elasticsearch,
@@ -395,6 +396,69 @@ class TestCreateIndexesServerless(unittest.TestCase):
mock_index.create.assert_called_once()
# ---------------------------------------------------------------------------
# migrate_indexes
# ---------------------------------------------------------------------------
class TestMigrateIndexes(unittest.TestCase):
"""migrate_indexes backfills dkim_results_combined/spf_results_combined
(issue #169) on pre-existing aggregate documents as a non-blocking
background task. It is guarded by a cheap count() query so repeated
startups against an already-backfilled index are a no-op, and any SDK
error is caught and logged rather than raised, so it never blocks
parsedmarc startup."""
def test_backfill_submitted_when_old_docs_exist(self):
with patch("parsedmarc.elastic.connections.get_connection") as mock_get_conn:
mock_client = MagicMock()
mock_client.count.return_value = {"count": 42}
mock_get_conn.return_value = mock_client
migrate_indexes(aggregate_indexes=["dmarc_aggregate"])
mock_client.update_by_query.assert_called_once()
kwargs = mock_client.update_by_query.call_args.kwargs
self.assertEqual(kwargs["index"], "dmarc_aggregate*")
self.assertEqual(kwargs["conflicts"], "proceed")
self.assertFalse(kwargs["wait_for_completion"])
self.assertEqual(kwargs["query"], elastic_module._COMBINED_BACKFILL_QUERY)
script_source = kwargs["script"]["source"]
self.assertIn("ctx._source.dkim_results_combined", script_source)
self.assertIn("ctx._source.spf_results_combined", script_source)
# The count() guard query also targets the date-suffixed pattern.
count_kwargs = mock_client.count.call_args.kwargs
self.assertEqual(count_kwargs["index"], "dmarc_aggregate*")
self.assertEqual(count_kwargs["query"], elastic_module._COMBINED_BACKFILL_QUERY)
def test_backfill_skipped_when_no_old_docs(self):
with patch("parsedmarc.elastic.connections.get_connection") as mock_get_conn:
mock_client = MagicMock()
mock_client.count.return_value = {"count": 0}
mock_get_conn.return_value = mock_client
migrate_indexes(aggregate_indexes=["dmarc_aggregate"])
mock_client.update_by_query.assert_not_called()
def test_backfill_skipped_when_no_aggregate_indexes(self):
with patch("parsedmarc.elastic.connections.get_connection") as mock_get_conn:
migrate_indexes()
migrate_indexes(aggregate_indexes=None)
mock_get_conn.assert_not_called()
def test_backfill_failure_does_not_raise(self):
with patch("parsedmarc.elastic.connections.get_connection") as mock_get_conn:
mock_client = MagicMock()
mock_client.count.side_effect = RuntimeError("cluster unreachable")
mock_get_conn.return_value = mock_client
with self.assertLogs("parsedmarc.log", level="WARNING") as cm:
migrate_indexes(aggregate_indexes=["dmarc_aggregate"])
self.assertTrue(any("cluster unreachable" in msg for msg in cm.output))
mock_client.update_by_query.assert_not_called()
# ---------------------------------------------------------------------------
# save_aggregate_report_to_elasticsearch
# ---------------------------------------------------------------------------
+102 -10
View File
@@ -376,16 +376,100 @@ class TestCreateIndexes(unittest.TestCase):
class TestMigrateIndexes(unittest.TestCase):
"""migrate_indexes backfills dkim_results_combined/spf_results_combined
(issue #169) on pre-existing aggregate documents as a non-blocking
background task. It is guarded by a cheap count() query so repeated
startups against an already-backfilled index are a no-op, and any SDK
error is caught and logged rather than raised, so it never blocks
parsedmarc startup."""
def test_backfill_submitted_when_old_docs_exist(self):
with (
patch("parsedmarc.opensearch.Index") as mock_index_cls,
patch("parsedmarc.opensearch.connections.get_connection") as mock_get_conn,
):
# The legacy fo migration that runs first sees no base index.
mock_index_cls.return_value.exists.return_value = False
mock_client = MagicMock()
mock_client.count.return_value = {"count": 42}
mock_get_conn.return_value = mock_client
migrate_indexes(aggregate_indexes=["dmarc_aggregate"])
mock_client.update_by_query.assert_called_once()
kwargs = mock_client.update_by_query.call_args.kwargs
self.assertEqual(kwargs["index"], "dmarc_aggregate*")
self.assertEqual(kwargs["conflicts"], "proceed")
self.assertFalse(kwargs["wait_for_completion"])
self.assertEqual(
kwargs["body"]["query"], opensearch_module._COMBINED_BACKFILL_QUERY
)
script_source = kwargs["body"]["script"]["source"]
self.assertIn("ctx._source.dkim_results_combined", script_source)
self.assertIn("ctx._source.spf_results_combined", script_source)
# The count() guard query also targets the date-suffixed pattern.
count_kwargs = mock_client.count.call_args.kwargs
self.assertEqual(count_kwargs["index"], "dmarc_aggregate*")
self.assertEqual(
count_kwargs["body"]["query"], opensearch_module._COMBINED_BACKFILL_QUERY
)
def test_backfill_skipped_when_no_old_docs(self):
with (
patch("parsedmarc.opensearch.Index") as mock_index_cls,
patch("parsedmarc.opensearch.connections.get_connection") as mock_get_conn,
):
mock_index_cls.return_value.exists.return_value = False
mock_client = MagicMock()
mock_client.count.return_value = {"count": 0}
mock_get_conn.return_value = mock_client
migrate_indexes(aggregate_indexes=["dmarc_aggregate"])
mock_client.update_by_query.assert_not_called()
def test_backfill_skipped_when_no_aggregate_indexes(self):
with patch("parsedmarc.opensearch.connections.get_connection") as mock_get_conn:
migrate_indexes()
migrate_indexes(aggregate_indexes=None)
mock_get_conn.assert_not_called()
def test_backfill_failure_does_not_raise(self):
with (
patch("parsedmarc.opensearch.Index") as mock_index_cls,
patch("parsedmarc.opensearch.connections.get_connection") as mock_get_conn,
):
mock_index_cls.return_value.exists.return_value = False
mock_client = MagicMock()
mock_client.count.side_effect = RuntimeError("cluster unreachable")
mock_get_conn.return_value = mock_client
with self.assertLogs("parsedmarc.log", level="WARNING") as cm:
migrate_indexes(aggregate_indexes=["dmarc_aggregate"])
self.assertTrue(any("cluster unreachable" in msg for msg in cm.output))
mock_client.update_by_query.assert_not_called()
class TestMigrateIndexesFoMigration(unittest.TestCase):
"""The legacy `published_policy.fo` field was mapped as `long` in
older indexes. migrate_indexes detects that and rebuilds the index
with the text/keyword shape. The branch is gnarly; a regression
would silently leave old data un-migrated."""
would silently leave old data un-migrated. Each test stubs the
combined-field backfill that now runs afterwards in the same call
(count 0 → no-op)."""
def test_no_indexes_is_noop(self):
migrate_indexes() # Should not raise
@staticmethod
def _noop_backfill_client():
client = MagicMock()
client.count.return_value = {"count": 0}
return client
def test_skips_non_existent_index(self):
with patch("parsedmarc.opensearch.Index") as mock_index_cls:
with (
patch("parsedmarc.opensearch.Index") as mock_index_cls,
patch("parsedmarc.opensearch.connections.get_connection") as mock_get_conn,
):
mock_get_conn.return_value = self._noop_backfill_client()
mock_index_cls.return_value.exists.return_value = False
migrate_indexes(aggregate_indexes=["missing"])
# exists() returned False — no field_mapping fetch.
@@ -394,12 +478,16 @@ class TestMigrateIndexes(unittest.TestCase):
def test_skips_when_doc_mapping_absent(self):
"""An index that has 'fo' but not under the 'doc' type
(e.g., empty index with default mapping) is left alone."""
with patch("parsedmarc.opensearch.Index") as mock_index_cls:
with (
patch("parsedmarc.opensearch.Index") as mock_index_cls,
patch("parsedmarc.opensearch.connections.get_connection") as mock_get_conn,
patch("parsedmarc.opensearch.reindex") as mock_reindex,
):
mock_get_conn.return_value = self._noop_backfill_client()
idx = mock_index_cls.return_value
idx.exists.return_value = True
idx.get_field_mapping.return_value = {"some_key": {"mappings": {}}}
with patch("parsedmarc.opensearch.reindex") as mock_reindex:
migrate_indexes(aggregate_indexes=["dmarc_aggregate-2023-01-01"])
migrate_indexes(aggregate_indexes=["dmarc_aggregate-2023-01-01"])
mock_reindex.assert_not_called()
def test_migrates_when_fo_is_long(self):
@@ -411,6 +499,8 @@ class TestMigrateIndexes(unittest.TestCase):
patch("parsedmarc.opensearch.reindex") as mock_reindex,
patch("parsedmarc.opensearch.connections.get_connection") as mock_get_conn,
):
mock_client = self._noop_backfill_client()
mock_get_conn.return_value = mock_client
idx = mock_index_cls.return_value
idx.exists.return_value = True
idx.get_field_mapping.return_value = {
@@ -423,16 +513,18 @@ class TestMigrateIndexes(unittest.TestCase):
}
}
migrate_indexes(aggregate_indexes=["dmarc_aggregate-2023-01-01"])
# reindex called from old → new (v2) index.
# reindex called from old → new (v2) index, with the client from
# connections.get_connection().
mock_reindex.assert_called_once()
# connections.get_connection consulted to get the ES client.
mock_get_conn.assert_called_once()
self.assertIs(mock_reindex.call_args.args[0], mock_client)
def test_skips_when_fo_already_text(self):
with (
patch("parsedmarc.opensearch.Index") as mock_index_cls,
patch("parsedmarc.opensearch.connections.get_connection") as mock_get_conn,
patch("parsedmarc.opensearch.reindex") as mock_reindex,
):
mock_get_conn.return_value = self._noop_backfill_client()
idx = mock_index_cls.return_value
idx.exists.return_value = True
idx.get_field_mapping.return_value = {