v1.3.2 Design Plan: Metrics/Stats Parity, Multi-Worker Streams, Record Shaping, and Batch Resilience
Status: Implemented on the 1.3.2 branch; archived as the consolidated design record
Target version: v1.3.2
Branch: observer
Scope: Standalone and manager modes. Worker mode is execution-only; it already reports stats upstream.
This archived plan now consolidates the later 1.3.2 follow-up design notes that were previously
tracked separately as:
structured-record-shaping-design.mdv1.3.2-batch-recovery-plan.md
Goal
v1.3.2 closes six backend gaps deferred from v1.3.0/v1.3.1 or discovered during the implementation/validation cycle:
- Standalone live stats parity (A) — standalone mode exposes live stats for active stream
pipelines via
StatsStoreand the existingGET /api/cluster/streamssurface, instead of only DB run-history and process-local Prometheus metrics. - Manager operational metrics (B) — Prometheus series on the manager process for dispatch decisions, reconciler actions, worker health, and callback receipt.
- UDP multi-worker streams (C) — lift the L008 linter block on
syslogandsnmp_trapsources in manager mode; provision K8s Services for UDP multi-worker streams via the same generickubernetes:block used by HTTP sources (kubernetes: enabled: truerequired in manager mode for all UDP push sources); validate the UDP path end-to-end in kind. - ASN.1 structured decode flattening (D) — support concatenated BER files that contain
multiple top-level records; add a generic
json_flattentransform for nested decoded records plus a generichex_decodetransform for text/hex heuristics and field-path decode overrides. - CDR structured record shaping (E) — add explicit nested-path shaping primitives and a
predictable
json_flattencontract for LTE/SGW/PGW-style decoded payloads, without relying on heuristic row expansion or vendor-specific flatteners. - Batch resilience for large CDR files (F) — reconcile lost worker-owned batch runs, process large decoded files in bounded record windows, and make serial record-safe file outputs safe across reruns and worker failure.
This release does not contain:
- UI revalidation (v1.3.3)
- Pipeline cloning, bulk actions, or live log streaming (backlog)
- Connector hardening (backlog)
- Blind schema-wide ASN.1 root-type auto-detection by default
Design Decisions
1. Standalone stats: stream-only, in-process loop, no HTTP hop
create_app() already creates StatsStore unconditionally and wires it into controller and
app.state.stats_store. The gap is that standalone _stream_worker() calls
executor.stream_run() with no stats= argument, so PipelineStats is never created and
StatsStore always stays empty.
Scope: stream pipelines only. Batch runs in standalone mode are finite and synchronous —
they complete in seconds to minutes and their final counts are already written to run_history
via manager.record_run(). There is no operational benefit to making batch runs visible in
StatsStore (no load scoring, no placement reconciler). Mid-run batch visibility is deferred.
In-process update, not HTTP. The standalone stats loop constructs PipelineStatsPayload
directly from PipelineStats.snapshot_and_reset_window() and calls stats_store.update()
in-process. No HTTP call, no /api/internal/pipeline-stats route involved.
Removal is explicit, never via is_final. The is_final flag is a worker-to-manager HTTP
protocol convention. Its only consumer is the POST /api/internal/pipeline-stats HTTP handler
in tram/api/routers/internal.py, which calls store.remove(run_id) when is_final=True.
There is no local consumer of is_final in the standalone process. Calling
stats_store.update(payload) with is_final=True in-process would leave completed entries in
the store until the 3× interval stale timeout — it does not trigger removal.
The standalone path therefore uses two explicit calls:
- The stats loop always calls
stats_store.update(payload)withis_final=False. _stream_worker()callsstats_store.remove(run_id)directly in itsfinallyblock when the stream exits.
Approximate parity, not structural parity. GET /api/cluster/streams in standalone reuses
the existing build_cluster_streams() fallback in _stream_views.py, which surfaces
StatsStore.all_active() entries that are not associated with a placement group. This path
already exists and requires no changes — but it produces synthetic views that are structurally
different from manager-mode placements in three ways:
placement_group_idisnull— no persisted group identity, no controller-owned slot.started_atisnull— no placement-backed timestamp; only stats timestamps are available.- No reconciliation: a standalone stream that stops updating will age out of
all_active()after3 × intervalseconds, but there is no reconciler to re-dispatch it.
The roadmap goal of “same live stats model” refers to counters and rates being live and observable, not to structural placement parity. True structural parity would require the controller to own synthetic placement objects — that is explicitly deferred beyond v1.3.2.
2. Manager metrics: instrument the server side, not the client side
The manager process owns POST /api/internal/run-complete and
POST /api/internal/pipeline-stats. Counters incremented in these handlers are genuinely
manager-side metrics. They record what the manager actually received, which is operationally
useful (e.g., “how many run-completes have I processed today?”).
Do NOT instrument _post_run_complete() or _post_stats() in tram/agent/server.py.
Those functions execute inside the worker process. Any counters placed there appear on
worker /metrics, not manager /metrics. They are the wrong instrumentation point for the
stated requirement.
There is no way for the manager to directly observe callback failures that occurred inside a
worker process. The inferred count tram_mgr_dispatch_total{result="accepted"} -
tram_mgr_run_complete_received_total is the closest approximation for batch dispatches. The
design acknowledges this gap explicitly and does not claim a direct “callback failures” counter.
3. Terminology: “broadcast” is a placement word, not a data word
The v1.3.0 feature is called “broadcast streams” and the placement table is named
broadcast_placements. “Broadcast” was chosen to describe the placement act — the pipeline
config is dispatched to all healthy workers simultaneously. It has no meaning about data routing.
In practice “broadcast” reads as a data-routing term to most engineers: “every packet goes to
every recipient.” Combined with count: all, the feature name implies data is being duplicated
across workers. It is not.
Decision: rename the user-facing feature to “multi-worker streams” in v1.3.2.
What changes:
- This document title and all docs going forward use “multi-worker streams”
GET /api/cluster/streamsresponse field names stay as-is (no API break)- Log messages: “broadcast stream dispatched” → “multi-worker stream dispatched”
roadmap.mdandarchitecture.mdupdated in this release
What does not change:
broadcast_placementsDB table name — internal, no user-visible surface, migration cost not worth itplacement_group_idfield — already in the API, renaming would be a breaking changeworkers.count: allYAML syntax —allandeachmean the same thing to a user; no gain- Internal code symbols (
_broadcast_placements,redispatch_broadcast_slot, etc.) — internal
4. UDP multi-worker streams: shared ingress vs. per-slot endpoints
workers.count: all for syslog and snmp_trap means N workers each independently run a
pipeline instance. No packet is replicated — each packet goes to exactly one worker.
For count: all + stateless sinks (Kafka, OpenSearch, REST, etc.) a shared UDP NodePort
Service works correctly. Kube-proxy ECMP hashes per flow (src_ip + src_port + dst_ip +
dst_port), so each upstream sender tends to consistently land on one worker, and it doesn’t
matter which worker processes which trap as long as one does. This is semantically equivalent
to how HTTP push sources work today with a shared LB — the distribution is automatic and no
per-pod addressing is needed.
Targeted Services are needed in two cases:
-
workers.list(constrained backend set): Only specific workers are dispatched. A Service selecting all workers would route traffic to non-participating pods. The implementation uses one Service with manual Endpoints pointing to only the listed workers. This constrains the backend set but does not guarantee a specific sender always reaches a specific worker — kube-proxy ECMP still distributes flows among the listed backends. True per-worker sender pinning (SNMP manager X always →worker-0) is not achievable with a Service; that would require the SNMP manager to send directly to the pod IP. This requirement is out of scope for v1.3.2. -
count: N(partial placement): Only N of the available workers are dispatched. A Service selecting all workers would route to non-participating pods. Same manual Endpoints approach asworkers.list.
Design choice for v1.3.2: Require kubernetes: enabled: true for all UDP push sources in
manager mode. There is no pre-existing worker-targeting UDP Service in the chart: the static
worker-ingress-service.yaml is TCP-only (:8767), and service.yaml selects the manager
pod in manager mode. A “shared NodePort for count: all with no kubernetes block” would have
no routing path to workers without introducing new Helm infrastructure. Rather than add that
infrastructure, the design reuses the same pipeline-specific Service mechanism already built
for HTTP workers.list: L012 is an error (not a warning) that blocks UDP push sources in
manager mode when kubernetes: enabled: true is absent, regardless of count value.
count: all + kubernetes: enabled: true gets a Service with a broad label selector (all
worker pods as backends). kube-proxy ECMP provides natural distribution. This is the
recommended model for snmp_trap → kafka and similar stateless-sink pipelines.
count: N and workers.list get a Service with manual Endpoints targeting only the
dispatched workers — the same path used by HTTP workers.list today.
Workstream A — Standalone Live Stats Parity
Objectives
GET /api/cluster/streamsreturns live stats for active stream pipelines in standalone mode.GET /api/pipelines/{name}/placementreturns a synthetic single-slot view for an active standalone stream (acknowledged approximation:placement_group_id: null,started_at: null). All other cases — non-stream pipelines, stopped streams, manager mode with no placement — keep existing404behavior.- Stale detection works identically: entries age out after
3 × intervalseconds if the loop stops updating. - All changes are gated on
self._worker_pool is None— zero effect in manager mode.
Runtime model
standalone stream pipeline running
│
├─ controller._stream_worker(config, stop_event)
│ run_id = str(uuid4())
│ stats = PipelineStats(run_id=run_id, pipeline_name=config.name, schedule_type="stream")
│ local_run = _LocalRun(run_id, pipeline_name, "stream", started_at=now, stats=stats)
│ _local_active_stats[run_id] = local_run ← protected by _local_stats_lock
│ executor.stream_run(config, stop_event, stats=stats)
│ [on exit/finally]
│ stats_store.remove(run_id) ← explicit removal, no is_final
│ _local_active_stats.pop(run_id)
│
└─ controller._local_stats_loop [daemon thread, standalone only]
every stats_interval seconds:
with _local_stats_lock:
runs = list(_local_active_stats.items())
for run_id, local_run in runs:
payload = PipelineStatsPayload(
worker_id=node_id,
pipeline_name=local_run.pipeline_name,
run_id=run_id,
schedule_type="stream",
uptime_seconds=(now - local_run.started_at).total_seconds(),
timestamp=now,
is_final=False,
**local_run.stats.snapshot_and_reset_window(),
)
stats_store.update(payload) ← in-process, no HTTP
Batch runs are excluded from StatsStore entirely. Their final counts reach run_history
via manager.record_run() as before. No stats_store.update() or stats_store.remove() call
is made for batch runs — not even with is_final=True. That flag has no local consumer and
would leave orphaned entries.
New data structures on PipelineController
@dataclass
class _LocalRun:
run_id: str
pipeline_name: str
schedule_type: str
started_at: datetime
stats: PipelineStats
# on PipelineController:
_local_active_stats: dict[str, _LocalRun] # guarded by _local_stats_lock
_local_stats_lock: threading.Lock
_local_stats_stop: threading.Event
_local_active_stats uses its own lock, separate from _stream_threads, because APScheduler
and the stats loop access different state.
run_id threading through _run_batch
Currently, the APScheduler path in _run_batch only passes run_id to the executor when
called from trigger_run(). In the standalone local path, run_id is generated inside
PipelineRunContext if not supplied, making it unobservable to the controller. A one-line
fix: always generate run_id = str(uuid.uuid4()) in _run_batch before calling the executor.
This has no observable effect on existing behavior and enables future batch stats work.
Placement endpoint behavior in standalone
GET /api/pipelines/{name}/placement currently returns 404 when no placement exists. The
v1.3.2 standalone path adds one new case only — active streams — and keeps all other cases as
404 to avoid broadening the API contract:
- Standalone, stream pipeline, active entry in
stats_store.for_pipeline(name): return a synthetic single-slot view withplacement_group_id: null,started_at: null,status: "running",target_count: 1, one slot with live stats. - All other cases (pipeline not found, pipeline not stream type, stream not yet reporting,
stream stopped, manager mode with no placement): keep existing
404behavior unchanged.
The synthetic view is documented as an approximation — it has no slot-identity stability and
no reconciliation backing. It exists solely to make GET /api/cluster/streams meaningful in
standalone; the placement endpoint is a secondary convenience.
Invariants
- Stats loop starts only when
self._worker_pool is None(standalone). - Loop is stopped via
_local_stats_stop.set()beforecontroller.stop()joins threads. stats_store.remove(run_id)is called in thefinallyblock of_stream_worker— always executes even on exception.- Loop is daemon=True, name=
tram-local-stats.
Code changes
| File | Change |
|---|---|
tram/pipeline/controller.py |
Add _LocalRun dataclass; _local_active_stats, _local_stats_lock, _local_stats_stop; _emit_local_stats_once(); _local_stats_loop() thread target; start/stop in start()/stop() guarded by self._worker_pool is None; _stream_worker() creates and registers _LocalRun, removes on exit; _run_batch() always generates run_id |
tram/api/routers/pipelines.py |
GET /{name}/placement: add standalone stream synthetic view; all other cases keep existing 404 behavior |
No changes needed to: executor.py, stats_store.py, _stream_views.py, internal.py,
health.py — the existing build_cluster_streams() fallback already handles non-placement
entries from stats_store.all_active().
Tests
| Test file | Cases |
|---|---|
tests/unit/test_controller_standalone_stats.py (new) |
_stream_worker registers _LocalRun; _emit_local_stats_once() calls stats_store.update() with correct payload fields; stream exit calls stats_store.remove(); _local_stats_loop exits on stop event; loop not started when worker_pool is not None |
tests/unit/test_pipelines_router.py (extend) |
Placement endpoint: standalone stream active → synthetic single-slot; standalone stream stopped → 404; non-stream pipeline → 404 |
Acceptance criteria
tram daemon(standalone): start a webhook stream pipeline; send at least one HTTP request into the webhook endpoint to produce traffic; pollGET /api/cluster/streamswithin 45 s; verify the pipeline appears in thestreamslist with non-zerorecords_in.GET /api/pipelines/{name}/placementreturns{"status": "running", "slots": [...]}withslot_count: 1andstale: falsestats.- Stop the pipeline; entry is absent from
GET /api/cluster/streamsimmediately (thefinallyblock in_stream_worker()callsstats_store.remove(run_id)explicitly — no waiting for the 3 × interval stale timeout). If the entry persists after a clean stop, that is a regression in explicit removal, not expected behavior. - All existing unit tests pass unchanged.
Workstream B — Manager Operational Metrics
Objectives
- Prometheus
/metricson the manager exposes dispatch decisions, reconciler actions, worker health, and callback receipt counters. - All series use the same no-op fallback as existing metrics when
prometheus_clientis absent. - Note in the metrics endpoint docstring that these are process-local; worker-side execution
metrics (
tram_records_*,tram_chunk_duration_seconds, etc.) require scraping worker pods.
New Prometheus series
| Name | Type | Labels | Description |
|---|---|---|---|
tram_mgr_dispatch_total |
Counter | pipeline, result (accepted/no_workers) |
Pipeline dispatch attempts by manager |
tram_mgr_redispatch_total |
Counter | pipeline |
Reconciler-triggered re-dispatches |
tram_mgr_reconcile_action_total |
Counter | pipeline, action (mark_stale/redispatch/resolve_running) |
Reconciler per-slot actions |
tram_mgr_placement_status |
Gauge | pipeline, status (running/degraded/reconciling/error) |
1 when placement is in that status, 0 otherwise |
tram_mgr_worker_healthy |
Gauge | — | Currently healthy worker count |
tram_mgr_worker_total |
Gauge | — | Total configured worker count |
tram_mgr_run_complete_received_total |
Counter | pipeline, status |
Run-complete callbacks received at manager |
tram_mgr_pipeline_stats_received_total |
Counter | — | Pipeline-stats callbacks received at manager |
tram_mgr_placement_status is a per-label gauge (1 = active, 0 = inactive). When a placement
transitions from degraded to running, the degraded label is set to 0 and running to 1.
Scope narrowing from roadmap: The v1.3.2 roadmap item says “callback failures.” This plan
narrows that to callback receipt counters plus inferred loss. There is no direct way for
the manager to count failures that occurred inside a worker process. What this release
delivers: tram_mgr_run_complete_received_total counts what the manager actually received.
Operators infer failures as tram_mgr_dispatch_total{result="accepted"} -
tram_mgr_run_complete_received_total for completed batch runs. This gap is documented in
the metrics description; it is a known boundary, not an oversight. The roadmap wording has
been updated to match.
Instrumentation points
| Call site | Series |
|---|---|
controller._run_batch() after worker_pool.dispatch() — accepted |
tram_mgr_dispatch_total{result="accepted"} |
controller._run_batch() when dispatch returns None |
tram_mgr_dispatch_total{result="no_workers"} |
controller._start_stream() single-dispatch path — accepted |
tram_mgr_dispatch_total{result="accepted"} |
controller._start_stream() multi-dispatch path — once per accepted slot |
tram_mgr_dispatch_total{result="accepted"} (N increments for N slots; counts individual worker dispatches, not logical placement decisions) |
controller._update_broadcast_placement_status() |
tram_mgr_placement_status (set 1 for new status, 0 for all others for that pipeline) |
reconciler.run_once() — slot["status"] set to "stale" |
tram_mgr_reconcile_action_total{action="mark_stale"} |
reconciler.run_once() — redispatch_broadcast_slot() returns True |
tram_mgr_redispatch_total, tram_mgr_reconcile_action_total{action="redispatch"} |
reconciler.run_once() — slot transitions to "running" |
tram_mgr_reconcile_action_total{action="resolve_running"} |
WorkerPool._poll_all() (at the end of each full poll cycle) |
tram_mgr_worker_healthy, tram_mgr_worker_total |
internal.run_complete() handler |
tram_mgr_run_complete_received_total{pipeline, status} |
internal.pipeline_stats() handler |
tram_mgr_pipeline_stats_received_total |
All instrumentation points are on the manager process. No worker-side code is touched.
Code changes
| File | Change |
|---|---|
tram/metrics/registry.py |
Add 8 new tram_mgr_* series with _NoOp* fallbacks |
tram/pipeline/controller.py |
Increment tram_mgr_dispatch_total, tram_mgr_placement_status |
tram/agent/reconciler.py |
Increment tram_mgr_redispatch_total, tram_mgr_reconcile_action_total |
tram/agent/worker_pool.py |
Set tram_mgr_worker_healthy, tram_mgr_worker_total after health poll |
tram/api/routers/internal.py |
Increment tram_mgr_run_complete_received_total, tram_mgr_pipeline_stats_received_total |
tram/api/routers/metrics_router.py |
Add docstring noting process-locality |
Tests
| Test file | Cases |
|---|---|
tests/unit/test_metrics_registry.py (extend) |
All new series importable; no-op when prometheus_client absent |
tests/unit/test_reconciler.py (extend) |
Stale detection increments mark_stale; re-dispatch increments redispatch and tram_mgr_redispatch_total |
tests/unit/test_controller_metrics.py (new) |
Batch dispatch accepted → result="accepted"; no-worker path → result="no_workers"; placement status gauge transitions correctly |
Acceptance criteria
tram daemon(manager mode, kind cluster): dispatch a batch pipeline;curl /metricsshowstram_mgr_dispatch_total{result="accepted"}andtram_mgr_worker_healthy/tram_mgr_worker_total. Labeled series that require specific events (tram_mgr_redispatch_total,tram_mgr_run_complete_received_total, etc.) will not appear until those events are exercised — verify each by triggering the relevant path explicitly:- Run a batch pipeline to completion →
tram_mgr_run_complete_received_totalappears. - Start a stream pipeline →
tram_mgr_pipeline_stats_received_totalappears after first stats interval. - Kill a worker and wait one reconcile cycle →
tram_mgr_redispatch_totalandtram_mgr_placement_status{status="degraded"}appear.
- Run a batch pipeline to completion →
- With
prometheus_clientuninstalled: daemon starts without error; notram_mgr_*series in response body.
Workstream C — UDP Multi-Worker Streams
Objectives
- Lift the L008 linter block:
syslogandsnmp_trapsources work withworkers.count: all,count: N, andworkers.listin manager mode whenkubernetes: enabled: trueis set. - Replace L008 with L012 (error): UDP push sources in manager mode require
kubernetes: enabled: true. Without it there is no worker-targeting UDP ingress path in the chart. This applies to all placement types (count: all,count: N,workers.list). - Fix
count: NService over-selection — a pre-existing gap affecting both HTTP and UDP push sources, and also fix the matching L006 contract: currently L006 blockscount: Neven with a kubernetes block, but the over-selection bug is being fixed here, so L006 should allowcount: N+kubernetes: enabled: true. - Extend
is_eligibleto allow UDP push sources to use a kubernetes Service. - End-to-end UDP path validated in kind before merge.
Semantics of workers.count: all for UDP sources
workers.count: all means: instantiate the pipeline on every healthy worker. Each worker runs
its own independent pipeline instance. No packet is replicated.
All UDP push sources in manager mode require kubernetes: enabled: true. The chart provides
no pre-existing worker-targeting UDP ingress; the per-pipeline Service created by
KubernetesServiceManager is the only ingress path. L012 enforces this as an error.
count: all + kubernetes: enabled: true: The Service uses a broad selector targeting all
worker pods. Kube-proxy ECMP hashes per flow (src_ip + src_port → one backend), so each
sender consistently lands on one worker. For snmp_trap → kafka this is the correct and
recommended model.
count: N or workers.list + kubernetes: enabled: true: Only specific workers are
dispatched. The Service uses manual Endpoints targeting only the dispatched worker pods — the
same path used today by HTTP workers.list.
The count: N over-selection fix (applies to HTTP and UDP equally)
Current behaviour: _build_selector() always returns a broad label selector matching all
worker pods — regardless of placement type. For workers.list this is already corrected by
manual Endpoints (_uses_manual_endpoints returns True when config.workers.worker_ids is
set). For count: N the Service over-selects.
Fix: ensure_service() gains an optional dispatched_worker_ids: list[str] | None
parameter. The controller derives these from placement slots for count: N and passes them
through. _uses_manual_endpoints() returns True when dispatched_worker_ids is provided,
triggering the same manual Endpoints path used by workers.list.
# In controller._activate_kubernetes_service:
dispatched_worker_ids = self._get_dispatched_worker_ids(config.name)
self._kubernetes_service_manager.ensure_service(config, dispatched_worker_ids=dispatched_worker_ids)
def _get_dispatched_worker_ids(self, pipeline_name: str) -> list[str] | None:
"""Return worker_ids for count:N placements only. None for count:all and workers.list."""
placement_group_id = self._active_placement_group.get(pipeline_name)
if placement_group_id is None:
return None
placement = self._broadcast_placements.get(placement_group_id)
if placement is None:
return None
target_count = placement.get("target_count")
if target_count == "all":
return None # shared selector is correct for count:all
state = self.manager.get(pipeline_name)
if state.config.workers and state.config.workers.worker_ids is not None:
return None # workers.list already uses config.workers.worker_ids path
return [slot["worker_id"] for slot in placement["slots"] if slot.get("worker_id")]
# In KubernetesServiceManager:
def _uses_manual_endpoints(
self, config: PipelineConfig, dispatched_worker_ids: list[str] | None = None
) -> bool:
if self._mode != "manager":
return False
if config.workers is not None and config.workers.worker_ids is not None:
return True # workers.list path
return dispatched_worker_ids is not None # count:N path
def _listed_worker_ids(
self, config: PipelineConfig, dispatched_worker_ids: list[str] | None = None
) -> list[str]:
if config.workers is not None and config.workers.worker_ids is not None:
return list(config.workers.worker_ids)
return list(dispatched_worker_ids or [])
UDP port binding model
HTTP push sources (webhook, prometheus_rw) share a single ingress listener on port :8767
and are routed by path (/webhooks/path-a). UDP sources work differently: each syslog or
snmp_trap pipeline binds its own raw socket directly on the pod using the source.port
value from the pipeline YAML. There is no shared UDP multiplexer.
Consequences:
- Two UDP pipelines on the same worker pod must use different
source.portvalues — the second bind will fail withAddress already in use. This is operator responsibility; the linter cannot catch runtime placement conflicts. - Helm does not need UDP container port declarations —
containerPortin the pod spec is informational only and not required for Service routing. - Privileged ports (< 1024, i.e., 162/514) require
NET_BIND_SERVICEcapability. Use non-privileged ports (≥ 1024) in kind/dev. Production clusters may need explicit security context configuration.
Generic kubernetes: config block
Rather than separate HTTP and UDP models, KubernetesServiceConfig is extended to be fully
generic. Protocol is auto-derived from source type — never user-configured. Port fields are
optional overrides.
Updated KubernetesServiceConfig model:
class KubernetesServiceConfig(BaseModel):
enabled: bool = True
service_type: Literal["ClusterIP", "NodePort", "LoadBalancer"] = "NodePort"
# Service cluster-side port. Default: target_port.
port: int | None = None
# Pod-side target port. Default: source.port (UDP) or worker ingress port (HTTP).
target_port: int | None = None
# Only valid when service_type=NodePort. Omit for auto-allocation.
node_port: int | None = None
# Only valid when service_type=LoadBalancer.
load_balancer_ip: str | None = None
# Merged into Service metadata.annotations.
annotations: dict[str, str] = Field(default_factory=dict)
# Optional name override (DNS-1123, max 63 chars).
service_name: str = ""
Example YAML — UDP snmp_trap pipeline, count:2:
source:
type: snmp_trap
port: 1162
kubernetes:
enabled: true
service_type: NodePort
port: 1162 # Service cluster-side port (same as pod port here)
target_port: 1162 # Explicit override; default would also resolve to source.port=1162
annotations:
metallb.universe.tf/address-pool: trap-pool
workers:
count: 2
Example YAML — HTTP webhook pipeline, LoadBalancer:
source:
type: webhook
path: /events
kubernetes:
enabled: true
service_type: LoadBalancer
port: 80 # Service cluster-side port
# target_port defaults to 8767 (shared worker ingress)
load_balancer_ip: 10.0.0.50
annotations:
service.beta.kubernetes.io/aws-load-balancer-type: nlb
_resolve_target_port() in KubernetesServiceManager:
UDP_PUSH_SOURCES = {"syslog", "snmp_trap"}
UDP_PORT_DEFAULTS = {"snmp_trap": 162, "syslog": 514}
def _resolve_target_port(self, config: PipelineConfig) -> int:
if config.kubernetes.target_port is not None:
return config.kubernetes.target_port
if config.source.type in UDP_PUSH_SOURCES:
return getattr(config.source, "port", None) or UDP_PORT_DEFAULTS[config.source.type]
return self._worker_ingress_port if self._mode == "manager" else self._standalone_port
def _resolve_protocol(self, config: PipelineConfig) -> str:
return "UDP" if config.source.type in UDP_PUSH_SOURCES else "TCP"
def _resolve_service_port(self, config: PipelineConfig) -> int:
if config.kubernetes.port is not None:
return config.kubernetes.port
return self._resolve_target_port(config)
_build_service_body() additions:
target_port = self._resolve_target_port(config)
service_port = self._resolve_service_port(config)
protocol = self._resolve_protocol(config)
port_spec = {
"name": "traffic",
"port": service_port,
"protocol": protocol,
"targetPort": target_port,
}
if config.kubernetes.service_type == "NodePort" and config.kubernetes.node_port is not None:
port_spec["nodePort"] = config.kubernetes.node_port
spec = {"type": config.kubernetes.service_type, "ports": [port_spec]}
if config.kubernetes.service_type == "LoadBalancer" and config.kubernetes.load_balancer_ip:
spec["loadBalancerIP"] = config.kubernetes.load_balancer_ip
metadata = {
"name": service_name,
"namespace": self._namespace,
"labels": {...},
"annotations": dict(config.kubernetes.annotations), # user annotations merged directly
}
is_eligible() is extended to cover ALL_PUSH_SOURCES = HTTP_PUSH_SOURCES | UDP_PUSH_SOURCES.
The PipelineConfig model validator that currently rejects non-HTTP sources with a kubernetes
block is updated to accept ALL_PUSH_SOURCES.
Linter rule changes
L006 update (HTTP push sources, kubernetes: enabled: true + count: N):
Currently L006 blocks count: N for HTTP push sources even when a kubernetes block is
present. The over-selection bug being fixed in this release makes count: N safe when
kubernetes: enabled: true. L006 is updated to allow count: N + kubernetes: enabled: true:
# Existing (unchanged): shared ingress path
no kubernetes block + workers.count: all → valid (shared HTTP ingress; v1.3.0 contract)
no kubernetes block + workers.list: [...] → L006 error (shared ingress cannot pin)
# Extended in v1.3.2: per-pipeline Service path
kubernetes: enabled: true + workers.count: all → no L006
kubernetes: enabled: true + workers.count: N → no L006 (over-selection bug now fixed)
kubernetes: enabled: true + workers.list: ... → no L006
# Still blocked in both paths:
no kubernetes block + workers.count: N → L006 error (shared ingress selects all workers)
L008 removed — the block is lifted now that UDP push sources have a proper Service path.
L012 added (error, manager mode only):
L012: UDP push source (
syslog,snmp_trap) in manager mode requireskubernetes: enabled: true. There is no pre-existing worker-targeting UDP ingress in the Helm chart; the per-pipeline Service is the only ingress path. This applies to allworkers.countvalues. Severity: Error (blocking).
L012 fires when: manager mode + UDP push source + no kubernetes: enabled: true block.
L012 does not fire for standalone mode (no Service provisioning needed).
Note: L011 is already taken by the risky-filename-partition rule
(_l011_risky_filename_partition_fields in linter.py:260). The new rule is L012.
Topology after fix
count: all + UDP (kubernetes: enabled: true required)
External sender → UDP NodePort (pipeline Service, broad selector) → kube-proxy ECMP → worker-0, 1, 2
One Service with broad label selector, all workers as backends. Correct for stateless sinks.
count: N or workers.list + UDP (kubernetes: enabled: true required)
External sender → UDP NodePort (pipeline Service) → manual Endpoints → dispatched workers only
Same manual Endpoints mechanism already used for HTTP workers.list.
kind validation checklist
Prerequisites:
- Manager + 3 workers deployed in kind.
- A single-broker Kafka deployment in the same kind cluster is required for steps 1–3
(e.g.,
helm install kafka oci://registry-1.docker.io/bitnamicharts/kafka --set replicaCount=1 --set zookeeper.enabled=false --set kraft.enabled=true). As an alternative, afileorrestsink can substitute for routing-level proof (steps 1–3) if Kafka is not yet available; swap back to Kafka for final merge validation.
- Register
snmp_trap → kafka(orsnmp_trap → file) withworkers.count: all, nokubernetesblock. Confirm L012 fires (error). Addkubernetes: enabled: true; confirm L012 clears and no L008. - Send test SNMP traps from at least two distinct source IPs to the cluster NodePort (from the pipeline Service). ECMP hashes per flow (src_ip + src_port), so a single sender lands on one backend consistently — multiple senders are needed to exercise distribution. Confirm records from both senders appear in the sink and that worker logs show each sender consistently reaching one worker (and no traps routed to a pod not running the pipeline).
- Register the same pipeline with
workers.count: 2+kubernetes: enabled: true. Confirm no L012. Confirm Service Endpoints target only the 2 dispatched workers, not all 3. - HTTP regression: register a webhook pipeline with
count: 2+kubernetes: enabled: true; confirm L006 does not fire and Service Endpoints target only the 2 dispatched workers.
Code changes
| File | Change |
|---|---|
tram/models/pipeline.py |
KubernetesServiceConfig: add port, target_port, load_balancer_ip, annotations fields; add ClusterIP to service_type literal; update model validator to accept ALL_PUSH_SOURCES |
tram/pipeline/linter.py |
Remove L008; add L012 (error); update L006 to allow count: N + kubernetes: enabled: true |
tram/pipeline/k8s_service_manager.py |
Add _resolve_target_port(), _resolve_protocol(), _resolve_service_port(); update _build_service_body() to use resolved values + annotations + loadBalancerIP; ensure_service(config, dispatched_worker_ids=None); extend _uses_manual_endpoints and _listed_worker_ids for dispatched_worker_ids; extend is_eligible to ALL_PUSH_SOURCES; add UDP_PUSH_SOURCES, UDP_PORT_DEFAULTS |
tram/pipeline/controller.py |
_activate_kubernetes_service calls _get_dispatched_worker_ids and passes to ensure_service; add _get_dispatched_worker_ids() |
No changes to: _stream_views.py, db.py, slots_json schema — no new slot fields needed.
Tests
| Test file | Cases |
|---|---|
tests/unit/test_linter.py (extend) |
L008 gone; L012 fires on UDP + manager + any count without kubernetes: enabled: true; clears when kubernetes block present; L006 no longer fires for HTTP push + count: N + kubernetes: enabled: true |
tests/unit/test_k8s_service_manager.py (extend) |
ensure_service(dispatched_worker_ids=[...]) uses manual Endpoints for those ids; UDP source produces protocol: UDP Service with correct targetPort; HTTP source keeps protocol: TCP and targetPort=8767; annotations merged into Service metadata; loadBalancerIP set when service_type=LoadBalancer; explicit port/target_port override defaults; is_eligible true for UDP + HTTP sources with kubernetes block |
tests/unit/test_pipeline_model.py (extend) |
KubernetesServiceConfig: annotations dict accepted; load_balancer_ip only valid with LoadBalancer; ClusterIP accepted as service_type; target_port and port optional; model validator accepts UDP push sources |
tests/unit/test_controller_udp.py (new) |
count: N dispatch passes dispatched_worker_ids to service manager; count: all passes None; workers.list passes None (handled by config path) |
| kind validation | See checklist above (manual, gates merge of PR-C2) |
Workstream D — ASN.1 Structured Decode Flattening
Objectives
serializer_in: {type: asn1}can split a BER file containing concatenated top-level objects and return one logical record per BER object instead of one record per file.- Root-type selection remains explicit by default (
message_class), with an optional ordered fallback list (message_classes) for cases where multiple top-level roots are valid. - A new generic
json_flattentransform can collapse nested decoded records into flat JSON / NDJSON / CSV friendly rows using configurable dict hoisting, list explosion, and list zipping. - A new generic
hex_decodetransform can turn byte-derived hex-string values into readable text where safe, while supporting explicit field-path overrides for semantic decode.
Design decisions
1. Multi-record BER split belongs in the ASN.1 serializer. BaseSerializer.parse() already
returns list[dict]; executor.py already counts and processes multiple logical records per
source chunk. The missing behavior is only in Asn1Serializer.parse(), which currently calls
compiled.decode(message_class, data) once on the whole file. The serializer is therefore the
correct layer to add split_records: true for BER files containing concatenated top-level
objects.
BER boundary walking is explicit and manual. asn1tools.decode() expects one complete BER
value at a time; it does not consume a concatenated BER stream for TRAM. D1 therefore adds a
small BER envelope walker inside asn1_serializer.py:
- parse outer tag bytes
- parse BER length bytes (short form and long form)
- for indefinite length BER, walk nested TLVs until end-of-contents (
00 00) - slice the exact top-level TLV bytes
- call
compiled.decode(root_type, record_bytes) - advance by the consumed byte count and continue until EOF
This is the same boundary model already proven in the existing external CDR decode scripts.
The walker is BER-only; split_records is rejected for der, per, uper, xer, and jer.
2. Root-type selection stays explicit by default. For schemas like PGW where the user can
provide a real top-level CHOICE root (GPRSRecord), ASN.1 already performs branch selection
from BER tags. TRAM should preserve that path. The new optional fallback is:
serializer_in:
type: asn1
schema_file: /data/schemas/asn1/ericsson/sgw-CDRFR9OLD.asn
message_classes: [CallEventRecord, GPRSRecord]
encoding: ber
split_records: true
Per top-level BER object, TRAM tries the configured root types in order until one decodes. This is intentionally narrower than “probe every type in every schema.” Blind schema-wide probing is deferred because it is expensive, ambiguous, and too easy to bias with keyword heuristics.
Failure path is explicit: decode failure raises SerializerError. If none of the configured
message_classes can decode a split BER object, the serializer raises SerializerError. The
executor already routes parse failures through the existing DLQ / error path. There is no silent
skip and no raw-record fallback in D1.
3. Flattening is generic, not ASN.1-branded. After decode, the data shape is nested JSON-like
dicts/lists/scalars. The transform should therefore be named json_flatten, not
asn1_flatten. ASN.1-specific patterns such as CHOICE objects represented as
{"type": ..., "value": ...} are supported through configuration rather than hardcoded as
the only behavior.
4. Hex interpretation is separate from structural flattening. json_flatten should not also
decode telecom encodings or guess text codecs. A separate hex_decode transform keeps the
structural and semantic concerns separate:
json_flattendecides record shapehex_decodedecides how to interpret byte-derived hex-string values
This keeps json_flatten reusable for non-ASN.1 nested JSON workflows as well.
asn1 serializer changes
Asn1SerializerConfig gains:
| Field | Default | Meaning |
|---|---|---|
split_records |
false |
For BER only: split concatenated top-level BER objects and decode each separately |
message_classes |
null |
Optional ordered fallback list of root types; mutually exclusive with message_class |
Rules:
split_recordsis valid only forencoding: ber- exactly one of
message_classormessage_classesmust be set - when
message_classesis present, TRAM tries each configured root type in order per split BER object and records the first success
Why a new json_flatten transform?
TRAM already has flatten, explode, and unnest. D2 does not replace them; it adds a
higher-level transform for cases where the current primitive chain becomes long, repetitive, and
schema-fragile.
The explicit gap relative to existing transforms is the combination of:
- recursive application of unnest/explode decisions across a nested tree
- configurable CHOICE-object handling (
{"type": ..., "value": ...}) - index-based list zipping (for patterns like labels[] + values[])
- one transform owning the flattening plan so pipelines do not need 8–12 ordered primitive transforms for common nested-record expansion
If a pipeline only needs one or two explicit unnest / explode steps, the existing transforms
remain the preferred tool. json_flatten is for the larger “nested structured payload to flat
row set” case.
json_flatten transform
json_flatten is a generic nested-record transform with heuristic defaults plus explicit
override controls.
Planned default behavior:
- recursively hoist nested dict wrappers where safe
- explode lists of dicts into multiple output records
- preserve parent scalar fields across explosions
- support CHOICE-like dicts via configurable
choice_mode - optionally zip paired lists by index
- optionally normalize output field names (
snake_case)
Planned controls:
| Field | Purpose |
|---|---|
explode_mode |
auto | off | paths |
explode_paths |
Explicit list paths to explode |
zip_lists |
auto | off | mappings |
zip_mappings |
Explicit list-pair mappings for index-based zipping |
choice_mode |
keep | unwrap_value | type_value |
rename_style |
none | snake_case |
drop_paths |
Remove unwanted subtrees after flattening |
keep_paths |
Whitelist specific subtrees |
max_depth |
Limit recursive flatten depth |
ambiguity_mode |
keep | error when auto rules detect more than one valid auto-expansion plan |
Example target usage:
transforms:
- type: json_flatten
explode_mode: auto
zip_lists: auto
choice_mode: type_value
rename_style: snake_case
This transform is intended to replace long hand-authored chains of unnest / explode /
add_field / drop / rename for common nested-record patterns, while still allowing a
pipeline author to fall back to the explicit primitive transforms when needed.
Ambiguity trigger in D2: “auto rules are not decisive” means json_flatten encounters more
than one candidate list-expansion plan at the same structural level and those plans would
produce materially different row sets. Example: two sibling list fields are both eligible for
auto-explode but are not declared as a zip pair and have incompatible lengths / semantics.
In D2 the allowed responses are:
keep: preserve the subtree without auto-expanding iterror: raiseTransformError
The earlier explode option is dropped from ambiguity_mode because it is too underspecified.
hex_decode transform
hex_decode is a second generic transform that operates after decode and after optional
structural flattening. Its goal is not “ASN.1 magic”; its goal is controlled interpretation
of byte-derived hex-string values.
Input model: by the time records reach transforms, _to_json_safe() in the ASN.1 serializer
has already converted bytes values to lowercase hex strings. D2 therefore does not operate
on Python bytes; it operates on scalar strings that may represent hex-encoded binary values.
Every override path first converts value -> bytes.fromhex(value) before applying a semantic
codec. Heuristic text decoding likewise means “decode the hex string as bytes, then try UTF-8 /
Latin-1.”
Scope boundary: hex_decode is not a second ASN.1 decoder. It does not re-run asn1tools
against nested BER fragments and it does not interpret structured wrapper objects such as
CHOICE dicts directly. Container selection / structural navigation belongs to json_flatten
and explicit path targeting. hex_decode is for leaf values that are already scalar hex
strings.
Planned generic behavior:
- default safe heuristics: decode printable UTF-8 / Latin-1 text; otherwise keep hex
- optional preservation of original raw value or raw hex shadow fields
- explicit field-path overrides for semantic decode
Proposed config shape:
transforms:
- type: hex_decode
mode: utf8_or_hex
preserve_original: false
overrides:
- path: recordOpeningTime
decode_as: timestamp
format: bcd_semi_octet
- path: pGWAddress.value.value
decode_as: ip
format: packed
- path: mSTimeZone
decode_as: timezone
format: tbcd_quarter_hour
- path: servedIMSI
decode_as: digits
format: tbcd
decode_as is semantic, not vendor-branded. format selects the wire codec. The code path is
therefore generic-but-extensible: hardcoded reusable codecs are acceptable; hardcoded vendor
field names are not.
Planned built-in semantic targets include:
textdigitstimestamptimezoneipenumbit_flagshex
Planned built-in wire formats include:
utf8latin1tbcdbcd_semi_octetpackedbit_string_bytestbcd_quarter_hour
These codecs are allowed to include IETF / 3GPP / ASN.1 reusable wire encodings. D2 does not promise support for vendor-proprietary binary layouts.
Non-goals
- No blind “try every type in every schema” root detection by default
- No vendor-specific flatteners in core (
mtas_flatten,pgw_flatten, etc.) - No implicit semantic reinterpretation of every
OCTET STRING; raw-safe output remains valid - No promise that
json_flattencan infer one universally correct flat row model for every arbitrary ASN.1 schema without hints - Built-in codecs may cover reusable standard wire formats (for example ASN.1 / BER / 3GPP-style
encodings such as
tbcd), but not vendor-proprietary binary layouts
Code changes
| File | Change |
|---|---|
tram/models/pipeline.py |
Extend Asn1SerializerConfig; add transform config models for json_flatten and hex_decode |
tram/serializers/asn1_serializer.py |
Add BER top-level splitter; add ordered message_classes fallback; preserve current single-message path when split_records=false |
tram/transforms/json_flatten.py |
New generic nested-record flatten transform |
tram/transforms/hex_decode.py |
New generic hex-string heuristic + override transform |
tram/transforms/__init__.py |
Register new transforms |
docs/connectors.md |
Document new ASN.1 serializer options and both transforms |
pipelines/ |
Add example ASN.1 pipelines for raw JSON, flat NDJSON, and flat CSV |
Tests
| Test file | Cases |
|---|---|
tests/unit/test_asn1_serializer.py |
BER split path; concatenated BER file returns N records; message_classes ordered fallback; invalid split_records on non-BER rejected |
tests/unit/test_json_flatten.py |
Dict hoist, list explode, choice handling, zip behavior, ambiguity modes |
tests/unit/test_hex_decode.py |
UTF-8 / Latin-1 heuristics, preserve-original behavior, semantic override dispatch by path |
tests/integration/test_asn1_pipeline.py |
ASN.1 raw export, flat NDJSON pipeline, flat CSV pipeline on sample BER input |
Workstream E — CDR Structured Record Shaping
Goal
The original ASN.1 flattening slice (D) solved generic BER multi-record decode and basic nested
row shaping, but the LTE/SGW/PGW pipelines exposed a second problem: structure alone was not
enough to express telecom CDR intent. Three list semantics needed to be represented explicitly:
- Alternate identifier lists — pick one or more entries by predicate without exploding the
record, such as
subscriptionID[]. - Multi-valued scalar attributes — choose or coalesce from several possible paths while still keeping one top-level CDR row.
- Child usage segments — explode only when the pipeline truly wants one child row per usage segment rather than one parent CDR row.
The delivered shaping path therefore had two layers:
- explicit
json_flattenfor ordered explode/zip/choice-unwrapping/final flattening on nested JSON-like records - smaller CDR-focused primitives (
select_from_list,coalesce_fields,project, conditionaldrop, dotted-path support, and narrow wildcard matching) so preserve-mode mediation pipelines could stay readable
Key design decisions
json_flattenmoved from heuristic row-shaping to an explicit ordered contract:explode_paths -> zip_groups -> choice_unwrap -> final flatten -> drop_paths- dotted-path support was extended to
unnest,explode,drop,rename,value_map, andcastthrough shared path helpers select_from_listshipped as the primary “pick by predicate, project up” primitive; it supports multi-select,first_item: true,on_no_match, and duplicate-output validationcoalesce_fieldsshipped as a small purpose-built primitive instead of repeated temp-field scaffoldingprojectbecame the final schema extraction/defaulting step rather than forcing repeatedjmespath + rename + add_fieldchainsdropabsorbed the conditional per-field form instead of introducing a separatedrop_emptytransform- wildcard path support was deliberately narrow: single-segment
*matching only, with exact path rules taking precedence - telecom semantic decode stayed outside
json_flatten;hex_decodehandles byte-derived scalar interpretation, includingbit_flagswith a companion bit-length field
Delivered surfaces
json_flattenredesigned around explicitexplode_paths,zip_groups,choice_unwrap,preserve_lists, anddrop_pathshex_decodeextended with structured overrides, semantic codecs, wildcard path support, andbit_flagsdecoding for values such asserviceConditionChange- new transforms:
select_from_list,coalesce_fields, andproject - shared path helpers:
get_path,set_path,delete_path,rename_path - dotted-path support across the core shaping transforms
- conditional
dropand narrow wildcard support onjson_flatten.drop_paths/hex_decode.overrides[].path
Outcome
The shipped CDR preserve-mode pipelines no longer depend on vendor-specific flatteners or heuristic list handling. They keep one top-level row per decoded CDR unless the pipeline author explicitly chooses an exploded child-row model, and the same shaping primitives remain reusable for non-telecom nested JSON payloads.
Workstream F — Batch Reconciliation, Incremental Record Processing, and Safe File Output
Goal
Large PGW batch runs in kind exposed three runtime gaps after the original 1.3.2 backend slices
were already underway:
- manager correctness when a worker-owned batch run disappears before callback
- unbounded memory use when a serializer eagerly materializes all decoded records from a large file
- unsafe append/reuse behavior for deterministic serial file outputs during rerun or failure
This workstream closed those gaps without introducing checkpoint/resume or automatic retry storms.
Key design decisions
- batch reconciliation uses manager dispatch tracking plus worker
/agent/status; Kubernetes termination reasons remain enrichment, not the core truth source - lost worker-owned batch runs are marked failed and cleared for rescheduling on the next schedule/manual trigger; no immediate automatic re-trigger by default
- adoption / stale cleanup after manager restart belongs in a dedicated
BatchReconciler, not in_boot_load() - incremental processing uses a bounded window config named
record_chunk_size, notbatch_size, to avoid colliding with the existing “stop after N records” behavior - serializers can optionally expose
parse_chunks(data, record_chunk_size); the executor uses that hook in the serial batch path, with a fallback slice path where appropriate - safe file output is implemented through a generic sink lifecycle hook
finalize_source(meta, success)and run-scoped staged temp files for record-safe serializers such asndjsonandcsv - heap cleanup after large batch runs is opt-in per pipeline through
post_batch_cleanup: true, not a global worker recycle or global trim policy
Delivered surfaces
BatchReconcilerrunning alongsidePlacementReconciler- active batch lease tracking in the controller plus adoption / mark-lost flows
BaseSerializer.parse_chunks(...)and incremental ASN.1 BER iteration for concatenated files- pipeline-level
record_chunk_size BaseSink.finalize_source(...)plus staged safe-finalize behavior in local/SFTP file sinks- stale temp-file cleanup for deterministic reruns
- per-pipeline
post_batch_cleanup
Validation outcome
- large PGW replay completed successfully with bounded serial record windows
- manager can now adopt or clear lost running batch state instead of leaving pipelines stuck in
running - serial
ndjson/csvoutputs now publish only on source-file success and remove staged temp files on failure
Recommended PR split
Historical note: the split below captures the initial A-D rollout plan. The later E/F
workstreams landed as follow-up slices during the same 1.3.2 cycle and are consolidated above
for archival completeness.
| PR | Branch | Content | Merge order |
|---|---|---|---|
| PR-A | observer-a-standalone-stats |
Workstream A: controller stats loop, _LocalRun dataclass, placement endpoint synthetic view |
1st |
| PR-B | observer-b-mgr-metrics |
Workstream B: registry series + instrumentation at controller/reconciler/worker_pool/internal router | 2nd (parallel with A) |
| PR-C1 | observer-c1-udp-model |
Workstream C: L008 removal, L012, count: N over-selection fix in service manager, UDP eligibility, unit tests |
3rd |
| PR-C2 | observer-c2-udp-e2e |
Workstream C: controller _get_dispatched_worker_ids wiring, kind validation evidence in PR body |
4th (gated on C1) |
| PR-D1 | observer-d1-asn1-split |
Workstream D: BER split, message_classes fallback, serializer tests |
5th |
| PR-D2 | observer-d2-json-flatten |
Workstream D: json_flatten, hex_decode, docs, example pipelines, integration tests |
6th (gated on D1) |
PR-A and PR-B share no code paths and can be reviewed in parallel. PR-C1 is a prerequisite for C2 because the kind validation requires the full stack wired together. PR-D1 should land before PR-D2 so the flattening examples and tests run against the real multi-record ASN.1 path instead of the legacy one-record-per-file limitation.
Open questions
-
tram_mgr_placement_statuscardinality: 4 label combinations (running/degraded/ reconciling/error) × number of active multi-worker pipelines. For deployments with >50 pipelines this could be high. Alternative: single gauge with value 0–3 encoded as an enum. Confirm expected pipeline count before deciding; 50+ is unusual for v1.3.2 target deployments. -
Batch stats mid-run visibility (deferred): If mid-run batch visibility is needed in standalone, the controller would need to snapshot periodically inside the thread pool. This adds complexity and is explicitly deferred. Re-evaluate in v1.3.3 after UI revalidation clarifies what stats the dashboard actually needs during a batch run.