fix(jobs): run delayed close and reminder deliveries

Wall-clock skip windows dropped the only weekly SQS attempt when Scheduler already fired in Eastern time. Dev schedules stay disabled.
This commit is contained in:
Adam Moussa 2026-09-21 14:48:59 -04:00
parent 3e9a1c60bf
commit d7ee6a3b74
No known key found for this signature in database
8 changed files with 99 additions and 178 deletions

View file

@ -1,10 +1,5 @@
from datetime import datetime
from zoneinfo import ZoneInfo
from shared.db import current_week, get_form_status, set_form_status from shared.db import current_week, get_form_status, set_form_status
EASTERN = ZoneInfo("America/New_York")
def _dispatch_job(payload: dict) -> None: def _dispatch_job(payload: dict) -> None:
from server.jobs import enqueue_job from server.jobs import enqueue_job
@ -13,18 +8,9 @@ def _dispatch_job(payload: dict) -> None:
def lambda_handler(event, context): def lambda_handler(event, context):
now_et = datetime.now(EASTERN) # Scheduler fires Thursday 23:59 America/New_York. Do not skip on wall
# EventBridge can fire slightly after midnight ET; accept Thu 23:xx or Fri 00–03 # clock: a delayed SQS delivery must still close the week. already-closed
# ET so a delayed cron still closes the form. Idempotency: already-closed is a no-op. # is the idempotency gate for redelivery.
in_close_window = (now_et.weekday() == 3 and now_et.hour == 23) or (
now_et.weekday() == 4 and now_et.hour < 4
)
if not in_close_window:
return {
"status": "skipped",
"reason": "outside close window (must be Thu 23:xx or Fri 00–03 ET)",
}
week = event.get("week", current_week()) week = event.get("week", current_week())
status = get_form_status(week) status = get_form_status(week)

View file

@ -1,8 +1,6 @@
import json import json
import os import os
from datetime import datetime
from decimal import Decimal from decimal import Decimal
from zoneinfo import ZoneInfo
from shared.db import current_week, get_orders, get_roster, get_settings, get_summary from shared.db import current_week, get_orders, get_roster, get_settings, get_summary
from shared.slack import post_channel_message, send_dm from shared.slack import post_channel_message, send_dm
@ -28,10 +26,6 @@ def lambda_handler(event, context):
elif event_type == "orders_aggregated": elif event_type == "orders_aggregated":
return handle_orders_aggregated(event) return handle_orders_aggregated(event)
elif event_type == "reminder": elif event_type == "reminder":
# Guard against duplicate triggers from dual EST/EDT schedules
now_et = datetime.now(ZoneInfo("America/New_York"))
if now_et.hour != 10 or now_et.weekday() != 3: # 10am Thursday
return {"status": "skipped", "reason": "outside reminder window"}
return handle_reminder(event) return handle_reminder(event)
elif event_type == "order_confirmed": elif event_type == "order_confirmed":
return handle_order_confirmed(event) return handle_order_confirmed(event)
@ -187,11 +181,9 @@ def handle_order_confirmed(event):
def handle_reminder(event): def handle_reminder(event):
# Guard against duplicate triggers from dual EST/EDT schedules # Scheduler fires Thursday 10:00 America/New_York. Delayed SQS delivery
now_et = datetime.now(ZoneInfo("America/New_York")) # must still DM anyone who has not ordered; already-ordered emails are
if now_et.hour != 10 or now_et.weekday() != 3: # 10am Thursday # the idempotency gate for redelivery.
return {"status": "skipped", "reason": "outside reminder window"}
week = event.get("week", current_week()) week = event.get("week", current_week())
form_url = os.environ.get("FORM_URL", "") form_url = os.environ.get("FORM_URL", "")

View file

@ -49,6 +49,11 @@ def main() -> None:
payload = json.loads(msg["Body"]) payload = json.loads(msg["Body"])
result = run_job(payload) result = run_job(payload)
logger.info("job result %s", result) logger.info("job result %s", result)
if result.get("status") == "skipped":
# Leave the message visible after timeout so a later
# receive can run it. Deleting here dropped the only
# weekly close/reminder when delivery landed late.
continue
sqs.delete_message(QueueUrl=queue_url, ReceiptHandle=receipt) sqs.delete_message(QueueUrl=queue_url, ReceiptHandle=receipt)
except Exception: except Exception:
logger.exception("job failed; leaving message for retry") logger.exception("job failed; leaving message for retry")

View file

@ -32,6 +32,7 @@ resource "aws_scheduler_schedule" "jobs" {
description = each.value.description description = each.value.description
schedule_expression = each.value.schedule schedule_expression = each.value.schedule
schedule_expression_timezone = "America/New_York" schedule_expression_timezone = "America/New_York"
state = local.is_prod ? "ENABLED" : "DISABLED"
flexible_time_window { flexible_time_window {
mode = "OFF" mode = "OFF"
} }

View file

@ -1,10 +1,8 @@
"""Unit tests for functions/close_form/handler.py — wall-clock guard.""" """Unit tests for src/server/jobs/close_form.py."""
import os import os
import sys import sys
from datetime import datetime
from unittest.mock import patch from unittest.mock import patch
from zoneinfo import ZoneInfo
os.environ.setdefault( os.environ.setdefault(
"AGGREGATE_FUNCTION_ARN", "AGGREGATE_FUNCTION_ARN",
@ -23,95 +21,35 @@ close_form_handler = importlib.util.module_from_spec(_spec)
sys.modules["close_form_handler"] = close_form_handler sys.modules["close_form_handler"] = close_form_handler
_spec.loader.exec_module(close_form_handler) _spec.loader.exec_module(close_form_handler)
ET = ZoneInfo("America/New_York")
def _make_datetime(year, month, day, hour, minute=0):
return datetime(year, month, day, hour, minute, tzinfo=ET)
class TestCloseFormGuard:
@patch("close_form_handler.datetime")
def test_skipped_on_wednesday(self, mock_dt):
"""Wednesday 11pm ET -> skipped (not Thursday)."""
mock_dt.now.return_value = _make_datetime(2026, 5, 13, 23) # Wednesday
result = close_form_handler.lambda_handler({}, None)
assert result["status"] == "skipped"
@patch("close_form_handler.datetime")
def test_skipped_on_friday_after_catchup_window(self, mock_dt):
"""Friday 4am ET -> skipped (past Thu 23 / Fri 00–03 catch-up window)."""
mock_dt.now.return_value = _make_datetime(2026, 5, 15, 4) # Friday 4am
result = close_form_handler.lambda_handler({}, None)
assert result["status"] == "skipped"
@patch("close_form_handler.datetime")
def test_skipped_thursday_before_11pm(self, mock_dt):
"""Thursday 10pm ET -> skipped (too early)."""
mock_dt.now.return_value = _make_datetime(2026, 5, 14, 22) # Thursday 10pm
result = close_form_handler.lambda_handler({}, None)
assert result["status"] == "skipped"
class TestCloseForm:
@patch("close_form_handler.set_form_status") @patch("close_form_handler.set_form_status")
@patch("close_form_handler.get_form_status", return_value="open") @patch("close_form_handler.get_form_status", return_value="open")
@patch("close_form_handler.current_week", return_value="2026-W19") @patch("close_form_handler.current_week", return_value="2026-W19")
@patch("close_form_handler._dispatch_job") @patch("close_form_handler._dispatch_job")
@patch("close_form_handler.datetime") def test_closes_open_form(self, mock_lam, mock_week, mock_status, mock_set):
def test_runs_thursday_at_11pm(
self, mock_dt, mock_lam, mock_week, mock_status, mock_set
):
"""Thursday 11pm ET -> proceeds to close."""
mock_dt.now.return_value = _make_datetime(2026, 5, 14, 23) # Thursday 11pm
result = close_form_handler.lambda_handler({}, None) result = close_form_handler.lambda_handler({}, None)
assert result["status"] == "closed" assert result["status"] == "closed"
mock_set.assert_called_once() mock_set.assert_called_once_with("2026-W19", "closed")
mock_lam.assert_called_once_with({"event": "aggregate", "week": "2026-W19"})
@patch("close_form_handler.set_form_status") @patch("close_form_handler.set_form_status")
@patch("close_form_handler.get_form_status", return_value="open") @patch("close_form_handler.get_form_status", return_value="open")
@patch("close_form_handler.current_week", return_value="2026-W19")
@patch("close_form_handler._dispatch_job") @patch("close_form_handler._dispatch_job")
@patch("close_form_handler.datetime") def test_uses_week_from_event(self, mock_lam, mock_status, mock_set):
def test_runs_thursday_at_1159pm( result = close_form_handler.lambda_handler({"week": "2026-W20"}, None)
self, mock_dt, mock_lam, mock_week, mock_status, mock_set
):
"""Thursday 11:59pm ET -> proceeds to close."""
mock_dt.now.return_value = _make_datetime(2026, 5, 14, 23, 59)
result = close_form_handler.lambda_handler({}, None)
assert result["status"] == "closed" assert result["status"] == "closed"
mock_set.assert_called_once_with("2026-W20", "closed")
@patch("close_form_handler.set_form_status") @patch("close_form_handler.set_form_status")
@patch("close_form_handler.get_form_status", return_value="open")
@patch("close_form_handler.current_week", return_value="2026-W19")
@patch("close_form_handler._dispatch_job")
@patch("close_form_handler.datetime")
def test_runs_friday_just_after_midnight(
self, mock_dt, mock_lam, mock_week, mock_status, mock_set
):
"""Friday 12:30am ET -> proceeds (delayed EventBridge past Thu 23:59)."""
mock_dt.now.return_value = _make_datetime(2026, 5, 15, 0, 30)
result = close_form_handler.lambda_handler({}, None)
assert result["status"] == "closed"
mock_set.assert_called_once()
@patch("close_form_handler.get_form_status", return_value="closed") @patch("close_form_handler.get_form_status", return_value="closed")
@patch("close_form_handler.current_week", return_value="2026-W19") @patch("close_form_handler.current_week", return_value="2026-W19")
@patch("close_form_handler.datetime") @patch("close_form_handler._dispatch_job")
def test_already_closed(self, mock_dt, mock_week, mock_status): def test_already_closed(self, mock_lam, mock_week, mock_status, mock_set):
"""Thursday 11pm but already closed -> returns already_closed."""
mock_dt.now.return_value = _make_datetime(2026, 5, 14, 23)
result = close_form_handler.lambda_handler({}, None) result = close_form_handler.lambda_handler({}, None)
assert result["status"] == "already_closed" assert result["status"] == "already_closed"
mock_set.assert_not_called()
mock_lam.assert_not_called()

View file

@ -2,14 +2,12 @@
import sys import sys
import os import os
from datetime import datetime
from decimal import Decimal from decimal import Decimal
from unittest.mock import patch from unittest.mock import patch
from zoneinfo import ZoneInfo
# --------------------------------------------------------------------------- # ---------------------------------------------------------------------------
# Ensure the handler module can be imported. The shared layer lives under # Ensure the handler module can be imported. The shared layer lives under
# src/shared/ and the handler lives under functions/slack_notifier/. # src/shared/ and the handler lives under src/server/jobs/.
# --------------------------------------------------------------------------- # ---------------------------------------------------------------------------
import importlib.util import importlib.util
@ -22,85 +20,34 @@ handler = importlib.util.module_from_spec(_spec)
sys.modules["slack_notifier_handler"] = handler sys.modules["slack_notifier_handler"] = handler
_spec.loader.exec_module(handler) _spec.loader.exec_module(handler)
ET = ZoneInfo("America/New_York")
# ──────────────────────────────────────────────────────────────────────────── # ────────────────────────────────────────────────────────────────────────────
# Helpers # Reminder delayed SQS delivery
# ──────────────────────────────────────────────────────────────────────────── # ────────────────────────────────────────────────────────────────────────────
def _make_datetime(year, month, day, hour, minute=0): class TestReminderDelayedDelivery:
"""Return a timezone-aware datetime in America/New_York.""" @patch("slack_notifier_handler.send_dm")
return datetime(year, month, day, hour, minute, tzinfo=ET) @patch("slack_notifier_handler.get_orders", return_value=[])
@patch("slack_notifier_handler.get_roster", return_value=[])
@patch("slack_notifier_handler.current_week", return_value="2026-W19")
def _thursday_10am(): def test_reminder_runs_after_scheduled_hour(
"""2026-05-14 is a Thursday.""" self, mock_week, mock_roster, mock_orders, mock_dm
return _make_datetime(2026, 5, 14, 10) ):
def _thursday_11am():
return _make_datetime(2026, 5, 14, 11)
def _wednesday_10am():
"""2026-05-13 is a Wednesday."""
return _make_datetime(2026, 5, 13, 10)
# ────────────────────────────────────────────────────────────────────────────
# Reminder Dedup Guard (Critical)
# ────────────────────────────────────────────────────────────────────────────
class TestReminderDedupGuard:
@patch("slack_notifier_handler.datetime")
def test_reminder_skipped_wrong_hour(self, mock_dt):
"""Invoked at 11am Thursday ET -> returns skipped."""
mock_dt.now.return_value = _thursday_11am()
result = handler.handle_reminder({"event": "reminder"}) result = handler.handle_reminder({"event": "reminder"})
assert result["status"] == "skipped", "Should skip when hour is not 10" assert result["status"] != "skipped"
assert "outside reminder window" in result["reason"]
@patch("slack_notifier_handler.datetime")
def test_reminder_skipped_wrong_day(self, mock_dt):
"""Invoked at 10am Wednesday ET -> returns skipped."""
mock_dt.now.return_value = _wednesday_10am()
result = handler.handle_reminder({"event": "reminder"})
assert result["status"] == "skipped", "Should skip when day is not Thursday"
assert "outside reminder window" in result["reason"]
@patch("slack_notifier_handler.send_dm") @patch("slack_notifier_handler.send_dm")
@patch("slack_notifier_handler.get_orders", return_value=[]) @patch("slack_notifier_handler.get_orders", return_value=[])
@patch("slack_notifier_handler.get_roster", return_value=[]) @patch("slack_notifier_handler.get_roster", return_value=[])
@patch("slack_notifier_handler.current_week", return_value="2026-W19") @patch("slack_notifier_handler.current_week", return_value="2026-W19")
@patch("slack_notifier_handler.datetime") def test_lambda_handler_reminder_does_not_skip(
def test_reminder_runs_at_correct_time( self, mock_week, mock_roster, mock_orders, mock_dm
self, mock_dt, mock_week, mock_roster, mock_orders, mock_dm
): ):
"""Invoked at 10am Thursday ET -> proceeds (does not skip)."""
mock_dt.now.return_value = _thursday_10am()
result = handler.handle_reminder({"event": "reminder"})
assert result["status"] != "skipped", "Should not skip at 10am Thursday"
@patch("slack_notifier_handler.datetime")
def test_lambda_handler_reminder_guard(self, mock_dt):
"""lambda_handler lines 32-34 also return skipped for wrong time."""
mock_dt.now.return_value = _thursday_11am()
result = handler.lambda_handler({"event": "reminder"}, None) result = handler.lambda_handler({"event": "reminder"}, None)
assert result["status"] == "skipped", ( assert result["status"] != "skipped"
"lambda_handler should short-circuit before calling handle_reminder"
)
assert "outside reminder window" in result["reason"]
# ──────────────────────────────────────────────────────────────────────────── # ────────────────────────────────────────────────────────────────────────────
@ -113,12 +60,10 @@ class TestReminderDMs:
@patch("slack_notifier_handler.get_orders") @patch("slack_notifier_handler.get_orders")
@patch("slack_notifier_handler.get_roster") @patch("slack_notifier_handler.get_roster")
@patch("slack_notifier_handler.current_week", return_value="2026-W19") @patch("slack_notifier_handler.current_week", return_value="2026-W19")
@patch("slack_notifier_handler.datetime")
def test_reminder_sends_to_non_ordered( def test_reminder_sends_to_non_ordered(
self, mock_dt, mock_week, mock_roster, mock_orders, mock_dm self, mock_week, mock_roster, mock_orders, mock_dm
): ):
"""Roster of 3, 1 has ordered -> DMs sent to 2 others.""" """Roster of 3, 1 has ordered -> DMs sent to 2 others."""
mock_dt.now.return_value = _thursday_10am()
mock_roster.return_value = [ mock_roster.return_value = [
{"email": "alice@x.com", "name": "Alice A", "slack_user_id": "U001"}, {"email": "alice@x.com", "name": "Alice A", "slack_user_id": "U001"},
{"email": "bob@x.com", "name": "Bob B", "slack_user_id": "U002"}, {"email": "bob@x.com", "name": "Bob B", "slack_user_id": "U002"},
@ -141,12 +86,10 @@ class TestReminderDMs:
@patch("slack_notifier_handler.get_orders", return_value=[]) @patch("slack_notifier_handler.get_orders", return_value=[])
@patch("slack_notifier_handler.get_roster") @patch("slack_notifier_handler.get_roster")
@patch("slack_notifier_handler.current_week", return_value="2026-W19") @patch("slack_notifier_handler.current_week", return_value="2026-W19")
@patch("slack_notifier_handler.datetime")
def test_reminder_skips_no_slack_id( def test_reminder_skips_no_slack_id(
self, mock_dt, mock_week, mock_roster, mock_orders, mock_dm self, mock_week, mock_roster, mock_orders, mock_dm
): ):
"""Employee without slack_user_id is counted as missing but not DM'd.""" """Employee without slack_user_id is counted as missing but not DM'd."""
mock_dt.now.return_value = _thursday_10am()
mock_roster.return_value = [ mock_roster.return_value = [
{"email": "dave@x.com", "name": "Dave D"}, # no slack_user_id {"email": "dave@x.com", "name": "Dave D"}, # no slack_user_id
{"email": "eve@x.com", "name": "Eve E", "slack_user_id": "U005"}, {"email": "eve@x.com", "name": "Eve E", "slack_user_id": "U005"},
@ -163,12 +106,10 @@ class TestReminderDMs:
@patch("slack_notifier_handler.get_orders") @patch("slack_notifier_handler.get_orders")
@patch("slack_notifier_handler.get_roster") @patch("slack_notifier_handler.get_roster")
@patch("slack_notifier_handler.current_week", return_value="2026-W19") @patch("slack_notifier_handler.current_week", return_value="2026-W19")
@patch("slack_notifier_handler.datetime")
def test_reminder_case_insensitive_email( def test_reminder_case_insensitive_email(
self, mock_dt, mock_week, mock_roster, mock_orders, mock_dm self, mock_week, mock_roster, mock_orders, mock_dm
): ):
"""'Adam@x.com' in roster matches 'adam@x.com' in orders.""" """'Adam@x.com' in roster matches 'adam@x.com' in orders."""
mock_dt.now.return_value = _thursday_10am()
mock_roster.return_value = [ mock_roster.return_value = [
{"email": "Adam@x.com", "name": "Adam M", "slack_user_id": "U010"}, {"email": "Adam@x.com", "name": "Adam M", "slack_user_id": "U010"},
] ]

View file

@ -36,3 +36,8 @@ def test_email_report_not_in_remaining_config():
assert "email_report" not in outputs assert "email_report" not in outputs
assert "payroll_email" not in variables assert "payroll_email" not in variables
assert "sender_email" not in variables assert "sender_email" not in variables
def test_job_schedules_disabled_outside_prod():
scheduler = (TERRAFORM / "scheduler.tf").read_text()
assert 'state = local.is_prod ? "ENABLED" : "DISABLED"' in scheduler

53
tests/test_worker.py Normal file
View file

@ -0,0 +1,53 @@
"""SQS worker deletes completed jobs and leaves skipped messages for retry."""
import json
from unittest.mock import MagicMock, patch
from server import worker
def _one_message_then_stop(body: dict):
def receive_message(**_kwargs):
worker._stop(None, None)
return {
"Messages": [
{
"ReceiptHandle": "rh-1",
"Body": json.dumps(body),
}
]
}
return receive_message
@patch("server.worker.boto3.client")
@patch("server.worker.run_job")
def test_worker_does_not_delete_skipped_jobs(mock_run, mock_client):
sqs = MagicMock()
mock_client.return_value = sqs
sqs.receive_message.side_effect = _one_message_then_stop({"event": "close"})
mock_run.return_value = {"status": "skipped", "reason": "outside window"}
with patch.dict("os.environ", {"JOBS_QUEUE_URL": "https://sqs.example/jobs"}):
worker._running = True
worker.main()
sqs.delete_message.assert_not_called()
@patch("server.worker.boto3.client")
@patch("server.worker.run_job")
def test_worker_deletes_completed_jobs(mock_run, mock_client):
sqs = MagicMock()
mock_client.return_value = sqs
sqs.receive_message.side_effect = _one_message_then_stop({"event": "close"})
mock_run.return_value = {"status": "closed", "week": "2026-W19"}
with patch.dict("os.environ", {"JOBS_QUEUE_URL": "https://sqs.example/jobs"}):
worker._running = True
worker.main()
sqs.delete_message.assert_called_once_with(
QueueUrl="https://sqs.example/jobs", ReceiptHandle="rh-1"
)