diff --git a/cdk/cdk.json b/cdk/cdk.json index f28e013..164a462 100644 --- a/cdk/cdk.json +++ b/cdk/cdk.json @@ -1,6 +1,10 @@ { "app": "python3 app.py", "context": { - "@aws-cdk/core:bootstrapQualifier": "hnb659fds" + "@aws-cdk/core:bootstrapQualifier": "hnb659fds", + "wildcardCertArn": "arn:aws:acm:us-east-1:328440206208:certificate/a66c0994-90d4-410a-a1d9-5595c2a3fae3", + "hostedZoneId": "Z06652411XKH89KTZD3XA", + "hostedZoneName": "seahaven.com", + "slackInteractionsDomain": "apm-wo.seahaven.com" } } diff --git a/cdk/stacks/pipeline_stack.py b/cdk/stacks/pipeline_stack.py index a258317..8b3a36e 100644 --- a/cdk/stacks/pipeline_stack.py +++ b/cdk/stacks/pipeline_stack.py @@ -2,8 +2,8 @@ Phase 1 — exports bucket + drop-folder uploader IAM user. Phase 2 — Glue database + S3-triggered classifier Lambda. -Phase 3 — partition-projection Glue table + Athena workgroup (this file). -Phase 4 — slack-post Lambda + scoped IAM (TODO). +Phase 3 — partition-projection Glue table + Athena workgroup. +Phase 4 — slack-post + interactions Lambdas, API Gateway, scoped IAM (this file). Lambdas: Python 3.12, ARM64, explicit LogGroup with 60-day retention. No DynamoDB — this is an S3 + Athena analytics workload (see CLAUDE.md). @@ -17,9 +17,18 @@ from aws_cdk import ( RemovalPolicy, Stack, ) +from aws_cdk import ( + aws_apigatewayv2 as apigwv2, +) +from aws_cdk import ( + aws_apigatewayv2_integrations as apigwv2_integrations, +) from aws_cdk import ( aws_athena as athena, ) +from aws_cdk import ( + aws_certificatemanager as acm, +) from aws_cdk import ( aws_glue as glue, ) @@ -32,6 +41,12 @@ from aws_cdk import ( from aws_cdk import ( aws_logs as logs, ) +from aws_cdk import ( + aws_route53 as route53, +) +from aws_cdk import ( + aws_route53_targets as route53_targets, +) from aws_cdk import ( aws_s3 as s3, ) @@ -41,12 +56,20 @@ from aws_cdk import ( from aws_cdk import ( aws_secretsmanager as secretsmanager, ) +from aws_cdk import ( + aws_ssm as ssm, +) from constructs import Construct GLUE_DATABASE = "apm_wo_analysis" GLUE_TABLE = "apm_wo_snapshots" ATHENA_WORKGROUP = "apm-wo-analysis" ANTHROPIC_SECRET = "apm-wo-analysis/anthropic-api-key" +SLACK_SECRET = "apm-wo-analysis/slack-credentials" +GRAFANA_URL_PARAM = "/apm-wo-analysis/grafana-base-url" +GRAFANA_DASHBOARD_URL = ( + "https://grafana.seahaven.com/d/apm-wo/apm-work-orders?from=now-30d&to=now" +) LAMBDAS_DIR = os.path.join(os.path.dirname(__file__), "..", "..", "lambdas") # Snapshot schema — mirrors the per-WO record written by the classifier @@ -241,4 +264,124 @@ class PipelineStack(Stack): self, "AnthropicKey", ANTHROPIC_SECRET ).grant_read(self.classifier_fn) - # Phase 4 — slack-post Lambda + scoped IAM; classifier async-invokes it. TODO + # ----- Phase 4 — Slack post + interactions Lambdas ----- + # Reused Slack app credentials (botToken/signingSecret/channelId) live in + # one Secrets Manager secret, created out-of-band like the Anthropic key. + slack_secret = secretsmanager.Secret.from_secret_name_v2( + self, "SlackCreds", SLACK_SECRET + ) + # Grafana dashboard URL is operational config — SSM so ops can repoint the + # 📊 button / modal links without a redeploy. The d/apm-wo slug is the + # forward contract Phase 5's dashboard must honor. + dashboard_param = ssm.StringParameter( + self, + "GrafanaUrlParam", + parameter_name=GRAFANA_URL_PARAM, + string_value=GRAFANA_DASHBOARD_URL, + ) + + slack_env = { + "SLACK_SECRET_NAME": SLACK_SECRET, + "DASHBOARD_URL_PARAM": GRAFANA_URL_PARAM, + "ANALYTICS_BUCKET": self.exports_bucket.bucket_name, + } + + def _slack_lambda(construct_id: str, fn_name: str, handler_path: str): + """A slack_post-package Lambda: read analytics/, the Slack secret, and + the dashboard param. slack_sdk is Docker-bundled for ARM64.""" + log_group = logs.LogGroup( + self, + f"{construct_id}Logs", + log_group_name=f"/aws/lambda/{fn_name}", + retention=logs.RetentionDays.TWO_MONTHS, + removal_policy=RemovalPolicy.DESTROY, + ) + fn = lambda_.Function( + self, + construct_id, + function_name=fn_name, + runtime=lambda_.Runtime.PYTHON_3_12, + architecture=lambda_.Architecture.ARM_64, + handler=handler_path, + memory_size=256, + timeout=Duration.seconds(30), + log_group=log_group, + environment=slack_env, + code=lambda_.Code.from_asset( + os.path.join(LAMBDAS_DIR, "slack_post"), + bundling=BundlingOptions( + image=lambda_.Runtime.PYTHON_3_12.bundling_image, + platform="linux/arm64", + command=[ + "bash", + "-c", + "pip install -r requirements.txt -t /asset-output " + "&& cp -au . /asset-output", + ], + ), + ), + ) + self.exports_bucket.grant_read(fn, "analytics/*") + slack_secret.grant_read(fn) + dashboard_param.grant_read(fn) + return fn + + slack_post_fn = _slack_lambda( + "SlackPost", "apm-wo-analysis-slack-post", "handler.handler" + ) + interactions_fn = _slack_lambda( + "SlackInteractions", + "apm-wo-analysis-slack-interactions", + "interactions.handler", + ) + + # Classifier async-invokes slack-post after writing the snapshot. + self.classifier_fn.add_environment( + "SLACK_POST_FUNCTION_NAME", slack_post_fn.function_name + ) + slack_post_fn.grant_invoke(self.classifier_fn) + + # HTTP API for Slack interactivity, on apm-wo.seahaven.com (the request URL + # registered in the Slack app manifest). The interactions Lambda verifies + # the Slack signature itself; the route is intentionally unauthenticated. + cert = acm.Certificate.from_certificate_arn( + self, "WildcardCert", self.node.try_get_context("wildcardCertArn") + ) + slack_domain_name = self.node.try_get_context("slackInteractionsDomain") + slack_domain = apigwv2.DomainName( + self, "SlackDomain", domain_name=slack_domain_name, certificate=cert + ) + slack_api = apigwv2.HttpApi( + self, + "SlackInteractionsApi", + default_domain_mapping=apigwv2.DomainMappingOptions( + domain_name=slack_domain + ), + ) + slack_api.add_routes( + path="/slack/interactions", + methods=[apigwv2.HttpMethod.POST], + integration=apigwv2_integrations.HttpLambdaIntegration( + "InteractionsIntegration", interactions_fn + ), + ) + + # Route53 alias apm-wo.seahaven.com → the API Gateway custom domain. + zone = route53.HostedZone.from_hosted_zone_attributes( + self, + "SeahavenZone", + hosted_zone_id=self.node.try_get_context("hostedZoneId"), + zone_name=self.node.try_get_context("hostedZoneName"), + ) + route53.ARecord( + self, + "SlackDomainAlias", + zone=zone, + record_name=slack_domain_name.split(".")[0], + target=route53.RecordTarget.from_alias( + route53_targets.ApiGatewayv2DomainProperties( + slack_domain.regional_domain_name, + slack_domain.regional_hosted_zone_id, + ) + ), + ) diff --git a/lambdas/classifier/handler.py b/lambdas/classifier/handler.py index ee8f32c..f356cac 100644 --- a/lambdas/classifier/handler.py +++ b/lambdas/classifier/handler.py @@ -51,6 +51,7 @@ COLUMN_MATCHERS = { } _s3 = boto3.client("s3") +_lambda = boto3.client("lambda") def _haiku_enabled() -> bool: @@ -177,6 +178,29 @@ def _build_summary(df: pd.DataFrame, dt: str, key: str, blank: int) -> dict: } +# Comment snippet length for the modal rows — Slack section text caps at 3000 +# chars; keep rows short so a full modal stays well under the block limit. +_SNIPPET_LEN = 280 + + +def _build_details(df: pd.DataFrame) -> list[dict]: + """Per-WO index for the Slack drill-down modals (read from S3 on demand).""" + return [ + { + "wo_number": r.wo_number, + "wo_description": r.wo_description, + "site": r.site, + "department": r.department, + "category": r.category, + "last_comment": r.last_comment[:_SNIPPET_LEN], + "is_escalation": bool(r.is_escalation), + "is_action": bool(r.is_action), + "mismatch": r.mismatch, + } + for r in df.itertuples() + ] + + def handler(event, context): """Classify each export dropped under raw/ into a daily Parquet snapshot.""" record = event["Records"][0] @@ -220,14 +244,38 @@ def handler(event, context): Body=json.dumps(summary, indent=2).encode("utf-8"), ContentType="application/json", ) + # details.json — per-WO index the slack-post Lambda reads for drill-down modals. + _s3.put_object( + Bucket=bucket, + Key=f"{ANALYTICS_PREFIX}/dt={dt}/details.json", + Body=json.dumps(_build_details(df)).encode("utf-8"), + ContentType="application/json", + ) - # Phase 4: async-invoke the slack-post Lambda here once it exists. print( f"Wrote {len(df)} rows, {summary['escalation_total']} escalations " f"({summary['third_escalation_count']} 3rd), {len(summary['mismatches'])} mismatches." ) + _invoke_slack_post(dt) return { "classified": int(len(df)), "dt": dt, "summary_key": f"dt={dt}/summary.json", } + + +def _invoke_slack_post(dt: str) -> None: + """Async-invoke the slack-post Lambda (name from env), if wired. Best-effort: + a Slack failure must not fail the classification that already landed in S3.""" + fn = os.environ.get("SLACK_POST_FUNCTION_NAME") + if not fn: + return + try: + _lambda.invoke( + FunctionName=fn, + InvocationType="Event", + Payload=json.dumps({"dt": dt}).encode("utf-8"), + ) + print(f"Invoked slack-post {fn} for dt={dt}") + except Exception as exc: # noqa: BLE001 — never let Slack break the pipeline + print(f"slack-post invoke failed (non-fatal): {exc}") diff --git a/lambdas/slack_post/blockkit.py b/lambdas/slack_post/blockkit.py index 84f7724..2823976 100644 --- a/lambdas/slack_post/blockkit.py +++ b/lambdas/slack_post/blockkit.py @@ -2,22 +2,432 @@ Two push surfaces, no App Home (see CLAUDE.md "Slack surfaces"): - daily summary: header, vs-yesterday deltas, escalation + action/routine - fields, top sites, mismatch callout, and a 📊 Open dashboard link button to + fields, top sites, mismatch callout, and a Open dashboard link button to Grafana. Long WO lists go in modals, never the channel (<100-block limit). - 3rd-escalation alert: standalone, batched, one @here — suppressed entirely on zero-3rd days (return None). Owned by the slack-blockkit-designer agent; built in Phase 4 (docs/BUILD.md). + +Slack constraints enforced here: + - Message / modal max 100 blocks. + - Actions block max 25 elements (kept to 5 per block for readability). + - section text max 3000 chars. + - Header text max 150 chars. + - No APM deep-links anywhere — Grafana only. """ from __future__ import annotations +# --------------------------------------------------------------------------- +# Constants +# --------------------------------------------------------------------------- -def build_daily_summary(today: dict, yesterday: dict | None) -> list[dict]: - """Build the daily summary blocks (with vs-yesterday deltas). (Phase 4)""" - raise NotImplementedError +ESCALATION_CATEGORIES: tuple[str, ...] = ( + "1st Escalation", + "2nd Escalation", + "3rd Escalation", + "SIM Ticket", + "Other Escalation", +) + +# Categories that drive the "drill" buttons — non-routine and worth calling out. +# Order determines button priority when we cap at 5 per actions block. +DRILL_CATEGORY_PRIORITY: tuple[str, ...] = ( + "3rd Escalation", + "2nd Escalation", + "1st Escalation", + "SIM Ticket", + "Vendor No-Show", + "Awaiting Scheduling", + "Report / Docs Needed", + "Awaiting Report / Invoice", + "Awaiting Vendor / Parts", + "Status Inquiry", + "Other Escalation", +) + +# Max WOs to render inline before the "+M more" overflow note in the alert. +ALERT_MAX_WOS = 20 + +# Max blocks reserved for WO rows in the modal (leaves headroom for +# title/footer blocks). Each WO gets 1 section block. +MODAL_HEADER_BLOCKS = 1 # title is not a block; we add 1 context block as preamble +MODAL_FOOTER_BLOCKS = 1 # overflow context block (may not be emitted) +MODAL_MAX_WO_BLOCKS = 100 - MODAL_HEADER_BLOCKS - MODAL_FOOTER_BLOCKS # = 98 -def build_escalation_alert(today: dict) -> list[dict] | None: - """Build the batched 3rd-escalation alert, or None when count is 0. (Phase 4)""" - raise NotImplementedError +# --------------------------------------------------------------------------- +# Helpers +# --------------------------------------------------------------------------- + + +def _divider() -> dict: + return {"type": "divider"} + + +def _header(text: str) -> dict: + # Header block text is plain_text, max 150 chars. + return { + "type": "header", + "text": {"type": "plain_text", "text": text[:150], "emoji": True}, + } + + +def _section(text: str) -> dict: + # Section text is mrkdwn, max 3000 chars. + return { + "type": "section", + "text": {"type": "mrkdwn", "text": text[:3000]}, + } + + +def _context(text: str) -> dict: + return { + "type": "context", + "elements": [{"type": "mrkdwn", "text": text[:3000]}], + } + + +def _fields_section(fields: list[str]) -> dict: + """Section with up to 10 fields (Slack limit), each <=2000 chars.""" + return { + "type": "section", + "fields": [{"type": "mrkdwn", "text": f[:2000]} for f in fields[:10]], + } + + +def _button(text: str, action_id: str, value: str) -> dict: + """Interactive button (opens modal via interactivity handler).""" + return { + "type": "button", + "text": {"type": "plain_text", "text": text, "emoji": False}, + "action_id": action_id, + "value": value, + } + + +def _link_button(text: str, url: str) -> dict: + """Link button — opens URL directly, no action_id.""" + return { + "type": "button", + "text": {"type": "plain_text", "text": text, "emoji": True}, + "url": url, + } + + +def _actions(elements: list[dict]) -> dict: + """Actions block; Slack limit 25 elements, we keep <=5 per block.""" + return {"type": "actions", "elements": elements[:25]} + + +def _delta_str(now: int, prev: int) -> str: + """Format a vs-yesterday delta as ▲+N or ▼-N or ~0.""" + diff = now - prev + if diff > 0: + return f"▲+{diff}" + if diff < 0: + return f"▼{diff}" + return "~0" + + +def _fmt_date(dt: str) -> str: + """Format 'YYYY-MM-DD' as 'May 28, 2026' for display.""" + try: + from datetime import date + + d = date.fromisoformat(dt) + # %-d is Linux-specific; strip the leading zero portably instead. + return f"{d.strftime('%b')} {d.day}, {d.year}" + except Exception: + return dt + + +# --------------------------------------------------------------------------- +# Public surface builders +# --------------------------------------------------------------------------- + + +def build_daily_summary( + today: dict, + yesterday: dict | None, + dashboard_url: str, +) -> list[dict]: + """Build the daily summary post blocks. + + Returns a list of Block Kit block dicts ready to pass to ``blocks=`` in a + ``chat.postMessage`` call. Always stays well under 100 blocks. + + Args: + today: The contents of today's ``summary.json``. + yesterday: The contents of yesterday's ``summary.json``, or None when + there is no previous snapshot (first run, etc.). + dashboard_url: The Grafana dashboard URL for the link button. + """ + blocks: list[dict] = [] + + dt = today.get("dt", "") + date_label = _fmt_date(dt) if dt else "Today" + + # ------------------------------------------------------------------ + # 1. Header + # ------------------------------------------------------------------ + blocks.append(_header(f"APM Work Orders — {date_label}")) + + # ------------------------------------------------------------------ + # 2. Context line: classified total + vs-yesterday deltas + # ------------------------------------------------------------------ + classified = today.get("classified_total", 0) + blank = today.get("blank_comment_rows", 0) + escalation = today.get("escalation_total", 0) + action = today.get("action_needed", 0) + routine = today.get("routine", 0) + + if yesterday is not None: + delta_classified = _delta_str(classified, yesterday.get("classified_total", 0)) + delta_escalation = _delta_str(escalation, yesterday.get("escalation_total", 0)) + delta_action = _delta_str(action, yesterday.get("action_needed", 0)) + context_text = ( + f"*{classified}* WOs classified ({blank} blank-comment rows excluded) " + f"| Escalations {delta_escalation} | Action-needed {delta_action} " + f"| vs yesterday: total {delta_classified}" + ) + else: + context_text = ( + f"*{classified}* WOs classified ({blank} blank-comment rows excluded) " + f"| No prior day for delta comparison" + ) + + blocks.append(_section(context_text)) + blocks.append(_divider()) + + # ------------------------------------------------------------------ + # 3. Escalation breakdown + # ------------------------------------------------------------------ + category_counts: dict[str, int] = today.get("category_counts", {}) + third_count = today.get( + "third_escalation_count", category_counts.get("3rd Escalation", 0) + ) + + esc_lines: list[str] = [] + for cat in ESCALATION_CATEGORIES: + count = category_counts.get(cat, 0) + if count == 0: + continue + prefix = "*" if cat == "3rd Escalation" else "" + suffix = "*" if cat == "3rd Escalation" else "" + esc_lines.append(f"{prefix}{cat}: {count}{suffix}") + + if esc_lines: + blocks.append(_section("*Escalations*\n" + "\n".join(esc_lines))) + else: + blocks.append(_section("*Escalations*\nNone today")) + + # ------------------------------------------------------------------ + # 4. Action-needed vs routine fields + # ------------------------------------------------------------------ + blocks.append( + _fields_section( + [ + f"*Action-needed*\n{action}", + f"*Routine*\n{routine}", + f"*3rd Escalations*\n{third_count}", + f"*Total classified*\n{classified}", + ] + ) + ) + blocks.append(_divider()) + + # ------------------------------------------------------------------ + # 5. Top 5 sites + # ------------------------------------------------------------------ + top_sites: list[dict] = today.get("top_sites", [])[:5] + if top_sites: + site_lines = "\n".join( + f"{i + 1}. *{s['site']}* — {s['count']} WOs" + for i, s in enumerate(top_sites) + ) + blocks.append(_section(f"*Top sites*\n{site_lines}")) + blocks.append(_divider()) + + # ------------------------------------------------------------------ + # 6. Mismatch callout (only when mismatches present) + # ------------------------------------------------------------------ + mismatches: list[dict] = today.get("mismatches", []) + if mismatches: + mm_count = len(mismatches) + blocks.append( + _section( + f":warning: *{mm_count} mismatch{'es' if mm_count != 1 else ''} " + f"flagged* — comment intent contradicts structured state. " + f"See the Grafana mismatch panel for details." + ) + ) + blocks.append(_divider()) + + # ------------------------------------------------------------------ + # 7. Drill buttons (action_id="drill_category") + dashboard link button + # ------------------------------------------------------------------ + # Pick the most-common non-routine categories up to 4 slots; the 5th slot + # is always the dashboard link button, keeping each actions block <=5. + drill_cats: list[tuple[str, int]] = [] + for cat in DRILL_CATEGORY_PRIORITY: + count = category_counts.get(cat, 0) + if count > 0: + drill_cats.append((cat, count)) + if len(drill_cats) >= 4: + break + + button_elements: list[dict] = [] + for cat, count in drill_cats: + short_label = cat.replace(" Escalation", " Esc.").replace("Awaiting ", "") + short_label = short_label[:20] # button text kept concise + button_elements.append( + _button(f"{short_label} ({count})", "drill_category", cat) + ) + + # Dashboard link button always present (no action_id — url button). + button_elements.append(_link_button("Open dashboard", dashboard_url)) + + blocks.append(_actions(button_elements)) + + # ------------------------------------------------------------------ + # 8. Footer context + # ------------------------------------------------------------------ + generated_at = today.get("generated_at", "") + footer_parts = [f"Generated {generated_at}" if generated_at else ""] + footer_parts.append(f"dt={dt}") + blocks.append(_context(" | ".join(p for p in footer_parts if p))) + + # Safety check — should never fire in practice given the bounded sections + # above, but surface a warning in the footer rather than silently truncate. + if len(blocks) > 95: + # Trim from the middle, preserve header and footer. + over = len(blocks) - 95 + blocks = blocks[:3] + blocks[3 + over :] + blocks.append( + _context(f"+{over} block(s) trimmed — view full breakdown in Grafana.") + ) + + return blocks + + +def build_escalation_alert(escalations: list[dict]) -> list[dict] | None: + """Build the batched 3rd-escalation alert blocks. + + Zero-3rd suppression: returns None when the list is empty. The caller must + check for None before posting — do NOT post a message with no content. + + Args: + escalations: Rows from ``details.json`` pre-filtered to + ``category == "3rd Escalation"``. + + Returns: + A list of Block Kit block dicts, or None. + """ + if not escalations: + return None + + blocks: list[dict] = [] + + # ------------------------------------------------------------------ + # 1. Header + # ------------------------------------------------------------------ + count = len(escalations) + blocks.append( + _header(f"3rd Escalation Alert — {count} WO{'s' if count != 1 else ''}") + ) + + # ------------------------------------------------------------------ + # 2. ONE @here mention (never one ping per WO) + # ------------------------------------------------------------------ + blocks.append( + _section( + f" *{count} work order{'s' if count != 1 else ''} " + f"at 3rd escalation* require immediate attention." + ) + ) + blocks.append(_divider()) + + # ------------------------------------------------------------------ + # 3. Compact WO list, capped at ALERT_MAX_WOS + # ------------------------------------------------------------------ + visible = escalations[:ALERT_MAX_WOS] + overflow = len(escalations) - len(visible) + + for wo in visible: + wo_num = wo.get("wo_number", "???") + site = wo.get("site", "—") + desc = wo.get("wo_description", "") + # Truncate description to keep the block tidy. + if len(desc) > 80: + desc = desc[:77] + "..." + blocks.append(_section(f"*{wo_num}* — {site} — {desc}")) + + if overflow > 0: + blocks.append( + _context(f"+{overflow} more — see the Grafana dashboard for the full list.") + ) + + return blocks + + +def build_wo_modal(title: str, wos: list[dict], dashboard_url: str) -> dict: + """Build a ``views.open`` payload (modal) listing work orders. + + Caps to Slack's 100-block modal limit by rendering the first chunk of WOs + and appending a "+M more" overflow context block when the list is too long. + No APM deep-links — Grafana only. + + Args: + title: Modal title (max 24 chars enforced by Slack; truncated here). + wos: List of per-WO dicts from ``details.json``. + dashboard_url: The Grafana dashboard URL for the overflow note link. + + Returns: + A dict suitable for passing to ``views.open`` as the ``view=`` arg. + """ + # Slack enforces a 24-char modal title limit. + title_text = title[:24] + + modal_blocks: list[dict] = [] + + # Each WO is one section block. Reserve space for the optional overflow + # context block. + max_wo = MODAL_MAX_WO_BLOCKS - 1 # one slot reserved for potential overflow + visible = wos[:max_wo] + overflow_count = len(wos) - len(visible) + + for wo in visible: + wo_num = wo.get("wo_number", "???") + site = wo.get("site", "—") + dept = wo.get("department", "—") + category = wo.get("category", "—") + snippet = wo.get("last_comment", "") + if len(snippet) > 150: + snippet = snippet[:147] + "..." + + text_parts = [ + f"*{wo_num}*", + f"_{site}_ | {dept} | {category}", + ] + if snippet: + text_parts.append(snippet) + + modal_blocks.append(_section("\n".join(text_parts))) + + if overflow_count > 0: + modal_blocks.append( + _context( + f"+{overflow_count} more — <{dashboard_url}|view the full list in Grafana>." + ) + ) + elif not modal_blocks: + modal_blocks.append(_section("No work orders to display.")) + + return { + "type": "modal", + "title": {"type": "plain_text", "text": title_text, "emoji": False}, + "close": {"type": "plain_text", "text": "Close", "emoji": False}, + "blocks": modal_blocks, + } diff --git a/lambdas/slack_post/handler.py b/lambdas/slack_post/handler.py index 9cba151..128b851 100644 --- a/lambdas/slack_post/handler.py +++ b/lambdas/slack_post/handler.py @@ -1,16 +1,59 @@ -"""Slack post + alert Lambda: daily summary and batched 3rd-escalation alert. +"""Slack post Lambda: daily summary + batched 3rd-escalation alert. -Invoked after the classifier finishes. Reads ``analytics/dt=today/summary.json`` -and the yesterday partition (for deltas), posts the daily summary to the WO -channel via the reused bot token (SSM, same pattern as payments-dashboard), and -— only if the 3rd-escalation count > 0 — posts the standalone batched alert. +Async-invoked by the classifier with ``{"dt": "YYYY-MM-DD"}`` once the day's +snapshot is written. Reads today's (and yesterday's, for deltas) ``summary.json`` +from ``analytics/``, posts the daily summary to the WO channel via the reused bot +token (Secrets Manager), and — only when the 3rd-escalation count > 0 — posts the +standalone batched alert (reading ``details.json`` for the 3rd-escalation rows). -Implemented in Phase 4 (docs/BUILD.md). +Runtime: Python 3.12, ARM64. No App Home — two push surfaces only (see CLAUDE.md). """ from __future__ import annotations +from datetime import datetime, timedelta, timezone + +import blockkit +import slackio + + +def _yesterday(dt: str) -> str: + return (datetime.strptime(dt, "%Y-%m-%d") - timedelta(days=1)).strftime("%Y-%m-%d") + def handler(event, context): - """Lambda entry point. (Phase 4)""" - raise NotImplementedError + """Post the daily summary and (conditionally) the 3rd-escalation alert.""" + dt = (event or {}).get("dt") or datetime.now(timezone.utc).strftime("%Y-%m-%d") + + today = slackio.read_analytics_json(dt, "summary.json") + if today is None: + print(f"No summary.json for dt={dt}; nothing to post.") + return {"posted": False, "reason": "no summary", "dt": dt} + yesterday = slackio.read_analytics_json(_yesterday(dt), "summary.json") + + client = slackio.web_client() + channel = slackio.channel_id() + dashboard_url = slackio.get_dashboard_url() + + summary_blocks = blockkit.build_daily_summary(today, yesterday, dashboard_url) + client.chat_postMessage( + channel=channel, blocks=summary_blocks, text=f"APM Work Orders — {dt}" + ) + posted = {"summary": True, "alert": False} + + # Standalone batched alert — only when there are 3rd escalations. Pull the + # rows from details.json so the alert can name the WOs. + if today.get("third_escalation_count", 0) > 0: + details = slackio.read_analytics_json(dt, "details.json") or [] + thirds = [d for d in details if d.get("category") == "3rd Escalation"] + alert_blocks = blockkit.build_escalation_alert(thirds) + if alert_blocks: + client.chat_postMessage( + channel=channel, + blocks=alert_blocks, + text=f"🚨 {len(thirds)} 3rd-escalation work orders", + ) + posted["alert"] = True + + print(f"Posted dt={dt}: {posted}") + return {"posted": True, "dt": dt, **posted} diff --git a/lambdas/slack_post/interactions.py b/lambdas/slack_post/interactions.py new file mode 100644 index 0000000..647f3fe --- /dev/null +++ b/lambdas/slack_post/interactions.py @@ -0,0 +1,69 @@ +"""Slack interactivity Lambda: drill-down modals behind API Gateway. + +Fronted by the ``apm-wo.seahaven.com`` HTTP API (Slack interactivity request URL). +Verifies the Slack request signature, then for a category/site drill button reads +the day's ``details.json``, filters it, and opens a WO-list modal via ``views.open`` +within Slack's 3-second ``trigger_id`` window. WO rows link to Grafana only — no +APM deep-links (see Phase 4 decisions). +""" + +from __future__ import annotations + +import base64 +import json +from datetime import datetime, timezone +from urllib.parse import parse_qs + +import blockkit +import slackio + +_FILTER_FIELD = {"drill_category": "category", "drill_site": "site"} + + +def _raw_body(event: dict) -> str: + body = event.get("body") or "" + if event.get("isBase64Encoded"): + body = base64.b64decode(body).decode("utf-8") + return body + + +def _headers(event: dict) -> dict: + # API Gateway HTTP API lowercases header names; be defensive either way. + return {k.lower(): v for k, v in (event.get("headers") or {}).items()} + + +def handler(event, context): + """Verify the Slack signature and open the requested drill-down modal.""" + body = _raw_body(event) + headers = _headers(event) + timestamp = headers.get("x-slack-request-timestamp", "") + signature = headers.get("x-slack-signature", "") + + if not slackio.verify_signature(body, timestamp, signature): + print("Rejected: invalid Slack signature.") + return {"statusCode": 401, "body": "invalid signature"} + + parsed = parse_qs(body) + payload = json.loads(parsed.get("payload", ["{}"])[0]) + if payload.get("type") != "block_actions": + return {"statusCode": 200, "body": ""} + + action = (payload.get("actions") or [{}])[0] + field = _FILTER_FIELD.get(action.get("action_id")) + value = action.get("value") + trigger_id = payload.get("trigger_id") + if not field or not value or not trigger_id: + return {"statusCode": 200, "body": ""} + + # The daily post is same-day; default the drill to today's snapshot. + dt = datetime.now(timezone.utc).strftime("%Y-%m-%d") + details = slackio.read_analytics_json(dt, "details.json") or [] + wos = [d for d in details if d.get(field) == value] + + view = blockkit.build_wo_modal( + title=f"{value} ({len(wos)})", + wos=wos, + dashboard_url=slackio.get_dashboard_url(), + ) + slackio.web_client().views_open(trigger_id=trigger_id, view=view) + return {"statusCode": 200, "body": ""} diff --git a/lambdas/slack_post/slackio.py b/lambdas/slack_post/slackio.py new file mode 100644 index 0000000..6b4b670 --- /dev/null +++ b/lambdas/slack_post/slackio.py @@ -0,0 +1,70 @@ +"""Shared Slack + AWS I/O for the post and interactions Lambdas. + +Keeps the Block Kit builders (blockkit.py) pure: everything that touches the +network or AWS lives here. Slack credentials are a single Secrets Manager secret +``{ botToken, signingSecret, channelId }``; the Grafana dashboard URL is +operational config in SSM (editable without a redeploy). Both are cached for the +life of the execution environment. +""" + +from __future__ import annotations + +import json +import os + +import boto3 +from slack_sdk import WebClient +from slack_sdk.signature import SignatureVerifier + +_secrets = boto3.client("secretsmanager") +_ssm = boto3.client("ssm") +_s3 = boto3.client("s3") + +_SECRET_NAME = os.environ["SLACK_SECRET_NAME"] +_DASHBOARD_PARAM = os.environ["DASHBOARD_URL_PARAM"] +_BUCKET = os.environ["ANALYTICS_BUCKET"] + +# Lazily-populated caches (warm across invocations in the same container). +_creds: dict | None = None +_dashboard_url: str | None = None + + +def get_credentials() -> dict: + """Return the Slack creds dict: ``botToken``, ``signingSecret``, ``channelId``.""" + global _creds + if _creds is None: + raw = _secrets.get_secret_value(SecretId=_SECRET_NAME)["SecretString"] + _creds = json.loads(raw) + return _creds + + +def get_dashboard_url() -> str: + """Grafana dashboard URL for the 📊 button / modal overflow links (SSM).""" + global _dashboard_url + if _dashboard_url is None: + _dashboard_url = _ssm.get_parameter(Name=_DASHBOARD_PARAM)["Parameter"]["Value"] + return _dashboard_url + + +def web_client() -> WebClient: + return WebClient(token=get_credentials()["botToken"]) + + +def channel_id() -> str: + return get_credentials()["channelId"] + + +def verify_signature(body: str, timestamp: str, signature: str) -> bool: + """Validate a Slack request signature (HMAC + 5-minute replay window).""" + verifier = SignatureVerifier(signing_secret=get_credentials()["signingSecret"]) + return verifier.is_valid(body=body, timestamp=timestamp, signature=signature) + + +def read_analytics_json(dt: str, name: str): + """Read ``analytics/dt=
/`` as JSON, or None if absent.""" + key = f"analytics/dt={dt}/{name}" + try: + obj = _s3.get_object(Bucket=_BUCKET, Key=key) + except _s3.exceptions.NoSuchKey: + return None + return json.loads(obj["Body"].read()) diff --git a/slack/manifest.yaml b/slack/manifest.yaml new file mode 100644 index 0000000..de37e10 --- /dev/null +++ b/slack/manifest.yaml @@ -0,0 +1,41 @@ +# Slack app manifest for "APM Work Orders". +# +# IMPORTANT: The request_url host (apm-wo.seahaven.com) must match the deployed +# API Gateway custom domain. If the domain changes, update interactivity.request_url +# here and redeploy the app via the Slack API or the App Configuration page. +# +# The bot user must be invited to the WO channel before it can post: +# /invite @APM Work Orders +# +# To apply this manifest: Slack App Configuration -> "App Manifest" tab -> paste. + +_metadata: + major_version: 1 + minor_version: 1 + +display_information: + name: APM Work Orders + description: Daily APM work-order analysis — escalation alerts and summary posts for Sea Haven facility ops. + background_color: "#1a1a2e" + +features: + bot_user: + display_name: APM Work Orders + always_online: true + +oauth_config: + scopes: + bot: + - chat:write + # chat:write.public allows posting to public channels without an explicit + # /invite. Remove if the WO channel is private and an invite is preferred. + - chat:write.public + +settings: + interactivity: + is_enabled: true + # Handler for drill-category buttons and modal opens. + request_url: https://apm-wo.seahaven.com/slack/interactions + org_deploy_enabled: false + socket_mode_enabled: false + token_rotation_enabled: false diff --git a/tests/test_blockkit.py b/tests/test_blockkit.py new file mode 100644 index 0000000..02eb1d7 --- /dev/null +++ b/tests/test_blockkit.py @@ -0,0 +1,439 @@ +"""Offline unit tests for the Block Kit surface builders. + +No AWS, no network, no slack_sdk. All fixtures are plain dicts. + +Run with the repo venv: + + ./.venv/bin/python -m pytest tests/test_blockkit.py -q +""" + +from __future__ import annotations + +import sys +from pathlib import Path + +sys.path.insert( + 0, str(Path(__file__).resolve().parent.parent / "lambdas" / "slack_post") +) + +import blockkit # noqa: E402 +from blockkit import ( # noqa: E402 + build_daily_summary, + build_escalation_alert, + build_wo_modal, +) + +# --------------------------------------------------------------------------- +# Shared fixtures +# --------------------------------------------------------------------------- + +DASHBOARD_URL = "https://grafana.seahaven.internal/d/apm-wo" + +CATEGORY_COUNTS_TYPICAL: dict[str, int] = { + "3rd Escalation": 4, + "2nd Escalation": 8, + "1st Escalation": 12, + "SIM Ticket": 2, + "Vendor No-Show": 3, + "Awaiting Scheduling": 25, + "Report / Docs Needed": 18, + "Awaiting Report / Invoice": 7, + "Awaiting Vendor / Parts": 5, + "Status Inquiry": 6, + "Other Escalation": 1, + "Schedule Confirmed": 90, + "Weekly WO Scheduled": 40, + "Completed / Pending Close": 22, + "Acknowledgement / No-op": 10, + "Cancelled": 5, + "On Hold": 3, + "Rescheduled": 4, + "Avetta Project Created": 2, + "Other": 15, +} + +TODAY_SUMMARY: dict = { + "dt": "2026-05-28", + "classified_total": 282, + "blank_comment_rows": 68, + "category_counts": CATEGORY_COUNTS_TYPICAL, + "escalation_total": 27, + "third_escalation_count": 4, + "action_needed": 91, + "routine": 191, + "top_sites": [ + {"site": "ABQ5", "count": 45}, + {"site": "ACY9", "count": 38}, + {"site": "BOS1", "count": 31}, + {"site": "DFW7", "count": 28}, + {"site": "LAX9", "count": 22}, + ], + "mismatches": [ + { + "wo_number": "WO-001", + "category": "Completed / Pending Close", + "mismatch": "comment claims completion while on REPORT hold", + }, + { + "wo_number": "WO-002", + "category": "Schedule Confirmed", + "mismatch": "comment confirms schedule while on SCHEDULING hold", + }, + ], + "generated_at": "2026-05-28T06:00:00Z", +} + +YESTERDAY_SUMMARY: dict = { + "dt": "2026-05-27", + "classified_total": 270, + "blank_comment_rows": 80, + "category_counts": {}, + "escalation_total": 30, + "third_escalation_count": 6, + "action_needed": 95, + "routine": 175, + "top_sites": [], + "mismatches": [], + "generated_at": "2026-05-27T06:00:00Z", +} + + +def _make_wo(i: int) -> dict: + return { + "wo_number": f"WO-{i:04d}", + "wo_description": f"Repair HVAC unit at dock {i}", + "site": "ABQ5", + "department": "SSP", + "category": "3rd Escalation", + "last_comment": f"3rd attempt to schedule vendor for dock {i} repair.", + "is_escalation": True, + "is_action": True, + "mismatch": None, + } + + +# --------------------------------------------------------------------------- +# build_daily_summary +# --------------------------------------------------------------------------- + + +class TestDailySummaryDeltas: + def test_with_yesterday_shows_deltas(self): + blocks = build_daily_summary(TODAY_SUMMARY, YESTERDAY_SUMMARY, DASHBOARD_URL) + # Collect all text from section/context blocks. + all_text = _extract_text(blocks) + # classified went 270 -> 282 (+12), escalation 30 -> 27 (-3). + assert "▲+12" in all_text, "expected classified delta ▲+12" + assert "▼-3" in all_text, "expected escalation delta ▼-3" + + def test_without_yesterday_omits_deltas(self): + """yesterday=None must not raise and must omit delta symbols.""" + blocks = build_daily_summary(TODAY_SUMMARY, None, DASHBOARD_URL) + all_text = _extract_text(blocks) + assert "▲" not in all_text + assert "▼" not in all_text + assert "No prior day" in all_text + + def test_without_yesterday_returns_blocks(self): + blocks = build_daily_summary(TODAY_SUMMARY, None, DASHBOARD_URL) + assert isinstance(blocks, list) + assert len(blocks) > 0 + + def test_zero_delta_shows_tilde(self): + yesterday_same = {**YESTERDAY_SUMMARY, "escalation_total": 27} + blocks = build_daily_summary(TODAY_SUMMARY, yesterday_same, DASHBOARD_URL) + all_text = _extract_text(blocks) + assert "~0" in all_text + + +class TestDailySummaryStructure: + def test_has_header_block(self): + blocks = build_daily_summary(TODAY_SUMMARY, YESTERDAY_SUMMARY, DASHBOARD_URL) + assert any(b.get("type") == "header" for b in blocks) + + def test_header_contains_date(self): + blocks = build_daily_summary(TODAY_SUMMARY, YESTERDAY_SUMMARY, DASHBOARD_URL) + header = next(b for b in blocks if b.get("type") == "header") + assert "2026" in header["text"]["text"] or "May" in header["text"]["text"] + + def test_mismatch_callout_present_when_mismatches_nonzero(self): + blocks = build_daily_summary(TODAY_SUMMARY, YESTERDAY_SUMMARY, DASHBOARD_URL) + all_text = _extract_text(blocks) + assert "mismatch" in all_text.lower() + + def test_mismatch_callout_absent_when_no_mismatches(self): + today_no_mm = {**TODAY_SUMMARY, "mismatches": []} + blocks = build_daily_summary(today_no_mm, YESTERDAY_SUMMARY, DASHBOARD_URL) + all_text = _extract_text(blocks) + assert "mismatch" not in all_text.lower() + + def test_dashboard_link_button_url(self): + blocks = build_daily_summary(TODAY_SUMMARY, YESTERDAY_SUMMARY, DASHBOARD_URL) + url = _find_link_button_url(blocks) + assert url == DASHBOARD_URL, f"expected {DASHBOARD_URL!r}, got {url!r}" + + def test_dashboard_button_has_no_action_id(self): + """The Open-dashboard button must be a url button, not an action button.""" + blocks = build_daily_summary(TODAY_SUMMARY, YESTERDAY_SUMMARY, DASHBOARD_URL) + for block in blocks: + if block.get("type") != "actions": + continue + for elem in block.get("elements", []): + if elem.get("url") == DASHBOARD_URL: + assert "action_id" not in elem, ( + "dashboard link button must not carry an action_id" + ) + + def test_drill_buttons_have_action_id(self): + blocks = build_daily_summary(TODAY_SUMMARY, YESTERDAY_SUMMARY, DASHBOARD_URL) + drill_found = False + for block in blocks: + if block.get("type") != "actions": + continue + for elem in block.get("elements", []): + if elem.get("action_id") == "drill_category": + drill_found = True + assert "value" in elem + assert drill_found, "expected at least one drill_category action button" + + def test_3rd_escalation_highlighted(self): + blocks = build_daily_summary(TODAY_SUMMARY, YESTERDAY_SUMMARY, DASHBOARD_URL) + all_text = _extract_text(blocks) + # 3rd Escalation should appear in bold mrkdwn. + assert "*3rd Escalation" in all_text or "3rd Escalation*" in all_text + + +class TestDailySummaryBlockLimit: + def _large_summary(self) -> dict: + """Summary with 25+ categories, 30 sites, and 200 mismatches.""" + big_cats: dict[str, int] = {f"Category {i}": i + 1 for i in range(30)} + big_cats.update(CATEGORY_COUNTS_TYPICAL) + big_sites = [{"site": f"S{i:03d}", "count": 30 - i} for i in range(30)] + big_mismatches = [ + {"wo_number": f"WO-{i:04d}", "category": "X", "mismatch": "test mm"} + for i in range(200) + ] + return { + **TODAY_SUMMARY, + "category_counts": big_cats, + "top_sites": big_sites, + "mismatches": big_mismatches, + } + + def test_stays_under_100_blocks_large_input(self): + big = self._large_summary() + blocks = build_daily_summary(big, YESTERDAY_SUMMARY, DASHBOARD_URL) + assert len(blocks) < 100, f"daily summary exceeded 100 blocks: {len(blocks)}" + + def test_reports_block_count(self, capsys): + big = self._large_summary() + blocks = build_daily_summary(big, YESTERDAY_SUMMARY, DASHBOARD_URL) + with capsys.disabled(): + print(f"\n[test] daily summary (large input): {len(blocks)} blocks") + + +# --------------------------------------------------------------------------- +# build_escalation_alert +# --------------------------------------------------------------------------- + + +class TestEscalationAlert: + def test_returns_none_on_empty_list(self): + assert build_escalation_alert([]) is None + + def test_returns_blocks_on_non_empty(self): + escalations = [_make_wo(i) for i in range(3)] + result = build_escalation_alert(escalations) + assert result is not None + assert isinstance(result, list) + assert len(result) > 0 + + def test_exactly_one_here_mention(self): + escalations = [_make_wo(i) for i in range(5)] + blocks = build_escalation_alert(escalations) + assert blocks is not None + all_text = _extract_text(blocks) + here_count = all_text.count("") + assert here_count == 1, f"expected exactly 1 , found {here_count}" + + def test_here_in_section_not_header(self): + """@here must be in a section block, not the header.""" + escalations = [_make_wo(i) for i in range(2)] + blocks = build_escalation_alert(escalations) + assert blocks is not None + for block in blocks: + if block.get("type") == "header": + header_text = block["text"]["text"] + assert "" not in header_text + + def test_no_here_on_empty(self): + result = build_escalation_alert([]) + assert result is None + + def test_overflow_note_when_many_wos(self): + """When list exceeds ALERT_MAX_WOS, a '+M more' context block is added.""" + many = [_make_wo(i) for i in range(blockkit.ALERT_MAX_WOS + 10)] + blocks = build_escalation_alert(many) + assert blocks is not None + all_text = _extract_text(blocks) + assert "more" in all_text.lower() + + def test_no_overflow_note_when_within_limit(self): + few = [_make_wo(i) for i in range(blockkit.ALERT_MAX_WOS - 1)] + blocks = build_escalation_alert(few) + assert blocks is not None + all_text = _extract_text(blocks) + # Should not have a "+N more" line. + assert " more —" not in all_text + + def test_single_wo_grammatically_correct(self): + """'WO' (singular) when count is 1.""" + blocks = build_escalation_alert([_make_wo(1)]) + assert blocks is not None + header = next(b for b in blocks if b.get("type") == "header") + assert "WOs" not in header["text"]["text"], "expected singular 'WO' not 'WOs'" + + def test_stays_under_100_blocks(self): + huge = [_make_wo(i) for i in range(500)] + blocks = build_escalation_alert(huge) + assert blocks is not None + assert len(blocks) < 100, f"alert exceeded 100 blocks: {len(blocks)}" + + def test_wo_number_appears_in_blocks(self): + escalations = [_make_wo(42)] + blocks = build_escalation_alert(escalations) + assert blocks is not None + all_text = _extract_text(blocks) + assert "WO-0042" in all_text + + +# --------------------------------------------------------------------------- +# build_wo_modal +# --------------------------------------------------------------------------- + + +class TestWoModal: + def test_returns_dict_with_modal_type(self): + modal = build_wo_modal("Test Modal", [_make_wo(1)], DASHBOARD_URL) + assert isinstance(modal, dict) + assert modal.get("type") == "modal" + + def test_modal_has_close_button(self): + modal = build_wo_modal("Test", [_make_wo(1)], DASHBOARD_URL) + assert "close" in modal + assert modal["close"]["type"] == "plain_text" + + def test_title_truncated_to_24_chars(self): + long_title = "A" * 50 + modal = build_wo_modal(long_title, [_make_wo(1)], DASHBOARD_URL) + assert len(modal["title"]["text"]) <= 24 + + def test_500_wos_truncated_to_under_100_blocks(self, capsys): + wos = [_make_wo(i) for i in range(500)] + modal = build_wo_modal("3rd Escalations", wos, DASHBOARD_URL) + block_count = len(modal["blocks"]) + with capsys.disabled(): + print(f"\n[test] 500-WO modal: {block_count} blocks") + assert block_count <= 100, ( + f"modal with 500 WOs has {block_count} blocks, exceeds 100" + ) + + def test_500_wos_has_overflow_note(self): + wos = [_make_wo(i) for i in range(500)] + modal = build_wo_modal("3rd Escalations", wos, DASHBOARD_URL) + all_text = _extract_text(modal["blocks"]) + assert "more" in all_text.lower(), "expected '+M more' overflow note" + + def test_overflow_note_contains_dashboard_url(self): + wos = [_make_wo(i) for i in range(500)] + modal = build_wo_modal("3rd Escalations", wos, DASHBOARD_URL) + all_text = _extract_text(modal["blocks"]) + assert DASHBOARD_URL in all_text, "overflow note must link to the dashboard URL" + + def test_overflow_count_correct(self): + """The overflow note's M value must equal len(wos) - rendered_count.""" + n = 500 + wos = [_make_wo(i) for i in range(n)] + modal = build_wo_modal("3rd Escalations", wos, DASHBOARD_URL) + blocks = modal["blocks"] + # Last block should be context with the overflow note. + last_text = _extract_text([blocks[-1]]) + # The overflow note pattern is "+M more". + import re + + match = re.search(r"\+(\d+) more", last_text) + assert match, f"no '+M more' pattern in last block: {last_text!r}" + reported_overflow = int(match.group(1)) + # Block count: total blocks = visible WO blocks + 1 overflow block. + rendered_wo_blocks = len(blocks) - 1 + assert reported_overflow == n - rendered_wo_blocks, ( + f"overflow note says +{reported_overflow} but " + f"{n} - {rendered_wo_blocks} = {n - rendered_wo_blocks}" + ) + + def test_empty_wo_list_returns_placeholder(self): + modal = build_wo_modal("Empty", [], DASHBOARD_URL) + all_text = _extract_text(modal["blocks"]) + assert "no work orders" in all_text.lower() + + def test_wo_number_appears_in_blocks(self): + modal = build_wo_modal("Test", [_make_wo(7)], DASHBOARD_URL) + all_text = _extract_text(modal["blocks"]) + assert "WO-0007" in all_text + + def test_snippet_truncated_to_150_chars(self): + wo = _make_wo(1) + wo["last_comment"] = "X" * 300 + modal = build_wo_modal("Test", [wo], DASHBOARD_URL) + all_text = _extract_text(modal["blocks"]) + # The snippet in the text should not be 300 X's. + assert "X" * 200 not in all_text + + def test_no_apm_deep_links(self): + """Modal must not contain any APM application deep-links.""" + wos = [_make_wo(i) for i in range(10)] + modal = build_wo_modal("3rd Escalations", wos, DASHBOARD_URL) + all_text = _extract_text(modal["blocks"]) + # APM URLs would start with typical patterns — confirm none present. + assert "apm://" not in all_text + assert "app.apm" not in all_text + + +# --------------------------------------------------------------------------- +# Helpers +# --------------------------------------------------------------------------- + + +def _extract_text(blocks: list[dict]) -> str: + """Collect all text strings from blocks for assertion convenience.""" + parts: list[str] = [] + for block in blocks: + btype = block.get("type") + if btype == "header": + parts.append(block.get("text", {}).get("text", "")) + elif btype == "section": + text_obj = block.get("text") + if text_obj: + parts.append(text_obj.get("text", "")) + for field in block.get("fields", []): + parts.append(field.get("text", "")) + elif btype == "context": + for elem in block.get("elements", []): + parts.append(elem.get("text", "")) + elif btype == "actions": + for elem in block.get("elements", []): + txt = elem.get("text", {}) + if isinstance(txt, dict): + parts.append(txt.get("text", "")) + return "\n".join(parts) + + +def _find_link_button_url(blocks: list[dict]) -> str | None: + """Return the URL from the first link button found in any actions block.""" + for block in blocks: + if block.get("type") != "actions": + continue + for elem in block.get("elements", []): + if "url" in elem: + return elem["url"] + return None diff --git a/tests/test_pipeline_synth.py b/tests/test_pipeline_synth.py index 19540c9..eec10ac 100644 --- a/tests/test_pipeline_synth.py +++ b/tests/test_pipeline_synth.py @@ -1,16 +1,19 @@ -"""Synth-level assertions for the Phase 3 analytics dataset. +"""Synth-level assertions for the analytics dataset (Phase 3) and Slack surfaces +(Phase 4). -Synthesizes ``apm-wo-analysis-pipeline`` and asserts the Glue table carries -partition projection, the column schema matches what the classifier writes, the -Athena workgroup enforces its result location, and the classifier role has zero -Glue access (projection means it never touches the catalog). No AWS, no Docker: -``aws:cdk:bundling-stacks=[]`` skips asset bundling so this is a fast offline gate. +Synthesizes ``apm-wo-analysis-pipeline`` and asserts: the Glue table carries +partition projection with the classifier's column schema; the Athena workgroup +enforces its result location; the classifier role has zero Glue access; and the +Phase 4 Slack post + interactions Lambdas exist behind an HTTP API with only the +scoped S3/Secrets/SSM permissions. No AWS, no Docker: ``aws:cdk:bundling-stacks=[]`` +skips asset bundling so this is a fast offline gate. Run with the repo venv: python -m pytest tests/test_pipeline_synth.py -q """ +import json import sys from pathlib import Path @@ -28,8 +31,16 @@ def _s3_path_ending(suffix: str): return {"Fn::Join": ["", Match.array_with([suffix])]} +def _cdk_context() -> dict: + """The real cdk.json context (cert ARN, hosted zone, Slack domain) — the + hand-built App below doesn't auto-load it the way `cdk synth` does.""" + ctx = json.loads((CDK_DIR / "cdk.json").read_text())["context"] + ctx["aws:cdk:bundling-stacks"] = [] + return ctx + + def _template() -> Template: - app = cdk.App(context={"aws:cdk:bundling-stacks": []}) + app = cdk.App(context=_cdk_context()) stack = PipelineStack( app, "apm-wo-analysis-pipeline", @@ -117,3 +128,59 @@ def test_classifier_role_has_no_glue_access(): a for a in actions if isinstance(a, str) and a.startswith("glue:") ] assert not offending, f"unexpected Glue access: {offending}" + + +# ----- Phase 4 — Slack surfaces ----- + + +def test_slack_lambdas_exist(): + t = _template() + for fn_name, handler in ( + ("apm-wo-analysis-slack-post", "handler.handler"), + ("apm-wo-analysis-slack-interactions", "interactions.handler"), + ): + t.has_resource_properties( + "AWS::Lambda::Function", + { + "FunctionName": fn_name, + "Handler": handler, + "Runtime": "python3.12", + "Architectures": ["arm64"], + }, + ) + + +def test_interactions_api_routes_post_to_slack_endpoint(): + t = _template() + t.resource_count_is("AWS::ApiGatewayV2::Api", 1) + t.has_resource_properties( + "AWS::ApiGatewayV2::Route", {"RouteKey": "POST /slack/interactions"} + ) + # Custom domain on apm-wo.seahaven.com + a Route53 alias for it. + t.has_resource_properties( + "AWS::ApiGatewayV2::DomainName", {"DomainName": "apm-wo.seahaven.com"} + ) + t.has_resource_properties("AWS::Route53::RecordSet", {"Type": "A"}) + + +def test_slack_roles_have_no_broad_or_write_access(): + # The Slack Lambdas should only read analytics/, the Slack secret, and the + # dashboard SSM param — never write S3, never s3:*/secretsmanager:* wildcards. + # Scope to the Slack policies by logical ID (the classifier/drop-uploader + # legitimately hold s3:PutObject). + t = _template() + slack_policies = { + lid: p + for lid, p in t.find_resources("AWS::IAM::Policy").items() + if lid.startswith(("SlackPost", "SlackInteractions")) + } + assert slack_policies, "expected scoped policies for the Slack roles" + for policy in slack_policies.values(): + for stmt in policy["Properties"]["PolicyDocument"]["Statement"]: + actions = stmt.get("Action", []) + actions = [actions] if isinstance(actions, str) else actions + for a in actions: + if not isinstance(a, str): + continue + assert a not in ("s3:*", "secretsmanager:*", "*"), f"too broad: {a}" + assert a != "s3:PutObject", "Slack roles must not write S3"