Ship one structured event model, propagate a correlation ID through every request, and send logs to the collection API outside the response path. TL;DR: for a healthtech experiment split across tenant cohorts, rollback safety depends on being able to join an outcome to the exact tenant, cohort, release, and request without turning any of those fields into uncontrolled metric labels.
The transport is not the hard part. The hard part is preserving evidence when the collector slows down, a retry duplicates a batch, or a sensitive field slips into an exception. This architecture decision record treats log delivery as bounded and lossy, while keeping the application decision deterministic: a cohort can advance only when its correlated error and latency signals remain inside a predeclared rollback envelope.
Keep it bounded.
What must remain true when delivery fails?
Four invariants drive the design.
- A valid inbound correlation ID is reused; otherwise the service creates one. The same value appears in the response header and every event produced for that request.
- Tenant and cohort are explicit fields, but patient identifiers, message bodies, OTP values, email addresses, and phone numbers never enter the event. GDPR Article 17 is one reason to minimize personal data before storage rather than depend on later deletion.
- The application response never waits for the remote log API. A full buffer drops the new diagnostic event and increments a low-cardinality counter; it does not turn a successful clinical workflow into an error.
- Every experiment event carries
release,experiment, andvariant, so rollback analysis does not infer deployment state from timestamps.
That last point matters during a staggered rollout. Two tenants can send the same request shape at the same instant while running different releases. A bare correlation ID reconstructs one request, but it cannot identify the cohort decision that exposed it.
The failure boundary is deliberate: process-local buffering can lose the final events during a crash. If audit records or legally required access trails must be durable, they belong in a transactional audit path with their own retention and access controls. Diagnostic logs cannot quietly inherit that responsibility. This is the central trade-off: response isolation is bought by accepting a measured gap in diagnostic delivery.
Compare the transport boundaries
| Option | Request latency | Failure behavior | Rollback evidence | Operational cost |
|---|---|---|---|---|
| Synchronous HTTP send | Includes collector latency | Collector failure can affect the caller | Immediate when available | Simple code, coupled systems |
| In-process bounded queue | Decoupled from collector latency | Drops are measurable when the queue is full or the process exits | Usually prompt; a small tail may be absent | Queue, worker, retry, and shutdown logic |
| Durable local or brokered queue | Decoupled | Survives more process failures | Delayed but recoverable | Another durable subsystem to operate |
For an experiment where rollback safety is the primary decision axis, a bounded queue is a reasonable default only when missing diagnostic events are visible and rollback can also rely on aggregate service metrics. Choose a durable queue when each event is itself evidence that must survive a restart. Choose synchronous delivery only for a control plane whose caller can safely absorb collector failure and latency.
Do not put tenant_id, correlation_id, or raw routes into metric labels. Prometheus instrumentation guidance warns that every unique label combination creates another time series. Keep metrics aggregated by bounded dimensions such as service, cohort, operation, and outcome; use structured logs to investigate a specific tenant or request.
How should a NestJS custom logger send structured logs over HTTP?
The example below is runnable with the Python standard library. It models the same boundaries a backend framework transport needs: request-context propagation, a bounded queue, batch delivery, redaction, timeouts, and a drop counter. In a NestJS service, the custom logger adapter should stay thin and hand each normalized event to the equivalent of emit(); lifecycle wiring belongs at the application boundary, while correlation context is established by request middleware.
from __future__ import annotations
import contextvars
import json
import queue
import re
import threading
import time
import urllib.request
import uuid
from dataclasses import asdict, dataclass
from typing import Any
correlation_id_var = contextvars.ContextVar("correlation_id", default="")
SENSITIVE_KEYS = re.compile(
r"^(authorization|cookie|email|phone|patient_id|otp|token)$", re.I
)
@dataclass(frozen=True)
class LogEvent:
timestamp_ms: int
level: str
message: str
correlation_id: str
tenant_id: str
cohort: str
release: str
experiment: str
variant: str
attributes: dict[str, Any]
def clean(value: Any) -> Any:
if isinstance(value, dict):
return {
key: "[REDACTED]" if SENSITIVE_KEYS.match(key) else clean(item)
for key, item in value.items()
}
if isinstance(value, list):
return [clean(item) for item in value]
return value
class HttpBatchTransport:
def __init__(self, endpoint: str, capacity: int = 2_000) -> None:
self.endpoint = endpoint
self.pending: queue.Queue[LogEvent] = queue.Queue(maxsize=capacity)
self.dropped = 0
self.stopping = threading.Event()
self.worker = threading.Thread(target=self._run, daemon=True)
self.worker.start()
def emit(self, event: LogEvent) -> None:
try:
self.pending.put_nowait(event)
except queue.Full:
self.dropped += 1
def _run(self) -> None:
while not self.stopping.is_set() or not self.pending.empty():
batch: list[LogEvent] = []
try:
batch.append(self.pending.get(timeout=0.25))
except queue.Empty:
continue
while len(batch) < 100:
try:
batch.append(self.pending.get_nowait())
except queue.Empty:
break
self._send(batch)
def _send(self, batch: list[LogEvent]) -> None:
body = json.dumps(
{"events": [clean(asdict(event)) for event in batch]},
separators=(",", ":"),
).encode("utf-8")
request = urllib.request.Request(
self.endpoint,
data=body,
headers={"Content-Type": "application/json"},
method="POST",
)
try:
with urllib.request.urlopen(request, timeout=2.0) as response:
response.read()
except OSError:
self.dropped += len(batch)
def close(self, timeout: float = 3.0) -> None:
self.stopping.set()
self.worker.join(timeout=timeout)
def begin_request(inbound_id: str | None) -> str:
candidate = (inbound_id or "").strip()
correlation_id = candidate if 1 <= len(candidate) <= 128 else str(uuid.uuid4())
correlation_id_var.set(correlation_id)
return correlation_id
def experiment_event(
*, tenant_id: str, cohort: str, outcome: str, duration_ms: int
) -> LogEvent:
return LogEvent(
timestamp_ms=int(time.time() * 1_000),
level="info",
message="experiment request completed",
correlation_id=correlation_id_var.get(),
tenant_id=tenant_id,
cohort=cohort,
release="release-42",
experiment="appointment-reminder-copy",
variant="treatment",
attributes={"outcome": outcome, "duration_ms": duration_ms},
)
Treat the inbound ID as untrusted input. Bound its length, reject characters your downstream format cannot represent, and create a UUID when it is absent or invalid. Also return the selected ID in the response so an operator can connect a support report to the server-side event. Do not accept tenant or cohort identity from a correlation header; derive those values from authenticated server-side context.
The concrete limits here are visible: the queue holds 2,000 events, each HTTP payload contains at most 100, the send timeout is 2.0 seconds, and shutdown waits up to 3.0 seconds. Those aren't universal defaults. They are a reviewable capacity contract. For example, if a cohort emits 400 events per second and collection is unavailable for ten seconds, this queue cannot preserve all 4,000 events; it will retain at most 2,000 pending events and count the rest as dropped. Raising capacity only changes how long the process can absorb the mismatch, and it consumes more memory while doing so. The correct values come from measured event size, burst rate, memory budget, and the maximum tolerable diagnostic gap.
The sample counts a failed batch as dropped instead of retrying forever. That is blunt on purpose. Production retry needs a maximum attempt count, exponential backoff with jitter, and a shutdown deadline. Without those bounds, an unreachable collector turns an observability helper into unbounded memory growth or a deployment that will not terminate.
It can still lose data.
Turn events into a rollback decision
Define the rule before exposure begins. For each bounded cohort, compare treatment with control over the same observation window and release. Useful service-side measures include request count, error ratio, and latency percentiles. If the experiment affects a web interaction, Core Web Vitals uses the 75th percentile as the assessment point for LCP, INP, and CLS; preserve that percentile convention when it is relevant rather than averaging away a slow tail.
A practical decision record states the minimum sample requirement, evaluation window, and rollback thresholds. It also names exclusions, such as synthetic checks and canceled requests. No magic number is universal, so the values must come from the workflow's clinical risk and existing service objective, not from the logging transport.
Rollback immediately when the predeclared safety boundary is crossed. Pause expansion when the log-drop counter rises, because missing evidence weakens the comparison even if aggregate metrics still look healthy. Continue only when treatment and control are comparable by release and observation window.
This separation is useful: metrics answer whether a cohort is unhealthy; correlated logs explain which request path failed. Trying to make one signal do both jobs creates either high-cardinality metrics or logs that are too vague to reconstruct an incident.
No single signal gets veto power.
Test the failure modes before exposure
Start the service with the collector unreachable and fill the queue. Requests should still complete, the drop counter should increase, memory should remain bounded, and shutdown should respect its deadline. Then send two concurrent requests with different correlation IDs and verify that context never crosses between them.
Run a redaction test with nested dictionaries and lists. Search the serialized body for a known email address, OTP, authorization value, and patient identifier. The correct count is zero.
Finally, replay one batch to the receiver. Because a timeout can occur after the receiver stores data but before the sender sees the response, duplicates are possible. Consumers should tolerate them, and analytical queries should avoid treating raw event count as unique request count. A stable event ID can support deduplication when that distinction matters.
Why reject synchronous delivery?
Sending each event inline looks attractive: there is no queue worker, and a successful call appears to prove delivery. It also couples patient-facing request latency and availability to the collector. For an appointment or OTP flow, that boundary is backwards; diagnostic plumbing must not widen a delivery gap.
The asynchronous design has a hard limitation: it is unsuitable when every record must be durably acknowledged before the business action commits. Its drop counter reports the evidence gap but cannot recover the missing events. A durable local journal or brokered queue is the better boundary for that requirement, at the cost of another storage lifecycle, backpressure policy, and recovery procedure.
Synchronous delivery still has a valid use case. A low-volume administrative action may require a separately designed, durable audit acknowledgement before the action succeeds. Call that an audit transaction, give it an explicit failure contract, and keep it out of the general application logger.
The resulting decision is narrow but defensible: use bounded asynchronous transport for diagnostic events, bounded metric dimensions for cohort health, and a separate durable path for records that cannot be lost. Correlation makes the evidence joinable. The rollback rule makes it actionable.
Top comments (0)