use faucet_core::Sink;
use faucet_sink_http::{HttpBatchMode, HttpSink, HttpSinkConfig};
use serde_json::json;
use wiremock::matchers::{body_json, method, path};
use wiremock::{Mock, MockServer, ResponseTemplate};
#[tokio::test]
async fn individual_mode_write_batch_partial_reports_only_the_failed_row() {
let server = MockServer::start().await;
for id in [0_u64, 2, 3] {
Mock::given(method("POST"))
.and(path("/ingest"))
.and(body_json(json!({ "id": id })))
.respond_with(ResponseTemplate::new(200))
.mount(&server)
.await;
}
Mock::given(method("POST"))
.and(path("/ingest"))
.and(body_json(json!({ "id": 1 })))
.respond_with(ResponseTemplate::new(400))
.mount(&server)
.await;
let config = HttpSinkConfig::new(format!("{}/ingest", server.uri()))
.batch_mode(HttpBatchMode::Individual)
.concurrency(4);
let sink = HttpSink::new(config);
let records: Vec<_> = (0..4).map(|i| json!({ "id": i })).collect();
let outcomes = sink
.write_batch_partial(&records)
.await
.expect("partial write must not surface an outer error in Individual mode");
assert_eq!(outcomes.len(), 4, "one outcome per record");
assert!(outcomes[0].is_ok());
assert!(
outcomes[1].is_err(),
"only the 400 record (id=1) is a failure"
);
assert!(outcomes[2].is_ok());
assert!(outcomes[3].is_ok());
let requests = server.received_requests().await.unwrap();
assert_eq!(requests.len(), 4, "every record is POSTed exactly once");
}
#[tokio::test]
async fn array_mode_write_batch_partial_surfaces_outer_error() {
let server = MockServer::start().await;
Mock::given(method("POST"))
.and(path("/ingest"))
.respond_with(ResponseTemplate::new(400))
.mount(&server)
.await;
let config = HttpSinkConfig::new(format!("{}/ingest", server.uri()))
.batch_mode(HttpBatchMode::Array)
.with_batch_size(0);
let sink = HttpSink::new(config);
let records: Vec<_> = (0..3).map(|i| json!({ "id": i })).collect();
let result = sink.write_batch_partial(&records).await;
assert!(
result.is_err(),
"array-mode failure must surface as an outer error, not per-row outcomes"
);
}
#[tokio::test]
async fn array_mode_multi_chunk_only_failed_chunk_rows_reported_failed() {
let server = MockServer::start().await;
Mock::given(method("POST"))
.and(path("/ingest"))
.and(body_json(json!([{ "id": 0 }, { "id": 1 }])))
.respond_with(ResponseTemplate::new(200))
.mount(&server)
.await;
Mock::given(method("POST"))
.and(path("/ingest"))
.and(body_json(json!([{ "id": 2 }, { "id": 3 }])))
.respond_with(ResponseTemplate::new(400))
.mount(&server)
.await;
let config = HttpSinkConfig::new(format!("{}/ingest", server.uri()))
.batch_mode(HttpBatchMode::Array)
.with_batch_size(2);
let sink = HttpSink::new(config);
let records: Vec<_> = (0..4).map(|i| json!({ "id": i })).collect();
let outcomes = sink
.write_batch_partial(&records)
.await
.expect("a later-chunk failure must report per-row outcomes, not an outer error");
assert_eq!(outcomes.len(), 4, "one outcome per record");
assert!(
outcomes[0].is_ok(),
"record 0 was in the delivered first chunk"
);
assert!(
outcomes[1].is_ok(),
"record 1 was in the delivered first chunk"
);
assert!(
outcomes[2].is_err(),
"record 2 was in the failed second chunk"
);
assert!(
outcomes[3].is_err(),
"record 3 was in the failed second chunk"
);
let requests = server.received_requests().await.unwrap();
assert_eq!(
requests.len(),
2,
"exactly two chunk POSTs (2 records each)"
);
}
#[tokio::test]
async fn array_mode_multi_chunk_first_chunk_failure_surfaces_outer_error() {
let server = MockServer::start().await;
Mock::given(method("POST"))
.and(path("/ingest"))
.respond_with(ResponseTemplate::new(400))
.mount(&server)
.await;
let config = HttpSinkConfig::new(format!("{}/ingest", server.uri()))
.batch_mode(HttpBatchMode::Array)
.with_batch_size(2);
let sink = HttpSink::new(config);
let records: Vec<_> = (0..4).map(|i| json!({ "id": i })).collect();
let result = sink.write_batch_partial(&records).await;
assert!(
result.is_err(),
"first-chunk failure (nothing delivered) must surface as an outer error"
);
}
#[tokio::test]
async fn array_mode_multi_chunk_unsent_chunks_after_failure_reported_failed() {
let server = MockServer::start().await;
Mock::given(method("POST"))
.and(path("/ingest"))
.and(body_json(json!([{ "id": 0 }])))
.respond_with(ResponseTemplate::new(200))
.mount(&server)
.await;
Mock::given(method("POST"))
.and(path("/ingest"))
.and(body_json(json!([{ "id": 1 }])))
.respond_with(ResponseTemplate::new(400))
.mount(&server)
.await;
let config = HttpSinkConfig::new(format!("{}/ingest", server.uri()))
.batch_mode(HttpBatchMode::Array)
.with_batch_size(1);
let sink = HttpSink::new(config);
let records: Vec<_> = (0..3).map(|i| json!({ "id": i })).collect();
let outcomes = sink
.write_batch_partial(&records)
.await
.expect("per-row outcomes once an earlier chunk was delivered");
assert_eq!(outcomes.len(), 3, "one outcome per record");
assert!(outcomes[0].is_ok(), "record 0 delivered");
assert!(outcomes[1].is_err(), "record 1 failed");
assert!(
outcomes[2].is_err(),
"record 2 never sent → reported failed"
);
let requests = server.received_requests().await.unwrap();
assert_eq!(
requests.len(),
2,
"only chunk 0 and chunk 1 are POSTed; chunk 2 is short-circuited"
);
}