1use crate::client::AxonFlowClient;
19use crate::error::AxonFlowError;
20use crate::types::hitl::{
21 HITLApprovalRequest, HITLCreateInput, HITLQueueListOptions, HITLQueueListResponse,
22 HITLReviewInput, HITLStats, HitlItemEnvelope, HitlListEnvelope, HitlStatsEnvelope,
23};
24use crate::PATH_SEGMENT;
25use percent_encoding::utf8_percent_encode;
26
27impl AxonFlowClient {
28 pub async fn list_hitl_queue(
49 &self,
50 opts: HITLQueueListOptions,
51 ) -> Result<HITLQueueListResponse, AxonFlowError> {
52 let mut url = format!("{}/api/v1/hitl/queue", self.endpoint());
53 let qs = build_list_query(&opts);
54 if !qs.is_empty() {
55 url.push('?');
56 url.push_str(&qs);
57 }
58
59 let resp = self.checked_get(&url).await?;
60 let body = resp.text().await?;
61 let envelope: HitlListEnvelope = serde_json::from_str(&body)?;
62 let total = envelope.meta.total;
63 let returned = envelope.data.len() as i64;
64 let offset = envelope.meta.offset;
65 Ok(HITLQueueListResponse {
66 items: envelope.data,
67 total,
68 has_more: offset + returned < total,
69 })
70 }
71
72 pub async fn get_hitl_request(
76 &self,
77 request_id: &str,
78 ) -> Result<HITLApprovalRequest, AxonFlowError> {
79 if request_id.is_empty() {
80 return Err(AxonFlowError::ConfigError(
81 "request_id is required".to_string(),
82 ));
83 }
84 let encoded = utf8_percent_encode(request_id, PATH_SEGMENT).to_string();
85 let url = format!("{}/api/v1/hitl/queue/{}", self.endpoint(), encoded);
86
87 let resp = self.checked_get(&url).await?;
88 let body = resp.text().await?;
89 let envelope: HitlItemEnvelope = serde_json::from_str(&body)?;
90 Ok(envelope.data)
91 }
92
93 pub async fn create_hitl_request(
144 &self,
145 input: HITLCreateInput,
146 ) -> Result<HITLApprovalRequest, AxonFlowError> {
147 if input.client_id.is_empty() {
148 return Err(AxonFlowError::ConfigError(
149 "client_id is required".to_string(),
150 ));
151 }
152 if input.original_query.is_empty() {
153 return Err(AxonFlowError::ConfigError(
154 "original_query is required".to_string(),
155 ));
156 }
157 if input.request_type.is_empty() {
158 return Err(AxonFlowError::ConfigError(
159 "request_type is required".to_string(),
160 ));
161 }
162
163 let url = format!("{}/api/v1/hitl/queue", self.endpoint());
164 let resp = self.checked_post_json(&url, &input).await?;
165 let body = resp.text().await?;
166 let envelope: HitlItemEnvelope = serde_json::from_str(&body)?;
167 Ok(envelope.data)
168 }
169
170 pub async fn approve_hitl_request(
174 &self,
175 request_id: &str,
176 review: HITLReviewInput,
177 ) -> Result<(), AxonFlowError> {
178 self.review_hitl_request(request_id, "approve", &review)
179 .await
180 }
181
182 pub async fn reject_hitl_request(
186 &self,
187 request_id: &str,
188 review: HITLReviewInput,
189 ) -> Result<(), AxonFlowError> {
190 self.review_hitl_request(request_id, "reject", &review)
191 .await
192 }
193
194 pub async fn get_hitl_stats(&self) -> Result<HITLStats, AxonFlowError> {
198 let url = format!("{}/api/v1/hitl/stats", self.endpoint());
199 let resp = self.checked_get(&url).await?;
200 let body = resp.text().await?;
201 let envelope: HitlStatsEnvelope = serde_json::from_str(&body)?;
202 Ok(envelope.data)
203 }
204
205 async fn review_hitl_request(
208 &self,
209 request_id: &str,
210 action: &str,
211 review: &HITLReviewInput,
212 ) -> Result<(), AxonFlowError> {
213 if request_id.is_empty() {
214 return Err(AxonFlowError::ConfigError(
215 "request_id is required".to_string(),
216 ));
217 }
218 let encoded = utf8_percent_encode(request_id, PATH_SEGMENT).to_string();
219 let url = format!(
220 "{}/api/v1/hitl/queue/{}/{}",
221 self.endpoint(),
222 encoded,
223 action
224 );
225 let _ = self.checked_post_json(&url, review).await?;
226 Ok(())
227 }
228}
229
230fn build_list_query(opts: &HITLQueueListOptions) -> String {
234 let mut pairs: Vec<(&str, String)> = Vec::with_capacity(4);
235 if let Some(status) = &opts.status {
236 pairs.push(("status", status.clone()));
237 }
238 if let Some(severity) = &opts.severity {
239 pairs.push(("severity", severity.clone()));
240 }
241 if let Some(limit) = opts.limit {
242 pairs.push(("limit", limit.to_string()));
243 }
244 if let Some(offset) = opts.offset {
245 pairs.push(("offset", offset.to_string()));
246 }
247 pairs
248 .into_iter()
249 .map(|(k, v)| {
250 let v = utf8_percent_encode(&v, PATH_SEGMENT).to_string();
251 format!("{k}={v}")
252 })
253 .collect::<Vec<_>>()
254 .join("&")
255}
256
257#[cfg(test)]
258mod tests {
259 use super::*;
260 use crate::{AxonFlowClient, AxonFlowConfig};
261 use serde_json::json;
262 use std::time::Duration;
263 use wiremock::matchers::{body_partial_json, method, path, query_param};
264 use wiremock::{Mock, MockServer, ResponseTemplate};
265
266 fn make_client(endpoint: String) -> AxonFlowClient {
267 let config = AxonFlowConfig {
268 endpoint,
269 timeout: Duration::from_secs(2),
270 ..Default::default()
271 };
272 AxonFlowClient::new(config).expect("client init")
273 }
274
275 fn sample_row() -> serde_json::Value {
276 json!({
277 "request_id": "hitl-req-runtime-001",
278 "org_id": "org-1",
279 "tenant_id": "tenant-1",
280 "client_id": "loan-desk",
281 "user_id": "cust-001",
282 "original_query": "disburse $50000 to cust-001",
283 "request_type": "adk-tool",
284 "request_context": {"tool_name": "disburse_payment"},
285 "triggered_policy_id": "loan-amount-cap",
286 "triggered_policy_name": "Loan amount cap",
287 "trigger_reason": "Disbursement above $10k requires manager approval",
288 "severity": "high",
289 "status": "pending",
290 "notify_url": "https://workflows.example.com/hooks/loan-approve",
291 "expires_at": "2026-05-23T11:00:00Z",
292 "created_at": "2026-05-23T10:00:00Z",
293 "updated_at": "2026-05-23T10:00:00Z",
294 })
295 }
296
297 #[tokio::test]
300 async fn list_happy_path_parses_payload_and_pagination() {
301 let server = MockServer::start().await;
302 Mock::given(method("GET"))
303 .and(path("/api/v1/hitl/queue"))
304 .respond_with(ResponseTemplate::new(200).set_body_json(json!({
305 "success": true,
306 "data": [sample_row()],
307 "meta": {"total": 1, "limit": 50, "offset": 0},
308 })))
309 .mount(&server)
310 .await;
311
312 let client = make_client(server.uri());
313 let page = client
314 .list_hitl_queue(HITLQueueListOptions::default())
315 .await
316 .unwrap();
317
318 assert_eq!(page.total, 1);
319 assert_eq!(page.items.len(), 1);
320 assert!(!page.has_more);
321 assert_eq!(page.items[0].request_id, "hitl-req-runtime-001");
322 assert_eq!(
323 page.items[0].notify_url.as_deref(),
324 Some("https://workflows.example.com/hooks/loan-approve")
325 );
326 }
327
328 #[tokio::test]
329 async fn list_passes_filters_via_query_string() {
330 let server = MockServer::start().await;
331 Mock::given(method("GET"))
332 .and(path("/api/v1/hitl/queue"))
333 .and(query_param("status", "pending"))
334 .and(query_param("severity", "critical"))
335 .and(query_param("limit", "5"))
336 .respond_with(ResponseTemplate::new(200).set_body_json(json!({
337 "success": true,
338 "data": [],
339 "meta": {"total": 0, "limit": 5, "offset": 0},
340 })))
341 .expect(1)
342 .mount(&server)
343 .await;
344
345 let client = make_client(server.uri());
346 let opts = HITLQueueListOptions {
347 status: Some("pending".into()),
348 severity: Some("critical".into()),
349 limit: Some(5),
350 offset: None,
351 };
352 let _ = client.list_hitl_queue(opts).await.unwrap();
353 }
354
355 #[tokio::test]
358 async fn get_happy_path_parses_full_row() {
359 let server = MockServer::start().await;
360 Mock::given(method("GET"))
361 .and(path("/api/v1/hitl/queue/hitl-req-runtime-001"))
362 .respond_with(ResponseTemplate::new(200).set_body_json(json!({
363 "success": true,
364 "data": sample_row(),
365 })))
366 .mount(&server)
367 .await;
368
369 let client = make_client(server.uri());
370 let got = client
371 .get_hitl_request("hitl-req-runtime-001")
372 .await
373 .unwrap();
374 assert_eq!(got.request_id, "hitl-req-runtime-001");
375 assert_eq!(got.severity, "high");
376 assert_eq!(
377 got.notify_url.as_deref(),
378 Some("https://workflows.example.com/hooks/loan-approve")
379 );
380 }
381
382 #[tokio::test]
383 async fn get_empty_id_returns_config_error() {
384 let client = make_client("http://127.0.0.1:1".into());
385 let err = client.get_hitl_request("").await.unwrap_err();
386 assert!(err.to_string().contains("request_id is required"));
387 }
388
389 #[tokio::test]
390 async fn get_404_surfaces_as_api_error() {
391 let server = MockServer::start().await;
392 Mock::given(method("GET"))
393 .and(path("/api/v1/hitl/queue/nope"))
394 .respond_with(ResponseTemplate::new(404).set_body_json(json!({"error": "not found"})))
395 .mount(&server)
396 .await;
397
398 let client = make_client(server.uri());
399 let err = client.get_hitl_request("nope").await.unwrap_err();
400 match err {
401 AxonFlowError::ApiError { status, .. } => assert_eq!(status, 404),
402 other => panic!("expected ApiError(404), got {other}"),
403 }
404 }
405
406 #[tokio::test]
409 async fn create_happy_path_round_trips_full_input() {
410 let server = MockServer::start().await;
411 Mock::given(method("POST"))
412 .and(path("/api/v1/hitl/queue"))
413 .and(body_partial_json(json!({
414 "client_id": "loan-desk",
415 "original_query": "disburse $50000 to cust-001",
416 "request_type": "adk-tool",
417 "notify_url": "https://workflows.example.com/hooks/loan-approve",
418 "severity": "high",
419 })))
420 .respond_with(ResponseTemplate::new(201).set_body_json(json!({
421 "success": true,
422 "data": sample_row(),
423 })))
424 .expect(1)
425 .mount(&server)
426 .await;
427
428 let client = make_client(server.uri());
429 let req = client
430 .create_hitl_request(HITLCreateInput {
431 client_id: "loan-desk".into(),
432 user_id: Some("cust-001".into()),
433 original_query: "disburse $50000 to cust-001".into(),
434 request_type: "adk-tool".into(),
435 triggered_policy_id: Some("loan-amount-cap".into()),
436 triggered_policy_name: Some("Loan amount cap".into()),
437 trigger_reason: Some("Disbursement above $10k requires manager approval".into()),
438 severity: Some("high".into()),
439 notify_url: Some("https://workflows.example.com/hooks/loan-approve".into()),
440 ..Default::default()
441 })
442 .await
443 .unwrap();
444 assert_eq!(req.request_id, "hitl-req-runtime-001");
445 assert_eq!(
446 req.notify_url.as_deref(),
447 Some("https://workflows.example.com/hooks/loan-approve")
448 );
449 }
450
451 #[tokio::test]
452 async fn create_minimal_required_fields_only() {
453 let server = MockServer::start().await;
454 Mock::given(method("POST"))
455 .and(path("/api/v1/hitl/queue"))
456 .respond_with(ResponseTemplate::new(201).set_body_json(json!({
457 "success": true,
458 "data": {
459 "request_id": "hitl-req-minimal",
460 "org_id": "org-1",
461 "tenant_id": "tenant-1",
462 "client_id": "c1",
463 "original_query": "q",
464 "request_type": "chat",
465 "triggered_policy_id": "",
466 "triggered_policy_name": "",
467 "trigger_reason": "",
468 "severity": "high",
469 "status": "pending",
470 "expires_at": "2026-05-23T11:00:00Z",
471 "created_at": "2026-05-23T10:00:00Z",
472 "updated_at": "2026-05-23T10:00:00Z",
473 },
474 })))
475 .mount(&server)
476 .await;
477
478 let client = make_client(server.uri());
479 let req = client
480 .create_hitl_request(HITLCreateInput {
481 client_id: "c1".into(),
482 original_query: "q".into(),
483 request_type: "chat".into(),
484 ..Default::default()
485 })
486 .await
487 .unwrap();
488 assert_eq!(req.request_id, "hitl-req-minimal");
489 assert_eq!(req.notify_url, None);
490 }
491
492 #[tokio::test]
493 async fn create_bad_notify_url_scheme_surfaces_400() {
494 let server = MockServer::start().await;
496 Mock::given(method("POST"))
497 .and(path("/api/v1/hitl/queue"))
498 .respond_with(ResponseTemplate::new(400).set_body_json(json!({
499 "success": false,
500 "error": "notify_url scheme \"javascript\" is not allowed (use https:// or http://)",
501 })))
502 .mount(&server)
503 .await;
504
505 let client = make_client(server.uri());
506 let err = client
507 .create_hitl_request(HITLCreateInput {
508 client_id: "loan-desk".into(),
509 original_query: "disburse $50000".into(),
510 request_type: "adk-tool".into(),
511 notify_url: Some("javascript:alert(1)".into()),
512 ..Default::default()
513 })
514 .await
515 .unwrap_err();
516 match err {
517 AxonFlowError::ApiError { status, .. } => assert_eq!(status, 400),
518 other => panic!("expected ApiError(400), got {other}"),
519 }
520 }
521
522 #[tokio::test]
523 async fn create_401_surfaces_as_api_error() {
524 let server = MockServer::start().await;
525 Mock::given(method("POST"))
526 .and(path("/api/v1/hitl/queue"))
527 .respond_with(ResponseTemplate::new(401).set_body_json(json!({
528 "success": false,
529 "error": "Invalid API key",
530 })))
531 .mount(&server)
532 .await;
533
534 let client = make_client(server.uri());
535 let err = client
536 .create_hitl_request(HITLCreateInput {
537 client_id: "loan-desk".into(),
538 original_query: "disburse $50000".into(),
539 request_type: "adk-tool".into(),
540 ..Default::default()
541 })
542 .await
543 .unwrap_err();
544 match err {
545 AxonFlowError::ApiError { status, .. } => assert_eq!(status, 401),
546 other => panic!("expected ApiError(401), got {other}"),
547 }
548 }
549
550 #[tokio::test]
551 async fn create_network_failure_surfaces_as_error() {
552 let server = MockServer::start().await;
554 let url = server.uri();
555 drop(server);
556
557 let client = make_client(url);
558 let err = client
559 .create_hitl_request(HITLCreateInput {
560 client_id: "loan-desk".into(),
561 original_query: "disburse $50000".into(),
562 request_type: "adk-tool".into(),
563 ..Default::default()
564 })
565 .await
566 .unwrap_err();
567 let _ = err;
570 }
571
572 #[tokio::test]
573 async fn create_missing_client_id_rejected() {
574 let client = make_client("http://127.0.0.1:1".into());
575 let err = client
576 .create_hitl_request(HITLCreateInput {
577 client_id: "".into(),
578 original_query: "q".into(),
579 request_type: "chat".into(),
580 ..Default::default()
581 })
582 .await
583 .unwrap_err();
584 assert!(err.to_string().contains("client_id is required"));
585 }
586
587 #[tokio::test]
588 async fn create_missing_original_query_rejected() {
589 let client = make_client("http://127.0.0.1:1".into());
590 let err = client
591 .create_hitl_request(HITLCreateInput {
592 client_id: "c1".into(),
593 original_query: "".into(),
594 request_type: "chat".into(),
595 ..Default::default()
596 })
597 .await
598 .unwrap_err();
599 assert!(err.to_string().contains("original_query is required"));
600 }
601
602 #[tokio::test]
603 async fn create_missing_request_type_rejected() {
604 let client = make_client("http://127.0.0.1:1".into());
605 let err = client
606 .create_hitl_request(HITLCreateInput {
607 client_id: "c1".into(),
608 original_query: "q".into(),
609 request_type: "".into(),
610 ..Default::default()
611 })
612 .await
613 .unwrap_err();
614 assert!(err.to_string().contains("request_type is required"));
615 }
616
617 #[tokio::test]
620 async fn approve_posts_review_input_to_correct_path() {
621 let server = MockServer::start().await;
622 Mock::given(method("POST"))
623 .and(path("/api/v1/hitl/queue/hitl-req-runtime-001/approve"))
624 .and(body_partial_json(json!({
625 "reviewer_id": "user_456",
626 "reviewer_email": "reviewer@example.com",
627 "comment": "Approved after review",
628 })))
629 .respond_with(ResponseTemplate::new(200).set_body_json(json!({"success": true})))
630 .expect(1)
631 .mount(&server)
632 .await;
633
634 let client = make_client(server.uri());
635 client
636 .approve_hitl_request(
637 "hitl-req-runtime-001",
638 HITLReviewInput {
639 reviewer_id: "user_456".into(),
640 reviewer_email: "reviewer@example.com".into(),
641 reviewer_role: None,
642 comment: Some("Approved after review".into()),
643 },
644 )
645 .await
646 .unwrap();
647 }
648
649 #[tokio::test]
650 async fn reject_posts_review_input_to_correct_path() {
651 let server = MockServer::start().await;
652 Mock::given(method("POST"))
653 .and(path("/api/v1/hitl/queue/hitl-req-runtime-001/reject"))
654 .respond_with(ResponseTemplate::new(200).set_body_json(json!({"success": true})))
655 .expect(1)
656 .mount(&server)
657 .await;
658
659 let client = make_client(server.uri());
660 client
661 .reject_hitl_request(
662 "hitl-req-runtime-001",
663 HITLReviewInput {
664 reviewer_id: "user_456".into(),
665 reviewer_email: "reviewer@example.com".into(),
666 reviewer_role: None,
667 comment: None,
668 },
669 )
670 .await
671 .unwrap();
672 }
673
674 #[tokio::test]
675 async fn approve_empty_id_rejected_before_http() {
676 let client = make_client("http://127.0.0.1:1".into());
677 let err = client
678 .approve_hitl_request(
679 "",
680 HITLReviewInput {
681 reviewer_id: "u".into(),
682 reviewer_email: "u@e".into(),
683 reviewer_role: None,
684 comment: None,
685 },
686 )
687 .await
688 .unwrap_err();
689 assert!(err.to_string().contains("request_id is required"));
690 }
691
692 #[tokio::test]
695 async fn stats_happy_path_parses_envelope() {
696 let server = MockServer::start().await;
697 Mock::given(method("GET"))
698 .and(path("/api/v1/hitl/stats"))
699 .respond_with(ResponseTemplate::new(200).set_body_json(json!({
700 "success": true,
701 "data": {
702 "total_pending": 12,
703 "high_priority": 4,
704 "critical_priority": 2,
705 "oldest_pending_hours": 9.5,
706 },
707 })))
708 .mount(&server)
709 .await;
710
711 let client = make_client(server.uri());
712 let stats = client.get_hitl_stats().await.unwrap();
713 assert_eq!(stats.total_pending, 12);
714 assert_eq!(stats.high_priority, 4);
715 assert_eq!(stats.critical_priority, 2);
716 assert_eq!(stats.oldest_pending_hours, Some(9.5));
717 }
718
719 #[test]
720 fn build_list_query_omits_none_fields() {
721 let qs = build_list_query(&HITLQueueListOptions::default());
722 assert_eq!(qs, "");
723 let qs = build_list_query(&HITLQueueListOptions {
724 status: Some("pending".into()),
725 severity: None,
726 limit: Some(20),
727 offset: None,
728 });
729 assert_eq!(qs, "status=pending&limit=20");
730 }
731}