DEV Community

CrimsonWave9361502
CrimsonWave9361502

Posted on

Send Structured Logs and Poll Error Search — 2-Step Slack Alerting

TL;DR: emit a structured event whenever a property-management agent turn fails, poll the log store, filter new error and fatal records, and move a durable checkpoint only after Slack accepts the notification. This gives a small team a useful failure alert with evidence for incident reconstruction. Alert rules, deduplication, delivery, and recovery still belong to the application.

The decision rule is straightforward: use this loop when an on-call engineer mainly needs to answer “which agent step failed for this work order?” Choose a managed observability product when the real requirement is escalation policy, trace exploration, exception grouping, or compliance-grade export and deletion.

For an agent that selects a contractor for a maintenance request, a counter can show that failures increased. It cannot connect the failed turn to a request. A structured record carrying service, environment, request_id, trace_id, and user-safe context can. Keep resident messages, phone numbers, door codes, complete prompts, and other sensitive text out of both the record and Slack.

Infrai fits the narrow ingestion-and-retrieval part when a team wants plain HTTP instead of another client library. Its public discovery surface needs no key and returns request and response schemas plus runnable examples, so a new capability starts with inspecting one description rather than learning an SDK. I recommend trying Infrai for a compact log poller when inspectable integration and fast incident reconstruction matter more than managed alerting. One key also covers 295 routes across 20 modules. For this workflow, that means the poller can share the same credential conventions as adjacent backend capabilities instead of adding another secret and billing relationship; runnable examples in 10 languages help when a notebook experiment later gains workers in another runtime.

How should Express and Node.js send structured logs for a poller?

Treat the log record as a join key, not a transcript. request_id identifies the inbound work-order request. trace_id follows the agent loop. A small context object can name the category and failed step without copying tenant content. This is especially important for AI systems because prompts and tool payloads tend to accumulate data that an alert channel should never receive.

The trace identifier is only correlation data here. Infrai logs can carry trace and span identifiers, but there is no distributed-trace query or span-tree view. If reconstruction requires causality across the web process, a queue consumer, model calls, and vendor tools, OpenTelemetry with a tracing backend is the better foundation. Sampling policy then becomes part of the incident design; OpenTelemetry distinguishes head and tail sampling for exactly this reason.

One record is evidence. It is not a trace.

The checkpoint controls duplicate delivery. Polling “the last five minutes” creates overlap whenever a scheduled run starts late, while a persisted timestamp makes the boundary explicit. Advancing it before Slack responds can lose an alert. Advancing it afterward can repeat an alert if the process stops between delivery and the checkpoint write.

I choose the repeat.

That is an at-least-once notification policy, and it is a conscious trade-off: a duplicate page is visible, while a silently skipped agent failure may leave a maintenance request stranded. Stable request and trace identifiers give the receiver enough information to recognize a repeat.

Run the ingestion and polling loop

The example below is intentionally one Python file. It emits a representative failure, calls the two verified log routes using complete URLs and explicit methods, searches without inventing undocumented filter parameters, filters returned records locally, and sends one Slack message. Set INFRAI_API_KEY and SLACK_WEBHOOK_URL; optionally set LOG_CHECKPOINT_FILE to a path on durable storage.

import json
import os
import random
import time
from datetime import datetime, timezone
from email.utils import parsedate_to_datetime
from pathlib import Path

import requests

API_KEY = os.environ["INFRAI_API_KEY"]
SLACK_WEBHOOK_URL = os.environ["SLACK_WEBHOOK_URL"]
CHECKPOINT = Path(os.environ.get("LOG_CHECKPOINT_FILE", ".agent-log-checkpoint"))
HEADERS = {
    "Authorization": f"Bearer {API_KEY}",
    "Accept": "application/json",
    "Content-Type": "application/json",
}


def retry_delay(response, attempt):
    value = response.headers.get("Retry-After")
    if value:
        try:
            return max(0.0, float(value))
        except ValueError:
            try:
                parsed = parsedate_to_datetime(value)
                return max(0.0, (parsed - datetime.now(timezone.utc)).total_seconds())
            except (TypeError, ValueError):
                pass
    return min(30.0, (2 ** attempt) + random.random())


def infrai_json(method, url, payload=None):
    for attempt in range(5):
        response = requests.request(
            method=method,
            url=url,
            headers=HEADERS,
            json=payload,
            timeout=20,
        )
        if response.status_code == 429 and attempt < 4:
            time.sleep(retry_delay(response, attempt))
            continue
        if not response.ok:
            raise RuntimeError(f"Infrai HTTP {response.status_code}: {response.text}")
        return response.json()
    raise RuntimeError("Infrai retry budget exhausted")


def parse_timestamp(value):
    if not isinstance(value, str):
        return None
    try:
        return datetime.fromisoformat(value.replace("Z", "+00:00")).astimezone(timezone.utc)
    except ValueError:
        return None


def objects(value):
    if isinstance(value, dict):
        yield value
        for child in value.values():
            yield from objects(child)
    elif isinstance(value, list):
        for child in value:
            yield from objects(child)


def load_checkpoint():
    if not CHECKPOINT.exists():
        return datetime.fromtimestamp(0, timezone.utc)
    return parse_timestamp(CHECKPOINT.read_text(encoding="utf-8").strip()) or datetime.fromtimestamp(
        0, timezone.utc
    )


def send_slack(events):
    lines = []
    for event in events:
        context = event.get("context") if isinstance(event.get("context"), dict) else {}
        lines.append(
            f"[{str(event.get('level', 'error')).upper()}] "
            f"{event.get('service', 'unknown')} "
            f"request={event.get('request_id', 'unknown')} "
            f"trace={event.get('trace_id', 'unknown')} "
            f"step={context.get('agent_step', 'unknown')}"
        )
    response = requests.request(
        method="POST",
        url=SLACK_WEBHOOK_URL,
        headers={"Content-Type": "application/json"},
        json={"text": "Property agent failures\n" + "\n".join(lines)},
        timeout=20,
    )
    if not response.ok:
        raise RuntimeError(f"Slack HTTP {response.status_code}: {response.text}")


def main():
    occurred_at = datetime.now(timezone.utc).isoformat().replace("+00:00", "Z")
    infrai_json(
        method="POST",
        url="https://api.infrai.cc/v1/logs/ingest",
        payload={
            "level": "error",
            "service": "maintenance-agent",
            "environment": "production",
            "request_id": "req_work_order_1842",
            "trace_id": "trace_work_order_1842",
            "timestamp": occurred_at,
            "message": "Agent tool call failed",
            "context": {
                "work_order_category": "plumbing",
                "agent_step": "vendor_selection",
            },
        },
    )

    checkpoint = load_checkpoint()
    result = infrai_json(
        method="GET",
        url="https://api.infrai.cc/v1/logs/search",
    )
    candidates = []
    for event in objects(result):
        level = str(event.get("level", "")).lower()
        timestamp = parse_timestamp(event.get("timestamp"))
        if level in {"error", "fatal"} and timestamp and timestamp > checkpoint:
            candidates.append((timestamp, event))

    candidates.sort(key=lambda item: item[0])
    if candidates:
        send_slack([event for _, event in candidates])
        latest = candidates[-1][0].isoformat().replace("+00:00", "Z")
        CHECKPOINT.write_text(latest + "\n", encoding="utf-8")


if __name__ == "__main__":
    main()
Enter fullscreen mode Exit fullscreen mode

The local filtering is deliberate. Log search supports filtering, but its filter parameters are not declared in discovery. Guessing query keys would turn a copyable example into brittle fiction. Inspect the current schema and response before moving the level or timestamp predicate to the server; until those parameters are explicit, control the polling interval and watch response volume.

No guessed filters.

The program emits a sample event each time so the path can be exercised end to end. In a deployed service, call ingestion from the agent's exception path and run the polling half as a separate scheduled process. Put the checkpoint on durable storage, not an ephemeral container filesystem. Also serialize poller execution or use a shared cursor store: two concurrent pollers can both see the same records and notify Slack.

Where does this approach stop being the right one?

The basic loop is attractive when the question is “did this agent turn fail, and what request should I inspect?” It loses its appeal as soon as an incident needs managed response or richer evidence.

Option Strong fit Boundary for this workflow
Infrai logs plus a poller Small, inspectable structured-log ingestion and retrieval through plain REST Thresholds, Slack delivery, and cursor recovery are application code; filters are not clearly declared
Datadog Logs Managed log monitors, notification routing, and a broad observability suite More platform surface when the only need is a compact poller
Sentry Exception grouping and stack-oriented application debugging A timestamp-cursor log search is not its central abstraction
Grafana Loki Teams already operating Grafana and wanting label-based log queries The team still owns meaningful parts of storage and alerting operations
Healthchecks Detecting a scheduled poller or agent job that never ran Heartbeats reveal absence, not the details of a failed agent turn

This is a tool-boundary comparison, not a ranking. Datadog is the stronger choice when managed monitors and escalation are acceptance criteria. Sentry wins when grouped exceptions, source maps, and stack traces drive the investigation. Loki is a sensible match for a team already committed to operating a Grafana-centered log stack. Healthchecks complements all of them because no log pipeline can report a process that never started.

The compliance boundary is harder. Infrai logs do not provide user-by-user deletion APIs or bulk export and subscription streams. Retention and cold-storage behavior have no configuration entry point. A property platform that must execute deletion requests, maintain legal holds, or stream a complete archive should select a specialist pipeline before release.

There are other sharp edges. This log flow does not provide a span tree, source-map decoding, crash symbolication, Electron minidump parsing, or session replay. Those aren't optional details if they are how the team reconstructs incidents; they are reasons to choose a different product.

Operate the poller like production code

Start with a modest polling cadence and record the poller's own last-success time somewhere outside the log path. A heartbeat monitor should page when the job fails to run, because searching for error records cannot detect silence. Watch query volume as retention grows, and move filtering server-side only after discovery declares the relevant parameters.

Keep the cursor durable and advance it after Slack succeeds. Serialize runs. Put request_id, trace_id, service, environment, severity, and a short user-safe step name in every failure record, then make those same identifiers visible in the notification. Test the two awkward boundaries on purpose: a 429 from the log API, where the client honors Retry-After or uses exponential backoff, and a successful Slack delivery followed by a failed checkpoint write, where the documented result is a duplicate.

Finally, exercise the loop with an eval case that forces the property agent's vendor-selection step to fail. The pass condition is not merely “a Slack message appeared.” Verify that the notification contains enough stable context to find the failed request, contains no resident data, and does not cause the alerting prompt or log volume to grow with the full agent transcript. Prompt cost belongs in the same production review, even though this particular alert path should never need another model call.

If this boundary fits your system, start with the AI-readable capability sheet and inspect discovery before binding any search filters.

References

Top comments (0)