Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -55,6 +55,10 @@ The bridge does not
inspect harness IDs, native error-code lists, or diagnostic text. Other
adapter-specific timeout codes are not automatically reclassified.

## Finalized Telemetry Artifacts

The task runner validates telemetry after the one-shot NeMo Fabric runtime shuts down. For Relay trajectories finalized during shutdown, it reads the runtime-owned plugin configuration and collects local ATIF files from that runtime's configured output directory. It validates the trajectory before promoting it to Harbor's `agent/trajectory.json`; ambiguous or malformed artifacts remain telemetry failures rather than changing the verifier reward. Remote-only output and nested filename templates are not collected by this fallback. Adapter invocation results and multi-turn behavior are unchanged.

## Invocation Accounting

The integration projects normalized `RunResult.usage` into Harbor's `AgentContext` without requiring NeMo Relay or ATIF. Harbor input counts include cache: inclusive NeMo Fabric input is copied unchanged, while exclusive input is combined with cache only when both counts are known. If input semantics are unknown, Harbor input remains unknown unless the adapter explicitly reports zero cache tokens. Legacy results remain valid, but their input count needs explicit cache semantics or an ATIF fallback to populate Harbor's inclusive count; the original values remain in result metadata. Unknown counts and costs remain `None`; estimates are not promoted to reported cost. ATIF metrics fill only missing fields and are never added to normalized usage. The same collection path applies to unsuccessful results that contain usage.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -78,6 +78,7 @@ async def run(payload: FabricRunPayload) -> RunResult:
Path(payload.logs_dir),
harbor_session_id=payload.request.context.get("harbor_session_id"),
harbor_context_id=payload.request.context.get("harbor_context_id"),
runtime_stopped=True,
)
return result

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@
import json
import re
import shutil
import tomllib
from collections import Counter
from datetime import datetime
from pathlib import Path
Expand All @@ -28,6 +29,7 @@ def publish_telemetry_evidence(
strict: bool = False,
harbor_session_id: str | None = None,
harbor_context_id: str | None = None,
runtime_stopped: bool = False,
) -> dict[str, Any]:
"""Validate telemetry, promote ATIF, and write a machine-readable summary.

Expand All @@ -43,8 +45,14 @@ def publish_telemetry_evidence(
logs_dir,
harbor_session_id=harbor_session_id,
harbor_context_id=harbor_context_id,
runtime_stopped=runtime_stopped,
)
except (OSError, json.JSONDecodeError, TelemetryValidationError, ValueError) as error:
except (
OSError,
json.JSONDecodeError,
TelemetryValidationError,
ValueError,
) as error:
summary = _base_summary(result)
summary.update(status="failed", error=str(error))
_write_summary(summary_path, summary, strict=strict)
Expand All @@ -71,6 +79,7 @@ def validate_telemetry(
*,
harbor_session_id: str | None = None,
harbor_context_id: str | None = None,
runtime_stopped: bool = False,
) -> dict[str, Any]:
"""Validate telemetry artifacts and promote exactly one valid ATIF trajectory."""

Expand All @@ -85,6 +94,10 @@ def validate_telemetry(
for artifact in artifacts
if artifact.kind == "atif"
]
if runtime_stopped:
atif_paths = sorted(
set(atif_paths).union(_finalized_relay_atif_paths(result, logs_dir))
)
summary = _base_summary(result)
summary["harbor_session_id"] = harbor_session_id
summary["harbor_context_id"] = harbor_context_id
Expand All @@ -93,16 +106,115 @@ def validate_telemetry(
atif_paths,
logs_dir,
)
summary["status"] = "not_emitted" if not atof_paths and not atif_paths else "succeeded"
summary["status"] = (
"not_emitted" if not atof_paths and not atif_paths else "succeeded"
)
return summary


def _finalized_relay_atif_paths(result: RunResult, logs_dir: Path) -> list[Path]:
"""Collect local Relay ATIF only after the one-shot runtime has shut down."""

paths: list[Path] = []
root = result.artifacts.root
if root is None or not any(ref.provider == "relay" for ref in result.telemetry):
return paths
run_root = _resolve_artifact_path(root, logs_dir).resolve()
if not run_root.is_relative_to(logs_dir.resolve()):
raise TelemetryValidationError(
"finalized ATIF root escapes the Harbor logs directory"
)
output = result.to_mapping()["output"]
relay_runtime = output.get("relay_runtime") if isinstance(output, dict) else None
config_path = (
relay_runtime.get("plugin_config_path")
if isinstance(relay_runtime, dict)
else None
)
if not isinstance(config_path, str):
return paths
plugin_path = _resolve_artifact_path(Path(config_path), logs_dir)
config_root = run_root / ".fabric" / result.runtime_id
if plugin_path.is_symlink() or not plugin_path.resolve().is_relative_to(
config_root
):
raise TelemetryValidationError("finalized Relay config is outside this runtime")
config = tomllib.loads(plugin_path.read_text(encoding="utf-8"))
components = config.get("components", [])
if not isinstance(components, list):
raise TelemetryValidationError("finalized Relay components must be an array")
for telemetry in result.telemetry:
if telemetry.provider != "relay":
continue
for component in components:
if not isinstance(component, dict):
raise TelemetryValidationError(
"finalized Relay component must be a table"
)
if (
component.get("kind") != "observability"
or component.get("enabled") is False
):
continue
component_config = component.get("config", {})
if not isinstance(component_config, dict):
raise TelemetryValidationError("finalized Relay config must be a table")
atif = component_config.get("atif", {})
if not isinstance(atif, dict):
raise TelemetryValidationError(
"finalized Relay ATIF config must be a table"
)
if atif.get("enabled") is not True or atif.get("storage"):
continue
Comment thread
AjayThorve marked this conversation as resolved.
if not isinstance(atif.get("output_directory"), str) or not isinstance(
atif.get("filename_template"), str
):
raise TelemetryValidationError(
"finalized ATIF requires an output directory and filename template"
)
directory = _resolve_artifact_path(Path(atif["output_directory"]), logs_dir)
# Relay owns the session ID, but Fabric owns the runtime-scoped directory.
if (
directory.is_symlink()
or directory.name != result.runtime_id
or not directory.resolve().is_relative_to(run_root)
):
raise TelemetryValidationError(
"finalized ATIF directory is outside this runtime"
)
if not directory.exists():
continue
template = atif["filename_template"]
if "/" in template or "\\" in template:
raise TelemetryValidationError(
"finalized ATIF filename template must be a basename"
)
# Adjacent placeholders need a minimum length, not overlapping wildcards.
pattern = "".join(
f"[^/]{{{part.count('{')},}}" if index % 2 else re.escape(part)
for index, part in enumerate(re.split(r"((?:\{[^{}]+\})+)", template))
)
for path in directory.iterdir():
if not re.fullmatch(pattern, path.name) or not path.is_file():
continue
if path.is_symlink() or not path.resolve().is_relative_to(
directory.resolve()
):
raise TelemetryValidationError(
"finalized ATIF artifact escapes its runtime directory"
)
paths.append(path)
return paths


def _resolve_artifact_path(path: Path, logs_dir: Path) -> Path:
"""Resolve collected task paths when validation runs on the Harbor host."""

task_logs = Path("/logs/agent")
if ".." in path.parts or not path.is_relative_to(task_logs):
raise TelemetryValidationError(f"telemetry artifact escapes /logs/agent: {path}")
raise TelemetryValidationError(
f"telemetry artifact escapes /logs/agent: {path}"
)
if logs_dir != task_logs:
relative = path.relative_to(task_logs)
return logs_dir / relative
Expand All @@ -123,10 +235,14 @@ def _validate_atof(paths: list[Path]) -> dict[str, Any]:
continue
value = json.loads(line)
if not isinstance(value, dict):
raise TelemetryValidationError(f"ATOF record must be an object: {path}:{line_number}")
raise TelemetryValidationError(
f"ATOF record must be an object: {path}:{line_number}"
)
missing = required.difference(value)
if missing:
raise TelemetryValidationError(f"ATOF record missing {sorted(missing)}: {path}:{line_number}")
raise TelemetryValidationError(
f"ATOF record missing {sorted(missing)}: {path}:{line_number}"
)
_validate_atof_record(value, path, line_number)
records += 1
counts[value["kind"]] += 1
Expand All @@ -141,27 +257,39 @@ def _validate_atof_record(value: dict[str, Any], path: Path, line_number: int) -
location = f"{path}:{line_number}"
for field in ("atof_version", "kind", "name", "timestamp", "uuid"):
if not isinstance(value[field], str) or not value[field]:
raise TelemetryValidationError(f"ATOF record field {field} must be a non-empty string: {location}")
raise TelemetryValidationError(
f"ATOF record field {field} must be a non-empty string: {location}"
)
if re.fullmatch(r"\d+\.\d+(?:\.\d+)?", value["atof_version"]) is None:
raise TelemetryValidationError(f"ATOF record atof_version is invalid: {location}")
raise TelemetryValidationError(
f"ATOF record atof_version is invalid: {location}"
)
try:
timestamp = datetime.fromisoformat(value["timestamp"].replace("Z", "+00:00"))
except ValueError as error:
raise TelemetryValidationError(f"ATOF record timestamp is invalid: {location}") from error
raise TelemetryValidationError(
f"ATOF record timestamp is invalid: {location}"
) from error
if timestamp.tzinfo is None:
raise TelemetryValidationError(f"ATOF record timestamp must include a timezone: {location}")
raise TelemetryValidationError(
f"ATOF record timestamp must include a timezone: {location}"
)
try:
UUID(value["uuid"])
except ValueError as error:
raise TelemetryValidationError(f"ATOF record uuid is invalid: {location}") from error
raise TelemetryValidationError(
f"ATOF record uuid is invalid: {location}"
) from error


def _validate_atif(
paths: list[Path],
logs_dir: Path,
) -> dict[str, Any]:
if len(paths) > 1:
raise TelemetryValidationError(f"expected at most one ATIF artifact, found {len(paths)}")
raise TelemetryValidationError(
f"expected at most one ATIF artifact, found {len(paths)}"
)
if not paths:
return {"files": [], "promoted": None}

Expand Down Expand Up @@ -196,30 +324,52 @@ def _validate_atif_structure(value: Any, path: Path) -> dict[str, Any]:
schema_version = value.get("schema_version")
supported_versions = {f"ATIF-v1.{minor}" for minor in range(8)}
if schema_version not in supported_versions:
raise TelemetryValidationError(f"unsupported ATIF schema_version {schema_version!r}: {path}")
raise TelemetryValidationError(
f"unsupported ATIF schema_version {schema_version!r}: {path}"
)
session_id = value.get("session_id")
if not isinstance(session_id, str) or not session_id:
raise TelemetryValidationError(f"ATIF session_id must be a non-empty string: {path}")
raise TelemetryValidationError(
f"ATIF session_id must be a non-empty string: {path}"
)
agent = value.get("agent")
if not isinstance(agent, dict) or not isinstance(agent.get("name"), str) or not agent["name"]:
raise TelemetryValidationError(f"ATIF agent.name must be a non-empty string: {path}")
if (
not isinstance(agent, dict)
or not isinstance(agent.get("name"), str)
or not agent["name"]
):
raise TelemetryValidationError(
f"ATIF agent.name must be a non-empty string: {path}"
)
if not isinstance(agent.get("version"), str) or not agent["version"]:
raise TelemetryValidationError(f"ATIF agent.version must be a non-empty string: {path}")
raise TelemetryValidationError(
f"ATIF agent.version must be a non-empty string: {path}"
)
steps = value.get("steps")
if not isinstance(steps, list) or not steps:
raise TelemetryValidationError(f"ATIF steps must be a non-empty array: {path}")
for index, step in enumerate(steps, 1):
if not isinstance(step, dict):
raise TelemetryValidationError(f"ATIF step {index} must be an object: {path}")
raise TelemetryValidationError(
f"ATIF step {index} must be an object: {path}"
)
if not isinstance(step.get("step_id"), int) or step["step_id"] < 1:
raise TelemetryValidationError(f"ATIF step {index} has an invalid step_id: {path}")
raise TelemetryValidationError(
f"ATIF step {index} has an invalid step_id: {path}"
)
if step.get("source") not in {"system", "user", "agent"}:
raise TelemetryValidationError(f"ATIF step {index} has an invalid source: {path}")
raise TelemetryValidationError(
f"ATIF step {index} has an invalid source: {path}"
)
if not isinstance(step.get("message"), str | list):
raise TelemetryValidationError(f"ATIF step {index} has an invalid message: {path}")
raise TelemetryValidationError(
f"ATIF step {index} has an invalid message: {path}"
)
final_metrics = value.get("final_metrics")
if final_metrics is not None and not isinstance(final_metrics, dict):
raise TelemetryValidationError(f"ATIF final_metrics must be an object or null: {path}")
raise TelemetryValidationError(
f"ATIF final_metrics must be an object or null: {path}"
)
return value


Expand All @@ -237,10 +387,14 @@ def _base_summary(result: RunResult) -> dict[str, Any]:

_SECRET_PATTERNS = (
re.compile(r"\b(?:sk|nvapi)-[A-Za-z0-9_-]{16,}\b"),
re.compile(r'(?i)["\'](?:api[_-]?key|access[_-]?token|authorization)["\']\s*:\s*["\'][^"\']{8,}["\']'),
re.compile(
r'(?i)["\'](?:api[_-]?key|access[_-]?token|authorization)["\']\s*:\s*["\'][^"\']{8,}["\']'
),
)


def _reject_obvious_secrets(text: str, path: Path) -> None:
if any(pattern.search(text) for pattern in _SECRET_PATTERNS):
raise TelemetryValidationError(f"telemetry contains a value that resembles a credential: {path}")
raise TelemetryValidationError(
f"telemetry contains a value that resembles a credential: {path}"
)
Loading
Loading