Skip to content

Commit 8e38457

Browse files
Merge pull request #674 from gsvigruha/llmobs-add-dataset-records-and-experiment-events
feat(llm-obs): add dataset-records reads + experiment eval-metric submit
2 parents 17d55fd + d27b1e2 commit 8e38457

4 files changed

Lines changed: 390 additions & 3 deletions

File tree

README.md

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -180,7 +180,7 @@ list of commands as built.
180180
| Data Streams (Kafka) || `kafka topic-configs`, `kafka broker-configs`, `kafka client-configs`, `kafka read-messages` | **Experimental** — Kafka cluster inspection via Datadog |
181181
| Restricted Datasets || `datasets list`, `datasets get`, `datasets create`, `datasets update`, `datasets delete` | Restricted dataset management for data access control |
182182
| Observability Pipelines || `obs-pipelines list`, `obs-pipelines get`, `obs-pipelines create`, `obs-pipelines update`, `obs-pipelines delete`, `obs-pipelines validate` | Full pipeline CRUD and validation |
183-
| LLM Observability || `llm-obs projects`, `llm-obs experiments`, `llm-obs datasets` | **New** — LLM Obs projects, experiments, and dataset management |
183+
| LLM Observability || `llm-obs projects`, `llm-obs experiments`, `llm-obs datasets` | **New** — LLM Obs projects, experiments (incl. `events submit`), and dataset management (incl. `datasets records` / `records-full`) |
184184
| Reference Tables || `reference-tables list`, `reference-tables get`, `reference-tables create`, `reference-tables batch-query` | **New** — Reference table management for log enrichment |
185185
| Miscellaneous || `misc ip-ranges`, `misc status` | IP ranges and status |
186186
| App Builder || `app-builder list`, `app-builder get`, `app-builder create`, `app-builder update`, `app-builder delete`, `app-builder publish` | Low-code app management with publish/unpublish and batch delete |

docs/COMMANDS.md

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -66,7 +66,7 @@ pup <domain> <subgroup> <action> [options] # Nested commands
6666
| data-deletion | requests (list, create, cancel) | src/commands/data_deletion.rs ||
6767
| data-governance | scanner-rules (list) | src/commands/data_governance.rs ||
6868
| obs-pipelines | list, get, create, update, delete, validate | src/commands/obs_pipelines.rs ||
69-
| llm-obs | projects (create, list), experiments (create, list, update, delete, summary, events (list, get), metric-values, dimension-values), datasets (create, list, batch-update, clone, restore), spans (search) | src/commands/llm_obs.rs ||
69+
| llm-obs | projects (create, list), experiments (create, list, update, delete, summary, events (list, get, submit), metric-values, dimension-values), datasets (create, list, batch-update, clone, restore, records, records-full), spans (search) | src/commands/llm_obs.rs ||
7070
| reference-tables | list, get, create, batch-query | src/commands/reference_tables.rs ||
7171
| network | flows list, devices (list, get, interfaces, tags), interfaces (list, update) | src/commands/network.rs ||
7272
| cloud | aws, gcp, azure, oci | src/commands/cloud.rs ||
@@ -300,7 +300,7 @@ steps) bypass `format_and_print` and do not honor `--jq`.
300300

301301
### v0.28.0 — New Command Groups and Full Pipeline Implementation
302302

303-
-**llm-obs** (new) — LLM Observability: projects (create, list), experiments (create, list, update, delete, summary, events (list, get), metric-values, dimension-values), datasets (create, list, batch-update, clone, restore), spans (search)
303+
-**llm-obs** (new) — LLM Observability: projects (create, list), experiments (create, list, update, delete, summary, events (list, get, submit), metric-values, dimension-values), datasets (create, list, batch-update, clone, restore, records, records-full), spans (search)
304304
-**reference-tables** (new) — Reference table management (list, get, create, batch-query)
305305
-**obs-pipelines** (upgraded from placeholder) — Full CRUD: list, get, create, update, delete, validate
306306
- **costs** — Added cloud cost configs: `aws-config`, `azure-config`, `gcp-config` (list, get, create, delete each)

src/commands/llm_obs.rs

Lines changed: 278 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -156,6 +156,71 @@ pub async fn datasets_restore(
156156
Ok(())
157157
}
158158

159+
// ---- Dataset records (no typed equivalent — unstable MCP endpoints) ----
160+
161+
#[allow(clippy::too_many_arguments)]
162+
pub async fn datasets_records(
163+
cfg: &Config,
164+
project_id: &str,
165+
dataset_id: &str,
166+
record_ids: Option<Vec<String>>,
167+
tags: Option<Vec<String>>,
168+
canonical_id: Option<String>,
169+
dataset_version: Option<i64>,
170+
limit: u32,
171+
cursor: Option<String>,
172+
compute_schema: Option<bool>,
173+
) -> Result<()> {
174+
let mut body = serde_json::json!({
175+
"project_id": project_id,
176+
"dataset_id": dataset_id,
177+
"limit": limit,
178+
});
179+
if let Some(ids) = record_ids {
180+
body["record_ids"] = serde_json::json!(ids);
181+
}
182+
if let Some(t) = tags {
183+
body["tags"] = serde_json::json!(t);
184+
}
185+
if let Some(c) = canonical_id {
186+
body["canonical_id"] = serde_json::json!(c);
187+
}
188+
if let Some(v) = dataset_version {
189+
body["dataset_version"] = serde_json::json!(v);
190+
}
191+
if let Some(c) = cursor {
192+
body["cursor"] = serde_json::json!(c);
193+
}
194+
if let Some(cs) = compute_schema {
195+
body["compute_schema"] = serde_json::json!(cs);
196+
}
197+
let resp = raw_client::raw_post(cfg, "/api/unstable/llm-obs-mcp/v1/dataset/records", body)
198+
.await
199+
.map_err(|e| anyhow::anyhow!("failed to get dataset records: {e:?}"))?;
200+
formatter::output(cfg, &resp)
201+
}
202+
203+
pub async fn datasets_records_full(
204+
cfg: &Config,
205+
project_id: &str,
206+
dataset_id: &str,
207+
record_ids: Vec<String>,
208+
) -> Result<()> {
209+
let body = serde_json::json!({
210+
"project_id": project_id,
211+
"dataset_id": dataset_id,
212+
"record_ids": record_ids,
213+
});
214+
let resp = raw_client::raw_post(
215+
cfg,
216+
"/api/unstable/llm-obs-mcp/v1/dataset/records-full",
217+
body,
218+
)
219+
.await
220+
.map_err(|e| anyhow::anyhow!("failed to get full dataset records: {e:?}"))?;
221+
formatter::output(cfg, &resp)
222+
}
223+
159224
// ---- Experiment analytics (no typed equivalent — unstable MCP endpoints) ----
160225

161226
pub async fn experiments_summary(cfg: &Config, experiment_id: &str) -> Result<()> {
@@ -214,6 +279,34 @@ pub async fn experiments_events_get(
214279
formatter::output(cfg, &resp)
215280
}
216281

282+
pub async fn experiments_events_submit(
283+
cfg: &Config,
284+
experiment_id: &str,
285+
metrics: &str,
286+
tags: Option<Vec<String>>,
287+
) -> Result<()> {
288+
// Mirrors the submit_llmobs_experiment_events MCP tool: experiment_id + metrics
289+
// (array) + optional tags. metrics is passed through as-is; the server validates it.
290+
let metrics: serde_json::Value = serde_json::from_str(metrics).map_err(|e| {
291+
anyhow::anyhow!("--metrics must be a JSON array of eval-metric events: {e}")
292+
})?;
293+
let mut body = serde_json::json!({
294+
"experiment_id": experiment_id,
295+
"metrics": metrics,
296+
});
297+
if let Some(t) = tags {
298+
body["tags"] = serde_json::json!(t);
299+
}
300+
let resp = raw_client::raw_post(
301+
cfg,
302+
"/api/unstable/llm-obs-mcp/v1/experiment/ingest-events",
303+
body,
304+
)
305+
.await
306+
.map_err(|e| anyhow::anyhow!("failed to submit experiment events: {e:?}"))?;
307+
formatter::output(cfg, &resp)
308+
}
309+
217310
pub async fn experiments_metric_values(
218311
cfg: &Config,
219312
experiment_id: &str,
@@ -2851,4 +2944,189 @@ mod tests {
28512944
cleanup_env();
28522945
std::env::remove_var("DD_TOKEN_STORAGE");
28532946
}
2947+
2948+
#[tokio::test]
2949+
async fn test_llm_obs_datasets_records() {
2950+
let _lock = lock_env().await;
2951+
let mut server = mockito::Server::new_async().await;
2952+
let cfg = test_config(&server.url());
2953+
2954+
let body = r#"{"status":"success","data":{"records":[{"id":"rec-1"}],"schema_summary":{},"returned":1,"truncated":false,"next_cursor":null}}"#;
2955+
let _mock = mock_post(
2956+
&mut server,
2957+
"/api/unstable/llm-obs-mcp/v1/dataset/records",
2958+
200,
2959+
body,
2960+
)
2961+
.await;
2962+
2963+
let result = super::datasets_records(
2964+
&cfg, "proj-1", "ds-1", None, None, None, None, 10, None, None,
2965+
)
2966+
.await;
2967+
assert!(
2968+
result.is_ok(),
2969+
"datasets_records failed: {:?}",
2970+
result.err()
2971+
);
2972+
cleanup_env();
2973+
}
2974+
2975+
#[tokio::test]
2976+
async fn test_llm_obs_datasets_records_filtered() {
2977+
let _lock = lock_env().await;
2978+
let mut server = mockito::Server::new_async().await;
2979+
let cfg = test_config(&server.url());
2980+
2981+
let body = r#"{"status":"success","data":{"records":[],"returned":0,"truncated":false,"next_cursor":null}}"#;
2982+
let _mock = mock_post(
2983+
&mut server,
2984+
"/api/unstable/llm-obs-mcp/v1/dataset/records",
2985+
200,
2986+
body,
2987+
)
2988+
.await;
2989+
2990+
let result = super::datasets_records(
2991+
&cfg,
2992+
"proj-1",
2993+
"ds-1",
2994+
Some(vec!["rec-1".into(), "rec-2".into()]),
2995+
Some(vec!["env:prod".into()]),
2996+
Some("canon-1".into()),
2997+
Some(3),
2998+
5,
2999+
Some("cursor-abc".into()),
3000+
Some(false),
3001+
)
3002+
.await;
3003+
assert!(
3004+
result.is_ok(),
3005+
"datasets_records filtered failed: {:?}",
3006+
result.err()
3007+
);
3008+
cleanup_env();
3009+
}
3010+
3011+
#[tokio::test]
3012+
async fn test_llm_obs_datasets_records_500() {
3013+
let _lock = lock_env().await;
3014+
let mut server = mockito::Server::new_async().await;
3015+
let cfg = test_config(&server.url());
3016+
3017+
let _mock = mock_post(
3018+
&mut server,
3019+
"/api/unstable/llm-obs-mcp/v1/dataset/records",
3020+
500,
3021+
r#"{"errors":["internal server error"]}"#,
3022+
)
3023+
.await;
3024+
3025+
let result = super::datasets_records(
3026+
&cfg, "proj-1", "ds-1", None, None, None, None, 10, None, None,
3027+
)
3028+
.await;
3029+
assert!(result.is_err(), "should fail on 500");
3030+
assert!(result.unwrap_err().to_string().contains("500"));
3031+
cleanup_env();
3032+
}
3033+
3034+
#[tokio::test]
3035+
async fn test_llm_obs_datasets_records_full() {
3036+
let _lock = lock_env().await;
3037+
let mut server = mockito::Server::new_async().await;
3038+
let cfg = test_config(&server.url());
3039+
3040+
let body = r#"{"status":"success","data":{"records":[{"id":"rec-1","input":{"prompt":"hello"},"expected_output":"world"}]}}"#;
3041+
let _mock = mock_post(
3042+
&mut server,
3043+
"/api/unstable/llm-obs-mcp/v1/dataset/records-full",
3044+
200,
3045+
body,
3046+
)
3047+
.await;
3048+
3049+
let result =
3050+
super::datasets_records_full(&cfg, "proj-1", "ds-1", vec!["rec-1".into()]).await;
3051+
assert!(
3052+
result.is_ok(),
3053+
"datasets_records_full failed: {:?}",
3054+
result.err()
3055+
);
3056+
cleanup_env();
3057+
}
3058+
3059+
#[tokio::test]
3060+
async fn test_llm_obs_experiments_events_submit() {
3061+
let _lock = lock_env().await;
3062+
std::env::set_var("DD_TOKEN_STORAGE", "file");
3063+
let mut server = mockito::Server::new_async().await;
3064+
let cfg = test_config(&server.url());
3065+
3066+
let resp_body = r#"{"status":"success","data":{"accepted":1}}"#;
3067+
let _mock = mock_post(
3068+
&mut server,
3069+
"/api/unstable/llm-obs-mcp/v1/experiment/ingest-events",
3070+
200,
3071+
resp_body,
3072+
)
3073+
.await;
3074+
3075+
let result = super::experiments_events_submit(
3076+
&cfg,
3077+
"exp-1",
3078+
r#"[{"label":"accuracy","metric_type":"score","score_value":0.9}]"#,
3079+
Some(vec!["run:1".to_string()]),
3080+
)
3081+
.await;
3082+
assert!(
3083+
result.is_ok(),
3084+
"experiments_events_submit failed: {:?}",
3085+
result.err()
3086+
);
3087+
cleanup_env();
3088+
std::env::remove_var("DD_TOKEN_STORAGE");
3089+
}
3090+
3091+
#[tokio::test]
3092+
async fn test_llm_obs_experiments_events_submit_400() {
3093+
let _lock = lock_env().await;
3094+
std::env::set_var("DD_TOKEN_STORAGE", "file");
3095+
let mut server = mockito::Server::new_async().await;
3096+
let cfg = test_config(&server.url());
3097+
3098+
let _mock = server
3099+
.mock("POST", mockito::Matcher::Any)
3100+
.match_query(mockito::Matcher::Any)
3101+
.with_status(400)
3102+
.with_header("content-type", "application/json")
3103+
.with_body(r#"{"errors":["bad request"]}"#)
3104+
.create_async()
3105+
.await;
3106+
3107+
let result = super::experiments_events_submit(
3108+
&cfg,
3109+
"exp-1",
3110+
r#"[{"label":"accuracy","metric_type":"score","score_value":0.9}]"#,
3111+
None,
3112+
)
3113+
.await;
3114+
assert!(result.is_err(), "should fail on 400");
3115+
cleanup_env();
3116+
std::env::remove_var("DD_TOKEN_STORAGE");
3117+
}
3118+
3119+
#[tokio::test]
3120+
async fn test_llm_obs_experiments_events_submit_invalid_json() {
3121+
let _lock = lock_env().await;
3122+
std::env::set_var("DD_TOKEN_STORAGE", "file");
3123+
let server = mockito::Server::new_async().await;
3124+
let cfg = test_config(&server.url());
3125+
3126+
// Malformed --metrics should fail locally before any request is made.
3127+
let result = super::experiments_events_submit(&cfg, "exp-1", "not-json", None).await;
3128+
assert!(result.is_err(), "should fail on invalid metrics JSON");
3129+
cleanup_env();
3130+
std::env::remove_var("DD_TOKEN_STORAGE");
3131+
}
28543132
}

0 commit comments

Comments
 (0)