Skip to content

Commit 440c8cb

Browse files
tillwfclaude
andcommitted
feat(llm-obs): add dataset-records reads + experiment eval-metric submit
Add three llm-obs subcommands backed by the unstable llm-obs-mcp v1 endpoints, closing the gap vs the LLM Obs MCP toolset: - experiments events submit -> POST /experiment/ingest-events (submit_llmobs_experiment_events; metrics/tags via --file JSON, experiment_id from the positional arg) - datasets records -> POST /dataset/records (get_llmobs_dataset_records; record-ids/tags/canonical-id/ dataset-version/limit/cursor/compute-schema flags) - datasets records-full -> POST /dataset/records-full (get_llmobs_full_dataset_records; 1-3 record ids) These are the tools the agent-observability auto-experiment skill needs that pup previously lacked. Docs (README coverage table, COMMANDS.md) updated; unit tests added for success + error paths. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
1 parent 17d55fd commit 440c8cb

4 files changed

Lines changed: 361 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: 257 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,29 @@ 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+
file: &str,
286+
) -> Result<()> {
287+
let mut body: serde_json::Value = util::read_json_file(file)?;
288+
if !body.is_object() {
289+
return Err(anyhow::anyhow!(
290+
"events file must contain a JSON object with a \"metrics\" array (and optional \"tags\")"
291+
));
292+
}
293+
// The experiment_id is taken from the positional arg; it overrides any value in the file.
294+
body["experiment_id"] = serde_json::json!(experiment_id);
295+
let resp = raw_client::raw_post(
296+
cfg,
297+
"/api/unstable/llm-obs-mcp/v1/experiment/ingest-events",
298+
body,
299+
)
300+
.await
301+
.map_err(|e| anyhow::anyhow!("failed to submit experiment events: {e:?}"))?;
302+
formatter::output(cfg, &resp)
303+
}
304+
217305
pub async fn experiments_metric_values(
218306
cfg: &Config,
219307
experiment_id: &str,
@@ -2851,4 +2939,173 @@ mod tests {
28512939
cleanup_env();
28522940
std::env::remove_var("DD_TOKEN_STORAGE");
28532941
}
2942+
2943+
#[tokio::test]
2944+
async fn test_llm_obs_datasets_records() {
2945+
let _lock = lock_env().await;
2946+
let mut server = mockito::Server::new_async().await;
2947+
let cfg = test_config(&server.url());
2948+
2949+
let body = r#"{"status":"success","data":{"records":[{"id":"rec-1"}],"schema_summary":{},"returned":1,"truncated":false,"next_cursor":null}}"#;
2950+
let _mock = mock_post(
2951+
&mut server,
2952+
"/api/unstable/llm-obs-mcp/v1/dataset/records",
2953+
200,
2954+
body,
2955+
)
2956+
.await;
2957+
2958+
let result = super::datasets_records(
2959+
&cfg, "proj-1", "ds-1", None, None, None, None, 10, None, None,
2960+
)
2961+
.await;
2962+
assert!(
2963+
result.is_ok(),
2964+
"datasets_records failed: {:?}",
2965+
result.err()
2966+
);
2967+
cleanup_env();
2968+
}
2969+
2970+
#[tokio::test]
2971+
async fn test_llm_obs_datasets_records_filtered() {
2972+
let _lock = lock_env().await;
2973+
let mut server = mockito::Server::new_async().await;
2974+
let cfg = test_config(&server.url());
2975+
2976+
let body = r#"{"status":"success","data":{"records":[],"returned":0,"truncated":false,"next_cursor":null}}"#;
2977+
let _mock = mock_post(
2978+
&mut server,
2979+
"/api/unstable/llm-obs-mcp/v1/dataset/records",
2980+
200,
2981+
body,
2982+
)
2983+
.await;
2984+
2985+
let result = super::datasets_records(
2986+
&cfg,
2987+
"proj-1",
2988+
"ds-1",
2989+
Some(vec!["rec-1".into(), "rec-2".into()]),
2990+
Some(vec!["env:prod".into()]),
2991+
Some("canon-1".into()),
2992+
Some(3),
2993+
5,
2994+
Some("cursor-abc".into()),
2995+
Some(false),
2996+
)
2997+
.await;
2998+
assert!(
2999+
result.is_ok(),
3000+
"datasets_records filtered failed: {:?}",
3001+
result.err()
3002+
);
3003+
cleanup_env();
3004+
}
3005+
3006+
#[tokio::test]
3007+
async fn test_llm_obs_datasets_records_500() {
3008+
let _lock = lock_env().await;
3009+
let mut server = mockito::Server::new_async().await;
3010+
let cfg = test_config(&server.url());
3011+
3012+
let _mock = mock_post(
3013+
&mut server,
3014+
"/api/unstable/llm-obs-mcp/v1/dataset/records",
3015+
500,
3016+
r#"{"errors":["internal server error"]}"#,
3017+
)
3018+
.await;
3019+
3020+
let result = super::datasets_records(
3021+
&cfg, "proj-1", "ds-1", None, None, None, None, 10, None, None,
3022+
)
3023+
.await;
3024+
assert!(result.is_err(), "should fail on 500");
3025+
assert!(result.unwrap_err().to_string().contains("500"));
3026+
cleanup_env();
3027+
}
3028+
3029+
#[tokio::test]
3030+
async fn test_llm_obs_datasets_records_full() {
3031+
let _lock = lock_env().await;
3032+
let mut server = mockito::Server::new_async().await;
3033+
let cfg = test_config(&server.url());
3034+
3035+
let body = r#"{"status":"success","data":{"records":[{"id":"rec-1","input":{"prompt":"hello"},"expected_output":"world"}]}}"#;
3036+
let _mock = mock_post(
3037+
&mut server,
3038+
"/api/unstable/llm-obs-mcp/v1/dataset/records-full",
3039+
200,
3040+
body,
3041+
)
3042+
.await;
3043+
3044+
let result =
3045+
super::datasets_records_full(&cfg, "proj-1", "ds-1", vec!["rec-1".into()]).await;
3046+
assert!(
3047+
result.is_ok(),
3048+
"datasets_records_full failed: {:?}",
3049+
result.err()
3050+
);
3051+
cleanup_env();
3052+
}
3053+
3054+
#[tokio::test]
3055+
async fn test_llm_obs_experiments_events_submit() {
3056+
let _lock = lock_env().await;
3057+
std::env::set_var("DD_TOKEN_STORAGE", "file");
3058+
let mut server = mockito::Server::new_async().await;
3059+
let cfg = test_config(&server.url());
3060+
3061+
let tmp = write_temp_json(
3062+
"pup_test_exp_ingest_events.json",
3063+
r#"{"metrics":[{"label":"accuracy","metric_type":"score","score_value":0.9}],"tags":["run:1"]}"#,
3064+
);
3065+
let resp_body = r#"{"status":"success","data":{"accepted":1}}"#;
3066+
let _mock = mock_post(
3067+
&mut server,
3068+
"/api/unstable/llm-obs-mcp/v1/experiment/ingest-events",
3069+
200,
3070+
resp_body,
3071+
)
3072+
.await;
3073+
3074+
let result = super::experiments_events_submit(&cfg, "exp-1", tmp.to_str().unwrap()).await;
3075+
assert!(
3076+
result.is_ok(),
3077+
"experiments_events_submit failed: {:?}",
3078+
result.err()
3079+
);
3080+
let _ = std::fs::remove_file(tmp);
3081+
cleanup_env();
3082+
std::env::remove_var("DD_TOKEN_STORAGE");
3083+
}
3084+
3085+
#[tokio::test]
3086+
async fn test_llm_obs_experiments_events_submit_400() {
3087+
let _lock = lock_env().await;
3088+
std::env::set_var("DD_TOKEN_STORAGE", "file");
3089+
let mut server = mockito::Server::new_async().await;
3090+
let cfg = test_config(&server.url());
3091+
3092+
let tmp = write_temp_json(
3093+
"pup_test_exp_ingest_events_400.json",
3094+
r#"{"metrics":[{"label":"accuracy","metric_type":"score","score_value":0.9}]}"#,
3095+
);
3096+
let _mock = server
3097+
.mock("POST", mockito::Matcher::Any)
3098+
.match_query(mockito::Matcher::Any)
3099+
.with_status(400)
3100+
.with_header("content-type", "application/json")
3101+
.with_body(r#"{"errors":["bad request"]}"#)
3102+
.create_async()
3103+
.await;
3104+
3105+
let result = super::experiments_events_submit(&cfg, "exp-1", tmp.to_str().unwrap()).await;
3106+
assert!(result.is_err(), "should fail on 400");
3107+
let _ = std::fs::remove_file(tmp);
3108+
cleanup_env();
3109+
std::env::remove_var("DD_TOKEN_STORAGE");
3110+
}
28543111
}

0 commit comments

Comments
 (0)