mirror of
https://github.com/Sea-Haven-Industries/procurement-ingest.git
synced 2026-10-04 10:11:57 +00:00
fix(shoc-emitter): skip empty-text comment webhooks (DEV-16) (#157)
* fix(shoc-emitter): skip empty-text comment webhooks Stop delivering #nocomment# WorkOrderComments rows as comment_added events so they no longer 400-park in the rejected queue. * style(shoc-emitter): ruff-format blank-comment handler test
This commit is contained in:
parent
67eaaa19d0
commit
8989dbde4a
6 changed files with 194 additions and 8 deletions
|
|
@ -16,6 +16,8 @@ Classification (docs/shoc-webhook-contract.md section 4):
|
||||||
NEW write_origin == "shoc-write-api" -> skip (echo guard; a no-op branch
|
NEW write_origin == "shoc-write-api" -> skip (echo guard; a no-op branch
|
||||||
today -- nothing writes that attribute yet -- kept so the phase-2
|
today -- nothing writes that attribute yet -- kept so the phase-2
|
||||||
write-back API never echoes SHOC's own writes back at it)
|
write-back API never echoes SHOC's own writes back at it)
|
||||||
|
comment_added with blank/missing text -> skip (email-processor #nocomment#
|
||||||
|
rows; SHOC requires non-empty text and would 400-park them)
|
||||||
Anything else (unknown table, comment MODIFYs) -> skip.
|
Anything else (unknown table, comment MODIFYs) -> skip.
|
||||||
"""
|
"""
|
||||||
|
|
||||||
|
|
@ -123,6 +125,33 @@ def _classify(table, event_name, new_image, old_image) -> str | None:
|
||||||
return None
|
return None
|
||||||
|
|
||||||
|
|
||||||
|
def _is_blank_comment_text(value) -> bool:
|
||||||
|
"""True when comment text is missing, empty, or whitespace-only."""
|
||||||
|
return value is None or (isinstance(value, str) and not value.strip())
|
||||||
|
|
||||||
|
|
||||||
|
def is_blank_comment_skip(record: dict) -> bool:
|
||||||
|
"""True when this stream record is a blank-text comment_added skip.
|
||||||
|
|
||||||
|
Used by the handler to emit a specific structured log for empty-comment
|
||||||
|
skips (distinct from REMOVE / echo guard / unmappable). Safe to call on
|
||||||
|
any record; returns False unless the record classifies as comment_added
|
||||||
|
with blank/missing text.
|
||||||
|
"""
|
||||||
|
event_name = record.get("eventName")
|
||||||
|
if event_name == "REMOVE":
|
||||||
|
return False
|
||||||
|
stream = record.get("dynamodb") or {}
|
||||||
|
new_image = _deserialize_image(stream.get("NewImage") or {})
|
||||||
|
if new_image.get("write_origin") == SHOC_WRITE_ORIGIN:
|
||||||
|
return False
|
||||||
|
table = _table_name(record.get("eventSourceARN") or "")
|
||||||
|
old_image = _deserialize_image(stream.get("OldImage") or {})
|
||||||
|
if _classify(table, event_name, new_image, old_image) != "work_order.comment_added":
|
||||||
|
return False
|
||||||
|
return _is_blank_comment_text(new_image.get("text"))
|
||||||
|
|
||||||
|
|
||||||
def build_event(record: dict) -> dict | None:
|
def build_event(record: dict) -> dict | None:
|
||||||
"""Map one DynamoDB stream record to a webhook envelope, or None to skip."""
|
"""Map one DynamoDB stream record to a webhook envelope, or None to skip."""
|
||||||
event_name = record.get("eventName")
|
event_name = record.get("eventName")
|
||||||
|
|
@ -138,6 +167,8 @@ def build_event(record: dict) -> dict | None:
|
||||||
if event_type is None:
|
if event_type is None:
|
||||||
return None
|
return None
|
||||||
if event_type == "work_order.comment_added":
|
if event_type == "work_order.comment_added":
|
||||||
|
if _is_blank_comment_text(new_image.get("text")):
|
||||||
|
return None
|
||||||
fields = COMMENT_DATA_FIELDS
|
fields = COMMENT_DATA_FIELDS
|
||||||
else:
|
else:
|
||||||
fields = WO_DATA_FIELDS
|
fields = WO_DATA_FIELDS
|
||||||
|
|
|
||||||
|
|
@ -99,11 +99,14 @@ def _park_rejected(event: dict, status_code: int):
|
||||||
def _batch_summary(records: int, emitted: int, skipped: int) -> str:
|
def _batch_summary(records: int, emitted: int, skipped: int) -> str:
|
||||||
"""Per-invocation emit/skip tally.
|
"""Per-invocation emit/skip tally.
|
||||||
|
|
||||||
Skips (REMOVE, echo guard, comment MODIFYs, unmappable records) return no
|
Skips (REMOVE, echo guard, blank-text comments, comment MODIFYs,
|
||||||
delivery and would otherwise be invisible: a systemic classification break
|
unmappable records) return no delivery and would otherwise be invisible:
|
||||||
-- e.g. a table rename desyncing the eventSourceARN parse -- would drop
|
a systemic classification break -- e.g. a table rename desyncing the
|
||||||
100%% of events while every batch still reports success. Logging the ratio
|
eventSourceARN parse -- would drop 100%% of events while every batch
|
||||||
makes that queryable in Logs Insights and alarmable.
|
still reports success. Logging the ratio makes that queryable in Logs
|
||||||
|
Insights and alarmable. Blank-text comment skips also emit a distinct
|
||||||
|
``shoc_skipped_empty_comment`` line so their volume stays observable
|
||||||
|
apart from other skip reasons.
|
||||||
"""
|
"""
|
||||||
return json.dumps(
|
return json.dumps(
|
||||||
{
|
{
|
||||||
|
|
@ -115,6 +118,25 @@ def _batch_summary(records: int, emitted: int, skipped: int) -> str:
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def _log_empty_comment_skip(record: dict) -> None:
|
||||||
|
"""Structured log for a blank-text comment_added skip (Logs Insights)."""
|
||||||
|
stream = record.get("dynamodb") or {}
|
||||||
|
# Attribute-value encoded NewImage; only the S forms are read for the
|
||||||
|
# log fields -- classification already decided this is a blank skip.
|
||||||
|
new_image = stream.get("NewImage") or {}
|
||||||
|
work_order_id = (new_image.get("work_order_id") or {}).get("S")
|
||||||
|
comment_id = (new_image.get("comment_id") or {}).get("S")
|
||||||
|
logger.info(
|
||||||
|
json.dumps(
|
||||||
|
{
|
||||||
|
"event": "shoc_skipped_empty_comment",
|
||||||
|
"work_order_id": work_order_id,
|
||||||
|
"comment_id": comment_id,
|
||||||
|
}
|
||||||
|
)
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
def handler(event, context):
|
def handler(event, context):
|
||||||
"""Lambda entry point. Triggered by the two WO-table stream ESMs."""
|
"""Lambda entry point. Triggered by the two WO-table stream ESMs."""
|
||||||
records = event.get("Records", [])
|
records = event.get("Records", [])
|
||||||
|
|
@ -124,6 +146,8 @@ def handler(event, context):
|
||||||
webhook_event = envelope.build_event(record)
|
webhook_event = envelope.build_event(record)
|
||||||
if webhook_event is None:
|
if webhook_event is None:
|
||||||
skipped += 1
|
skipped += 1
|
||||||
|
if envelope.is_blank_comment_skip(record):
|
||||||
|
_log_empty_comment_skip(record)
|
||||||
continue
|
continue
|
||||||
# Capture the sequence number BEFORE any delivery attempt so the
|
# Capture the sequence number BEFORE any delivery attempt so the
|
||||||
# failure paths below can never KeyError on the subscript (Open SWE
|
# failure paths below can never KeyError on the subscript (Open SWE
|
||||||
|
|
|
||||||
|
|
@ -256,6 +256,12 @@ def build_state_event(item):
|
||||||
|
|
||||||
|
|
||||||
def build_comment_event(item):
|
def build_comment_event(item):
|
||||||
|
# Skip blank/missing text the same way the live emitter does -- otherwise
|
||||||
|
# a --since / --work-order-id replay would re-POST #nocomment# rows and
|
||||||
|
# re-park them on the rejected queue.
|
||||||
|
text = item.get("text")
|
||||||
|
if text is None or (isinstance(text, str) and not text.strip()):
|
||||||
|
return None
|
||||||
occurred_at = item.get("ingested_at") or datetime.now(timezone.utc).isoformat()
|
occurred_at = item.get("ingested_at") or datetime.now(timezone.utc).isoformat()
|
||||||
delivery_id = replay_delivery_id(
|
delivery_id = replay_delivery_id(
|
||||||
COMMENTS_TABLE, item["work_order_id"], item["comment_id"], occurred_at
|
COMMENTS_TABLE, item["work_order_id"], item["comment_id"], occurred_at
|
||||||
|
|
@ -294,8 +300,21 @@ def scan_since(table, attr, since):
|
||||||
kwargs["ExclusiveStartKey"] = lek
|
kwargs["ExclusiveStartKey"] = lek
|
||||||
|
|
||||||
|
|
||||||
|
def _append_comment_events(envelopes, items):
|
||||||
|
"""Build comment envelopes, dropping blank-text skips; return skip count."""
|
||||||
|
skipped = 0
|
||||||
|
for item in items:
|
||||||
|
event = build_comment_event(item)
|
||||||
|
if event is None:
|
||||||
|
skipped += 1
|
||||||
|
continue
|
||||||
|
envelopes.append(event)
|
||||||
|
return skipped
|
||||||
|
|
||||||
|
|
||||||
def select_by_ids(dynamodb, work_order_ids, want_state, want_comments):
|
def select_by_ids(dynamodb, work_order_ids, want_state, want_comments):
|
||||||
envelopes = []
|
envelopes = []
|
||||||
|
skipped_empty = 0
|
||||||
if want_state:
|
if want_state:
|
||||||
table = dynamodb.Table(WO_TABLE)
|
table = dynamodb.Table(WO_TABLE)
|
||||||
for wid in work_order_ids:
|
for wid in work_order_ids:
|
||||||
|
|
@ -306,12 +325,17 @@ def select_by_ids(dynamodb, work_order_ids, want_state, want_comments):
|
||||||
if want_comments:
|
if want_comments:
|
||||||
table = dynamodb.Table(COMMENTS_TABLE)
|
table = dynamodb.Table(COMMENTS_TABLE)
|
||||||
for wid in work_order_ids:
|
for wid in work_order_ids:
|
||||||
envelopes.extend(build_comment_event(i) for i in query_comments(table, wid))
|
skipped_empty += _append_comment_events(
|
||||||
|
envelopes, query_comments(table, wid)
|
||||||
|
)
|
||||||
|
if skipped_empty:
|
||||||
|
print(f"skipped {skipped_empty} empty-text comment(s)")
|
||||||
return envelopes
|
return envelopes
|
||||||
|
|
||||||
|
|
||||||
def select_since(dynamodb, since, want_state, want_comments):
|
def select_since(dynamodb, since, want_state, want_comments):
|
||||||
envelopes = []
|
envelopes = []
|
||||||
|
skipped_empty = 0
|
||||||
if want_state:
|
if want_state:
|
||||||
table = dynamodb.Table(WO_TABLE)
|
table = dynamodb.Table(WO_TABLE)
|
||||||
items = list(scan_since(table, "updated_at", since))
|
items = list(scan_since(table, "updated_at", since))
|
||||||
|
|
@ -326,7 +350,9 @@ def select_since(dynamodb, since, want_state, want_comments):
|
||||||
f"[{COMMENTS_TABLE}] full table Scan (ingested_at >= {since}): "
|
f"[{COMMENTS_TABLE}] full table Scan (ingested_at >= {since}): "
|
||||||
f"{len(items)} row(s)"
|
f"{len(items)} row(s)"
|
||||||
)
|
)
|
||||||
envelopes.extend(build_comment_event(i) for i in items)
|
skipped_empty += _append_comment_events(envelopes, items)
|
||||||
|
if skipped_empty:
|
||||||
|
print(f"skipped {skipped_empty} empty-text comment(s)")
|
||||||
return envelopes
|
return envelopes
|
||||||
|
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -190,6 +190,20 @@ def test_replay_envelopes_carry_replay_true_and_exclude_s3_key():
|
||||||
assert set(comment["data"]) == set(replay.COMMENT_DATA_FIELDS)
|
assert set(comment["data"]) == set(replay.COMMENT_DATA_FIELDS)
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.parametrize(
|
||||||
|
"text",
|
||||||
|
["", " \t", None],
|
||||||
|
ids=["empty", "whitespace-only", "missing"],
|
||||||
|
)
|
||||||
|
def test_blank_comment_text_returns_none(text):
|
||||||
|
item = dict(COMMENT_ITEM)
|
||||||
|
if text is None:
|
||||||
|
del item["text"]
|
||||||
|
else:
|
||||||
|
item["text"] = text
|
||||||
|
assert replay.build_comment_event(item) is None
|
||||||
|
|
||||||
|
|
||||||
def test_delivery_id_is_stable_across_builds():
|
def test_delivery_id_is_stable_across_builds():
|
||||||
first = replay.build_comment_event(dict(COMMENT_ITEM))
|
first = replay.build_comment_event(dict(COMMENT_ITEM))
|
||||||
second = replay.build_comment_event(dict(COMMENT_ITEM))
|
second = replay.build_comment_event(dict(COMMENT_ITEM))
|
||||||
|
|
|
||||||
|
|
@ -150,9 +150,38 @@ def test_comment_insert_is_comment_added(envelope):
|
||||||
assert set(event["data"]) == set(envelope.COMMENT_DATA_FIELDS)
|
assert set(event["data"]) == set(envelope.COMMENT_DATA_FIELDS)
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.parametrize(
|
||||||
|
"overrides",
|
||||||
|
[
|
||||||
|
{"text": {"S": ""}},
|
||||||
|
{"text": {"S": " \t\n"}},
|
||||||
|
{"_drop": "text"},
|
||||||
|
],
|
||||||
|
ids=["empty", "whitespace-only", "missing"],
|
||||||
|
)
|
||||||
|
def test_blank_comment_text_is_skipped(envelope, overrides):
|
||||||
|
image = _comment_image()
|
||||||
|
if overrides.get("_drop") == "text":
|
||||||
|
del image["text"]
|
||||||
|
else:
|
||||||
|
image.update(overrides)
|
||||||
|
record = _record(COMMENTS_STREAM_ARN, "INSERT", image)
|
||||||
|
assert envelope.build_event(record) is None
|
||||||
|
assert envelope.is_blank_comment_skip(record) is True
|
||||||
|
|
||||||
|
|
||||||
|
def test_non_blank_comment_is_not_blank_skip(envelope):
|
||||||
|
record = _record(COMMENTS_STREAM_ARN, "INSERT", _comment_image())
|
||||||
|
assert envelope.build_event(record) is not None
|
||||||
|
assert envelope.is_blank_comment_skip(record) is False
|
||||||
|
|
||||||
|
|
||||||
def test_remove_is_skipped(envelope):
|
def test_remove_is_skipped(envelope):
|
||||||
assert envelope.build_event(_record(WO_STREAM_ARN, "REMOVE")) is None
|
assert envelope.build_event(_record(WO_STREAM_ARN, "REMOVE")) is None
|
||||||
assert envelope.build_event(_record(COMMENTS_STREAM_ARN, "REMOVE")) is None
|
assert envelope.build_event(_record(COMMENTS_STREAM_ARN, "REMOVE")) is None
|
||||||
|
assert (
|
||||||
|
envelope.is_blank_comment_skip(_record(COMMENTS_STREAM_ARN, "REMOVE")) is False
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
def test_comment_modify_is_skipped(envelope):
|
def test_comment_modify_is_skipped(envelope):
|
||||||
|
|
|
||||||
|
|
@ -6,7 +6,7 @@ report_batch_item_failures (so the ESM retries from it, in order) and later
|
||||||
records are never attempted; earlier in-batch successes are not re-delivered.
|
records are never attempted; earlier in-batch successes are not re-delivered.
|
||||||
Non-retryable rejections park the full envelope on the rejected queue and the
|
Non-retryable rejections park the full envelope on the rejected queue and the
|
||||||
loop CONTINUES (a contract bug must never block the shard). Skip records
|
loop CONTINUES (a contract bug must never block the shard). Skip records
|
||||||
(REMOVE / echo guard) produce no delivery attempts at all.
|
(REMOVE / echo guard / blank-text comments) produce no delivery attempts at all.
|
||||||
|
|
||||||
delivery.deliver is monkeypatched on the handler's own sibling instance (the
|
delivery.deliver is monkeypatched on the handler's own sibling instance (the
|
||||||
handler imports it by bare name); envelope.build_event runs for real, so
|
handler imports it by bare name); envelope.build_event runs for real, so
|
||||||
|
|
@ -171,6 +171,68 @@ def test_skip_records_produce_no_delivery_calls(monkeypatch, handler_mod):
|
||||||
assert fake_sqs.sent == []
|
assert fake_sqs.sent == []
|
||||||
|
|
||||||
|
|
||||||
|
COMMENTS_STREAM_ARN = (
|
||||||
|
"arn:aws:dynamodb:us-east-1:011934824531:table/WorkOrderComments"
|
||||||
|
"/stream/2026-07-23T00:00:00.000"
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def _comment_record(event_id, sequence_number, **image_overrides):
|
||||||
|
image = {
|
||||||
|
"work_order_id": {"S": f"wo-{event_id}"},
|
||||||
|
"comment_id": {"S": f"wo-{event_id}#nocomment#abc123"},
|
||||||
|
"record_type": {"S": "new_work_order"},
|
||||||
|
"commenter": {"S": ""},
|
||||||
|
"text": {"S": ""},
|
||||||
|
"created_at": {"S": "2026-07-16T14:05:00+00:00"},
|
||||||
|
"ingested_at": {"S": "2026-07-16T14:05:00+00:00"},
|
||||||
|
}
|
||||||
|
image.update(image_overrides)
|
||||||
|
return {
|
||||||
|
"eventID": event_id,
|
||||||
|
"eventName": "INSERT",
|
||||||
|
"eventSource": "aws:dynamodb",
|
||||||
|
"eventSourceARN": COMMENTS_STREAM_ARN,
|
||||||
|
"dynamodb": {
|
||||||
|
"ApproximateCreationDateTime": 1784642602.0,
|
||||||
|
"SequenceNumber": sequence_number,
|
||||||
|
"StreamViewType": "NEW_AND_OLD_IMAGES",
|
||||||
|
"NewImage": image,
|
||||||
|
},
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
def test_blank_comment_skip_logs_and_produces_no_delivery(
|
||||||
|
monkeypatch, handler_mod, caplog
|
||||||
|
):
|
||||||
|
event = {
|
||||||
|
"Records": [
|
||||||
|
_comment_record("evt-empty", "501"),
|
||||||
|
_wo_record("evt-wo", "502"),
|
||||||
|
]
|
||||||
|
}
|
||||||
|
attempted, fake_sqs = _wire(
|
||||||
|
monkeypatch, handler_mod, {"evt-wo": ("delivered", 200)}
|
||||||
|
)
|
||||||
|
|
||||||
|
with caplog.at_level("INFO"):
|
||||||
|
result = handler_mod.handler(event, None)
|
||||||
|
|
||||||
|
assert result == {"batchItemFailures": []}
|
||||||
|
assert attempted == ["evt-wo"]
|
||||||
|
assert fake_sqs.sent == []
|
||||||
|
|
||||||
|
skip_logs = [
|
||||||
|
json.loads(rec.message)
|
||||||
|
for rec in caplog.records
|
||||||
|
if rec.message.startswith("{") and '"shoc_skipped_empty_comment"' in rec.message
|
||||||
|
]
|
||||||
|
assert len(skip_logs) == 1
|
||||||
|
assert skip_logs[0]["event"] == "shoc_skipped_empty_comment"
|
||||||
|
assert skip_logs[0]["work_order_id"] == "wo-evt-empty"
|
||||||
|
assert skip_logs[0]["comment_id"] == "wo-evt-empty#nocomment#abc123"
|
||||||
|
|
||||||
|
|
||||||
def test_unexpected_error_reports_record_not_whole_batch(monkeypatch, handler_mod):
|
def test_unexpected_error_reports_record_not_whole_batch(monkeypatch, handler_mod):
|
||||||
# An unexpected exception (here: SQS park failure on a 4xx rejection) must
|
# An unexpected exception (here: SQS park failure on a 4xx rejection) must
|
||||||
# NOT escape the loop -- that would fail the invocation and make the ESM
|
# NOT escape the loop -- that would fail the invocation and make the ESM
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue