diff --git a/config/config.md b/config/config.md
index 91cb4fe766..992296d482 100644
--- a/config/config.md
+++ b/config/config.md
@@ -76,6 +76,7 @@
| `otlp.enable` | Bool | `true` | Whether to enable OpenTelemetry protocol in HTTP API. |
| `otlp.experimental_enable_exponential_histogram` | Bool | `false` | Experimental: enable cumulative OTLP exponential histogram ingestion. |
| `otlp.trace_ingest_chunk_size` | Integer | `512` | Maximum spans per trace ingest chunk. Set to 0 to disable splitting. |
+| `otlp.experimental_enable_resource_info` | Bool | `false` | Whether to synthesize the `greptime_otel_resource_info` table from OTLP metric
resource attributes, so metrics-only services reach the semantic graph. |
| `prom_store` | -- | -- | Prometheus remote storage options |
| `prom_store.enable` | Bool | `true` | Whether to enable Prometheus remote write and read in HTTP API. |
| `prom_store.with_metric_engine` | Bool | `true` | Whether to store the data from Prometheus remote write in metric engine. |
@@ -319,6 +320,7 @@
| `otlp.enable` | Bool | `true` | Whether to enable OpenTelemetry protocol in HTTP API. |
| `otlp.experimental_enable_exponential_histogram` | Bool | `false` | Experimental: enable cumulative OTLP exponential histogram ingestion. |
| `otlp.trace_ingest_chunk_size` | Integer | `512` | Maximum spans per trace ingest chunk. Set to 0 to disable splitting. |
+| `otlp.experimental_enable_resource_info` | Bool | `false` | Whether to synthesize the `greptime_otel_resource_info` table from OTLP metric
resource attributes, so metrics-only services reach the semantic graph. |
| `prom_store` | -- | -- | Prometheus remote storage options |
| `prom_store.enable` | Bool | `true` | Whether to enable Prometheus remote write and read in HTTP API. |
| `prom_store.with_metric_engine` | Bool | `true` | Whether to store the data from Prometheus remote write in metric engine. |
diff --git a/config/frontend.example.toml b/config/frontend.example.toml
index 3635a63f06..8053485572 100644
--- a/config/frontend.example.toml
+++ b/config/frontend.example.toml
@@ -241,6 +241,9 @@ enable = true
experimental_enable_exponential_histogram = false
## Maximum spans per trace ingest chunk. Set to 0 to disable splitting.
trace_ingest_chunk_size = 512
+## Whether to synthesize the `greptime_otel_resource_info` table from OTLP metric
+## resource attributes, so metrics-only services reach the semantic graph.
+experimental_enable_resource_info = false
## Prometheus remote storage options
[prom_store]
diff --git a/config/standalone.example.toml b/config/standalone.example.toml
index 72ba992cce..70f74c7f94 100644
--- a/config/standalone.example.toml
+++ b/config/standalone.example.toml
@@ -208,6 +208,9 @@ enable = true
experimental_enable_exponential_histogram = false
## Maximum spans per trace ingest chunk. Set to 0 to disable splitting.
trace_ingest_chunk_size = 512
+## Whether to synthesize the `greptime_otel_resource_info` table from OTLP metric
+## resource attributes, so metrics-only services reach the semantic graph.
+experimental_enable_resource_info = false
## Prometheus remote storage options
[prom_store]
diff --git a/docs/rfcs/2026-06-25-entity-relationships-and-graph-query.md b/docs/rfcs/2026-06-25-entity-relationships-and-graph-query.md
index 995cf6cab7..1c78e91f25 100644
--- a/docs/rfcs/2026-06-25-entity-relationships-and-graph-query.md
+++ b/docs/rfcs/2026-06-25-entity-relationships-and-graph-query.md
@@ -149,7 +149,7 @@ The decisions behind the shape:
- **`provenance` is part of edge identity**, so a declared edge and a derived edge between the same pair coexist, and a user-declared edge survives when derivation is off.
- **`confidence` expresses derivation certainty, not statistical completeness**: `1.0` for a successfully paired or declared edge, `< 1.0` for virtual-node or agent-inferred edges. It does not correct for trace sampling.
- **RED metrics describe the observed span-pair population** — the `service_graph` connector semantics. Under sampling, counts understate true traffic, and ratios (error rate, mean duration) are representative only when sampling is unbiased with respect to status and latency — tail sampling that keeps errors and slow traces skews both.
-- **Endpoint encoding (v1).** A single-attribute identity is the value verbatim; a composite identity is the attributes sorted by key, rendered `k1=v1,k2=v2` — a human-readable traversal key, not a collision-free canonical encoding. Nodes carry the structured `entity_id_attrs`; edges carry only the string key. Known v1 limitation: composite components containing unescaped `,` or `=` are not guaranteed to be distinguished.
+- **Endpoint encoding.** An identity renders as its values in declared order (broad to narrow), joined by `,` with `,` and `\` escaped in each value. The identifying attribute *keys* are deliberately absent: the same identity reaches the graph under different column names per signal — a trace table's `service_name` against a metric table's `job` — and encoding the keys would split one entity into one entity per signal. Where a source carries a qualifier separately (`service.namespace` beside `service.name`) the convention folds it in as `/`, reproducing the identity a source that pre-composed it (Prometheus `job`) already emits. Nodes carry the structured `entity_id_attrs`; edges carry only the string key.
- **Hand-declared edges** (`provenance = 'declared'`) live in a physical table of the same shape (plus a business validity window) in `greptime_private` and are unioned into the computed `semantic_relationships`.
The relationship vocabulary is small, typed, and inverse-paired; stored direction is `src -> dst`, the inverse is a query concern:
@@ -197,7 +197,7 @@ GROUP BY 1, src_id, dst_id;
Edge endpoints are built from each trace table's `service` entity declaration (the same id expression the registry uses), so edges land on exactly the entity ids the registry emits, self-calls are judged on the full endpoint id, and a table whose declaration is unusable contributes no edges rather than silently switching identity. The join is bounded by a time-proximity window and the results of **all** trace tables combine into one edge per window with merged RED metrics. One assumption to note: a deployment can route traces to multiple tables (`x-greptime-trace-table-name`), so one distributed trace is not guaranteed to live in a single table — pairing must account for that, and how is an implementation concern. Uninstrumented peers become **virtual nodes**: a client span with no matching server span yields an edge to a synthetic node named from `peer.service`/`db.name`/`server.address`.
-**The entity registry.** Per declaring table: filter to the window, project (time bin, casts, JSON assembly), `DISTINCT`. The per-table results are unioned without any cross-source merge: no descriptive-conflict or JSON-merging policy lives in the derivation. Rows whose identity columns are NULL identify nothing and are filtered out; single-column ids are used verbatim, composite ids render the sorted `k=v` form plus the `entity_id_attrs` JSON. Consumers deduplicate with `SELECT DISTINCT entity_type, entity_id`; a unique node set is the snapshot relation's job (see *Querying the graph*).
+**The entity registry.** Per declaring table: filter to the window, project (time bin, casts, JSON assembly), `DISTINCT`. The per-table results are unioned without any cross-source merge: no descriptive-conflict or JSON-merging policy lives in the derivation. Rows whose identity columns are NULL identify nothing and are filtered out; ids render as their escaped values joined in declared order, with the `entity_id_attrs` JSON alongside once an id is assembled from more than one column. Consumers deduplicate with `SELECT DISTINCT entity_type, entity_id`; a unique node set is the snapshot relation's job (see *Querying the graph*).
**Attribute edges.** Shared attributes provide join keys; co-declaration or a relationship template provides the relationship semantics. A shared value alone determines neither direction nor edge type, so the derivation is rule-based, not join-everything. A table declaring both a `service.instance` and a `host` identity on the same rows derives `runs_on` between them: the row itself witnesses the relationship, and the built-in vocabulary fixes the direction. Cross-table edges require a declared relationship template (endpoint mappings, `rel_type`, direction): a derivation *rule* declared on schema, not per-instance edge data.
@@ -284,7 +284,7 @@ Drawbacks:
- Read-time derivation costs a self-join over trace tables on every scan of the computed tables. The default one-hour window bounds it, but very large deployments will want materialisation (Future Work).
- `semantic_entities` exposes per-observation rows; deduplication falls on the consumer (`DISTINCT`) until the snapshot relation lands in M2.
-- The v1 `k=v` id encoding can collide on hostile values (unescaped `,`/`=`); hardening is an open question.
+- The id encoding is readable rather than canonical: escaping keeps components distinguishable, but a value containing the qualifier separator (`/`) is still ambiguous against a qualified name.
- The graph is only as connected as the identity columns tables actually share: the `job`/`service_name` misalignment above is the canonical example, and enrichment discipline gates connectivity.
# Implementation Plan
@@ -306,7 +306,7 @@ Drawbacks:
# Open Questions
-1. **Endpoint encoding hardening.** The v1 `k=v` rendering (with `entity_id_attrs` JSON as the source of truth) is readable but collides when values contain `,`/`=`. Open: a fixed-width `entity_id_hash` join key (New-Relic-style); typed value encoding; the criteria for introducing the `semantic_key=column` mapping form; `schema_url`-gated entity merge. Related and unresolved: whether the identifying attribute *key* should participate in entity equality — a single-attribute id keeps only the value, so `job=api` and `service_name=api` currently name the same entity even when the keys differ semantically.
+1. **Endpoint encoding hardening.** The rendering (with `entity_id_attrs` JSON as the source of truth) is readable rather than canonical. Open: a fixed-width `entity_id_hash` join key (New-Relic-style); typed value encoding; the criteria for introducing the `semantic_key=column` mapping form; `schema_url`-gated entity merge. Settled: identifying attribute *keys* do not participate in entity equality, so `job=api` and `service_name=api` name the same entity — the property that lets one service reach the graph from several signals.
2. **Identity reconciliation across signals.** Trace edges key on the trace table's service declaration; host edges key on `host.id`/`k8s.pod.uid`. Linking them relies on tables carrying both columns; how aggressively to enrich (a `k8sattributes`-equivalent) gates how connected the graph is.
3. **`GRAPH_TABLE` surface scope.** How much of SQL/PGQ to implement vs. a minimal single-hop subset. Start minimal and grow with demand.
4. **Visibility beyond table grants.** The derivation contract filters by table-read permission; whether topology needs finer policies (row-level, tenant-scoped visibility of cross-schema edges) is open.
diff --git a/src/common/catalog/src/consts.rs b/src/common/catalog/src/consts.rs
index b976cbc256..63aa103b83 100644
--- a/src/common/catalog/src/consts.rs
+++ b/src/common/catalog/src/consts.rs
@@ -220,6 +220,11 @@ pub const SEMANTIC_RELATIONSHIPS_TABLE_NAME: &str = "semantic_relationships";
/// `semantic_relationships`; see [`is_ddl_reserved_table`] for its lifecycle.
pub const SEMANTIC_RELATIONSHIPS_DECLARED_TABLE_NAME: &str = "semantic_relationships_declared";
+/// Width of the window the graph derivation bins observations into. Sources
+/// synthesizing observations must land a row in every window they describe,
+/// so this is shared rather than restated per crate.
+pub const SEMANTIC_GRAPH_WINDOW_NANOS: i64 = 60 * 1_000_000_000;
+
// Column names of the graph tables: `catalog` exposes these schemas and the
// read-time plans in `operator` must project exactly them.
pub const OBSERVED_AT_COLUMN: &str = "observed_at";
diff --git a/src/frontend/src/instance.rs b/src/frontend/src/instance.rs
index 78b701a344..412ff495f6 100644
--- a/src/frontend/src/instance.rs
+++ b/src/frontend/src/instance.rs
@@ -134,6 +134,7 @@ pub struct Instance {
slow_query_options: SlowQueryOptions,
influxdb_default_merge_mode: InfluxdbMergeMode,
trace_ingest_chunk_size: usize,
+ otlp_resource_info: bool,
suspend: Arc,
// cache for otlp metrics
diff --git a/src/frontend/src/instance/builder.rs b/src/frontend/src/instance/builder.rs
index c49a8394be..0f87dc424f 100644
--- a/src/frontend/src/instance/builder.rs
+++ b/src/frontend/src/instance/builder.rs
@@ -367,6 +367,7 @@ impl FrontendBuilder {
slow_query_options: self.options.slow_query.clone(),
influxdb_default_merge_mode: self.options.influxdb.default_merge_mode,
trace_ingest_chunk_size: self.options.otlp.trace_ingest_chunk_size,
+ otlp_resource_info: self.options.otlp.experimental_enable_resource_info,
suspend: Arc::new(AtomicBool::new(false)),
})
}
diff --git a/src/frontend/src/instance/entity_graph.rs b/src/frontend/src/instance/entity_graph.rs
index 14f18ed1b7..a0817cf280 100644
--- a/src/frontend/src/instance/entity_graph.rs
+++ b/src/frontend/src/instance/entity_graph.rs
@@ -23,7 +23,7 @@
//! the query engine. Injected into the catalog manager after the engine is built,
//! breaking the `catalog -> query` cycle.
-use std::collections::HashMap;
+use std::collections::{BTreeMap, HashMap};
use std::sync::Weak;
use async_trait::async_trait;
@@ -60,8 +60,9 @@ use table::TableRef;
use table::metadata::TableInfo;
use table::predicate::{TimeRangeExtraction, extract_time_range_strict};
use table::requests::{
- EntityRole, SEMANTIC_SIGNAL_TYPE, SEMANTIC_SOURCE, SIGNAL_TYPE_METRIC, SOURCE_PROMETHEUS,
- is_trace_v1_table, parse_entity_columns, parse_entity_option_key,
+ EntityRole, SEMANTIC_METRIC_TYPE, SEMANTIC_SIGNAL_TYPE, SEMANTIC_SOURCE, SIGNAL_TYPE_METRIC,
+ SOURCE_OPENTELEMETRY, SOURCE_PROMETHEUS, is_trace_v1_table, parse_entity_columns,
+ parse_entity_option_key,
};
use crate::error;
@@ -177,6 +178,7 @@ impl EntityGraphProviderImpl {
time_index: time_index.clone(),
entity_type,
id_columns,
+ id_qualifier: None,
descriptive_columns,
scope_columns,
})
@@ -188,10 +190,11 @@ impl EntityGraphProviderImpl {
/// All entity declarations of one table: the explicit options plus the
/// zero-configuration conventions (`otlp_trace_entities` for trace-v1
/// tables — including the `service` identity of tables created before the
- /// ingest-side auto-stamp — and the Prometheus descriptor whitelist). An
- /// explicit declaration of a type always suppresses the implicit one, even
- /// when the explicit declaration is invalid and skipped: silently falling
- /// back would change entity identity behind the user's back.
+ /// ingest-side auto-stamp — and the Prometheus/OTel descriptor
+ /// whitelists). An explicit declaration of a type always suppresses the
+ /// implicit one, even when the explicit declaration is invalid and
+ /// skipped: silently falling back would change entity identity behind the
+ /// user's back.
fn declarations_for(
table_info: &TableInfo,
conventions: &Conventions,
@@ -204,7 +207,20 @@ impl EntityGraphProviderImpl {
&mut declarations,
);
}
- Self::extend_with_prometheus_conventions(table_info, conventions, &mut declarations);
+ Self::extend_with_info_metric_conventions(
+ table_info,
+ &conventions.prometheus_info_metrics,
+ SOURCE_PROMETHEUS,
+ None,
+ &mut declarations,
+ );
+ Self::extend_with_info_metric_conventions(
+ table_info,
+ &conventions.otel_info_metrics,
+ SOURCE_OPENTELEMETRY,
+ Some(servers::semantic::METRIC_TYPE_INFO),
+ &mut declarations,
+ );
declarations
}
@@ -216,29 +232,34 @@ impl EntityGraphProviderImpl {
})
}
- /// Implicit declarations of the well-known Prometheus entity-descriptor
- /// metrics (the `prometheus_info_metrics` whitelist of `conventions.yaml`),
- /// gated on the ingest-stamped `signal_type=metric` + `source=prometheus`
- /// options and keyed by table name. The metric engine's physical table
+ /// Implicit declarations of the well-known entity-descriptor metrics
+ /// (the `prometheus_info_metrics` / `otel_info_metrics` whitelists of
+ /// `conventions.yaml`), gated on the ingest-stamped `signal_type=metric`
+ /// option plus the whitelist's expected `source`; OTel descriptors also
+ /// require `metric.type=info`. The metric engine's physical table
/// aggregates every logical table's columns and must not contribute a
/// duplicate source.
- fn extend_with_prometheus_conventions(
+ fn extend_with_info_metric_conventions(
table_info: &TableInfo,
- conventions: &Conventions,
+ whitelist: &BTreeMap>,
+ expected_source: &str,
+ expected_metric_type: Option<&str>,
declarations: &mut Vec,
) {
- let Some(implicit_entities) = conventions.prometheus_info_metrics.get(&table_info.name)
- else {
+ let Some(implicit_entities) = whitelist.get(&table_info.name) else {
return;
};
let options = &table_info.meta.options.extra_options;
if options.get(SEMANTIC_SIGNAL_TYPE).map(String::as_str) != Some(SIGNAL_TYPE_METRIC)
- || options.get(SEMANTIC_SOURCE).map(String::as_str) != Some(SOURCE_PROMETHEUS)
+ || options.get(SEMANTIC_SOURCE).map(String::as_str) != Some(expected_source)
+ || expected_metric_type.is_some_and(|expected| {
+ options.get(SEMANTIC_METRIC_TYPE).map(String::as_str) != Some(expected)
+ })
|| table_info.is_physical_table()
{
debug!(
- "Table `{}` matches the info-metric whitelist but is not a prometheus logical \
- metric table; skipping its implicit declarations",
+ "Table `{}` matches the info-metric whitelist but is not an eligible \
+ `{expected_source}` info-metric source; skipping its implicit declarations",
table_info.name
);
return;
@@ -296,12 +317,19 @@ impl EntityGraphProviderImpl {
.cloned()
.collect()
};
+ // A table predating the qualifier column keeps the unqualified
+ // identity rather than losing the declaration.
+ let id_qualifier = implicit
+ .qualified_by
+ .clone()
+ .filter(|c| schema.column_schema_by_name(c).is_some());
declarations.push(EntityDeclaration {
schema: table_info.schema_name.clone(),
table: table_info.name.clone(),
time_index: time_index.clone(),
entity_type: implicit.entity.clone(),
id_columns: implicit.id.clone(),
+ id_qualifier,
descriptive_columns,
scope_columns: vec![],
});
@@ -832,6 +860,24 @@ mod tests {
declarations[2].id_columns,
vec!["service_name", "resource_attributes.service.instance.id"]
);
+ assert_eq!(declarations[1].id_qualifier, None);
+
+ let namespaced = table_info(
+ &[
+ "service_name",
+ "resource_attributes.service.namespace",
+ "resource_attributes.service.instance.id",
+ ],
+ &[(TABLE_DATA_MODEL, TABLE_DATA_MODEL_TRACE_V1)],
+ );
+ for declaration in sorted_declarations(&namespaced) {
+ assert_eq!(
+ declaration.id_qualifier.as_deref(),
+ Some("resource_attributes.service.namespace"),
+ "{} must qualify its identity like the metric side's job",
+ declaration.entity_type
+ );
+ }
// Missing uid column: no pod entity synthesized, no name-based guess.
let no_uid = table_info(
@@ -856,6 +902,47 @@ mod tests {
assert!(sorted_declarations(&invalid_explicit).is_empty());
}
+ #[test]
+ fn trace_table_host_and_container_require_stable_ids() {
+ let with_ids = table_info(
+ &[
+ "service_name",
+ "resource_attributes.host.id",
+ "resource_attributes.host.name",
+ "resource_attributes.container.id",
+ ],
+ &[(TABLE_DATA_MODEL, TABLE_DATA_MODEL_TRACE_V1)],
+ );
+ let declarations = sorted_declarations(&with_ids);
+ let types: Vec<&str> = declarations
+ .iter()
+ .map(|d| d.entity_type.as_str())
+ .collect();
+ assert_eq!(types, vec!["container", "host", "service"]);
+ assert_eq!(
+ declarations[1].id_columns,
+ vec!["resource_attributes.host.id"]
+ );
+ assert_eq!(
+ declarations[1].descriptive_columns,
+ vec!["resource_attributes.host.name"]
+ );
+
+ let names_only = table_info(
+ &[
+ "service_name",
+ "resource_attributes.host.name",
+ "resource_attributes.container.name",
+ ],
+ &[(TABLE_DATA_MODEL, TABLE_DATA_MODEL_TRACE_V1)],
+ );
+ let types: Vec = sorted_declarations(&names_only)
+ .into_iter()
+ .map(|d| d.entity_type)
+ .collect();
+ assert_eq!(types, vec!["service"]);
+ }
+
#[test]
fn prometheus_implicit_declarations_are_gated() {
let labels: &[&str] = &["namespace", "pod", "node"];
@@ -900,6 +987,111 @@ mod tests {
assert_eq!(declarations[1].id_columns, vec!["pod"]);
}
+ const OTEL_STAMPS: &[(&str, &str)] = &[
+ (SEMANTIC_SIGNAL_TYPE, SIGNAL_TYPE_METRIC),
+ (SEMANTIC_SOURCE, SOURCE_OPENTELEMETRY),
+ (SEMANTIC_METRIC_TYPE, servers::semantic::METRIC_TYPE_INFO),
+ ];
+
+ #[test]
+ fn greptime_otel_resource_info_gets_implicit_declarations() {
+ let info = prom_table_info(
+ "greptime_otel_resource_info",
+ &[
+ "job",
+ "instance",
+ "service.name",
+ "service.namespace",
+ "host.id",
+ "host.name",
+ "container.id",
+ "container.name",
+ "k8s.pod.uid",
+ "k8s.pod.name",
+ "k8s.namespace.name",
+ ],
+ OTEL_STAMPS,
+ );
+ let declarations = sorted_declarations(&info);
+ let types: Vec<&str> = declarations
+ .iter()
+ .map(|d| d.entity_type.as_str())
+ .collect();
+ assert_eq!(
+ types,
+ vec![
+ "container",
+ "host",
+ "k8s.pod",
+ "service",
+ "service.instance"
+ ]
+ );
+ assert_eq!(declarations[1].id_columns, vec!["host.id"]);
+ assert_eq!(declarations[1].descriptive_columns, vec!["host.name"]);
+ assert_eq!(declarations[3].id_columns, vec!["job"]);
+ assert_eq!(
+ declarations[3].descriptive_columns,
+ vec!["service.name", "service.namespace"]
+ );
+ assert_eq!(declarations[4].id_columns, vec!["job", "instance"]);
+ assert!(declarations[4].descriptive_columns.is_empty());
+
+ let partial = prom_table_info(
+ "greptime_otel_resource_info",
+ &["job", "service.name", "host.id"],
+ OTEL_STAMPS,
+ );
+ let types: Vec = sorted_declarations(&partial)
+ .into_iter()
+ .map(|d| d.entity_type)
+ .collect();
+ assert_eq!(types, vec!["host", "service"]);
+ }
+
+ #[test]
+ fn otel_implicit_declarations_are_gated() {
+ let labels: &[&str] = &["job", "instance", "host.id"];
+ assert!(
+ sorted_declarations(&prom_table_info(
+ "greptime_otel_resource_info",
+ labels,
+ PROM_STAMPS
+ ))
+ .is_empty()
+ );
+ let mut stamps = OTEL_STAMPS.to_vec();
+ stamps.push((PHYSICAL_TABLE_METADATA_KEY, "true"));
+ assert!(
+ sorted_declarations(&prom_table_info(
+ "greptime_otel_resource_info",
+ labels,
+ &stamps
+ ))
+ .is_empty()
+ );
+ assert!(
+ conventions()
+ .unwrap()
+ .otel_info_metrics
+ .contains_key(servers::otlp::metrics::OTEL_RESOURCE_INFO_TABLE_NAME)
+ );
+
+ let wrong_type = [
+ (SEMANTIC_SIGNAL_TYPE, SIGNAL_TYPE_METRIC),
+ (SEMANTIC_SOURCE, SOURCE_OPENTELEMETRY),
+ (SEMANTIC_METRIC_TYPE, servers::semantic::METRIC_TYPE_GAUGE),
+ ];
+ assert!(
+ sorted_declarations(&prom_table_info(
+ "greptime_otel_resource_info",
+ labels,
+ &wrong_type
+ ))
+ .is_empty()
+ );
+ }
+
#[test]
fn target_info_descriptive_rest_covers_remaining_tags() {
let info = prom_table_info(
diff --git a/src/frontend/src/instance/otlp.rs b/src/frontend/src/instance/otlp.rs
index 0808619853..ccc09c9978 100644
--- a/src/frontend/src/instance/otlp.rs
+++ b/src/frontend/src/instance/otlp.rs
@@ -50,7 +50,7 @@ use table::requests::{
use self::trace_ingest::trace_conventions;
use crate::instance::Instance;
-use crate::metrics::{OTLP_LOGS_ROWS, OTLP_METRICS_ROWS};
+use crate::metrics::{OTLP_LOGS_ROWS, OTLP_METRICS_ROWS, OTLP_RESOURCE_INFO_WRITE_ERRORS};
fn trace_permission_targets(
table_name: &str,
@@ -119,9 +119,15 @@ impl OpenTelemetryProtocolHandler for Instance {
.cloned()
.unwrap_or_default();
metric_ctx.is_legacy = is_legacy;
+ metric_ctx.resource_info = self.otlp_resource_info;
- let (requests, rows, semantic_index, mut outcome) =
- otlp::metrics::to_grpc_insert_requests(request, &mut metric_ctx)?;
+ let otlp::metrics::MetricsConversion {
+ requests,
+ rows,
+ semantic_index,
+ resource_info,
+ mut outcome,
+ } = otlp::metrics::to_grpc_insert_requests(request, &mut metric_ctx)?;
if outcome.rejected_data_points > 0 {
warn!(
"Rejected {} OTLP exponential histogram data points: {}",
@@ -156,7 +162,7 @@ impl OpenTelemetryProtocolHandler for Instance {
// OTLP tables have one sample field in both the legacy and physical paths.
let output = if metric_ctx.is_legacy || !metric_ctx.with_metric_engine {
- self.handle_row_inserts(requests, ctx, false, true)
+ self.handle_row_inserts(requests, ctx.clone(), false, true)
.await
.map_err(BoxedError::new)
.context(error::ExecuteGrpcQuerySnafu)
@@ -165,12 +171,43 @@ impl OpenTelemetryProtocolHandler for Instance {
.extension(PHYSICAL_TABLE_PARAM)
.unwrap_or(GREPTIME_PHYSICAL_TABLE)
.to_string();
- self.handle_metric_row_inserts(requests, ctx, physical_table)
+ self.handle_metric_row_inserts(requests, ctx.clone(), physical_table)
.await
.map_err(BoxedError::new)
.context(error::ExecuteGrpcQuerySnafu)
}?;
outcome.write_cost = output.meta.cost;
+
+ // Derived enrichment, written after the metric data is committed:
+ // failing here would make the client retry data the server already
+ // accepted, so every failure degrades to a warning instead.
+ if let Some(resource_info) = resource_info {
+ let written = match self.check_row_insert_permission(
+ &resource_info,
+ &ctx,
+ PermissionReq::Action(OTLP_WRITE),
+ ) {
+ Ok(_) => self
+ .handle_row_inserts(resource_info, ctx, false, false)
+ .await
+ .map_err(BoxedError::new)
+ .map_err(|e| e.to_string()),
+ Err(e) => Err(e.to_string()),
+ };
+ match written {
+ Ok(descriptor_output) => outcome.write_cost += descriptor_output.meta.cost,
+ Err(e) => {
+ OTLP_RESOURCE_INFO_WRITE_ERRORS.inc();
+ warn!("Failed to write the OTLP resource descriptor table: {e}");
+ outcome.error_message.get_or_insert(format!(
+ "metric data was accepted, but writing the resource \
+ descriptor table `{}` failed: {e}",
+ otlp::metrics::OTEL_RESOURCE_INFO_TABLE_NAME
+ ));
+ }
+ }
+ }
+
Ok(outcome)
}
diff --git a/src/frontend/src/metrics.rs b/src/frontend/src/metrics.rs
index aba33637cf..f649db0b46 100644
--- a/src/frontend/src/metrics.rs
+++ b/src/frontend/src/metrics.rs
@@ -45,6 +45,14 @@ lazy_static! {
)
.unwrap();
+ /// Failed writes of the synthesized OTLP resource descriptor table; these
+ /// surface as OTLP partial-success warnings, not request failures.
+ pub static ref OTLP_RESOURCE_INFO_WRITE_ERRORS: IntCounter = register_int_counter!(
+ "greptime_frontend_otlp_resource_info_write_errors",
+ "frontend otlp resource descriptor write errors"
+ )
+ .unwrap();
+
/// The number of OpenTelemetry traces send by frontend node.
pub static ref OTLP_TRACES_ROWS: IntCounter = register_int_counter!(
"greptime_frontend_otlp_traces_rows",
diff --git a/src/frontend/src/service_config/otlp.rs b/src/frontend/src/service_config/otlp.rs
index 5dcab2c8ce..d6a36f702e 100644
--- a/src/frontend/src/service_config/otlp.rs
+++ b/src/frontend/src/service_config/otlp.rs
@@ -23,6 +23,11 @@ pub struct OtlpOptions {
pub experimental_enable_exponential_histogram: bool,
/// Maximum spans per trace ingest chunk. Set to 0 to disable splitting.
pub trace_ingest_chunk_size: usize,
+ /// Whether to synthesize the `greptime_otel_resource_info` descriptor
+ /// table from the resource attributes of OTLP metrics, so metrics-only
+ /// services reach the semantic graph. Off by default: it creates and
+ /// writes a table the user did not send.
+ pub experimental_enable_resource_info: bool,
}
impl Default for OtlpOptions {
@@ -31,6 +36,7 @@ impl Default for OtlpOptions {
enable: true,
experimental_enable_exponential_histogram: false,
trace_ingest_chunk_size: DEFAULT_TRACE_INGEST_CHUNK_SIZE,
+ experimental_enable_resource_info: false,
}
}
}
@@ -45,6 +51,7 @@ mod tests {
assert!(default.enable);
assert!(!default.experimental_enable_exponential_histogram);
assert_eq!(default.trace_ingest_chunk_size, 512);
+ assert!(!default.experimental_enable_resource_info);
let options: OtlpOptions = toml::from_str("enable = false").unwrap();
assert!(!options.enable);
diff --git a/src/operator/src/statement/semantic_graph.rs b/src/operator/src/statement/semantic_graph.rs
index 7014a94f8a..86827cf8c7 100644
--- a/src/operator/src/statement/semantic_graph.rs
+++ b/src/operator/src/statement/semantic_graph.rs
@@ -41,8 +41,9 @@ use common_catalog::consts::{
ENTITY_DESCRIPTIVE_COLUMN, ENTITY_ID_ATTRS_COLUMN, ENTITY_ID_COLUMN, ENTITY_SCOPE_COLUMN,
ENTITY_TYPE_COLUMN, ERROR_COUNT_COLUMN, FRESH_UNTIL_COLUMN, GENERATION_ID_COLUMN,
OBSERVED_AT_COLUMN, PROVENANCE_COLUMN, REL_TYPE_COLUMN, REQUEST_COUNT_COLUMN,
- SEMANTIC_RELATIONSHIPS_DECLARED_TABLE_NAME, SOURCE_TABLES_COLUMN, SRC_ID_COLUMN,
- SRC_TYPE_COLUMN, VALID_FROM_COLUMN, VALID_UNTIL_COLUMN, WINDOW_END_COLUMN, WINDOW_START_COLUMN,
+ SEMANTIC_GRAPH_WINDOW_NANOS, SEMANTIC_RELATIONSHIPS_DECLARED_TABLE_NAME, SOURCE_TABLES_COLUMN,
+ SRC_ID_COLUMN, SRC_TYPE_COLUMN, VALID_FROM_COLUMN, VALID_UNTIL_COLUMN, WINDOW_END_COLUMN,
+ WINDOW_START_COLUMN,
};
use common_function::function::FunctionContext;
use common_function::function_registry::FUNCTION_REGISTRY;
@@ -60,7 +61,7 @@ use datafusion::dataframe::DataFrame;
use datafusion::functions::{core as core_fns, datetime as datetime_fns, string as string_fns};
use datafusion::functions_nested::expr_fn::make_array;
use datafusion_common::{Column, Result as DfResult, ScalarValue};
-use datafusion_expr::{Expr, LogicalPlan, ScalarUDF, cast, ident, lit};
+use datafusion_expr::{Case, Expr, LogicalPlan, ScalarUDF, cast, ident, lit};
pub use relationships::{
CallsSource, CoDeclaredSource, DeclaredSource, RelationshipSources, build_relationships_plan,
};
@@ -122,9 +123,10 @@ pub fn declared_relationships_schema_matches(table_info: &table::metadata::Table
})
}
-/// Bin width for the temporal window of derived rows: 60s buckets, matching the
-/// service-graph convention.
-const BIN_NANOS: i64 = 60 * 1_000_000_000;
+/// Bin width for the temporal window of derived rows, matching the
+/// service-graph convention. Shared so ingestion-synthesized observations
+/// land in the same buckets.
+const BIN_NANOS: i64 = SEMANTIC_GRAPH_WINDOW_NANOS;
/// Default retention for the declared-edge table; expiry slides the topology window.
const DEFAULT_DECLARED_RELATIONSHIPS_TTL: &str = "90d";
@@ -284,8 +286,12 @@ pub struct EntityDeclaration {
/// The table's time index column, used for the temporal window filter.
pub time_index: String,
pub entity_type: String,
- /// Identifying columns (>= 1). One column → id verbatim; several → composite.
+ /// Identifying columns (>= 1), ordered broad to narrow.
pub id_columns: Vec,
+ /// Optional column qualifying the first identity component, so a source
+ /// carrying the parts separately (a trace table's `service.namespace`)
+ /// yields the same id as one carrying them pre-composed (`job`).
+ pub id_qualifier: Option,
/// Descriptive columns snapshotted into the `descriptive` JSON (may be empty).
pub descriptive_columns: Vec,
/// Scope columns (namespace/environment). One column → scope verbatim;
@@ -417,28 +423,91 @@ fn cast_string_or_empty(column: &str) -> Expr {
core_fns::coalesce().call(vec![cast(ident(column), DataType::Utf8), lit("")])
}
-/// The canonical entity-id expression for `id_columns`: the value verbatim for
-/// a single column, the sorted `k=v,k=v` rendering for a composite. `col`
-/// constructs the column reference (unqualified for registry branches,
-/// join-side-qualified for the calls derivation).
/// A row identifies an entity only when every identity component is present
/// and non-empty: kube-state-metrics descriptors emit empty-string labels (an
/// unscheduled pod's `node`, an owner-less pod's `owner_*`), and an empty
-/// string is never a meaningful entity id.
+/// string is never a meaningful entity id. The qualifier is optional by
+/// construction and so is not guarded.
pub(crate) fn identifies(column: &str) -> Expr {
ident(column)
.is_not_null()
.and(cast(ident(column), DataType::Utf8).not_eq(lit("")))
}
-fn entity_id_expr(id_columns: &[String], col: &dyn Fn(&str) -> Expr) -> Expr {
- if let [id] = id_columns {
- cast(col(id), DataType::Utf8)
- } else {
- let mut cols = id_columns.to_vec();
- cols.sort();
- sorted_kv_expr_with(&cols, false, col)
+const ID_SEPARATOR: &str = ",";
+const ID_ESCAPE: &str = "\\";
+/// Left unescaped: `/` is how Prometheus renders `job`.
+const ID_QUALIFIER_SEPARATOR: &str = "/";
+
+/// Escapes so a composite id decodes back to its components. The escape
+/// character goes first, or it would double the escapes the separator pass
+/// introduces.
+fn escaped_id_value(value: Expr) -> Expr {
+ let escape_escape = string_fns::replace().call(vec![
+ value,
+ lit(ID_ESCAPE),
+ lit(format!("{ID_ESCAPE}{ID_ESCAPE}")),
+ ]);
+ string_fns::replace().call(vec![
+ escape_escape,
+ lit(ID_SEPARATOR),
+ lit(format!("{ID_ESCAPE}{ID_SEPARATOR}")),
+ ])
+}
+
+/// The identity values in declared order (broad to narrow), escaped and
+/// joined.
+///
+/// The identifying *column names* are deliberately absent: the same identity
+/// reaches us under different names per source — a trace table's
+/// `service_name` against a metric table's `job` — and encoding them would
+/// split one entity into one per signal. `entity_id_attrs` keeps the
+/// structured form.
+///
+/// `col` constructs the column reference (unqualified for registry branches,
+/// join-side-qualified for the calls derivation).
+fn entity_id_expr(
+ id_columns: &[String],
+ qualifier: Option<&str>,
+ col: &dyn Fn(&str) -> Expr,
+) -> Expr {
+ let mut parts = Vec::with_capacity(id_columns.len() * 2);
+ for (index, column) in id_columns.iter().enumerate() {
+ if index > 0 {
+ parts.push(lit(ID_SEPARATOR));
+ }
+ let value = escaped_id_value(cast(col(column), DataType::Utf8));
+ parts.push(match qualifier {
+ Some(qualifier) if index == 0 => qualified_id_expr(qualifier, value, col),
+ _ => value,
+ });
}
+ if let [single] = parts.as_slice() {
+ single.clone()
+ } else {
+ concat_expr(parts)
+ }
+}
+
+/// `/`, or bare `value` when the qualifier is empty on this
+/// row — the spec's rule composing `job` from `service.namespace` and
+/// `service.name`.
+fn qualified_id_expr(qualifier: &str, value: Expr, col: &dyn Fn(&str) -> Expr) -> Expr {
+ let qualifier = escaped_id_value(
+ core_fns::coalesce().call(vec![cast(col(qualifier), DataType::Utf8), lit("")]),
+ );
+ Expr::Case(Case::new(
+ None,
+ vec![(
+ Box::new(qualifier.clone().eq(lit(""))),
+ Box::new(value.clone()),
+ )],
+ Some(Box::new(concat_expr(vec![
+ qualifier,
+ lit(ID_QUALIFIER_SEPARATOR),
+ value,
+ ]))),
+ ))
}
/// The `parse_json` UDF, shared by all derivation plans. Resolved from the
@@ -583,14 +652,22 @@ fn registry_source(
let mut rows = Vec::with_capacity(1 + rest.len());
for decl in std::iter::once(first).chain(rest) {
// CAST even a single-column id: id columns need not be strings, and
- // the computed table declares entity_id STRING. Composite ids
- // additionally carry a JSON object of the id columns in
- // entity_id_attrs.
- let entity_id = entity_id_expr(&decl.id_columns, &|c| ident(c));
- let entity_id_attrs = if decl.id_columns.len() == 1 {
+ // the computed table declares entity_id STRING. An id assembled from
+ // more than one column additionally carries its parts as a JSON
+ // object in entity_id_attrs, the structured form of the id string.
+ let entity_id = entity_id_expr(&decl.id_columns, decl.id_qualifier.as_deref(), &|c| {
+ ident(c)
+ });
+ let id_parts = decl
+ .id_qualifier
+ .iter()
+ .chain(&decl.id_columns)
+ .cloned()
+ .collect::>();
+ let entity_id_attrs = if id_parts.len() == 1 {
null_json()
} else {
- json_object_expr(&decl.id_columns)
+ json_object_expr(&id_parts)
};
let scope = match decl.scope_columns.as_slice() {
@@ -811,11 +888,74 @@ mod tests {
time_index: "ts".to_string(),
entity_type: entity_type.to_string(),
id_columns: id_columns.iter().map(|s| s.to_string()).collect(),
+ id_qualifier: None,
descriptive_columns: vec![],
scope_columns: vec![],
}
}
+ /// The id string must decode back to its components, and a qualifier must
+ /// reproduce the identity a source that pre-composed it already emits.
+ #[tokio::test]
+ async fn registry_id_escaping_and_qualifier() {
+ let schema = Arc::new(Schema::new(vec![
+ Field::new(
+ "ts",
+ DataType::Timestamp(TimeUnit::Millisecond, None),
+ false,
+ ),
+ Field::new("service_name", DataType::Utf8, false),
+ Field::new("instance", DataType::Utf8, false),
+ Field::new("namespace", DataType::Utf8, true),
+ ]));
+ let batch = RecordBatch::try_new(
+ schema.clone(),
+ vec![
+ Arc::new(TimestampMillisecondArray::from(vec![1_000, 2_000, 3_000])) as ArrayRef,
+ Arc::new(StringArray::from(vec!["cart", "a,b", "we\\ird"])),
+ Arc::new(StringArray::from(vec!["i-1", "c", "i-3"])),
+ Arc::new(StringArray::from(vec![Some("shop"), None, Some("")])),
+ ],
+ )
+ .unwrap();
+ let ctx = SessionContext::new();
+ ctx.register_table(
+ "svc",
+ Arc::new(MemTable::try_new(schema, vec![vec![batch]]).unwrap()),
+ )
+ .unwrap();
+
+ let mut declaration = decl("service.instance", &["service_name", "instance"]);
+ declaration.id_qualifier = Some("namespace".to_string());
+ declaration.time_index = "ts".to_string();
+ let plan = build_registry_plan(
+ vec![RegistrySource {
+ declarations: vec![declaration],
+ scan: ctx.table("svc").await.unwrap(),
+ }],
+ &test_window(),
+ )
+ .unwrap()
+ .unwrap();
+
+ let mut ids: Vec = collect(&ctx, plan)
+ .await
+ .iter()
+ .flat_map(|batch| strings(batch, 5))
+ .collect();
+ ids.sort();
+ assert_eq!(
+ ids,
+ vec![
+ // an absent qualifier leaves the identity bare, and a value
+ // holding the separator stays distinguishable from two values
+ "a\\,b,c".to_string(),
+ "shop/cart,i-1".to_string(),
+ "we\\\\ird,i-3".to_string(),
+ ]
+ );
+ }
+
#[tokio::test]
async fn registry_single_column_identity() {
let ctx = metric_table_ctx();
@@ -900,24 +1040,24 @@ mod tests {
.collect();
rows.sort();
- // Composite id -> sorted `k=v,k=v` plus a JSON object of the id columns;
- // descriptive JSON keeps `\`, `"` and control characters intact in
- // runtime values, NULL -> "".
+ // Composite id -> values in declared order plus a JSON object of the
+ // id columns; descriptive JSON keeps `\`, `"` and control characters
+ // intact in runtime values, NULL -> "".
assert_eq!(
rows,
vec![
(
- "pid=42,service_name=cart".to_string(),
+ "cart,42".to_string(),
Some(r#"{"pid":"42","service_name":"cart"}"#.to_string()),
Some(r#"{"host":""}"#.to_string()),
),
(
- "pid=42,service_name=cart".to_string(),
+ "cart,42".to_string(),
Some(r#"{"pid":"42","service_name":"cart"}"#.to_string()),
Some(r#"{"host":"h2"}"#.to_string()),
),
(
- "pid=42,service_name=cart".to_string(),
+ "cart,42".to_string(),
Some(r#"{"pid":"42","service_name":"cart"}"#.to_string()),
Some(r#"{"host":"we\"ird\\\nhost"}"#.to_string()),
),
diff --git a/src/operator/src/statement/semantic_graph/conventions.rs b/src/operator/src/statement/semantic_graph/conventions.rs
index 36d058921b..9ec6f52e15 100644
--- a/src/operator/src/statement/semantic_graph/conventions.rs
+++ b/src/operator/src/statement/semantic_graph/conventions.rs
@@ -19,7 +19,7 @@
//! with the binary, not an operator-editable configuration surface; explicit
//! `greptime.semantic.entity.*` declarations always override it.
-use std::collections::{BTreeMap, HashSet};
+use std::collections::{BTreeMap, HashMap, HashSet};
use std::sync::LazyLock;
use serde::Deserialize;
@@ -43,15 +43,20 @@ pub struct VirtualDstCandidate {
pub connection_type: String,
}
-/// One implicit entity declaration: of a whitelisted Prometheus info metric,
-/// or of a trace-v1 table's flattened resource attributes.
+/// One implicit entity declaration: of a whitelisted Prometheus or OTel info
+/// metric, or of a trace-v1 table's flattened resource attributes.
#[derive(Debug, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct ImplicitEntity {
pub entity: String,
- /// Identifying label columns; every one must exist on the table for the
- /// declaration to apply.
+ /// Identifying label columns, ordered broad to narrow; every one must
+ /// exist on the table for the declaration to apply.
pub id: Vec,
+ /// Column qualifying `id[0]` as `/`, so a source
+ /// carrying the parts separately matches one carrying them pre-composed.
+ /// Skipped when the column is absent or empty.
+ #[serde(default)]
+ pub qualified_by: Option,
/// Descriptive label columns, filtered to those present (kube-state-metrics
/// label sets vary across versions).
#[serde(default)]
@@ -70,7 +75,12 @@ pub struct Conventions {
pub trace_co_declared_edges: Vec,
pub virtual_dst_candidates: Vec,
pub otlp_trace_entities: Vec,
+ /// Table name -> the entities that table declares, for Prometheus-sourced
+ /// descriptor metrics (`source = prometheus`).
pub prometheus_info_metrics: BTreeMap>,
+ /// Table name -> the entities that table declares, for OTLP-sourced
+ /// descriptor tables (`source = opentelemetry`).
+ pub otel_info_metrics: BTreeMap>,
}
/// The built-in entity-type vocabulary. User-declared types are open-ended;
@@ -78,6 +88,7 @@ pub struct Conventions {
pub const ENTITY_TYPE_SERVICE: &str = "service";
pub const ENTITY_TYPE_SERVICE_INSTANCE: &str = "service.instance";
pub const ENTITY_TYPE_HOST: &str = "host";
+pub const ENTITY_TYPE_CONTAINER: &str = "container";
pub const ENTITY_TYPE_PROCESS: &str = "process";
pub const ENTITY_TYPE_K8S_POD: &str = "k8s.pod";
pub const ENTITY_TYPE_K8S_NODE: &str = "k8s.node";
@@ -88,10 +99,11 @@ pub const ENTITY_TYPE_GEN_AI_AGENT: &str = "gen_ai.agent";
pub const ENTITY_TYPE_GEN_AI_MODEL: &str = "gen_ai.model";
pub const ENTITY_TYPE_GEN_AI_TOOL: &str = "gen_ai.tool";
-const ENTITY_TYPES: [&str; 12] = [
+const ENTITY_TYPES: [&str; 13] = [
ENTITY_TYPE_SERVICE,
ENTITY_TYPE_SERVICE_INSTANCE,
ENTITY_TYPE_HOST,
+ ENTITY_TYPE_CONTAINER,
ENTITY_TYPE_PROCESS,
ENTITY_TYPE_K8S_POD,
ENTITY_TYPE_K8S_NODE,
@@ -200,11 +212,16 @@ fn validate(conventions: &Conventions) -> Result<(), String> {
let per_table = conventions
.prometheus_info_metrics
.iter()
+ .chain(&conventions.otel_info_metrics)
.map(|(table, entities)| (table.as_str(), entities))
.chain(std::iter::once((
"otlp traces",
&conventions.otlp_trace_entities,
)));
+ // An entity type declared with a different number of id columns by two
+ // sources yields ids that can never match, silently splitting one entity
+ // in two.
+ let mut arity = HashMap::new();
for (table, entities) in per_table {
let mut seen_types = HashSet::new();
for implicit in entities {
@@ -226,6 +243,12 @@ fn validate(conventions: &Conventions) -> Result<(), String> {
implicit.entity
));
}
+ if implicit.qualified_by.as_ref().is_some_and(String::is_empty) {
+ return Err(format!(
+ "entity `{}` of info metric `{table}` has an empty qualified_by",
+ implicit.entity
+ ));
+ }
if implicit.descriptive_rest && !implicit.descriptive.is_empty() {
return Err(format!(
"entity `{}` of info metric `{table}` sets both descriptive and \
@@ -233,6 +256,17 @@ fn validate(conventions: &Conventions) -> Result<(), String> {
implicit.entity
));
}
+ match arity.insert(implicit.entity.as_str(), implicit.id.len()) {
+ Some(previous) if previous != implicit.id.len() => {
+ return Err(format!(
+ "entity `{}` is declared with {previous} id columns elsewhere but \
+ {} for `{table}`",
+ implicit.entity,
+ implicit.id.len()
+ ));
+ }
+ _ => {}
+ }
}
}
Ok(())
@@ -249,25 +283,28 @@ mod tests {
/// Each case must fail on exactly the rule it names, so the shared
/// boilerplate is valid and the mutated part is inside the vocabulary.
- fn broken(edges: &str, info_metrics: &str) -> String {
+ fn broken(edges: &str, info_metrics: &str, otel_metrics: &str) -> String {
format!(
"co_declared_edges: [{edges}]\ntrace_co_declared_edges: []\n\
virtual_dst_candidates: []\notlp_trace_entities: []\n\
- prometheus_info_metrics: {{{info_metrics}}}"
+ prometheus_info_metrics: {{{info_metrics}}}\n\
+ otel_info_metrics: {{{otel_metrics}}}"
)
}
#[test]
fn validation_rejects_broken_conventions() {
- let err = |edges, info| parse(&broken(edges, info)).unwrap_err();
+ let err = |edges, info, otel| parse(&broken(edges, info, otel)).unwrap_err();
- assert!(err("{src: host, dst: service, rel: pets}", "").contains("unknown rel_type"));
+ assert!(err("{src: host, dst: service, rel: pets}", "", "").contains("unknown rel_type"));
assert!(
- err("{src: k8s.pods, dst: k8s.node, rel: runs_on}", "").contains("unknown entity type")
+ err("{src: k8s.pods, dst: k8s.node, rel: runs_on}", "", "")
+ .contains("unknown entity type")
);
assert!(
err(
"{src: host, dst: service, rel: uses}, {src: host, dst: service, rel: uses}",
+ "",
""
)
.contains("duplicate edge rule")
@@ -275,14 +312,27 @@ mod tests {
assert!(
err(
"",
- "t: [{entity: host, id: [x], descriptive: [y], descriptive_rest: true}]"
+ "t: [{entity: host, id: [x], descriptive: [y], descriptive_rest: true}]",
+ ""
)
.contains("descriptive_rest")
);
+ // the otel map runs through the same per-table validation
+ assert!(err("", "", "t: [{entity: hosts, id: [x]}]").contains("unknown entity type"));
+ // ids of a different arity for one type can never match each other
+ assert!(
+ err(
+ "",
+ "a: [{entity: host, id: [x]}]",
+ "b: [{entity: host, id: [x, y]}]"
+ )
+ .contains("id columns")
+ );
// Unknown YAML keys are rejected, catching typos in the embedded file.
assert!(
parse(&broken(
"{src: host, dst: service, rel: uses, direction: down}",
+ "",
""
))
.is_err()
diff --git a/src/operator/src/statement/semantic_graph/conventions.yaml b/src/operator/src/statement/semantic_graph/conventions.yaml
index 9fc37fe2a0..808fe0f265 100644
--- a/src/operator/src/statement/semantic_graph/conventions.yaml
+++ b/src/operator/src/statement/semantic_graph/conventions.yaml
@@ -2,6 +2,10 @@
# versioned with the binary (include_str!), not an operator-editable surface;
# non-standard deployments override per table with explicit
# greptime.semantic.entity.* declarations, which always win.
+#
+# Every `id` list is ordered broad to narrow, and that order is part of the
+# identity: the id renders as the values joined in this order, so all
+# declarations of one entity type must agree on it or they name two entities.
co_declared_edges:
- { src: service.instance, dst: host, rel: runs_on }
@@ -11,6 +15,8 @@ co_declared_edges:
- { src: service.instance, dst: service, rel: part_of }
- { src: k8s.pod, dst: k8s.workload, rel: part_of }
- { src: service.instance, dst: k8s.pod, rel: runs_on }
+ - { src: service.instance, dst: container, rel: runs_on }
+ - { src: container, dst: host, rel: runs_on }
# Applied only to trace sources; the endpoints are Greptime entity types
# derived from GenAI semantic-convention attributes.
@@ -30,12 +36,25 @@ virtual_dst_candidates:
- { column: span_attributes.server.address, connection_type: virtual_node }
# Implicit declarations for greptime_trace_v1 tables; the pod UID bridges
-# trace-side pods onto the kube-state-metrics entities below.
+# trace-side pods onto the kube-state-metrics entities below. host.name and
+# container.name vary by SDK and resource detector, so identity is the id.
+# qualified_by folds the namespace into the service identity the same way the
+# metric sources' `job` already carries it, so both signals name one entity.
otlp_trace_entities:
- entity: service
id: [service_name]
+ qualified_by: resource_attributes.service.namespace
- entity: service.instance
id: [service_name, resource_attributes.service.instance.id]
+ qualified_by: resource_attributes.service.namespace
+ - entity: host
+ id: [resource_attributes.host.id]
+ descriptive:
+ - resource_attributes.host.name
+ - entity: container
+ id: [resource_attributes.container.id]
+ descriptive:
+ - resource_attributes.container.name
- entity: k8s.pod
id: [resource_attributes.k8s.pod.uid]
descriptive:
@@ -113,3 +132,32 @@ prometheus_info_metrics:
- entity: service.instance
id: [job, instance]
descriptive_rest: true
+
+# Implicit declarations for the OTLP resource descriptor synthesized at
+# metrics ingestion, gated on source=opentelemetry and metric.type=info.
+# job/instance are the derived Prometheus-compatible identity, so services
+# line up with target_info. Unlike target_info the attribute set is known, so
+# descriptive columns are listed instead of swept up by descriptive_rest,
+# which would misfile other entities' ids under service.instance.
+otel_info_metrics:
+ greptime_otel_resource_info:
+ - entity: service
+ id: [job]
+ descriptive:
+ - service.name
+ - service.namespace
+ - entity: service.instance
+ id: [job, instance]
+ - entity: host
+ id: [host.id]
+ descriptive:
+ - host.name
+ - entity: container
+ id: [container.id]
+ descriptive:
+ - container.name
+ - entity: k8s.pod
+ id: [k8s.pod.uid]
+ descriptive:
+ - k8s.pod.name
+ - k8s.namespace.name
diff --git a/src/operator/src/statement/semantic_graph/relationships.rs b/src/operator/src/statement/semantic_graph/relationships.rs
index 3b5c6d37b7..257c844a9a 100644
--- a/src/operator/src/statement/semantic_graph/relationships.rs
+++ b/src/operator/src/statement/semantic_graph/relationships.rs
@@ -200,9 +200,9 @@ fn co_declared_branch(
bin.clone() + bin_interval(),
bin.clone() + bin_interval(),
lit(src.entity_type.as_str()),
- entity_id_expr(&src.id_columns, &|c| ident(c)),
+ entity_id_expr(&src.id_columns, src.id_qualifier.as_deref(), &|c| ident(c)),
lit(dst.entity_type.as_str()),
- entity_id_expr(&dst.id_columns, &|c| ident(c)),
+ entity_id_expr(&dst.id_columns, dst.id_qualifier.as_deref(), &|c| ident(c)),
lit(*rel_type),
lit(*provenance),
lit(1.0_f64),
@@ -501,6 +501,13 @@ fn span_predicate(service: &EntityDeclaration, window: &GraphQueryWindow, strict
}
/// One trace table's client spans, normalized to
+/// A trace table's `duration_nano` is UInt64 on tables created before the
+/// signed-integer ingest change and Int64 after it. The per-table selects are
+/// unioned, and those two have no common integer type.
+fn duration_nano_expr() -> Expr {
+ cast(ident(DURATION_NANO_COLUMN), DataType::Int64).alias(DURATION_NANO_COLUMN)
+}
+
/// `(timestamp, trace_id, span_id, src_id, status_code, duration_nano,
/// virtual_dst, virtual_conn)`. `src_id` is built from the table's `service`
/// declaration, so edges land on exactly the entity ids the registry emits (a
@@ -553,9 +560,12 @@ fn client_spans(
ident(SPAN_ID_COLUMN),
// The cast inside entity_id_expr also normalizes tag columns, which
// come out of the storage engine dictionary-encoded.
- entity_id_expr(&service.id_columns, &|c| ident(c)).alias(SRC_ID_COLUMN),
+ entity_id_expr(&service.id_columns, service.id_qualifier.as_deref(), &|c| {
+ ident(c)
+ })
+ .alias(SRC_ID_COLUMN),
ident(SPAN_STATUS_CODE_COLUMN),
- ident(DURATION_NANO_COLUMN),
+ duration_nano_expr(),
virtual_dst.alias("virtual_dst"),
virtual_conn.alias("virtual_conn"),
])
@@ -578,9 +588,12 @@ fn server_spans(
ident(TRACE_TIMESTAMP_COLUMN),
ident(TRACE_ID_COLUMN),
ident(PARENT_SPAN_ID_COLUMN),
- entity_id_expr(&service.id_columns, &|c| ident(c)).alias(DST_ID_COLUMN),
+ entity_id_expr(&service.id_columns, service.id_qualifier.as_deref(), &|c| {
+ ident(c)
+ })
+ .alias(DST_ID_COLUMN),
ident(SPAN_STATUS_CODE_COLUMN),
- ident(DURATION_NANO_COLUMN),
+ duration_nano_expr(),
])
}
@@ -609,7 +622,10 @@ fn agent_calls_branch(
ident(TRACE_TIMESTAMP_COLUMN),
ident(TRACE_ID_COLUMN),
ident(SPAN_ID_COLUMN),
- entity_id_expr(&agent.id_columns, &|c| ident(c)).alias(SRC_ID_COLUMN),
+ entity_id_expr(&agent.id_columns, agent.id_qualifier.as_deref(), &|c| {
+ ident(c)
+ })
+ .alias(SRC_ID_COLUMN),
])?;
let child = trace
.scan
@@ -619,9 +635,12 @@ fn agent_calls_branch(
ident(TRACE_TIMESTAMP_COLUMN),
ident(TRACE_ID_COLUMN),
ident(PARENT_SPAN_ID_COLUMN),
- entity_id_expr(&agent.id_columns, &|c| ident(c)).alias(DST_ID_COLUMN),
+ entity_id_expr(&agent.id_columns, agent.id_qualifier.as_deref(), &|c| {
+ ident(c)
+ })
+ .alias(DST_ID_COLUMN),
ident(SPAN_STATUS_CODE_COLUMN),
- ident(DURATION_NANO_COLUMN),
+ duration_nano_expr(),
])?;
parents = union_all(parents, parent)?;
children = union_all(children, child)?;
@@ -732,6 +751,16 @@ mod tests {
name: &str,
extra: &[(&str, &[Option<&str>])],
spans: &[Span<'_>],
+ ) {
+ register_typed_trace_table(ctx, name, extra, spans, DataType::UInt64)
+ }
+
+ fn register_typed_trace_table(
+ ctx: &SessionContext,
+ name: &str,
+ extra: &[(&str, &[Option<&str>])],
+ spans: &[Span<'_>],
+ duration_type: DataType,
) {
let mut fields = vec![
Field::new(
@@ -745,7 +774,7 @@ mod tests {
Field::new(SPAN_KIND_COLUMN, DataType::Utf8, false),
Field::new(SPAN_STATUS_CODE_COLUMN, DataType::Utf8, false),
Field::new(SERVICE_NAME_COLUMN, DataType::Utf8, false),
- Field::new(DURATION_NANO_COLUMN, DataType::UInt64, false),
+ Field::new(DURATION_NANO_COLUMN, duration_type.clone(), false),
];
for (column, _) in extra {
fields.push(Field::new(*column, DataType::Utf8, true));
@@ -774,9 +803,14 @@ mod tests {
Arc::new(StringArray::from(
spans.iter().map(|s| s.6).collect::>(),
)),
- Arc::new(UInt64Array::from(
- spans.iter().map(|s| s.7).collect::>(),
- )),
+ match duration_type {
+ DataType::Int64 => Arc::new(Int64Array::from(
+ spans.iter().map(|s| s.7 as i64).collect::>(),
+ )) as ArrayRef,
+ _ => Arc::new(UInt64Array::from(
+ spans.iter().map(|s| s.7).collect::>(),
+ )),
+ },
];
for (_, values) in extra {
columns.push(Arc::new(StringArray::from(values.to_vec())));
@@ -843,6 +877,7 @@ mod tests {
time_index: TRACE_TIMESTAMP_COLUMN.to_string(),
entity_type: "service".to_string(),
id_columns: id_columns.iter().map(|s| s.to_string()).collect(),
+ id_qualifier: None,
descriptive_columns: vec![],
scope_columns: vec![],
}
@@ -972,7 +1007,7 @@ mod tests {
),
],
);
- register_trace_table(
+ register_typed_trace_table(
&ctx,
"trace_b",
&[],
@@ -989,6 +1024,7 @@ mod tests {
1_500_000_000,
),
],
+ DataType::Int64,
);
let a = ctx.table("trace_a").await.unwrap();
let b = ctx.table("trace_b").await.unwrap();
@@ -1294,6 +1330,7 @@ mod tests {
time_index: "ts".to_string(),
entity_type: entity_type.to_string(),
id_columns: id_columns.iter().map(|s| s.to_string()).collect(),
+ id_qualifier: None,
descriptive_columns: vec![],
scope_columns: vec![],
}
@@ -1544,16 +1581,10 @@ mod tests {
let total: usize = batches.iter().map(|b| b.num_rows()).sum();
assert_eq!(total, 1);
let batch = &batches[0];
- // Composite service identity renders the same sorted `k=v` form the
- // registry emits, so edges land on registry entity ids.
- assert_eq!(
- strings(batch, 5),
- vec!["service_name=frontend,service_namespace=ns1"]
- );
- assert_eq!(
- strings(batch, 7),
- vec!["service_name=cart,service_namespace=ns1"]
- );
+ // Composite service identity renders exactly as the registry emits it,
+ // so edges land on registry entity ids.
+ assert_eq!(strings(batch, 5), vec!["frontend,ns1"]);
+ assert_eq!(strings(batch, 7), vec!["cart,ns1"]);
}
fn build_relationships_plan_for_test(
diff --git a/src/servers/src/http/otlp.rs b/src/servers/src/http/otlp.rs
index d04022cdf8..3ecaba74a4 100644
--- a/src/servers/src/http/otlp.rs
+++ b/src/servers/src/http/otlp.rs
@@ -118,8 +118,9 @@ pub async fn metrics(
promote_scope_attrs: http_opts.promote_scope_attrs,
with_metric_engine,
experimental_enable_exponential_histogram,
- // set is_legacy later
+ // set by the frontend from its config
is_legacy: false,
+ resource_info: false,
metric_type: MetricType::Init,
metric_translation_strategy: http_opts.metric_translation_strategy,
}));
@@ -128,7 +129,7 @@ pub async fn metrics(
handler.metrics(request, query_ctx).await.map(|outcome| {
if outcome.accepted_data_points == 0 && outcome.rejected_data_points > 0 {
OtlpMetricsResponse::Failure(outcome)
- } else if outcome.rejected_data_points > 0 {
+ } else if outcome.rejected_data_points > 0 || outcome.error_message.is_some() {
OtlpMetricsResponse::PartialSuccess(outcome)
} else {
OtlpMetricsResponse::FullSuccess(outcome)
diff --git a/src/servers/src/otlp/metrics.rs b/src/servers/src/otlp/metrics.rs
index bd0bfb2ce7..d65e117e07 100644
--- a/src/servers/src/otlp/metrics.rs
+++ b/src/servers/src/otlp/metrics.rs
@@ -25,6 +25,7 @@ use common_query::native_histogram::{
};
use common_query::prelude::{GREPTIME_COUNT, greptime_timestamp, greptime_value};
use common_query::prometheus::PROMETHEUS_STALE_NAN_BITS;
+use common_telemetry::warn;
use lazy_static::lazy_static;
use otel_arrow_rust::proto::opentelemetry::collector::metrics::v1::ExportMetricsServiceRequest;
use otel_arrow_rust::proto::opentelemetry::common::v1::{AnyValue, KeyValue, any_value};
@@ -36,7 +37,7 @@ use table::requests::{
};
use crate::error::{self, Result};
-use crate::otlp::trace::{KEY_SERVICE_INSTANCE_ID, KEY_SERVICE_NAME};
+use crate::otlp::trace::{KEY_SERVICE_INSTANCE_ID, KEY_SERVICE_NAME, KEY_SERVICE_NAMESPACE};
use crate::query_handler::MetricsIngestOutcome;
use crate::row_writer::{self, MultiTableData, TableData};
pub use crate::semantic::SemanticIndex;
@@ -45,8 +46,11 @@ use crate::semantic::{
METRIC_TYPE_UPDOWN_COUNTER,
};
+mod resource_info;
mod translator;
+pub use resource_info::OTEL_RESOURCE_INFO_TABLE_NAME;
+use resource_info::ResourceInfoData;
pub use translator::legacy_normalize_otlp_name;
pub(crate) use translator::ucum_to_openmetrics_unit;
use translator::{translate_label_name, translate_metric_name};
@@ -96,28 +100,43 @@ const MIN_EXPONENTIAL_HISTOGRAM_SCALE: i32 = -4;
const MAX_EXPONENTIAL_HISTOGRAM_SCALE: i32 = 8;
const MAX_REJECTION_MESSAGE_BYTES: usize = 512;
+/// Result of converting one OTLP metrics request.
+#[derive(Debug)]
+pub struct MetricsConversion {
+ pub requests: RowInsertRequests,
+ /// Row count of `requests`; the resource descriptor is not counted.
+ pub rows: usize,
+ /// Per-table semantic index for the auto-create path to stamp as table
+ /// options; covers the descriptor table too.
+ pub semantic_index: SemanticIndex,
+ /// The synthesized resource descriptor, written separately from the
+ /// metric data. See [`resource_info`].
+ pub resource_info: Option,
+ pub outcome: MetricsIngestOutcome,
+}
+
/// Convert OpenTelemetry metrics to GreptimeDB insert requests
///
/// See
///
/// for data structure of OTLP metrics.
-///
-/// Returns `InsertRequests`, total number of rows to ingest, and the per-table
-/// semantic index for the auto-create path to stamp as table options.
pub fn to_grpc_insert_requests(
request: ExportMetricsServiceRequest,
metric_ctx: &mut OtlpMetricCtx,
-) -> Result<(
- RowInsertRequests,
- usize,
- SemanticIndex,
- MetricsIngestOutcome,
-)> {
+) -> Result {
let mut table_writer = MultiTableData::default();
let mut semantic_index = SemanticIndex::default();
let mut outcome = MetricsIngestOutcome::default();
+ let mut resource_info = ResourceInfoData::default();
for resource in &request.resource_metrics {
+ if metric_ctx.resource_info
+ && !metric_ctx.is_legacy
+ && let Some(r) = resource.resource.as_ref()
+ {
+ resource_info.observe(&r.attributes, resource, metric_ctx);
+ }
+
let resource_attrs = resource.resource.as_ref().map(|r| {
let mut attrs = r.attributes.clone();
process_resource_attrs(&mut attrs, metric_ctx);
@@ -149,8 +168,46 @@ pub fn to_grpc_insert_requests(
}
let (requests, rows) = table_writer.into_row_insert_requests();
+
validate_sample_kinds(&requests)?;
- Ok((requests, rows, semantic_index, outcome))
+
+ // The metric and the descriptor would fight over the same table across
+ // two engines. The user's metric wins.
+ let resource_info = if !metric_ctx.resource_info {
+ None
+ } else if requests
+ .inserts
+ .iter()
+ .any(|r| r.table_name == OTEL_RESOURCE_INFO_TABLE_NAME)
+ {
+ warn!(
+ "Skipping OTLP resource descriptor synthesis: the request writes \
+ a metric named `{OTEL_RESOURCE_INFO_TABLE_NAME}`"
+ );
+ None
+ } else {
+ resource_info.into_row_insert_requests()?
+ };
+ if resource_info.is_some() {
+ semantic_index.record_scalar(
+ OTEL_RESOURCE_INFO_TABLE_NAME,
+ SEMANTIC_METRIC_TYPE,
+ crate::semantic::METRIC_TYPE_INFO,
+ );
+ semantic_index.record_scalar(
+ OTEL_RESOURCE_INFO_TABLE_NAME,
+ SEMANTIC_METRIC_METADATA_QUALITY,
+ METADATA_QUALITY_DECLARED,
+ );
+ }
+
+ Ok(MetricsConversion {
+ requests,
+ rows,
+ semantic_index,
+ outcome,
+ resource_info,
+ })
}
fn validate_sample_kinds(requests: &RowInsertRequests) -> Result<()> {
@@ -287,30 +344,61 @@ fn from_metric_type(data: &metric::Data) -> MetricType {
}
}
+/// Non-scalar values (bool, arrays, maps, bytes) are not representable as tags.
+fn scalar_value_string(value: Option<&AnyValue>) -> Option {
+ match value.and_then(|v| v.value.as_ref())? {
+ any_value::Value::StringValue(s) => Some(s.clone()),
+ any_value::Value::IntValue(v) => Some(v.to_string()),
+ any_value::Value::DoubleValue(v) => Some(v.to_string()),
+ _ => None,
+ }
+}
+
+/// Prometheus-style `(job, instance)` identity. Per the OTel Prometheus
+/// compatibility spec `job` folds in `service.namespace` when present; a
+/// resource without `service.name` gets no job rather than a fabricated one.
+/// Both fields are optional strings, so a named type keeps a swapped
+/// destructuring from type-checking at the call sites.
+pub(crate) struct ServiceIdentity {
+ pub job: Option,
+ pub instance: Option,
+}
+
+pub(crate) fn service_identity(attrs: &[KeyValue]) -> ServiceIdentity {
+ let mut name = None;
+ let mut namespace = None;
+ let mut instance = None;
+ for kv in attrs {
+ match kv.key.as_str() {
+ KEY_SERVICE_NAME => name = scalar_value_string(kv.value.as_ref()),
+ KEY_SERVICE_NAMESPACE => namespace = scalar_value_string(kv.value.as_ref()),
+ KEY_SERVICE_INSTANCE_ID => instance = scalar_value_string(kv.value.as_ref()),
+ _ => {}
+ }
+ }
+ let job = name.map(|name| match namespace {
+ Some(ns) if !ns.is_empty() => format!("{ns}/{name}"),
+ _ => name,
+ });
+ ServiceIdentity { job, instance }
+}
+
+fn string_key_value(key: &str, value: String) -> KeyValue {
+ KeyValue {
+ key: key.to_string(),
+ value: Some(AnyValue {
+ value: Some(any_value::Value::StringValue(value)),
+ }),
+ }
+}
+
fn process_resource_attrs(attrs: &mut Vec, metric_ctx: &OtlpMetricCtx) {
if metric_ctx.is_legacy {
return;
}
- // remap service.name and service.instance.id to job and instance
- let mut tmp = Vec::with_capacity(2);
- for kv in attrs.iter() {
- match &kv.key as &str {
- KEY_SERVICE_NAME => {
- tmp.push(KeyValue {
- key: JOB_KEY.to_string(),
- value: kv.value.clone(),
- });
- }
- KEY_SERVICE_INSTANCE_ID => {
- tmp.push(KeyValue {
- key: INSTANCE_KEY.to_string(),
- value: kv.value.clone(),
- });
- }
- _ => {}
- }
- }
+ // remap the service identity attributes to job and instance
+ let ServiceIdentity { job, instance } = service_identity(attrs);
// if promote all, then exclude the list, else, include the list
if metric_ctx.promote_all_resource_attrs {
@@ -322,7 +410,12 @@ fn process_resource_attrs(attrs: &mut Vec, metric_ctx: &OtlpMetricCtx)
});
}
- attrs.extend(tmp);
+ if let Some(job) = job {
+ attrs.push(string_key_value(JOB_KEY, job));
+ }
+ if let Some(instance) = instance {
+ attrs.push(string_key_value(INSTANCE_KEY, instance));
+ }
}
fn process_scope_attrs(scope: &ScopeMetrics, metric_ctx: &OtlpMetricCtx) -> Option> {
@@ -536,35 +629,13 @@ fn encode_exponential_histogram(
metric_ctx: &OtlpMetricCtx,
outcome: &mut MetricsIngestOutcome,
) -> Result {
- if !metric_ctx.experimental_enable_exponential_histogram {
+ if let Err(rejection) = exponential_histogram_gate(histogram, metric_ctx) {
reject_data_points(outcome, histogram.data_points.len(), || {
- format!(
- "metric `{name}` uses OTLP exponential histograms; set otlp.experimental_enable_exponential_histogram = true to enable ingestion"
- )
+ rejection.message(name)
})?;
return Ok(false);
}
- match AggregationTemporality::try_from(histogram.aggregation_temporality) {
- Ok(AggregationTemporality::Cumulative) => {}
- Ok(AggregationTemporality::Delta) => {
- reject_data_points(outcome, histogram.data_points.len(), || {
- format!(
- "metric `{name}` uses delta OTLP exponential histograms; only cumulative temporality is supported"
- )
- })?;
- return Ok(false);
- }
- _ => {
- reject_data_points(outcome, histogram.data_points.len(), || {
- format!(
- "metric `{name}` has unspecified OTLP exponential histogram temporality; cumulative temporality is required"
- )
- })?;
- return Ok(false);
- }
- }
-
let column_schema = native_histogram_column_schema().map_err(|error| {
error::InternalSnafu {
err_msg: error.to_string(),
@@ -611,7 +682,46 @@ fn encode_exponential_histogram(
Ok(emitted)
}
-fn exponential_histogram_value(
+pub(crate) enum ExponentialHistogramRejection {
+ Disabled,
+ DeltaTemporality,
+ UnspecifiedTemporality,
+}
+
+impl ExponentialHistogramRejection {
+ fn message(&self, name: &str) -> String {
+ match self {
+ Self::Disabled => format!(
+ "metric `{name}` uses OTLP exponential histograms; set otlp.experimental_enable_exponential_histogram = true to enable ingestion"
+ ),
+ Self::DeltaTemporality => format!(
+ "metric `{name}` uses delta OTLP exponential histograms; only cumulative temporality is supported"
+ ),
+ Self::UnspecifiedTemporality => format!(
+ "metric `{name}` has unspecified OTLP exponential histogram temporality; cumulative temporality is required"
+ ),
+ }
+ }
+}
+
+/// Whole-metric acceptance, decided once for the encoder and for the resource
+/// descriptor, which must not describe a resource whose data was rejected.
+/// Individual points can still fail [`exponential_histogram_value`].
+pub(crate) fn exponential_histogram_gate(
+ histogram: &ExponentialHistogram,
+ metric_ctx: &OtlpMetricCtx,
+) -> std::result::Result<(), ExponentialHistogramRejection> {
+ if !metric_ctx.experimental_enable_exponential_histogram {
+ return Err(ExponentialHistogramRejection::Disabled);
+ }
+ match AggregationTemporality::try_from(histogram.aggregation_temporality) {
+ Ok(AggregationTemporality::Cumulative) => Ok(()),
+ Ok(AggregationTemporality::Delta) => Err(ExponentialHistogramRejection::DeltaTemporality),
+ _ => Err(ExponentialHistogramRejection::UnspecifiedTemporality),
+ }
+}
+
+pub(crate) fn exponential_histogram_value(
data_point: &ExponentialHistogramDataPoint,
) -> std::result::Result<(ValueData, i64), String> {
if data_point.start_time_unix_nano > data_point.time_unix_nano {
@@ -823,29 +933,21 @@ fn write_attributes(
};
let tags = attrs.iter().filter_map(|attr| {
- attr.value
- .as_ref()
- .and_then(|v| v.value.as_ref())
- .and_then(|val| {
- let key = match attribute_type {
- AttributeType::Resource | AttributeType::DataPoint => {
- translate_label_name(&attr.key, metric_ctx.metric_translation_strategy)
- }
- AttributeType::Scope => {
- format!(
- "otel_scope_{}",
- translate_label_name(&attr.key, metric_ctx.metric_translation_strategy)
- )
- }
- AttributeType::Legacy => legacy_normalize_otlp_name(&attr.key),
- };
- match val {
- any_value::Value::StringValue(s) => Some((key, s.clone())),
- any_value::Value::IntValue(v) => Some((key, v.to_string())),
- any_value::Value::DoubleValue(v) => Some((key, v.to_string())),
- _ => None, // TODO(sunng87): allow different type of values
- }
- })
+ // TODO(sunng87): allow different type of values
+ let value = scalar_value_string(attr.value.as_ref())?;
+ let key = match attribute_type {
+ AttributeType::Resource | AttributeType::DataPoint => {
+ translate_label_name(&attr.key, metric_ctx.metric_translation_strategy)
+ }
+ AttributeType::Scope => {
+ format!(
+ "otel_scope_{}",
+ translate_label_name(&attr.key, metric_ctx.metric_translation_strategy)
+ )
+ }
+ AttributeType::Legacy => legacy_normalize_otlp_name(&attr.key),
+ };
+ Some((key, value))
});
row_writer::write_tags(writer, tags, row)?;
@@ -1295,6 +1397,169 @@ mod tests {
}
}
+ fn descriptor_ctx() -> OtlpMetricCtx {
+ OtlpMetricCtx {
+ resource_info: true,
+ ..Default::default()
+ }
+ }
+
+ fn attr_value(attrs: &[KeyValue], key: &str) -> Option {
+ attrs
+ .iter()
+ .find(|kv| kv.key == key)
+ .and_then(|kv| scalar_value_string(kv.value.as_ref()))
+ }
+
+ fn gauge_request(
+ resource_attrs: Vec,
+ metric_name: &str,
+ ) -> ExportMetricsServiceRequest {
+ use otel_arrow_rust::proto::opentelemetry::resource::v1::Resource;
+ ExportMetricsServiceRequest {
+ resource_metrics: vec![ResourceMetrics {
+ resource: Some(Resource {
+ attributes: resource_attrs,
+ ..Default::default()
+ }),
+ scope_metrics: vec![ScopeMetrics {
+ metrics: vec![Metric {
+ name: metric_name.to_string(),
+ data: Some(metric::Data::Gauge(Gauge {
+ data_points: vec![NumberDataPoint {
+ time_unix_nano: 1_000_000,
+ value: Some(Value::AsInt(1)),
+ ..Default::default()
+ }],
+ })),
+ ..Default::default()
+ }],
+ ..Default::default()
+ }],
+ ..Default::default()
+ }],
+ }
+ }
+
+ fn column_names(request: &RowInsertRequests, table: &str) -> Vec {
+ request
+ .inserts
+ .iter()
+ .find(|r| r.table_name == table)
+ .unwrap_or_else(|| panic!("missing table {table}"))
+ .rows
+ .as_ref()
+ .unwrap()
+ .schema
+ .iter()
+ .map(|c| c.column_name.clone())
+ .collect()
+ }
+
+ #[test]
+ fn test_conversion_synthesizes_resource_descriptor() {
+ set_default_prefix(None).unwrap();
+ let request = gauge_request(
+ vec![keyvalue("service.name", "api"), keyvalue("host.id", "h-1")],
+ "my_gauge",
+ );
+ let conversion = to_grpc_insert_requests(request, &mut descriptor_ctx()).unwrap();
+
+ // descriptor keeps raw OTel keys while the metric table's labels went
+ // through the (default underscore-escaping) translation strategy
+ let resource_info = conversion.resource_info.expect("descriptor synthesized");
+ let descriptor_cols = column_names(&resource_info, OTEL_RESOURCE_INFO_TABLE_NAME);
+ assert!(descriptor_cols.contains(&"host.id".to_string()));
+ assert!(descriptor_cols.contains(&"service.name".to_string()));
+ assert!(descriptor_cols.contains(&"job".to_string()));
+ let metric_cols = column_names(&conversion.requests, "my_gauge");
+ assert!(metric_cols.contains(&"service_name".to_string()));
+ assert!(!metric_cols.contains(&"service.name".to_string()));
+ // host.id is not in the promote list: only the descriptor keeps it
+ assert!(!metric_cols.contains(&"host_id".to_string()));
+
+ let decoded = decode(&conversion.semantic_index);
+ let t = &decoded[OTEL_RESOURCE_INFO_TABLE_NAME];
+ assert_eq!(
+ t.get(SEMANTIC_METRIC_TYPE).map(String::as_str),
+ Some("info")
+ );
+ assert_eq!(
+ t.get(SEMANTIC_METRIC_METADATA_QUALITY).map(String::as_str),
+ Some("declared")
+ );
+ }
+
+ #[test]
+ fn test_conversion_skips_descriptor_for_legacy_mode() {
+ set_default_prefix(None).unwrap();
+ let request = gauge_request(
+ vec![keyvalue("service.name", "api"), keyvalue("host.id", "h-1")],
+ "my_gauge",
+ );
+ let mut ctx = OtlpMetricCtx {
+ is_legacy: true,
+ ..descriptor_ctx()
+ };
+ let conversion = to_grpc_insert_requests(request, &mut ctx).unwrap();
+ assert!(conversion.resource_info.is_none());
+
+ // legacy tables predate job/instance and the promote filter; adding
+ // either would alter the schema of tables already in use
+ let cols = column_names(&conversion.requests, "my_gauge");
+ assert!(!cols.contains(&"job".to_string()));
+ assert!(cols.contains(&"host_id".to_string()));
+ }
+
+ #[test]
+ fn test_conversion_skips_descriptor_on_metric_name_collision() {
+ set_default_prefix(None).unwrap();
+ let request = gauge_request(
+ vec![keyvalue("service.name", "api")],
+ OTEL_RESOURCE_INFO_TABLE_NAME,
+ );
+ let conversion = to_grpc_insert_requests(request, &mut descriptor_ctx()).unwrap();
+ assert!(conversion.resource_info.is_none());
+ // the metric itself still goes through the main path
+ assert!(
+ conversion
+ .requests
+ .inserts
+ .iter()
+ .any(|r| r.table_name == OTEL_RESOURCE_INFO_TABLE_NAME)
+ );
+ }
+
+ #[test]
+ fn test_job_composition_follows_service_namespace() {
+ let mut attrs = vec![
+ keyvalue("service.name", "api"),
+ keyvalue("service.namespace", "shop"),
+ ];
+ process_resource_attrs(&mut attrs, &OtlpMetricCtx::default());
+ assert_eq!(attr_value(&attrs, "job").as_deref(), Some("shop/api"));
+
+ let mut attrs = vec![keyvalue("service.name", "api")];
+ process_resource_attrs(&mut attrs, &OtlpMetricCtx::default());
+ assert_eq!(attr_value(&attrs, "job").as_deref(), Some("api"));
+
+ let mut attrs = vec![
+ keyvalue("service.name", "api"),
+ keyvalue("service.namespace", ""),
+ ];
+ process_resource_attrs(&mut attrs, &OtlpMetricCtx::default());
+ assert_eq!(attr_value(&attrs, "job").as_deref(), Some("api"));
+
+ // no service.name: no job is fabricated
+ let mut attrs = vec![
+ keyvalue("service.namespace", "shop"),
+ keyvalue("service.instance.id", "inst-1"),
+ ];
+ process_resource_attrs(&mut attrs, &OtlpMetricCtx::default());
+ assert_eq!(attr_value(&attrs, "job"), None);
+ assert_eq!(attr_value(&attrs, "instance").as_deref(), Some("inst-1"));
+ }
+
#[test]
fn test_encode_gauge() {
let mut tables = MultiTableData::default();
@@ -2104,8 +2369,12 @@ mod tests {
AggregationTemporality::Cumulative,
),
]);
- let (requests, _, semantic_index, outcome) =
- to_grpc_insert_requests(request, &mut OtlpMetricCtx::default()).unwrap();
+ let MetricsConversion {
+ requests,
+ semantic_index,
+ outcome,
+ ..
+ } = to_grpc_insert_requests(request, &mut OtlpMetricCtx::default()).unwrap();
assert_eq!(outcome.accepted_data_points, 1);
assert_eq!(outcome.rejected_data_points, 1);
@@ -2127,8 +2396,9 @@ mod tests {
vec![],
AggregationTemporality::Cumulative,
)]);
- let (_, _, _, outcome) =
- to_grpc_insert_requests(empty, &mut OtlpMetricCtx::default()).unwrap();
+ let outcome = to_grpc_insert_requests(empty, &mut OtlpMetricCtx::default())
+ .unwrap()
+ .outcome;
assert_eq!(outcome.rejected_data_points, 0);
assert_eq!(outcome.error_message, None);
}
@@ -2208,8 +2478,12 @@ mod tests {
.collect(),
};
- let (requests, rows, _, outcome) =
- to_grpc_insert_requests(request, &mut OtlpMetricCtx::default()).unwrap();
+ let MetricsConversion {
+ requests,
+ rows,
+ outcome,
+ ..
+ } = to_grpc_insert_requests(request, &mut OtlpMetricCtx::default()).unwrap();
assert_eq!(outcome.accepted_data_points, 2);
assert_eq!(rows, 6);
@@ -2243,8 +2517,13 @@ mod tests {
experimental_enable_exponential_histogram: true,
..Default::default()
};
- let (requests, rows, semantic_index, outcome) =
- to_grpc_insert_requests(request, &mut ctx).unwrap();
+ let MetricsConversion {
+ requests,
+ rows,
+ semantic_index,
+ outcome,
+ ..
+ } = to_grpc_insert_requests(request, &mut ctx).unwrap();
assert_eq!(outcome.accepted_data_points, 0);
assert_eq!(outcome.rejected_data_points, 1);
@@ -2269,14 +2548,17 @@ mod tests {
experimental_enable_exponential_histogram: true,
..Default::default()
};
- let (new_requests, _, _, _) =
- to_grpc_insert_requests(request.clone(), &mut new_ctx).unwrap();
+ let new_requests = to_grpc_insert_requests(request.clone(), &mut new_ctx)
+ .unwrap()
+ .requests;
let mut legacy_ctx = OtlpMetricCtx {
experimental_enable_exponential_histogram: true,
is_legacy: true,
..Default::default()
};
- let (legacy_requests, _, _, _) = to_grpc_insert_requests(request, &mut legacy_ctx).unwrap();
+ let legacy_requests = to_grpc_insert_requests(request, &mut legacy_ctx)
+ .unwrap()
+ .requests;
let new_insert = &new_requests.inserts[0];
let legacy_insert = &legacy_requests.inserts[0];
@@ -2325,8 +2607,9 @@ mod tests {
vec![ExponentialHistogramDataPoint::default()],
AggregationTemporality::Cumulative,
)]);
- let (_, _, _, outcome) =
- to_grpc_insert_requests(request, &mut OtlpMetricCtx::default()).unwrap();
+ let outcome = to_grpc_insert_requests(request, &mut OtlpMetricCtx::default())
+ .unwrap()
+ .outcome;
assert_eq!(outcome.rejected_data_points, 1);
assert!(outcome.error_message.unwrap().len() <= MAX_REJECTION_MESSAGE_BYTES);
diff --git a/src/servers/src/otlp/metrics/resource_info.rs b/src/servers/src/otlp/metrics/resource_info.rs
new file mode 100644
index 0000000000..cd3427f716
--- /dev/null
+++ b/src/servers/src/otlp/metrics/resource_info.rs
@@ -0,0 +1,393 @@
+// Copyright 2023 Greptime Team
+//
+// Licensed under the Apache License, Version 2.0 (the "License");
+// you may not use this file except in compliance with the License.
+// You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing, software
+// distributed under the License is distributed on an "AS IS" BASIS,
+// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+// See the License for the specific language governing permissions and
+// limitations under the License.
+
+//! Resource descriptor synthesized from OTLP metrics requests, so the entity
+//! graph reads one table instead of scanning every logical metric table for
+//! attributes the promote filter may have dropped.
+//!
+//! Columns are a fixed allowlist under the raw OTel attribute names: the
+//! conventions whitelist matches fixed names, so they must not follow the
+//! per-request label translation strategy or the promote/ignore headers.
+
+use std::collections::BTreeMap;
+
+use api::v1::RowInsertRequests;
+use common_catalog::consts::SEMANTIC_GRAPH_WINDOW_NANOS;
+use common_grpc::precision::Precision;
+use common_query::prelude::{greptime_timestamp, greptime_value};
+use otel_arrow_rust::proto::opentelemetry::common::v1::KeyValue;
+use otel_arrow_rust::proto::opentelemetry::metrics::v1::{ResourceMetrics, metric};
+use session::protocol_ctx::OtlpMetricCtx;
+
+use crate::error::Result;
+use crate::otlp::metrics::{
+ INSTANCE_KEY, JOB_KEY, ServiceIdentity, exponential_histogram_gate,
+ exponential_histogram_value, scalar_value_string, service_identity,
+};
+use crate::otlp::trace::{
+ KEY_CONTAINER_ID, KEY_CONTAINER_NAME, KEY_HOST_ID, KEY_HOST_NAME, KEY_K8S_NAMESPACE_NAME,
+ KEY_K8S_POD_NAME, KEY_K8S_POD_UID, KEY_SERVICE_NAME, KEY_SERVICE_NAMESPACE,
+};
+use crate::row_writer::{self, MultiTableData};
+
+/// Prefixed like the other engine-managed tables, so a user metric is
+/// unlikely to claim the name.
+pub const OTEL_RESOURCE_INFO_TABLE_NAME: &str = "greptime_otel_resource_info";
+
+/// Attributes projected under their raw OTel keys. `service.instance.id` is
+/// absent on purpose: it lands in `instance`.
+///
+/// Matched instead of scanned: this runs for every attribute of every
+/// resource, and the compiler turns it into a length-and-prefix dispatch.
+fn is_projected_attr(key: &str) -> bool {
+ matches!(
+ key,
+ KEY_SERVICE_NAME
+ | KEY_SERVICE_NAMESPACE
+ | KEY_HOST_ID
+ | KEY_HOST_NAME
+ | KEY_CONTAINER_ID
+ | KEY_CONTAINER_NAME
+ | KEY_K8S_POD_UID
+ | KEY_K8S_POD_NAME
+ | KEY_K8S_NAMESPACE_NAME
+ )
+}
+
+/// Upper bound of [`is_projected_attr`] plus the derived `job`/`instance`,
+/// used to size the per-row buffers.
+const MAX_PROJECTED_TAGS: usize = 11;
+
+/// Projected attributes (sorted `(name, value)` pairs) -> graph window ->
+/// the newest data-point time seen in that window, which is what the row for
+/// that window is stamped with.
+///
+/// Windows are an inner map so the attributes are stored, and moved, once per
+/// resource. They are keyed separately because one request may carry data for
+/// several of them: a single row per resource would describe only the newest
+/// window, leaving the earlier ones with metric rows but no entities.
+#[derive(Debug, Default)]
+pub struct ResourceInfoData {
+ rows: BTreeMap, BTreeMap>,
+}
+
+impl ResourceInfoData {
+ /// Takes the raw attributes, before the promote filter runs on them.
+ pub fn observe(
+ &mut self,
+ raw_attrs: &[KeyValue],
+ resource: &ResourceMetrics,
+ metric_ctx: &OtlpMetricCtx,
+ ) {
+ let mut tags = Vec::with_capacity(MAX_PROJECTED_TAGS);
+ let ServiceIdentity { job, instance } = service_identity(raw_attrs);
+ if let Some(job) = job {
+ tags.push((JOB_KEY.to_string(), job));
+ }
+ if let Some(instance) = instance {
+ tags.push((INSTANCE_KEY.to_string(), instance));
+ }
+ for kv in raw_attrs {
+ if is_projected_attr(&kv.key)
+ && let Some(value) = scalar_value_string(kv.value.as_ref())
+ {
+ tags.push((kv.key.clone(), value));
+ }
+ }
+ if tags.is_empty() {
+ return;
+ }
+ // Sorted so equal attribute sets share a key, and so the emitted
+ // columns keep a stable order.
+ tags.sort_unstable();
+
+ let mut observed: BTreeMap = BTreeMap::new();
+ for_each_encoded_time(resource, metric_ctx, |ts| {
+ let window = ts - ts.rem_euclid(SEMANTIC_GRAPH_WINDOW_NANOS);
+ observed
+ .entry(window)
+ .and_modify(|newest| *newest = (*newest).max(ts))
+ .or_insert(ts);
+ });
+ if observed.is_empty() {
+ return;
+ }
+
+ let windows = self.rows.entry(tags).or_default();
+ for (window, newest) in observed {
+ windows
+ .entry(window)
+ .and_modify(|seen| *seen = (*seen).max(newest))
+ .or_insert(newest);
+ }
+ }
+
+ /// Every projected attribute becomes a tag, so auto-create puts it in the
+ /// primary key where the conventions expect it.
+ pub fn into_row_insert_requests(self) -> Result