diff --git a/adapters/common/src/nemo_fabric_adapters/common/utils.py b/adapters/common/src/nemo_fabric_adapters/common/utils.py index e91a1309e..e8dad5d3b 100644 --- a/adapters/common/src/nemo_fabric_adapters/common/utils.py +++ b/adapters/common/src/nemo_fabric_adapters/common/utils.py @@ -5,24 +5,12 @@ from __future__ import annotations -import copy import json import os import sys from pathlib import Path -from typing import TYPE_CHECKING from typing import Any -if TYPE_CHECKING: - from nemo_relay import plugin - from nemo_relay.observability import AtifConfig - from nemo_relay.observability import AtofConfig - from nemo_relay.observability import AtofFileSinkConfig - from nemo_relay.observability import AtofStreamSinkConfig - from nemo_relay.observability import HttpStorageConfig - from nemo_relay.observability import OtlpConfig - from nemo_relay.observability import S3StorageConfig - def current_virtualenv() -> Path | None: """Return the current virtual environment, if Python is running in one.""" @@ -196,10 +184,6 @@ def merge_unique(*values: Any) -> list[str]: return merged -def without_none(mapping: dict[str, Any]) -> dict[str, Any]: - return {key: value for key, value in mapping.items() if value is not None} - - def dump_yaml(value: dict[str, Any]) -> str: try: import yaml @@ -226,7 +210,7 @@ def load_relay_plugin_config(payload: dict[str, Any]) -> dict[str, Any]: { "kind": "observability", "enabled": True, - "config": plugin_config or {"version": 1}, + "config": plugin_config or {"version": 2}, } ], } @@ -243,214 +227,41 @@ def normalize_relay_output_dirs(plugin_config: dict[str, Any], payload: dict[str if component.get("kind") != "observability": continue config = component.setdefault("config", {}) - config.setdefault("version", 1) - for section_name in ("atof", "atif"): - section = config.get(section_name) - if not isinstance(section, dict) or not section.get("enabled"): - continue - - output_directory = section.get("output_directory") - if output_directory: - path = Path(output_directory) - if not path.is_absolute(): - path = base / path - else: - path = base / "artifacts" / "relay" - - section["output_directory"] = str(path / str(runtime_id)) - Path(section["output_directory"]).mkdir(parents=True, exist_ok=True) - if section_name == "atof": - section.setdefault("filename", "events.atof.jsonl") - section.setdefault("mode", "overwrite") - if section_name == "atif": - section.setdefault("filename_template", "trajectory-{session_id}.atif.json") - section.setdefault("agent_name", agent_name(payload)) - section.setdefault("model_name", relay_model_name(payload)) - - -def relay_api_plugin_config(plugin_config: dict[str, Any]) -> plugin.PluginConfig: - from nemo_relay import plugin - from nemo_relay.observability import ComponentSpec - from nemo_relay.observability import ConfigPolicy - from nemo_relay.observability import ObservabilityConfig - - components: list[Any] = [] - for component in plugin_config.get("components", []): - if not isinstance(component, dict): - continue - enabled = bool(component.get("enabled", True)) - config = component.get("config") or {} - if component.get("kind") == "observability" and isinstance(config, dict): - policy = config.get("policy") if isinstance(config.get("policy"), dict) else {} - components.append( - ComponentSpec( - ObservabilityConfig( - # Relay 0.6 only accepts the v2 observability API model. - # Fabric still accepts its existing flat/v1 configuration - # below and translates it at this API boundary. - version=2, - atof=_relay_api_atof_config(config.get("atof")), - atif=_relay_api_atif_config( - config.get("atif"), - ), - opentelemetry=_relay_api_otlp_config(config.get("opentelemetry")), - openinference=_relay_api_otlp_config(config.get("openinference")), - policy=ConfigPolicy( - unknown_component=policy.get("unknown_component", "warn"), - unknown_field=policy.get("unknown_field", "warn"), - unsupported_value=policy.get("unsupported_value", "error"), - ), - ), - enabled=enabled, - ) - ) - continue - components.append( - plugin.ComponentSpec( - kind=str(component.get("kind") or ""), - enabled=enabled, - config=config if isinstance(config, dict) else {}, - ) - ) - - policy = plugin_config.get("policy") if isinstance(plugin_config.get("policy"), dict) else {} - plugin_config = plugin.PluginConfig( - version=int(plugin_config.get("version", 1)), - components=components, - policy=plugin.ConfigPolicy( - unknown_component=policy.get("unknown_component", "warn"), - unknown_field=policy.get("unknown_field", "warn"), - unsupported_value=policy.get("unsupported_value", "error"), - ), - ) - - report = plugin.validate(plugin_config) - if any(diagnostic["level"] == "error" for diagnostic in report["diagnostics"]): - raise RuntimeError(report["diagnostics"]) + config.setdefault("version", 2) - return plugin_config - - -def _relay_api_atof_config(value: Any) -> AtofConfig | None: - if not isinstance(value, dict): - return None - from nemo_relay.observability import AtofConfig - from nemo_relay.observability import AtofFileSinkConfig - from nemo_relay.observability import AtofStreamSinkConfig - - sinks: list[AtofFileSinkConfig | AtofStreamSinkConfig] = [] - has_explicit_file_sink = False - for sink in value.get("sinks") or []: - if not isinstance(sink, dict): + atof = config.get("atof") + if isinstance(atof, dict) and atof.get("enabled"): + for sink in atof.get("sinks") or []: + if not isinstance(sink, dict) or sink.get("type") != "file": + continue + output_directory = sink.get("output_directory") + if output_directory: + path = Path(output_directory) + if not path.is_absolute(): + path = base / path + else: + path = base / "artifacts" / "relay" + sink["output_directory"] = str(path / str(runtime_id)) + Path(sink["output_directory"]).mkdir(parents=True, exist_ok=True) + sink.setdefault("filename", "events.atof.jsonl") + sink.setdefault("mode", "overwrite") + + atif = config.get("atif") + if not isinstance(atif, dict) or not atif.get("enabled"): continue - if sink.get("type") == "file": - has_explicit_file_sink = True - sinks.append(_relay_api_atof_file_sink_config(sink)) - elif sink.get("type") == "stream": - sinks.append(_relay_api_atof_stream_sink_config(sink)) - - if not has_explicit_file_sink and any( - key in value for key in ("output_directory", "filename", "mode") - ): - sinks.append(_relay_api_atof_file_sink_config(value)) - - for endpoint in value.get("endpoints") or []: - if isinstance(endpoint, dict): - sinks.append(_relay_api_atof_stream_sink_config(endpoint)) - - return AtofConfig( - enabled=bool(value.get("enabled", False)), - sinks=sinks, - ) - - -def _relay_api_atof_file_sink_config(value: dict[str, Any]) -> AtofFileSinkConfig: - from nemo_relay.observability import AtofFileSinkConfig - - return AtofFileSinkConfig( - output_directory=value.get("output_directory"), - filename=value.get("filename"), - mode=value.get("mode", "append"), - ) - + output_directory = atif.get("output_directory") + if output_directory: + path = Path(output_directory) + if not path.is_absolute(): + path = base / path + else: + path = base / "artifacts" / "relay" -def _relay_api_atof_stream_sink_config(value: dict[str, Any]) -> AtofStreamSinkConfig: - from nemo_relay.observability import AtofStreamSinkConfig - - return AtofStreamSinkConfig( - url=str(value.get("url", "")), - transport=value.get("transport", "http_post"), - headers=value.get("headers", {}), - header_env=value.get("header_env", {}), - timeout_millis=int(value.get("timeout_millis", 3000)), - field_name_policy=value.get("field_name_policy", "preserve"), - name=value.get("name"), - ) - - -def _relay_api_atif_config(value: Any) -> AtifConfig | None: - if not isinstance(value, dict): - return None - from nemo_relay.observability import AtifConfig - - storage_configs = value.get("storage") - storage = None - if isinstance(storage_configs, list): - storage = [_relay_api_storage_config(item) for item in storage_configs if isinstance(item, dict)] - return AtifConfig( - enabled=bool(value.get("enabled", False)), - agent_name=value.get("agent_name", "NeMo Relay"), - agent_version=value.get("agent_version"), - model_name=value.get("model_name", "unknown"), - tool_definitions=value.get("tool_definitions"), - extra=value.get("extra"), - output_directory=value.get("output_directory"), - filename_template=value.get("filename_template", "nemo-relay-atif-{session_id}.json"), - storage=storage, - ) - - -def _relay_api_storage_config(value: dict[str, Any]) -> HttpStorageConfig | S3StorageConfig: - if value.get("type") == "s3": - from nemo_relay.observability import S3StorageConfig - - return S3StorageConfig( - bucket=value.get("bucket", ""), - key_prefix=value.get("key_prefix"), - access_key_id=value.get("access_key_id"), - secret_access_key_var=value.get("secret_access_key_var"), - session_token_var=value.get("session_token_var"), - region=value.get("region"), - endpoint_url=value.get("endpoint_url"), - allow_http=value.get("allow_http"), - ) - from nemo_relay.observability import HttpStorageConfig - - return HttpStorageConfig( - endpoint=value.get("endpoint", ""), - headers=value.get("headers", {}), - header_env=value.get("header_env", {}), - timeout_millis=int(value.get("timeout_millis", 3000)), - ) - - -def _relay_api_otlp_config(value: Any) -> OtlpConfig | None: - if not isinstance(value, dict): - return None - from nemo_relay.observability import OtlpConfig - - return OtlpConfig( - enabled=bool(value.get("enabled", False)), - transport=value.get("transport", "http_binary"), - endpoint=value.get("endpoint"), - headers=value.get("headers", {}), - resource_attributes=value.get("resource_attributes", {}), - service_name=value.get("service_name", "nemo-relay"), - service_namespace=value.get("service_namespace"), - service_version=value.get("service_version"), - instrumentation_scope=value.get("instrumentation_scope"), - timeout_millis=int(value.get("timeout_millis", 3000)), - ) + atif["output_directory"] = str(path / str(runtime_id)) + Path(atif["output_directory"]).mkdir(parents=True, exist_ok=True) + atif.setdefault("filename_template", "trajectory-{session_id}.atif.json") + atif.setdefault("agent_name", agent_name(payload)) + atif.setdefault("model_name", relay_model_name(payload)) def collect_relay_artifacts(plugin_config: dict[str, Any]) -> list[dict[str, str]]: @@ -459,84 +270,39 @@ def collect_relay_artifacts(plugin_config: dict[str, Any]) -> list[dict[str, str if component.get("kind") != "observability": continue config = component.get("config") or {} - for section_name, pattern in ( - ("atof", "*.jsonl"), - ("atif", "*.json"), - ): - section = config.get(section_name) - if not isinstance(section, dict) or not section.get("enabled"): - continue - directory = Path(section.get("output_directory") or ".") - if not directory.exists(): - continue - for path in sorted(directory.glob(pattern)): - artifacts.append({"kind": section_name, "path": str(path)}) - return artifacts - - -def relay_cli_plugin_config( - plugin_config: dict[str, Any], *, observability_version: int -) -> dict[str, Any]: - """Render normalized Relay intent for the current external CLI contract.""" - - rendered = copy.deepcopy(plugin_config) - if observability_version == 1: - return rendered - if observability_version != 2: - raise ValueError( - f"unsupported NeMo Relay observability config version {observability_version}" - ) - for component in rendered.get("components", []): - if not isinstance(component, dict) or component.get("kind") != "observability": + atof = config.get("atof") + if isinstance(atof, dict) and atof.get("enabled"): + for sink in atof.get("sinks") or []: + if not isinstance(sink, dict) or sink.get("type") != "file": + continue + output_directory = sink.get("output_directory") + if not output_directory: + continue + directory = Path(output_directory) + if not directory.exists(): + continue + for path in sorted(directory.glob("*.jsonl")): + artifacts.append({"kind": "atof", "path": str(path)}) + + atif = config.get("atif") + if not isinstance(atif, dict) or not atif.get("enabled"): continue - config = component.get("config") - if not isinstance(config, dict) or int(config.get("version", 1)) != 1: + output_directory = atif.get("output_directory") + if not output_directory: continue - - atof = config.get("atof") - if isinstance(atof, dict): - sinks = list(atof.get("sinks") or []) - if atof.get("enabled"): - file_sink = without_none( - { - "type": "file", - "output_directory": atof.get("output_directory"), - "filename": atof.get("filename"), - "mode": atof.get("mode", "append"), - } - ) - sinks.append(file_sink) - for endpoint in atof.get("endpoints") or []: - if not isinstance(endpoint, dict): - continue - sinks.append( - without_none( - { - "type": "stream", - "url": endpoint.get("url"), - "transport": endpoint.get("transport", "http_post"), - "headers": endpoint.get("headers", {}), - "header_env": endpoint.get("header_env", {}), - "timeout_millis": endpoint.get("timeout_millis", 3000), - "field_name_policy": endpoint.get( - "field_name_policy", "preserve" - ), - } - ) - ) - config["atof"] = { - "enabled": bool(atof.get("enabled", False)), - "sinks": sinks, - } - config["version"] = 2 - return rendered + directory = Path(output_directory) + if not directory.exists(): + continue + for path in sorted(directory.glob("*.json")): + artifacts.append({"kind": "atif", "path": str(path)}) + return artifacts def write_relay_configs( *, relay_config: dict[str, Any] | None = None, plugin_config: dict[str, Any] | None = None, - observability_version: int = 1, + observability_version: int = 2, ) -> tuple[Path | None, Path | None]: try: import tomli_w @@ -556,14 +322,13 @@ def write_relay_configs( relay_config_path.write_text(tomli_w.dumps(relay_config), encoding="utf-8") if plugin_config is not None: + if observability_version != 2: + raise ValueError( + f"unsupported NeMo Relay observability config version {observability_version}" + ) plugin_config_path = config_dir / "plugins.toml" plugin_config_path.write_text( - tomli_w.dumps( - relay_cli_plugin_config( - plugin_config, - observability_version=observability_version, - ) - ), + tomli_w.dumps(plugin_config), encoding="utf-8", ) diff --git a/adapters/deepagents/README.md b/adapters/deepagents/README.md index d3b697aa1..3b105e020 100644 --- a/adapters/deepagents/README.md +++ b/adapters/deepagents/README.md @@ -181,6 +181,7 @@ CLI flags are involved: from nemo_fabric import ( RelayAtifConfig, RelayAtofConfig, + RelayAtofFileSinkConfig, RelayObservabilityConfig, ) from examples.code_review_agent import deepagents_config @@ -192,9 +193,13 @@ config.enable_relay( observability=RelayObservabilityConfig( atof=RelayAtofConfig( enabled=True, - output_directory="./artifacts/relay", - filename="events.atof.jsonl", - mode="overwrite", + sinks=[ + RelayAtofFileSinkConfig( + output_directory="./artifacts/relay", + filename="events.atof.jsonl", + mode="overwrite", + ) + ], ), atif=RelayAtifConfig( enabled=True, diff --git a/adapters/deepagents/src/nemo_fabric_adapters/deepagents/adapter.py b/adapters/deepagents/src/nemo_fabric_adapters/deepagents/adapter.py index 0a3523e7b..2afc77939 100644 --- a/adapters/deepagents/src/nemo_fabric_adapters/deepagents/adapter.py +++ b/adapters/deepagents/src/nemo_fabric_adapters/deepagents/adapter.py @@ -480,7 +480,7 @@ def __init__(self) -> None: self._relay_plugin: Any = None self._relay_scope: Any = None self._relay_scope_type: Any = None - self._relay_api_config: Any = None + self._relay_plugin_config: dict[str, Any] | None = None self._callback_handler_type: Any = None async def start(self, payload: dict[str, Any]) -> None: @@ -549,9 +549,7 @@ def _configure_observability(self, agent_kwargs: dict[str, Any]) -> dict[str, An raise _relay_dependency_error() from exc assert self._observability is not None - self._relay_api_config = common_utils.relay_api_plugin_config( - self._observability.plugin_config - ) + self._relay_plugin_config = self._observability.plugin_config self._relay_plugin = plugin self._relay_scope = scope self._relay_scope_type = ScopeType @@ -590,7 +588,7 @@ async def invoke(self, invocation: dict[str, Any]) -> dict[str, Any]: try: if self._observability is not None: callback_handler = self._callback_handler_type() - async with self._relay_plugin.plugin(self._relay_api_config): + async with self._relay_plugin.plugin(self._relay_plugin_config): with self._relay_scope.scope( "deepagents-request", self._relay_scope_type.Agent ): @@ -663,7 +661,7 @@ async def stop(self) -> None: self._relay_plugin = None self._relay_scope = None self._relay_scope_type = None - self._relay_api_config = None + self._relay_plugin_config = None self._callback_handler_type = None self._started = False if checkpointer is not None: diff --git a/adapters/hermes/src/nemo_fabric_adapters/hermes/adapter.py b/adapters/hermes/src/nemo_fabric_adapters/hermes/adapter.py index 07fd53aa3..b89876eff 100755 --- a/adapters/hermes/src/nemo_fabric_adapters/hermes/adapter.py +++ b/adapters/hermes/src/nemo_fabric_adapters/hermes/adapter.py @@ -238,12 +238,9 @@ async def start(self, payload: dict[str, Any]) -> None: self._relay_plugin_config = common_utils.load_relay_plugin_config( payload ) - relay_api_config = common_utils.relay_api_plugin_config( - self._relay_plugin_config - ) from nemo_relay import plugin - self._relay_context = plugin.plugin(relay_api_config) + self._relay_context = plugin.plugin(self._relay_plugin_config) await self._relay_context.__aenter__() self._relay_context_entered = True diff --git a/crates/fabric-core/src/config.rs b/crates/fabric-core/src/config.rs index 5316c4337..376ac0925 100644 --- a/crates/fabric-core/src/config.rs +++ b/crates/fabric-core/src/config.rs @@ -625,43 +625,59 @@ pub struct RelayAtofConfig { /// Whether ATOF export is enabled. #[serde(default)] pub enabled: bool, - /// Directory used for ATOF files. - #[serde(default, skip_serializing_if = "Option::is_none")] - pub output_directory: Option, - /// ATOF file name. - #[serde(default, skip_serializing_if = "Option::is_none")] - pub filename: Option, - /// File write mode. - #[serde(default)] - pub mode: RelayAtofMode, - /// Optional remote ATOF endpoints. + /// ATOF file and stream sinks. #[serde(default, skip_serializing_if = "Vec::is_empty")] - pub endpoints: Vec, + pub sinks: Vec, /// Additive ATOF fields. #[serde(default, flatten)] pub extensions: BTreeMap, } -/// Relay ATOF endpoint configuration. +/// Relay ATOF sink configuration. #[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)] -pub struct RelayAtofEndpointConfig { - /// Endpoint URL. - pub url: String, - /// Endpoint transport. - #[serde(default)] - pub transport: RelayAtofEndpointTransport, - /// Endpoint headers. - #[serde(default, skip_serializing_if = "BTreeMap::is_empty")] - pub headers: BTreeMap, - /// Request timeout in milliseconds. - #[serde(default = "default_relay_timeout_millis")] - pub timeout_millis: u64, - /// Field-name handling policy. - #[serde(default)] - pub field_name_policy: RelayAtofEndpointFieldNamePolicy, - /// Additive endpoint fields. - #[serde(default, flatten)] - pub extensions: BTreeMap, +#[serde(tag = "type", rename_all = "snake_case")] +pub enum RelayAtofSinkConfig { + /// Write ATOF records to a local file. + File { + /// Directory used for ATOF files. + #[serde(default, skip_serializing_if = "Option::is_none")] + output_directory: Option, + /// ATOF file name. + #[serde(default, skip_serializing_if = "Option::is_none")] + filename: Option, + /// File write mode. + #[serde(default)] + mode: RelayAtofMode, + /// Additive file sink fields. + #[serde(default, flatten)] + extensions: BTreeMap, + }, + /// Send ATOF records to a remote stream. + Stream { + /// Stream URL. + url: String, + /// Stream transport. + #[serde(default)] + transport: RelayAtofStreamTransport, + /// Static stream headers. + #[serde(default, skip_serializing_if = "BTreeMap::is_empty")] + headers: BTreeMap, + /// Environment-variable-backed stream headers. + #[serde(default, skip_serializing_if = "BTreeMap::is_empty")] + header_env: BTreeMap, + /// Request timeout in milliseconds. + #[serde(default = "default_relay_timeout_millis")] + timeout_millis: u64, + /// Field-name handling policy. + #[serde(default)] + field_name_policy: RelayAtofStreamFieldNamePolicy, + /// Optional stream sink name. + #[serde(default, skip_serializing_if = "Option::is_none")] + name: Option, + /// Additive stream sink fields. + #[serde(default, flatten)] + extensions: BTreeMap, + }, } /// Relay ATIF export configuration. @@ -874,10 +890,10 @@ pub enum RelayAtofMode { Overwrite, } -/// Relay ATOF endpoint transport. +/// Relay ATOF stream transport. #[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize, JsonSchema)] #[serde(rename_all = "snake_case")] -pub enum RelayAtofEndpointTransport { +pub enum RelayAtofStreamTransport { /// HTTP POST transport. #[default] HttpPost, @@ -887,10 +903,10 @@ pub enum RelayAtofEndpointTransport { Ndjson, } -/// Relay ATOF endpoint field-name policy. +/// Relay ATOF stream field-name policy. #[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize, JsonSchema)] #[serde(rename_all = "snake_case")] -pub enum RelayAtofEndpointFieldNamePolicy { +pub enum RelayAtofStreamFieldNamePolicy { /// Preserve field names. #[default] Preserve, @@ -933,7 +949,7 @@ impl TelemetryProvider { } fn default_relay_config_version() -> u32 { - 1 + 2 } fn default_enabled() -> bool { @@ -1643,6 +1659,36 @@ mod tests { PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("../..") } + #[test] + fn relay_observability_uses_v2_typed_atof_sinks() { + let observability: RelayObservabilityConfig = serde_json::from_value(serde_json::json!({ + "atof": { + "enabled": true, + "sinks": [ + { + "type": "file", + "output_directory": "artifacts/relay", + "filename": "events.atof.jsonl", + "mode": "overwrite" + }, + { + "type": "stream", + "url": "http://localhost:4319/events", + "transport": "ndjson", + "header_env": {"authorization": "RELAY_AUTHORIZATION"}, + "name": "live-events" + } + ] + } + })) + .expect("Relay v2 observability config"); + + let value = serde_json::to_value(observability).expect("serialized observability"); + assert_eq!(value["version"], 2); + assert_eq!(value["atof"]["sinks"][0]["type"], "file"); + assert_eq!(value["atof"]["sinks"][1]["type"], "stream"); + } + #[test] fn resolves_complete_typed_config_with_explicit_base_dir() { let base_dir = repository_root(); diff --git a/docs/reference/api/python-library-reference/index.md b/docs/reference/api/python-library-reference/index.md index 8d580c39c..b4096061c 100644 --- a/docs/reference/api/python-library-reference/index.md +++ b/docs/reference/api/python-library-reference/index.md @@ -31,7 +31,8 @@ SPDX-License-Identifier: Apache-2.0 */} - [`models.ModelConfig`](./nemo_fabric.models.md#class-modelconfig): Model alias configuration. - [`models.RelayAtifConfig`](./nemo_fabric.models.md#class-relayatifconfig): NeMo Relay ATIF export configuration. - [`models.RelayAtofConfig`](./nemo_fabric.models.md#class-relayatofconfig): NeMo Relay ATOF export configuration. -- [`models.RelayAtofEndpointConfig`](./nemo_fabric.models.md#class-relayatofendpointconfig): NeMo Relay ATOF remote endpoint configuration. +- [`models.RelayAtofFileSinkConfig`](./nemo_fabric.models.md#class-relayatoffilesinkconfig): NeMo Relay ATOF file sink configuration. +- [`models.RelayAtofStreamSinkConfig`](./nemo_fabric.models.md#class-relayatofstreamsinkconfig): NeMo Relay ATOF stream sink configuration. - [`models.RelayComponentConfig`](./nemo_fabric.models.md#class-relaycomponentconfig): Generic NeMo Relay plugin component configuration. - [`models.RelayConfig`](./nemo_fabric.models.md#class-relayconfig): First-class NeMo Relay integration configuration. - [`models.RelayConfigPolicy`](./nemo_fabric.models.md#class-relayconfigpolicy): NeMo Relay config validation policy. diff --git a/docs/reference/api/python-library-reference/nemo_fabric.models.md b/docs/reference/api/python-library-reference/nemo_fabric.models.md index 2736be7ce..ee2a955fd 100644 --- a/docs/reference/api/python-library-reference/nemo_fabric.models.md +++ b/docs/reference/api/python-library-reference/nemo_fabric.models.md @@ -668,8 +668,68 @@ Return a detached JSON-compatible mapping for Rust/core calls. --- -## class `RelayAtofEndpointConfig` -NeMo Relay ATOF remote endpoint configuration. +## class `RelayAtofFileSinkConfig` +NeMo Relay ATOF file sink configuration. + + +--- + +### property extra_fields + +Return fields preserved by the extension point for this model. + +--- + +### property model_extra + +Get extra fields set during validation. + + + +**Returns:** + A dictionary of extra fields, or `None` if `config.extra` is not set to `"allow"`. + +--- + +### property model_fields_set + +Returns the set of fields that have been explicitly set on this model instance. + + + +**Returns:** + A set of strings representing the fields that have been set, i.e. that were not filled from defaults. + + + +--- + + +### classmethod `from_mapping` + +```python +from_mapping(value: 'Mapping[str, Any]') → Self +``` + +Validate a mapping using this Pydantic model. + +--- + + +### method `to_mapping` + +```python +to_mapping() → dict[str, Any] +``` + +Return a detached JSON-compatible mapping for Rust/core calls. + + +--- + + +## class `RelayAtofStreamSinkConfig` +NeMo Relay ATOF stream sink configuration. --- diff --git a/docs/reference/api/rust-library-reference/nemo-fabric-core/config/enum-capabilitykind.mdx b/docs/reference/api/rust-library-reference/nemo-fabric-core/config/enum-capabilitykind.mdx index b0d72b888..b22fc853a 100644 --- a/docs/reference/api/rust-library-reference/nemo-fabric-core/config/enum-capabilitykind.mdx +++ b/docs/reference/api/rust-library-reference/nemo-fabric-core/config/enum-capabilitykind.mdx @@ -2,7 +2,7 @@ title: "Enum Capability Kind" sidebar-title: "CapabilityKind" description: "Capability kind." -position: 39 +position: 38 --- {/* SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved. SPDX-License-Identifier: Apache-2.0 */} diff --git a/docs/reference/api/rust-library-reference/nemo-fabric-core/config/enum-capabilitytarget.mdx b/docs/reference/api/rust-library-reference/nemo-fabric-core/config/enum-capabilitytarget.mdx index 2b82413a1..cfd744fbe 100644 --- a/docs/reference/api/rust-library-reference/nemo-fabric-core/config/enum-capabilitytarget.mdx +++ b/docs/reference/api/rust-library-reference/nemo-fabric-core/config/enum-capabilitytarget.mdx @@ -2,7 +2,7 @@ title: "Enum Capability Target" sidebar-title: "CapabilityTarget" description: "Capability routing target." -position: 40 +position: 39 --- {/* SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved. SPDX-License-Identifier: Apache-2.0 */} diff --git a/docs/reference/api/rust-library-reference/nemo-fabric-core/config/enum-relayatifstorageconfig.mdx b/docs/reference/api/rust-library-reference/nemo-fabric-core/config/enum-relayatifstorageconfig.mdx index 1b2589f2c..aed4582ca 100644 --- a/docs/reference/api/rust-library-reference/nemo-fabric-core/config/enum-relayatifstorageconfig.mdx +++ b/docs/reference/api/rust-library-reference/nemo-fabric-core/config/enum-relayatifstorageconfig.mdx @@ -2,7 +2,7 @@ title: "Enum Relay Atif Storage Config" sidebar-title: "RelayAtifStorageConfig" description: "Relay ATIF remote storage configuration." -position: 44 +position: 43 --- {/* SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved. SPDX-License-Identifier: Apache-2.0 */} diff --git a/docs/reference/api/rust-library-reference/nemo-fabric-core/config/enum-relayatofmode.mdx b/docs/reference/api/rust-library-reference/nemo-fabric-core/config/enum-relayatofmode.mdx index 010b88501..6c6d8a7e9 100644 --- a/docs/reference/api/rust-library-reference/nemo-fabric-core/config/enum-relayatofmode.mdx +++ b/docs/reference/api/rust-library-reference/nemo-fabric-core/config/enum-relayatofmode.mdx @@ -2,7 +2,7 @@ title: "Enum Relay Atof Mode" sidebar-title: "RelayAtofMode" description: "Relay ATOF file mode." -position: 47 +position: 44 --- {/* SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved. SPDX-License-Identifier: Apache-2.0 */} diff --git a/docs/reference/api/rust-library-reference/nemo-fabric-core/config/struct-relayatofendpointconfig.mdx b/docs/reference/api/rust-library-reference/nemo-fabric-core/config/enum-relayatofsinkconfig.mdx similarity index 52% rename from docs/reference/api/rust-library-reference/nemo-fabric-core/config/struct-relayatofendpointconfig.mdx rename to docs/reference/api/rust-library-reference/nemo-fabric-core/config/enum-relayatofsinkconfig.mdx index 74f698421..5af1bf71f 100644 --- a/docs/reference/api/rust-library-reference/nemo-fabric-core/config/struct-relayatofendpointconfig.mdx +++ b/docs/reference/api/rust-library-reference/nemo-fabric-core/config/enum-relayatofsinkconfig.mdx @@ -1,77 +1,117 @@ --- -title: "Struct Relay Atof Endpoint Config" -sidebar-title: "RelayAtofEndpointConfig" -description: "Relay ATOF endpoint configuration." -position: 20 +title: "Enum Relay Atof Sink Config" +sidebar-title: "RelayAtofSinkConfig" +description: "Relay ATOF sink configuration." +position: 45 --- {/* SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved. SPDX-License-Identifier: Apache-2.0 */} Generated from `cargo doc --no-deps -p nemo-fabric-core`. -
String,\n    pub transport: RelayAtofEndpointTransport,\n    pub headers: BTreeMap<String, String>,\n    pub timeout_millis: u64,\n    pub field_name_policy: RelayAtofEndpointFieldNamePolicy,\n    pub extensions: BTreeMap<String, Value>,\n}"}} />
+
Option<PathBuf>,\n        filename: Option<String>,\n        mode: RelayAtofMode,\n        extensions: BTreeMap<String, Value>,\n    },\n    Stream {\n        url: String,\n        transport: RelayAtofStreamTransport,\n        headers: BTreeMap<String, String>,\n        header_env: BTreeMap<String, String>,\n        timeout_millis: u64,\n        field_name_policy: RelayAtofStreamFieldNamePolicy,\n        name: Option<String>,\n        extensions: BTreeMap<String, Value>,\n    },\n}"}} />
-Relay ATOF endpoint configuration. +Relay ATOF sink configuration. -## Fields +## Variants + +### `File` + +
+ +Write ATOF records to a local file. + +#### Fields + +### `output_directory: Option` + +Directory used for ATOF files. + +### `filename: Option` + +ATOF file name. + +### `mode: RelayAtofMode` + +File write mode. + +### `extensions: BTreeMap` + +Additive file sink fields. + +### `Stream` + +
+ +Send ATOF records to a remote stream. + +#### Fields ### `url: String` -Endpoint URL. +Stream URL. -### `transport: RelayAtofEndpointTransport` +### `transport: RelayAtofStreamTransport` -Endpoint transport. +Stream transport. ### `headers: BTreeMap` -Endpoint headers. +Static stream headers. + +### `header_env: BTreeMap` + +Environment-variable-backed stream headers. ### `timeout_millis: u64` Request timeout in milliseconds. -### `field_name_policy: RelayAtofEndpointFieldNamePolicy` +### `field_name_policy: RelayAtofStreamFieldNamePolicy` Field-name handling policy. +### `name: Option` + +Optional stream sink name. + ### `extensions: BTreeMap` -Additive endpoint fields. +Additive stream sink fields. ## Trait Implementations -### `impl Clone for RelayAtofEndpointConfig` +### `impl Clone for RelayAtofSinkConfig` -
Clone for RelayAtofEndpointConfig"}} />
+
Clone for RelayAtofSinkConfig"}} />
#### `clone` -
clone(&self) -> RelayAtofEndpointConfig"}} />
+
clone(&self) -> RelayAtofSinkConfig"}} />
#### `clone_from`
clone_from(&mut self, source: &Self)"}} />
-### `impl Debug for RelayAtofEndpointConfig` +### `impl Debug for RelayAtofSinkConfig` -
Debug for RelayAtofEndpointConfig"}} />
+
Debug for RelayAtofSinkConfig"}} />
#### `fmt`
fmt(&self, f: &mut Formatter<'_>) -> Result"}} />
-### `impl<'de> Deserialize<'de> for RelayAtofEndpointConfig` +### `impl<'de> Deserialize<'de> for RelayAtofSinkConfig` -
Deserialize<'de> for RelayAtofEndpointConfig"}} />
+
Deserialize<'de> for RelayAtofSinkConfig"}} />
#### `deserialize`
deserialize<__D>(__deserializer: __D) -> Result<Self, __D::Error>where\n    __D: Deserializer<'de>,"}} />
-### `impl JsonSchema for RelayAtofEndpointConfig` +### `impl JsonSchema for RelayAtofSinkConfig` -
RelayAtofEndpointConfig"}} />
+
RelayAtofSinkConfig"}} />
#### `schema_name` @@ -89,26 +129,26 @@ Additive endpoint fields.
bool"}} />
-### `impl PartialEq for RelayAtofEndpointConfig` +### `impl PartialEq for RelayAtofSinkConfig` -
PartialEq for RelayAtofEndpointConfig"}} />
+
PartialEq for RelayAtofSinkConfig"}} />
#### `eq` -
eq(&self, other: &RelayAtofEndpointConfig) -> bool"}} />
+
eq(&self, other: &RelayAtofSinkConfig) -> bool"}} />
#### `ne`
ne(&self, other: &Rhs) -> bool"}} />
-### `impl Serialize for RelayAtofEndpointConfig` +### `impl Serialize for RelayAtofSinkConfig` -
Serialize for RelayAtofEndpointConfig"}} />
+
Serialize for RelayAtofSinkConfig"}} />
#### `serialize`
serialize<__S>(&self, __serializer: __S) -> Result<__S::Ok, __S::Error>where\n    __S: Serializer,"}} />
-### `impl StructuralPartialEq for RelayAtofEndpointConfig` +### `impl StructuralPartialEq for RelayAtofSinkConfig` -
StructuralPartialEq for RelayAtofEndpointConfig"}} />
+
StructuralPartialEq for RelayAtofSinkConfig"}} />
diff --git a/docs/reference/api/rust-library-reference/nemo-fabric-core/config/enum-relayatofendpointfieldnamepolicy.mdx b/docs/reference/api/rust-library-reference/nemo-fabric-core/config/enum-relayatofstreamfieldnamepolicy.mdx similarity index 72% rename from docs/reference/api/rust-library-reference/nemo-fabric-core/config/enum-relayatofendpointfieldnamepolicy.mdx rename to docs/reference/api/rust-library-reference/nemo-fabric-core/config/enum-relayatofstreamfieldnamepolicy.mdx index e500f6aa1..9beca42d2 100644 --- a/docs/reference/api/rust-library-reference/nemo-fabric-core/config/enum-relayatofendpointfieldnamepolicy.mdx +++ b/docs/reference/api/rust-library-reference/nemo-fabric-core/config/enum-relayatofstreamfieldnamepolicy.mdx @@ -1,8 +1,8 @@ --- -title: "Enum Relay Atof Endpoint Field Name Policy" -sidebar-title: "RelayAtofEndpointFieldNamePolicy" -description: "Relay ATOF endpoint field-name policy." -position: 45 +title: "Enum Relay Atof Stream Field Name Policy" +sidebar-title: "RelayAtofStreamFieldNamePolicy" +description: "Relay ATOF stream field-name policy." +position: 46 --- {/* SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved. SPDX-License-Identifier: Apache-2.0 */} @@ -10,13 +10,13 @@ SPDX-License-Identifier: Apache-2.0 */} Generated from `cargo doc --no-deps -p nemo-fabric-core`. ```rust -pub enum RelayAtofEndpointFieldNamePolicy { +pub enum RelayAtofStreamFieldNamePolicy { Preserve, ReplaceDots, } ``` -Relay ATOF endpoint field-name policy. +Relay ATOF stream field-name policy. ## Variants @@ -34,45 +34,45 @@ Replace dots in field names. ## Trait Implementations -### `impl Clone for RelayAtofEndpointFieldNamePolicy` +### `impl Clone for RelayAtofStreamFieldNamePolicy` -
Clone for RelayAtofEndpointFieldNamePolicy"}} />
+
Clone for RelayAtofStreamFieldNamePolicy"}} />
#### `clone` -
clone(&self) -> RelayAtofEndpointFieldNamePolicy"}} />
+
clone(&self) -> RelayAtofStreamFieldNamePolicy"}} />
#### `clone_from`
clone_from(&mut self, source: &Self)"}} />
-### `impl Debug for RelayAtofEndpointFieldNamePolicy` +### `impl Debug for RelayAtofStreamFieldNamePolicy` -
Debug for RelayAtofEndpointFieldNamePolicy"}} />
+
Debug for RelayAtofStreamFieldNamePolicy"}} />
#### `fmt`
fmt(&self, f: &mut Formatter<'_>) -> Result"}} />
-### `impl Default for RelayAtofEndpointFieldNamePolicy` +### `impl Default for RelayAtofStreamFieldNamePolicy` -
Default for RelayAtofEndpointFieldNamePolicy"}} />
+
Default for RelayAtofStreamFieldNamePolicy"}} />
#### `default` -
default() -> RelayAtofEndpointFieldNamePolicy"}} />
+
default() -> RelayAtofStreamFieldNamePolicy"}} />
-### `impl<'de> Deserialize<'de> for RelayAtofEndpointFieldNamePolicy` +### `impl<'de> Deserialize<'de> for RelayAtofStreamFieldNamePolicy` -
Deserialize<'de> for RelayAtofEndpointFieldNamePolicy"}} />
+
Deserialize<'de> for RelayAtofStreamFieldNamePolicy"}} />
#### `deserialize`
deserialize<__D>(__deserializer: __D) -> Result<Self, __D::Error>where\n    __D: Deserializer<'de>,"}} />
-### `impl JsonSchema for RelayAtofEndpointFieldNamePolicy` +### `impl JsonSchema for RelayAtofStreamFieldNamePolicy` -
RelayAtofEndpointFieldNamePolicy"}} />
+
RelayAtofStreamFieldNamePolicy"}} />
#### `schema_name` @@ -90,34 +90,34 @@ Replace dots in field names.
bool"}} />
-### `impl PartialEq for RelayAtofEndpointFieldNamePolicy` +### `impl PartialEq for RelayAtofStreamFieldNamePolicy` -
PartialEq for RelayAtofEndpointFieldNamePolicy"}} />
+
PartialEq for RelayAtofStreamFieldNamePolicy"}} />
#### `eq` -
eq(&self, other: &RelayAtofEndpointFieldNamePolicy) -> bool"}} />
+
eq(&self, other: &RelayAtofStreamFieldNamePolicy) -> bool"}} />
#### `ne`
ne(&self, other: &Rhs) -> bool"}} />
-### `impl Serialize for RelayAtofEndpointFieldNamePolicy` +### `impl Serialize for RelayAtofStreamFieldNamePolicy` -
Serialize for RelayAtofEndpointFieldNamePolicy"}} />
+
Serialize for RelayAtofStreamFieldNamePolicy"}} />
#### `serialize`
serialize<__S>(&self, __serializer: __S) -> Result<__S::Ok, __S::Error>where\n    __S: Serializer,"}} />
-### `impl Copy for RelayAtofEndpointFieldNamePolicy` +### `impl Copy for RelayAtofStreamFieldNamePolicy` -
Copy for RelayAtofEndpointFieldNamePolicy"}} />
+
Copy for RelayAtofStreamFieldNamePolicy"}} />
-### `impl Eq for RelayAtofEndpointFieldNamePolicy` +### `impl Eq for RelayAtofStreamFieldNamePolicy` -
Eq for RelayAtofEndpointFieldNamePolicy"}} />
+
Eq for RelayAtofStreamFieldNamePolicy"}} />
-### `impl StructuralPartialEq for RelayAtofEndpointFieldNamePolicy` +### `impl StructuralPartialEq for RelayAtofStreamFieldNamePolicy` -
StructuralPartialEq for RelayAtofEndpointFieldNamePolicy"}} />
+
StructuralPartialEq for RelayAtofStreamFieldNamePolicy"}} />
diff --git a/docs/reference/api/rust-library-reference/nemo-fabric-core/config/enum-relayatofendpointtransport.mdx b/docs/reference/api/rust-library-reference/nemo-fabric-core/config/enum-relayatofstreamtransport.mdx similarity index 75% rename from docs/reference/api/rust-library-reference/nemo-fabric-core/config/enum-relayatofendpointtransport.mdx rename to docs/reference/api/rust-library-reference/nemo-fabric-core/config/enum-relayatofstreamtransport.mdx index 089169e4e..4bf15bdef 100644 --- a/docs/reference/api/rust-library-reference/nemo-fabric-core/config/enum-relayatofendpointtransport.mdx +++ b/docs/reference/api/rust-library-reference/nemo-fabric-core/config/enum-relayatofstreamtransport.mdx @@ -1,8 +1,8 @@ --- -title: "Enum Relay Atof Endpoint Transport" -sidebar-title: "RelayAtofEndpointTransport" -description: "Relay ATOF endpoint transport." -position: 46 +title: "Enum Relay Atof Stream Transport" +sidebar-title: "RelayAtofStreamTransport" +description: "Relay ATOF stream transport." +position: 47 --- {/* SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved. SPDX-License-Identifier: Apache-2.0 */} @@ -10,14 +10,14 @@ SPDX-License-Identifier: Apache-2.0 */} Generated from `cargo doc --no-deps -p nemo-fabric-core`. ```rust -pub enum RelayAtofEndpointTransport { +pub enum RelayAtofStreamTransport { HttpPost, Websocket, Ndjson, } ``` -Relay ATOF endpoint transport. +Relay ATOF stream transport. ## Variants @@ -41,45 +41,45 @@ NDJSON transport. ## Trait Implementations -### `impl Clone for RelayAtofEndpointTransport` +### `impl Clone for RelayAtofStreamTransport` -
Clone for RelayAtofEndpointTransport"}} />
+
Clone for RelayAtofStreamTransport"}} />
#### `clone` -
clone(&self) -> RelayAtofEndpointTransport"}} />
+
clone(&self) -> RelayAtofStreamTransport"}} />
#### `clone_from`
clone_from(&mut self, source: &Self)"}} />
-### `impl Debug for RelayAtofEndpointTransport` +### `impl Debug for RelayAtofStreamTransport` -
Debug for RelayAtofEndpointTransport"}} />
+
Debug for RelayAtofStreamTransport"}} />
#### `fmt`
fmt(&self, f: &mut Formatter<'_>) -> Result"}} />
-### `impl Default for RelayAtofEndpointTransport` +### `impl Default for RelayAtofStreamTransport` -
Default for RelayAtofEndpointTransport"}} />
+
Default for RelayAtofStreamTransport"}} />
#### `default` -
default() -> RelayAtofEndpointTransport"}} />
+
default() -> RelayAtofStreamTransport"}} />
-### `impl<'de> Deserialize<'de> for RelayAtofEndpointTransport` +### `impl<'de> Deserialize<'de> for RelayAtofStreamTransport` -
Deserialize<'de> for RelayAtofEndpointTransport"}} />
+
Deserialize<'de> for RelayAtofStreamTransport"}} />
#### `deserialize`
deserialize<__D>(__deserializer: __D) -> Result<Self, __D::Error>where\n    __D: Deserializer<'de>,"}} />
-### `impl JsonSchema for RelayAtofEndpointTransport` +### `impl JsonSchema for RelayAtofStreamTransport` -
RelayAtofEndpointTransport"}} />
+
RelayAtofStreamTransport"}} />
#### `schema_name` @@ -97,34 +97,34 @@ NDJSON transport.
bool"}} />
-### `impl PartialEq for RelayAtofEndpointTransport` +### `impl PartialEq for RelayAtofStreamTransport` -
PartialEq for RelayAtofEndpointTransport"}} />
+
PartialEq for RelayAtofStreamTransport"}} />
#### `eq` -
eq(&self, other: &RelayAtofEndpointTransport) -> bool"}} />
+
eq(&self, other: &RelayAtofStreamTransport) -> bool"}} />
#### `ne`
ne(&self, other: &Rhs) -> bool"}} />
-### `impl Serialize for RelayAtofEndpointTransport` +### `impl Serialize for RelayAtofStreamTransport` -
Serialize for RelayAtofEndpointTransport"}} />
+
Serialize for RelayAtofStreamTransport"}} />
#### `serialize`
serialize<__S>(&self, __serializer: __S) -> Result<__S::Ok, __S::Error>where\n    __S: Serializer,"}} />
-### `impl Copy for RelayAtofEndpointTransport` +### `impl Copy for RelayAtofStreamTransport` -
Copy for RelayAtofEndpointTransport"}} />
+
Copy for RelayAtofStreamTransport"}} />
-### `impl Eq for RelayAtofEndpointTransport` +### `impl Eq for RelayAtofStreamTransport` -
Eq for RelayAtofEndpointTransport"}} />
+
Eq for RelayAtofStreamTransport"}} />
-### `impl StructuralPartialEq for RelayAtofEndpointTransport` +### `impl StructuralPartialEq for RelayAtofStreamTransport` -
StructuralPartialEq for RelayAtofEndpointTransport"}} />
+
StructuralPartialEq for RelayAtofStreamTransport"}} />
diff --git a/docs/reference/api/rust-library-reference/nemo-fabric-core/config/index.mdx b/docs/reference/api/rust-library-reference/nemo-fabric-core/config/index.mdx index 4de4cef4d..d82bf8c58 100644 --- a/docs/reference/api/rust-library-reference/nemo-fabric-core/config/index.mdx +++ b/docs/reference/api/rust-library-reference/nemo-fabric-core/config/index.mdx @@ -32,7 +32,6 @@ Fabric config models and loading helpers. - [ModelConfig](/reference/api/rust-library-reference/nemo-fabric-core/config/struct-modelconfig): Model configuration. - [RelayAtifConfig](/reference/api/rust-library-reference/nemo-fabric-core/config/struct-relayatifconfig): Relay ATIF export configuration. - [RelayAtofConfig](/reference/api/rust-library-reference/nemo-fabric-core/config/struct-relayatofconfig): Relay ATOF export configuration. -- [RelayAtofEndpointConfig](/reference/api/rust-library-reference/nemo-fabric-core/config/struct-relayatofendpointconfig): Relay ATOF endpoint configuration. - [RelayComponentConfig](/reference/api/rust-library-reference/nemo-fabric-core/config/struct-relaycomponentconfig): Generic NeMo Relay plugin component configuration. - [RelayConfig](/reference/api/rust-library-reference/nemo-fabric-core/config/struct-relayconfig): NeMo Relay integration configuration. - [RelayConfigPolicy](/reference/api/rust-library-reference/nemo-fabric-core/config/struct-relayconfigpolicy): Relay validation policy. @@ -60,9 +59,10 @@ Fabric config models and loading helpers. - [EnvironmentOwnership](/reference/api/rust-library-reference/nemo-fabric-core/config/enum-environmentownership): Whether Fabric owns the underlying environment resource. - [McpExposure](/reference/api/rust-library-reference/nemo-fabric-core/config/enum-mcpexposure): MCP exposure strategy. - [RelayAtifStorageConfig](/reference/api/rust-library-reference/nemo-fabric-core/config/enum-relayatifstorageconfig): Relay ATIF remote storage configuration. -- [RelayAtofEndpointFieldNamePolicy](/reference/api/rust-library-reference/nemo-fabric-core/config/enum-relayatofendpointfieldnamepolicy): Relay ATOF endpoint field-name policy. -- [RelayAtofEndpointTransport](/reference/api/rust-library-reference/nemo-fabric-core/config/enum-relayatofendpointtransport): Relay ATOF endpoint transport. - [RelayAtofMode](/reference/api/rust-library-reference/nemo-fabric-core/config/enum-relayatofmode): Relay ATOF file mode. +- [RelayAtofSinkConfig](/reference/api/rust-library-reference/nemo-fabric-core/config/enum-relayatofsinkconfig): Relay ATOF sink configuration. +- [RelayAtofStreamFieldNamePolicy](/reference/api/rust-library-reference/nemo-fabric-core/config/enum-relayatofstreamfieldnamepolicy): Relay ATOF stream field-name policy. +- [RelayAtofStreamTransport](/reference/api/rust-library-reference/nemo-fabric-core/config/enum-relayatofstreamtransport): Relay ATOF stream transport. - [RelayOtlpTransport](/reference/api/rust-library-reference/nemo-fabric-core/config/enum-relayotlptransport): Relay OTLP transport. - [RelayUnsupportedBehavior](/reference/api/rust-library-reference/nemo-fabric-core/config/enum-relayunsupportedbehavior): Relay unsupported/unknown config handling. - [ResolutionStrategy](/reference/api/rust-library-reference/nemo-fabric-core/config/enum-resolutionstrategy): Adapter install or availability strategy. diff --git a/docs/reference/api/rust-library-reference/nemo-fabric-core/config/struct-relayatofconfig.mdx b/docs/reference/api/rust-library-reference/nemo-fabric-core/config/struct-relayatofconfig.mdx index 0d91c211a..6d30976bf 100644 --- a/docs/reference/api/rust-library-reference/nemo-fabric-core/config/struct-relayatofconfig.mdx +++ b/docs/reference/api/rust-library-reference/nemo-fabric-core/config/struct-relayatofconfig.mdx @@ -9,7 +9,7 @@ SPDX-License-Identifier: Apache-2.0 */} Generated from `cargo doc --no-deps -p nemo-fabric-core`. -
bool,\n    pub output_directory: Option<PathBuf>,\n    pub filename: Option<String>,\n    pub mode: RelayAtofMode,\n    pub endpoints: Vec<RelayAtofEndpointConfig>,\n    pub extensions: BTreeMap<String, Value>,\n}"}} />
+
bool,\n    pub sinks: Vec<RelayAtofSinkConfig>,\n    pub extensions: BTreeMap<String, Value>,\n}"}} />
Relay ATOF export configuration. @@ -19,21 +19,9 @@ Relay ATOF export configuration. Whether ATOF export is enabled. -### `output_directory: Option` +### `sinks: Vec` -Directory used for ATOF files. - -### `filename: Option` - -ATOF file name. - -### `mode: RelayAtofMode` - -File write mode. - -### `endpoints: Vec` - -Optional remote ATOF endpoints. +ATOF file and stream sinks. ### `extensions: BTreeMap` diff --git a/docs/reference/api/rust-library-reference/nemo-fabric-core/config/struct-relaycomponentconfig.mdx b/docs/reference/api/rust-library-reference/nemo-fabric-core/config/struct-relaycomponentconfig.mdx index 44df6f4e6..308f7ec67 100644 --- a/docs/reference/api/rust-library-reference/nemo-fabric-core/config/struct-relaycomponentconfig.mdx +++ b/docs/reference/api/rust-library-reference/nemo-fabric-core/config/struct-relaycomponentconfig.mdx @@ -2,7 +2,7 @@ title: "Struct Relay Component Config" sidebar-title: "RelayComponentConfig" description: "Generic NeMo Relay plugin component configuration." -position: 21 +position: 20 --- {/* SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved. SPDX-License-Identifier: Apache-2.0 */} diff --git a/docs/reference/api/rust-library-reference/nemo-fabric-core/config/struct-relayconfig.mdx b/docs/reference/api/rust-library-reference/nemo-fabric-core/config/struct-relayconfig.mdx index 852abee85..d00e15d59 100644 --- a/docs/reference/api/rust-library-reference/nemo-fabric-core/config/struct-relayconfig.mdx +++ b/docs/reference/api/rust-library-reference/nemo-fabric-core/config/struct-relayconfig.mdx @@ -2,7 +2,7 @@ title: "Struct Relay Config" sidebar-title: "RelayConfig" description: "NeMo Relay integration configuration." -position: 22 +position: 21 --- {/* SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved. SPDX-License-Identifier: Apache-2.0 */} diff --git a/docs/reference/api/rust-library-reference/nemo-fabric-core/config/struct-relayconfigpolicy.mdx b/docs/reference/api/rust-library-reference/nemo-fabric-core/config/struct-relayconfigpolicy.mdx index 66b07af70..33ff49c38 100644 --- a/docs/reference/api/rust-library-reference/nemo-fabric-core/config/struct-relayconfigpolicy.mdx +++ b/docs/reference/api/rust-library-reference/nemo-fabric-core/config/struct-relayconfigpolicy.mdx @@ -2,7 +2,7 @@ title: "Struct Relay Config Policy" sidebar-title: "RelayConfigPolicy" description: "Relay validation policy." -position: 23 +position: 22 --- {/* SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved. SPDX-License-Identifier: Apache-2.0 */} diff --git a/docs/reference/api/rust-library-reference/nemo-fabric-core/config/struct-relayobservabilityconfig.mdx b/docs/reference/api/rust-library-reference/nemo-fabric-core/config/struct-relayobservabilityconfig.mdx index 1fba378f6..2deeb06be 100644 --- a/docs/reference/api/rust-library-reference/nemo-fabric-core/config/struct-relayobservabilityconfig.mdx +++ b/docs/reference/api/rust-library-reference/nemo-fabric-core/config/struct-relayobservabilityconfig.mdx @@ -2,7 +2,7 @@ title: "Struct Relay Observability Config" sidebar-title: "RelayObservabilityConfig" description: "NeMo Relay observability component configuration." -position: 24 +position: 23 --- {/* SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved. SPDX-License-Identifier: Apache-2.0 */} diff --git a/docs/reference/api/rust-library-reference/nemo-fabric-core/config/struct-relayotlpconfig.mdx b/docs/reference/api/rust-library-reference/nemo-fabric-core/config/struct-relayotlpconfig.mdx index 0d9cfc977..d51f995d0 100644 --- a/docs/reference/api/rust-library-reference/nemo-fabric-core/config/struct-relayotlpconfig.mdx +++ b/docs/reference/api/rust-library-reference/nemo-fabric-core/config/struct-relayotlpconfig.mdx @@ -2,7 +2,7 @@ title: "Struct Relay Otlp Config" sidebar-title: "RelayOtlpConfig" description: "Relay OpenTelemetry/OpenInference export configuration." -position: 25 +position: 24 --- {/* SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved. SPDX-License-Identifier: Apache-2.0 */} diff --git a/docs/reference/api/rust-library-reference/nemo-fabric-core/config/struct-toolsconfig.mdx b/docs/reference/api/rust-library-reference/nemo-fabric-core/config/struct-toolsconfig.mdx index ccedcf918..94985e2cc 100644 --- a/docs/reference/api/rust-library-reference/nemo-fabric-core/config/struct-toolsconfig.mdx +++ b/docs/reference/api/rust-library-reference/nemo-fabric-core/config/struct-toolsconfig.mdx @@ -2,7 +2,7 @@ title: "Struct Tools Config" sidebar-title: "ToolsConfig" description: "Harness-neutral tool capability configuration." -position: 35 +position: 34 --- {/* SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved. SPDX-License-Identifier: Apache-2.0 */} diff --git a/docs/reference/api/rust-library-reference/nemo-fabric-core/config/struct-toolsplan.mdx b/docs/reference/api/rust-library-reference/nemo-fabric-core/config/struct-toolsplan.mdx index 948ab6ecf..204811185 100644 --- a/docs/reference/api/rust-library-reference/nemo-fabric-core/config/struct-toolsplan.mdx +++ b/docs/reference/api/rust-library-reference/nemo-fabric-core/config/struct-toolsplan.mdx @@ -2,7 +2,7 @@ title: "Struct Tools Plan" sidebar-title: "ToolsPlan" description: "Normalized tool policy for a run." -position: 36 +position: 35 --- {/* SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved. SPDX-License-Identifier: Apache-2.0 */} diff --git a/docs/sdk/python.mdx b/docs/sdk/python.mdx index 75c7ac7aa..a8cc4f7e8 100644 --- a/docs/sdk/python.mdx +++ b/docs/sdk/python.mdx @@ -213,22 +213,42 @@ variant = review_agent_config(config, github_mcp=True, relay=True) ``` Relay observability is represented directly in the SDK config's top-level -`relay` block, and additional Relay plugin components can be supplied -generically when their component package is available in the runtime -environment: +`relay` block. ATOF uses Relay 0.6 file and stream sinks: ```python -from nemo_fabric import RelayComponentConfig +from nemo_fabric import ( + RelayAtofConfig, + RelayAtofFileSinkConfig, + RelayAtofStreamSinkConfig, + RelayObservabilityConfig, +) relay_config = config.model_copy(deep=True) relay_config.enable_relay( output_dir="./artifacts/relay", - components=[ - RelayComponentConfig(kind="switchyard", config={"route": "canary"}), - ], + observability=RelayObservabilityConfig( + atof=RelayAtofConfig( + enabled=True, + sinks=[ + RelayAtofFileSinkConfig( + output_directory="./artifacts/relay", + filename="events.atof.jsonl", + mode="overwrite", + ), + RelayAtofStreamSinkConfig( + url="http://localhost:4319/events", + transport="ndjson", + ), + ], + ), + ), ) ``` +Additional Relay plugin components can be supplied generically with +`RelayComponentConfig` when their component package is available in the runtime +environment. + The repository's [code-review example](https://github.com/NVIDIA/NeMo-Fabric/tree/main/examples/code_review_agent) uses this pattern for complete Hermes Agent, Codex, Deep Agents, diff --git a/examples/code_review_agent/config.py b/examples/code_review_agent/config.py index 127f545cd..0add94980 100644 --- a/examples/code_review_agent/config.py +++ b/examples/code_review_agent/config.py @@ -14,6 +14,7 @@ from nemo_fabric import ModelConfig from nemo_fabric import RelayAtifConfig from nemo_fabric import RelayAtofConfig +from nemo_fabric import RelayAtofFileSinkConfig from nemo_fabric import RelayObservabilityConfig from nemo_fabric import RelayOtlpConfig from nemo_fabric import RuntimeConfig @@ -239,9 +240,13 @@ def with_relay(base: FabricConfig) -> FabricConfig: ), atof=RelayAtofConfig( enabled=True, - output_directory="./artifacts/relay", - filename="events.atof.jsonl", - mode="overwrite", + sinks=[ + RelayAtofFileSinkConfig( + output_directory="./artifacts/relay", + filename="events.atof.jsonl", + mode="overwrite", + ) + ], ), ), ) @@ -292,7 +297,9 @@ def with_relay_openinference(base: FabricConfig) -> FabricConfig: if isinstance(observability.atif, RelayAtifConfig): observability.atif.output_directory = "./artifacts/relay-openinference" if isinstance(observability.atof, RelayAtofConfig): - observability.atof.output_directory = "./artifacts/relay-openinference" + for sink in observability.atof.sinks or []: + if isinstance(sink, RelayAtofFileSinkConfig): + sink.output_directory = "./artifacts/relay-openinference" return config diff --git a/examples/notebooks/02_variations.ipynb b/examples/notebooks/02_variations.ipynb index 991fb3795..785beedfa 100644 --- a/examples/notebooks/02_variations.ipynb +++ b/examples/notebooks/02_variations.ipynb @@ -104,6 +104,7 @@ " HarnessConfig,\n", " ModelConfig,\n", " RelayAtofConfig,\n", + " RelayAtofFileSinkConfig,\n", " RelayObservabilityConfig,\n", ")\n", "\n", @@ -361,8 +362,11 @@ " output_dir=str(relay_dir),\n", " observability=RelayObservabilityConfig(\n", " atof=RelayAtofConfig(\n", - " enabled=True, output_directory=str(relay_dir),\n", - " filename=\"events.atof.jsonl\", mode=\"overwrite\",\n", + " enabled=True,\n", + " sinks=[RelayAtofFileSinkConfig(\n", + " output_directory=str(relay_dir),\n", + " filename=\"events.atof.jsonl\", mode=\"overwrite\",\n", + " )],\n", " ),\n", " ),\n", " )\n", diff --git a/python/src/nemo_fabric/__init__.py b/python/src/nemo_fabric/__init__.py index df761ccf0..e9836d2b7 100644 --- a/python/src/nemo_fabric/__init__.py +++ b/python/src/nemo_fabric/__init__.py @@ -20,7 +20,8 @@ from nemo_fabric.models import ModelConfig from nemo_fabric.models import RelayAtifConfig from nemo_fabric.models import RelayAtofConfig -from nemo_fabric.models import RelayAtofEndpointConfig +from nemo_fabric.models import RelayAtofFileSinkConfig +from nemo_fabric.models import RelayAtofStreamSinkConfig from nemo_fabric.models import RelayComponentConfig from nemo_fabric.models import RelayConfig from nemo_fabric.models import RelayConfigPolicy @@ -72,7 +73,8 @@ "ModelConfig", "RelayAtifConfig", "RelayAtofConfig", - "RelayAtofEndpointConfig", + "RelayAtofFileSinkConfig", + "RelayAtofStreamSinkConfig", "RelayComponentConfig", "RelayConfigPolicy", "RelayHttpStorageConfig", diff --git a/python/src/nemo_fabric/integrations/harbor/fabric_agent.py b/python/src/nemo_fabric/integrations/harbor/fabric_agent.py index 7f78e70ce..6d21d9316 100644 --- a/python/src/nemo_fabric/integrations/harbor/fabric_agent.py +++ b/python/src/nemo_fabric/integrations/harbor/fabric_agent.py @@ -23,6 +23,7 @@ from nemo_fabric import ModelConfig from nemo_fabric import RelayAtifConfig from nemo_fabric import RelayAtofConfig +from nemo_fabric import RelayAtofFileSinkConfig from nemo_fabric import RelayObservabilityConfig from nemo_fabric import RunRequest from nemo_fabric import RunResult @@ -409,9 +410,13 @@ def build_harbor_config( ), atof=RelayAtofConfig( enabled=True, - output_directory=relay_output, - filename="events.atof.jsonl", - mode="overwrite", + sinks=[ + RelayAtofFileSinkConfig( + output_directory=relay_output, + filename="events.atof.jsonl", + mode="overwrite", + ) + ], ), ), ) diff --git a/python/src/nemo_fabric/models.py b/python/src/nemo_fabric/models.py index 3b7dd813d..b93d67143 100644 --- a/python/src/nemo_fabric/models.py +++ b/python/src/nemo_fabric/models.py @@ -240,24 +240,42 @@ class RelayConfigPolicy(FabricBaseModel): unsupported_value: Literal["ignore", "warn", "error"] = "error" -class RelayAtofEndpointConfig(FabricBaseModel): - """NeMo Relay ATOF remote endpoint configuration.""" +class RelayAtofFileSinkConfig(FabricBaseModel): + """NeMo Relay ATOF file sink configuration.""" + type: Literal["file"] = "file" + output_directory: str | Path | None = None + filename: str | None = None + mode: Literal["append", "overwrite"] = "append" + + +class RelayAtofStreamSinkConfig(FabricBaseModel): + """NeMo Relay ATOF stream sink configuration.""" + + type: Literal["stream"] = "stream" url: str transport: Literal["http_post", "websocket", "ndjson"] = "http_post" - headers: dict[str, str] = Field(default_factory=dict) + headers: dict[str, str] = Field(default_factory=dict, exclude_if=lambda value: not value) + header_env: dict[str, str] = Field(default_factory=dict, exclude_if=lambda value: not value) timeout_millis: int = 3000 field_name_policy: Literal["preserve", "replace_dots"] = "preserve" + name: str | None = None class RelayAtofConfig(FabricBaseModel): """NeMo Relay ATOF export configuration.""" enabled: bool = False - output_directory: str | Path | None = None - filename: str | None = None - mode: Literal["append", "overwrite"] = "append" - endpoints: list[RelayAtofEndpointConfig | dict[str, Any]] | None = None + sinks: ( + list[ + Annotated[ + RelayAtofFileSinkConfig | RelayAtofStreamSinkConfig, + Field(discriminator="type"), + ] + | dict[str, Any] + ] + | None + ) = None class RelayS3StorageConfig(FabricBaseModel): @@ -325,7 +343,7 @@ class RelayOtlpConfig(FabricBaseModel): class RelayObservabilityConfig(FabricBaseModel): """NeMo Relay observability component configuration.""" - version: int = 1 + version: int = 2 atof: RelayAtofConfig | dict[str, Any] | None = None atif: RelayAtifConfig | dict[str, Any] | None = None opentelemetry: RelayOtlpConfig | dict[str, Any] | None = None diff --git a/schemas/agent.schema.json b/schemas/agent.schema.json index 8e0be3c05..17b0bf863 100644 --- a/schemas/agent.schema.json +++ b/schemas/agent.schema.json @@ -410,75 +410,128 @@ "description": "Whether ATOF export is enabled.", "type": "boolean" }, - "endpoints": { - "description": "Optional remote ATOF endpoints.", + "sinks": { + "description": "ATOF file and stream sinks.", "items": { - "$ref": "#/$defs/RelayAtofEndpointConfig" + "$ref": "#/$defs/RelayAtofSinkConfig" }, "type": "array" - }, - "filename": { - "description": "ATOF file name.", - "type": [ - "string", - "null" - ] - }, - "mode": { - "$ref": "#/$defs/RelayAtofMode", - "default": "append", - "description": "File write mode." - }, - "output_directory": { - "description": "Directory used for ATOF files.", - "type": [ - "string", - "null" - ] } }, "type": "object" }, - "RelayAtofEndpointConfig": { - "additionalProperties": true, - "description": "Relay ATOF endpoint configuration.", - "properties": { - "field_name_policy": { - "$ref": "#/$defs/RelayAtofEndpointFieldNamePolicy", - "default": "preserve", - "description": "Field-name handling policy." + "RelayAtofMode": { + "description": "Relay ATOF file mode.", + "oneOf": [ + { + "const": "append", + "description": "Append to an existing ATOF file.", + "type": "string" }, - "headers": { - "additionalProperties": { - "type": "string" + { + "const": "overwrite", + "description": "Overwrite an existing ATOF file.", + "type": "string" + } + ] + }, + "RelayAtofSinkConfig": { + "description": "Relay ATOF sink configuration.", + "oneOf": [ + { + "additionalProperties": true, + "description": "Write ATOF records to a local file.", + "properties": { + "filename": { + "description": "ATOF file name.", + "type": [ + "string", + "null" + ] + }, + "mode": { + "$ref": "#/$defs/RelayAtofMode", + "default": "append", + "description": "File write mode." + }, + "output_directory": { + "description": "Directory used for ATOF files.", + "type": [ + "string", + "null" + ] + }, + "type": { + "const": "file", + "type": "string" + } }, - "description": "Endpoint headers.", + "required": [ + "type" + ], "type": "object" }, - "timeout_millis": { - "default": 3000, - "description": "Request timeout in milliseconds.", - "format": "uint64", - "minimum": 0, - "type": "integer" - }, - "transport": { - "$ref": "#/$defs/RelayAtofEndpointTransport", - "default": "http_post", - "description": "Endpoint transport." - }, - "url": { - "description": "Endpoint URL.", - "type": "string" + { + "additionalProperties": true, + "description": "Send ATOF records to a remote stream.", + "properties": { + "field_name_policy": { + "$ref": "#/$defs/RelayAtofStreamFieldNamePolicy", + "default": "preserve", + "description": "Field-name handling policy." + }, + "header_env": { + "additionalProperties": { + "type": "string" + }, + "description": "Environment-variable-backed stream headers.", + "type": "object" + }, + "headers": { + "additionalProperties": { + "type": "string" + }, + "description": "Static stream headers.", + "type": "object" + }, + "name": { + "description": "Optional stream sink name.", + "type": [ + "string", + "null" + ] + }, + "timeout_millis": { + "default": 3000, + "description": "Request timeout in milliseconds.", + "format": "uint64", + "minimum": 0, + "type": "integer" + }, + "transport": { + "$ref": "#/$defs/RelayAtofStreamTransport", + "default": "http_post", + "description": "Stream transport." + }, + "type": { + "const": "stream", + "type": "string" + }, + "url": { + "description": "Stream URL.", + "type": "string" + } + }, + "required": [ + "type", + "url" + ], + "type": "object" } - }, - "required": [ - "url" - ], - "type": "object" + ] }, - "RelayAtofEndpointFieldNamePolicy": { - "description": "Relay ATOF endpoint field-name policy.", + "RelayAtofStreamFieldNamePolicy": { + "description": "Relay ATOF stream field-name policy.", "oneOf": [ { "const": "preserve", @@ -492,8 +545,8 @@ } ] }, - "RelayAtofEndpointTransport": { - "description": "Relay ATOF endpoint transport.", + "RelayAtofStreamTransport": { + "description": "Relay ATOF stream transport.", "oneOf": [ { "const": "http_post", @@ -512,21 +565,6 @@ } ] }, - "RelayAtofMode": { - "description": "Relay ATOF file mode.", - "oneOf": [ - { - "const": "append", - "description": "Append to an existing ATOF file.", - "type": "string" - }, - { - "const": "overwrite", - "description": "Overwrite an existing ATOF file.", - "type": "string" - } - ] - }, "RelayComponentConfig": { "additionalProperties": true, "description": "Generic NeMo Relay plugin component configuration.", @@ -682,7 +720,7 @@ "description": "Relay config validation policy." }, "version": { - "default": 1, + "default": 2, "description": "Relay observability config version.", "format": "uint32", "minimum": 0, diff --git a/schemas/run-plan.schema.json b/schemas/run-plan.schema.json index db8bb4ea4..41ad885ef 100644 --- a/schemas/run-plan.schema.json +++ b/schemas/run-plan.schema.json @@ -938,75 +938,128 @@ "description": "Whether ATOF export is enabled.", "type": "boolean" }, - "endpoints": { - "description": "Optional remote ATOF endpoints.", + "sinks": { + "description": "ATOF file and stream sinks.", "items": { - "$ref": "#/$defs/RelayAtofEndpointConfig" + "$ref": "#/$defs/RelayAtofSinkConfig" }, "type": "array" - }, - "filename": { - "description": "ATOF file name.", - "type": [ - "string", - "null" - ] - }, - "mode": { - "$ref": "#/$defs/RelayAtofMode", - "default": "append", - "description": "File write mode." - }, - "output_directory": { - "description": "Directory used for ATOF files.", - "type": [ - "string", - "null" - ] } }, "type": "object" }, - "RelayAtofEndpointConfig": { - "additionalProperties": true, - "description": "Relay ATOF endpoint configuration.", - "properties": { - "field_name_policy": { - "$ref": "#/$defs/RelayAtofEndpointFieldNamePolicy", - "default": "preserve", - "description": "Field-name handling policy." + "RelayAtofMode": { + "description": "Relay ATOF file mode.", + "oneOf": [ + { + "const": "append", + "description": "Append to an existing ATOF file.", + "type": "string" }, - "headers": { - "additionalProperties": { - "type": "string" + { + "const": "overwrite", + "description": "Overwrite an existing ATOF file.", + "type": "string" + } + ] + }, + "RelayAtofSinkConfig": { + "description": "Relay ATOF sink configuration.", + "oneOf": [ + { + "additionalProperties": true, + "description": "Write ATOF records to a local file.", + "properties": { + "filename": { + "description": "ATOF file name.", + "type": [ + "string", + "null" + ] + }, + "mode": { + "$ref": "#/$defs/RelayAtofMode", + "default": "append", + "description": "File write mode." + }, + "output_directory": { + "description": "Directory used for ATOF files.", + "type": [ + "string", + "null" + ] + }, + "type": { + "const": "file", + "type": "string" + } }, - "description": "Endpoint headers.", + "required": [ + "type" + ], "type": "object" }, - "timeout_millis": { - "default": 3000, - "description": "Request timeout in milliseconds.", - "format": "uint64", - "minimum": 0, - "type": "integer" - }, - "transport": { - "$ref": "#/$defs/RelayAtofEndpointTransport", - "default": "http_post", - "description": "Endpoint transport." - }, - "url": { - "description": "Endpoint URL.", - "type": "string" + { + "additionalProperties": true, + "description": "Send ATOF records to a remote stream.", + "properties": { + "field_name_policy": { + "$ref": "#/$defs/RelayAtofStreamFieldNamePolicy", + "default": "preserve", + "description": "Field-name handling policy." + }, + "header_env": { + "additionalProperties": { + "type": "string" + }, + "description": "Environment-variable-backed stream headers.", + "type": "object" + }, + "headers": { + "additionalProperties": { + "type": "string" + }, + "description": "Static stream headers.", + "type": "object" + }, + "name": { + "description": "Optional stream sink name.", + "type": [ + "string", + "null" + ] + }, + "timeout_millis": { + "default": 3000, + "description": "Request timeout in milliseconds.", + "format": "uint64", + "minimum": 0, + "type": "integer" + }, + "transport": { + "$ref": "#/$defs/RelayAtofStreamTransport", + "default": "http_post", + "description": "Stream transport." + }, + "type": { + "const": "stream", + "type": "string" + }, + "url": { + "description": "Stream URL.", + "type": "string" + } + }, + "required": [ + "type", + "url" + ], + "type": "object" } - }, - "required": [ - "url" - ], - "type": "object" + ] }, - "RelayAtofEndpointFieldNamePolicy": { - "description": "Relay ATOF endpoint field-name policy.", + "RelayAtofStreamFieldNamePolicy": { + "description": "Relay ATOF stream field-name policy.", "oneOf": [ { "const": "preserve", @@ -1020,8 +1073,8 @@ } ] }, - "RelayAtofEndpointTransport": { - "description": "Relay ATOF endpoint transport.", + "RelayAtofStreamTransport": { + "description": "Relay ATOF stream transport.", "oneOf": [ { "const": "http_post", @@ -1040,21 +1093,6 @@ } ] }, - "RelayAtofMode": { - "description": "Relay ATOF file mode.", - "oneOf": [ - { - "const": "append", - "description": "Append to an existing ATOF file.", - "type": "string" - }, - { - "const": "overwrite", - "description": "Overwrite an existing ATOF file.", - "type": "string" - } - ] - }, "RelayComponentConfig": { "additionalProperties": true, "description": "Generic NeMo Relay plugin component configuration.", @@ -1210,7 +1248,7 @@ "description": "Relay config validation policy." }, "version": { - "default": 1, + "default": 2, "description": "Relay observability config version.", "format": "uint32", "minimum": 0, diff --git a/skills/nemo-fabric-integrate/references/config-mapping.md b/skills/nemo-fabric-integrate/references/config-mapping.md index 98d9a3596..9a93b6784 100644 --- a/skills/nemo-fabric-integrate/references/config-mapping.md +++ b/skills/nemo-fabric-integrate/references/config-mapping.md @@ -65,6 +65,10 @@ def with_relay(base: FabricConfig) -> FabricConfig: Use this function-and-copy pattern for every variant; keep all variation in ordinary Python. +For ATOF, author the Relay 0.6 sink model directly. Put +`RelayAtofFileSinkConfig` and `RelayAtofStreamSinkConfig` instances in +`RelayAtofConfig.sinks`, and set `RelayAtofConfig.enabled=True`. + ## Relative Paths If the config uses relative paths for skills, workspaces, or artifacts, pass diff --git a/tests/adapters/test_adapaters_common_utils.py b/tests/adapters/test_adapaters_common_utils.py index 7b322722c..2a8c498aa 100644 --- a/tests/adapters/test_adapaters_common_utils.py +++ b/tests/adapters/test_adapaters_common_utils.py @@ -2,7 +2,6 @@ # SPDX-License-Identifier: Apache-2.0 import builtins -import dataclasses import json import os import sys @@ -272,10 +271,6 @@ def test_normalize_list(value: object, expected: list[str]): assert common_utils.normalize_list(value) == expected -def test_without_none(): - assert common_utils.without_none({"a": 1, "b": None, "c": False}) == {"a": 1, "c": False} - - def test_load_relay_plugin_config_wraps_and_normalizes_bare_observability_config( tmp_path: Path, ): @@ -287,7 +282,16 @@ def test_load_relay_plugin_config_wraps_and_normalizes_bare_observability_config "config": { "atof": { "enabled": True, - "output_directory": "custom-relay", + "sinks": [ + { + "type": "file", + "output_directory": "custom-relay", + }, + { + "type": "stream", + "url": "https://example.test/events", + }, + ], }, "atif": {"enabled": True}, } @@ -320,25 +324,24 @@ def test_load_relay_plugin_config_wraps_and_normalizes_bare_observability_config assert plugin_config["version"] == 1 assert plugin_config["components"][0]["kind"] == "observability" - assert observability["atof"]["output_directory"] == str( - tmp_path / "custom-relay" / "runtime-current" - ) - assert observability["atof"]["filename"] == "events.atof.jsonl" - assert observability["atof"]["mode"] == "overwrite" - assert Path(observability["atof"]["output_directory"]).is_dir() - assert observability["atif"]["output_directory"] == str( - tmp_path / "artifacts" / "relay" / "runtime-current" - ) + assert observability["version"] == 2 + file_sink, stream_sink = observability["atof"]["sinks"] + assert file_sink["output_directory"] == str(tmp_path / "custom-relay" / "runtime-current") + assert file_sink["filename"] == "events.atof.jsonl" + assert file_sink["mode"] == "overwrite" + assert Path(file_sink["output_directory"]).is_dir() + assert stream_sink == { + "type": "stream", + "url": "https://example.test/events", + } + assert observability["atif"]["output_directory"] == str(tmp_path / "artifacts" / "relay" / "runtime-current") assert observability["atif"]["filename_template"] == "trajectory-{session_id}.atif.json" assert observability["atif"]["agent_name"] == "review-agent" assert observability["atif"]["model_name"] == "nvidia/review-model" assert Path(observability["atif"]["output_directory"]).is_dir() - atof_file = Path(observability["atof"]["output_directory"]) / "events.atof.jsonl" - atif_file = ( - Path(observability["atif"]["output_directory"]) - / "trajectory-current.atif.json" - ) + atof_file = Path(file_sink["output_directory"]) / "events.atof.jsonl" + atif_file = Path(observability["atif"]["output_directory"]) / "trajectory-current.atif.json" atof_file.write_text("{}", encoding="utf-8") atif_file.write_text("{}", encoding="utf-8") @@ -364,7 +367,19 @@ def test_collect_relay_artifacts(tmp_path: Path): { "kind": "observability", "config": { - "atof": {"enabled": True, "output_directory": str(atof_dir)}, + "atof": { + "enabled": True, + "sinks": [ + { + "type": "file", + "output_directory": str(atof_dir), + }, + { + "type": "stream", + "url": "https://example.test/events", + }, + ], + }, "atif": {"enabled": True, "output_directory": str(atif_dir)}, }, } @@ -377,7 +392,34 @@ def test_collect_relay_artifacts(tmp_path: Path): ] -def test_relay_api_plugin_config_translates_flat_atof_to_relay_v06_sinks(): +def test_collect_relay_artifacts_ignores_missing_output_directories( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +): + (tmp_path / "unrelated.jsonl").write_text("{}", encoding="utf-8") + (tmp_path / "unrelated.json").write_text("{}", encoding="utf-8") + monkeypatch.chdir(tmp_path) + plugin_config = { + "components": [ + { + "kind": "observability", + "config": { + "atof": { + "enabled": True, + "sinks": [{"type": "file"}], + }, + "atif": {"enabled": True}, + }, + } + ] + } + + assert common_utils.collect_relay_artifacts(plugin_config) == [] + + +def test_relay_validates_raw_v06_plugin_config(): + from nemo_relay import plugin + os.environ["TOKEN"] = "test-token" plugin_config = { "version": 1, @@ -386,21 +428,26 @@ def test_relay_api_plugin_config_translates_flat_atof_to_relay_v06_sinks(): "kind": "observability", "enabled": True, "config": { - "version": 1, + "version": 2, "atof": { "enabled": True, - "output_directory": "/tmp/atof", - "filename": "events.jsonl", - "mode": "overwrite", - "endpoints": [ + "sinks": [ { + "type": "file", + "output_directory": "/tmp/atof", + "filename": "events.jsonl", + "mode": "overwrite", + }, + { + "type": "stream", "url": "https://example.test/events", + "transport": "ndjson", "headers": {"x-test": "value"}, "header_env": {"authorization": "TOKEN"}, "timeout_millis": 1000, "field_name_policy": "replace_dots", "name": "phoenix", - } + }, ], }, }, @@ -408,30 +455,41 @@ def test_relay_api_plugin_config_translates_flat_atof_to_relay_v06_sinks(): ], } - rendered = common_utils.relay_api_plugin_config(plugin_config) - observability = rendered.components[0].config + assert plugin.validate(plugin_config)["diagnostics"] == [] - assert observability.version == 2 - assert dataclasses.asdict(observability.atof) == { - "enabled": True, - "sinks": [ - { - "output_directory": "/tmp/atof", - "filename": "events.jsonl", - "mode": "overwrite", - }, + +def test_relay_validates_unknown_atof_sink_type(): + from nemo_relay import plugin + + plugin_config = { + "version": 1, + "components": [ { - "url": "https://example.test/events", - "transport": "http_post", - "headers": {"x-test": "value"}, - "header_env": {"authorization": "TOKEN"}, - "timeout_millis": 1000, - "field_name_policy": "replace_dots", - "name": "phoenix", - }, + "kind": "observability", + "enabled": True, + "config": { + "version": 2, + "atof": { + "enabled": True, + "sinks": [{"type": "unknown"}], + }, + }, + } ], } + assert plugin.validate(plugin_config)["diagnostics"] == [ + { + "code": "observability.invalid_plugin_config", + "component": "observability", + "level": "error", + "message": ( + "invalid config: invalid observability plugin config: " + "unknown variant `unknown`, expected `file` or `stream`" + ), + } + ] + @pytest.mark.parametrize( ("relay_config", "plugin_config", "expected_names"), @@ -466,7 +524,7 @@ def test_write_relay_configs( assert tomllib.load(stream) == config -def test_write_relay_configs_migrates_atof_to_current_cli_contract(tmp_path: Path): +def test_write_relay_configs_preserves_current_cli_contract(tmp_path: Path): os.environ["FABRIC_RELAY_CONFIG_PATH"] = str(tmp_path / "relay.json") plugin_config = { "version": 1, @@ -475,21 +533,25 @@ def test_write_relay_configs_migrates_atof_to_current_cli_contract(tmp_path: Pat "kind": "observability", "enabled": True, "config": { - "version": 1, + "version": 2, "atof": { "enabled": True, - "output_directory": "/tmp/atof", - "filename": "events.jsonl", - "mode": "overwrite", - "endpoints": [ + "sinks": [ + { + "type": "file", + "output_directory": "/tmp/atof", + "filename": "events.jsonl", + "mode": "overwrite", + }, { + "type": "stream", "url": "https://example.test/events", "transport": "http_post", "headers": {"x-test": "value"}, "header_env": {"authorization": "TOKEN"}, "timeout_millis": 1000, "field_name_policy": "replace_dots", - } + }, ], }, "atif": {"enabled": True, "output_directory": "/tmp/atif"}, @@ -532,5 +594,4 @@ def test_write_relay_configs_migrates_atof_to_current_cli_contract(tmp_path: Pat "enabled": True, "output_directory": "/tmp/atif", } - assert plugin_config["components"][0]["config"]["version"] == 1 - assert "sinks" not in plugin_config["components"][0]["config"]["atof"] + assert rendered == plugin_config diff --git a/tests/adapters/test_deepagents.py b/tests/adapters/test_deepagents.py index e0373fe3a..d718050b4 100644 --- a/tests/adapters/test_deepagents.py +++ b/tests/adapters/test_deepagents.py @@ -15,6 +15,7 @@ import os import sys import types +from collections.abc import AsyncIterator from collections.abc import Iterator from pathlib import Path from typing import Any @@ -187,9 +188,10 @@ def add_nemo_relay_integration(kwargs, **_): return merged @contextlib.asynccontextmanager - async def plugin_ctx(_config): + async def plugin_ctx(config: object) -> AsyncIterator[None]: calls["plugin_open"] = True calls["plugin_enters"] = calls.get("plugin_enters", 0) + 1 + calls.setdefault("plugin_configs", []).append(config) try: yield finally: @@ -331,13 +333,11 @@ async def test_relay_telemetry_wraps_agent_and_reports_artifacts( tmp_path, make_payload, monkeypatch, fake_sdks, fake_relay ): artifacts = [{"kind": "atof", "path": str(tmp_path / "events.atof.jsonl")}] + plugin_config = {"version": 1, "components": []} monkeypatch.setattr( adapter.common_utils, "load_relay_plugin_config", - lambda _p: {"version": 1, "components": []}, - ) - monkeypatch.setattr( - adapter.common_utils, "relay_api_plugin_config", lambda _c: object() + lambda _p: plugin_config, ) monkeypatch.setattr( adapter.common_utils, "collect_relay_artifacts", lambda _c: artifacts @@ -357,6 +357,7 @@ async def test_relay_telemetry_wraps_agent_and_reports_artifacts( assert fake_relay["wrapped"] assert fake_relay["plugin_open"] + assert fake_relay["plugin_configs"] == [plugin_config] assert output["telemetry"] == { "enabled": True, "provider": "relay", @@ -378,10 +379,6 @@ async def test_relay_telemetry_wraps_agent_and_reports_artifacts( async def test_native_telemetry_exports_without_artifacts( tmp_path, make_payload, monkeypatch, fake_sdks, fake_relay ): - monkeypatch.setattr( - adapter.common_utils, "relay_api_plugin_config", lambda _c: object() - ) - payload = make_payload(tmp_path) payload["telemetry_plan"] = { "providers": ["native"], @@ -412,6 +409,9 @@ async def test_native_telemetry_exports_without_artifacts( assert fake_relay["wrapped"] assert fake_relay["plugin_open"] + assert fake_relay["plugin_configs"] == [ + payload["telemetry_plan"]["native_config"] + ] assert output["telemetry"] == { "enabled": True, "provider": "native", @@ -794,13 +794,11 @@ async def test_persistent_runtime_scopes_relay_per_invocation( tmp_path, make_payload, monkeypatch, fake_sdks, fake_relay ): artifacts = [{"kind": "atif", "path": str(tmp_path / "trajectory.json")}] + plugin_config = {"version": 1, "components": []} monkeypatch.setattr( adapter.common_utils, "load_relay_plugin_config", - lambda _payload: {"version": 1, "components": []}, - ) - monkeypatch.setattr( - adapter.common_utils, "relay_api_plugin_config", lambda _config: object() + lambda _payload: plugin_config, ) monkeypatch.setattr( adapter.common_utils, @@ -829,6 +827,7 @@ async def test_persistent_runtime_scopes_relay_per_invocation( assert fake_relay["integration_adds"] == 1 assert fake_relay["plugin_enters"] == 2 assert fake_relay["plugin_exits"] == 2 + assert fake_relay["plugin_configs"] == [plugin_config, plugin_config] assert fake_relay["scopes"] == [ ("deepagents-request", "agent"), ("deepagents-request", "agent"), diff --git a/tests/adapters/test_hermes_adapter.py b/tests/adapters/test_hermes_adapter.py index 710b4af9c..5c1186aac 100644 --- a/tests/adapters/test_hermes_adapter.py +++ b/tests/adapters/test_hermes_adapter.py @@ -288,7 +288,15 @@ def test_hermes_config_variation_matrix_surfaces_supported_capabilities( { "relay": { "config": { - "atof": {"enabled": True, "output_directory": "relay/atof"}, + "atof": { + "enabled": True, + "sinks": [ + { + "type": "file", + "output_directory": "relay/atof", + } + ], + }, "atif": {"enabled": True, "output_directory": "relay/atif"}, } } @@ -369,7 +377,7 @@ def test_hermes_config_variation_matrix_surfaces_supported_capabilities( } assert config["platform_toolsets"] == {"cli": ["git", "shell"]} assert config["plugins"]["enabled"] == ["observability/nemo_relay"] - assert observability["atof"]["output_directory"] == str( + assert observability["atof"]["sinks"][0]["output_directory"] == str( tmp_path / "relay" / "atof" / "runtime-matrix" ) assert observability["atif"]["output_directory"] == str( diff --git a/tests/e2e/test_claude.py b/tests/e2e/test_claude.py index 5ec382ab8..986705b47 100644 --- a/tests/e2e/test_claude.py +++ b/tests/e2e/test_claude.py @@ -21,6 +21,7 @@ ModelConfig, RelayAtifConfig, RelayAtofConfig, + RelayAtofFileSinkConfig, RelayObservabilityConfig, RuntimeConfig, ) @@ -123,7 +124,10 @@ def fabric_config( if relay: config.enable_relay( observability=RelayObservabilityConfig( - atof=RelayAtofConfig(enabled=True), + atof=RelayAtofConfig( + enabled=True, + sinks=[RelayAtofFileSinkConfig()], + ), atif=RelayAtifConfig(enabled=True), ) ) diff --git a/tests/python/test_sdk_contract.py b/tests/python/test_sdk_contract.py index 2b44272c5..e37ee1487 100644 --- a/tests/python/test_sdk_contract.py +++ b/tests/python/test_sdk_contract.py @@ -30,6 +30,8 @@ from nemo_fabric import MetadataConfig from nemo_fabric import RelayAtifConfig from nemo_fabric import RelayAtofConfig +from nemo_fabric import RelayAtofFileSinkConfig +from nemo_fabric import RelayAtofStreamSinkConfig from nemo_fabric import RelayComponentConfig from nemo_fabric import RelayConfigPolicy from nemo_fabric import RelayObservabilityConfig @@ -220,9 +222,19 @@ def test_fabric_config_authors_first_class_relay_observability(): observability=RelayObservabilityConfig( atof=RelayAtofConfig( enabled=True, - output_directory="./artifacts/relay", - filename="events.atof.jsonl", - mode="overwrite", + sinks=[ + RelayAtofFileSinkConfig( + output_directory="./artifacts/relay", + filename="events.atof.jsonl", + mode="overwrite", + ), + RelayAtofStreamSinkConfig( + url="http://localhost:4319/events", + transport="ndjson", + header_env={"authorization": "RELAY_AUTHORIZATION"}, + name="live-events", + ), + ], ), atif=RelayAtifConfig( enabled=True, @@ -243,12 +255,26 @@ def test_fabric_config_authors_first_class_relay_observability(): assert config.to_mapping()["relay"] == { "output_dir": "./artifacts/relay", "observability": { - "version": 1, + "version": 2, "atof": { "enabled": True, - "output_directory": "./artifacts/relay", - "filename": "events.atof.jsonl", - "mode": "overwrite", + "sinks": [ + { + "type": "file", + "output_directory": "./artifacts/relay", + "filename": "events.atof.jsonl", + "mode": "overwrite", + }, + { + "type": "stream", + "url": "http://localhost:4319/events", + "transport": "ndjson", + "header_env": {"authorization": "RELAY_AUTHORIZATION"}, + "timeout_millis": 3000, + "field_name_policy": "preserve", + "name": "live-events", + }, + ], }, "atif": { "enabled": True, @@ -273,6 +299,18 @@ def test_fabric_config_authors_first_class_relay_observability(): } +def test_relay_atof_stream_sink_omits_empty_header_maps(): + sink = RelayAtofStreamSinkConfig(url="https://example.test/events") + + assert sink.to_mapping() == { + "type": "stream", + "url": "https://example.test/events", + "transport": "http_post", + "timeout_millis": 3000, + "field_name_policy": "preserve", + } + + def test_fabric_config_enable_relay_preserves_omitted_fields(): config = _fabric_config()