diff --git a/examples/v1/logs/SubmitLog_2329487120.java b/examples/v1/logs/SubmitLog_2329487120.java new file mode 100644 index 00000000000..f802d1a99da --- /dev/null +++ b/examples/v1/logs/SubmitLog_2329487120.java @@ -0,0 +1,33 @@ +// Send deflate logs with explicit compression returns "Response from server (always 200 empty +// JSON)." response + +import com.datadog.api.client.ApiClient; +import com.datadog.api.client.ApiException; +import com.datadog.api.client.v1.api.LogsApi; +import com.datadog.api.client.v1.api.LogsApi.SubmitLogOptionalParameters; +import com.datadog.api.client.v1.model.ContentEncoding; +import com.datadog.api.client.v1.model.HTTPLogItem; +import java.util.Collections; +import java.util.List; + +public class Example { + public static void main(String[] args) { + ApiClient defaultClient = ApiClient.getDefaultApiClient(); + LogsApi apiInstance = new LogsApi(defaultClient); + + List body = + Collections.singletonList( + new HTTPLogItem().message("Example-Log").ddtags("host:ExampleLog")); + + try { + apiInstance.submitLog( + body, new SubmitLogOptionalParameters().contentEncoding(ContentEncoding.DEFLATE)); + } catch (ApiException e) { + System.err.println("Exception when calling LogsApi#submitLog"); + System.err.println("Status code: " + e.getCode()); + System.err.println("Reason: " + e.getResponseBody()); + System.err.println("Response headers: " + e.getResponseHeaders()); + e.printStackTrace(); + } + } +} diff --git a/examples/v1/logs/SubmitLog_348923381.java b/examples/v1/logs/SubmitLog_348923381.java new file mode 100644 index 00000000000..ebb31abfc30 --- /dev/null +++ b/examples/v1/logs/SubmitLog_348923381.java @@ -0,0 +1,33 @@ +// Send gzip logs with explicit compression returns "Response from server (always 200 empty JSON)." +// response + +import com.datadog.api.client.ApiClient; +import com.datadog.api.client.ApiException; +import com.datadog.api.client.v1.api.LogsApi; +import com.datadog.api.client.v1.api.LogsApi.SubmitLogOptionalParameters; +import com.datadog.api.client.v1.model.ContentEncoding; +import com.datadog.api.client.v1.model.HTTPLogItem; +import java.util.Collections; +import java.util.List; + +public class Example { + public static void main(String[] args) { + ApiClient defaultClient = ApiClient.getDefaultApiClient(); + LogsApi apiInstance = new LogsApi(defaultClient); + + List body = + Collections.singletonList( + new HTTPLogItem().message("Example-Log").ddtags("host:ExampleLog")); + + try { + apiInstance.submitLog( + body, new SubmitLogOptionalParameters().contentEncoding(ContentEncoding.GZIP)); + } catch (ApiException e) { + System.err.println("Exception when calling LogsApi#submitLog"); + System.err.println("Status code: " + e.getCode()); + System.err.println("Reason: " + e.getResponseBody()); + System.err.println("Response headers: " + e.getResponseHeaders()); + e.printStackTrace(); + } + } +} diff --git a/examples/v2/logs/SubmitLog_1695361854.java b/examples/v2/logs/SubmitLog_1695361854.java new file mode 100644 index 00000000000..094198d312e --- /dev/null +++ b/examples/v2/logs/SubmitLog_1695361854.java @@ -0,0 +1,38 @@ +// Send deflate logs with explicit compression returns "Request accepted for processing (always 202 +// empty JSON)." response + +import com.datadog.api.client.ApiClient; +import com.datadog.api.client.ApiException; +import com.datadog.api.client.v2.api.LogsApi; +import com.datadog.api.client.v2.api.LogsApi.SubmitLogOptionalParameters; +import com.datadog.api.client.v2.model.ContentEncoding; +import com.datadog.api.client.v2.model.HTTPLogItem; +import java.util.Collections; +import java.util.List; + +public class Example { + public static void main(String[] args) { + ApiClient defaultClient = ApiClient.getDefaultApiClient(); + LogsApi apiInstance = new LogsApi(defaultClient); + + List body = + Collections.singletonList( + new HTTPLogItem() + .ddsource("nginx") + .ddtags("env:staging,version:5.1") + .hostname("i-012345678") + .message("2019-11-19T14:37:58,995 INFO [process.name][20081] Hello World") + .service("payment")); + + try { + apiInstance.submitLog( + body, new SubmitLogOptionalParameters().contentEncoding(ContentEncoding.DEFLATE)); + } catch (ApiException e) { + System.err.println("Exception when calling LogsApi#submitLog"); + System.err.println("Status code: " + e.getCode()); + System.err.println("Reason: " + e.getResponseBody()); + System.err.println("Response headers: " + e.getResponseHeaders()); + e.printStackTrace(); + } + } +} diff --git a/examples/v2/logs/SubmitLog_3804768803.java b/examples/v2/logs/SubmitLog_3804768803.java new file mode 100644 index 00000000000..4b99f44d434 --- /dev/null +++ b/examples/v2/logs/SubmitLog_3804768803.java @@ -0,0 +1,38 @@ +// Send gzip logs with explicit compression returns "Request accepted for processing (always 202 +// empty JSON)." response + +import com.datadog.api.client.ApiClient; +import com.datadog.api.client.ApiException; +import com.datadog.api.client.v2.api.LogsApi; +import com.datadog.api.client.v2.api.LogsApi.SubmitLogOptionalParameters; +import com.datadog.api.client.v2.model.ContentEncoding; +import com.datadog.api.client.v2.model.HTTPLogItem; +import java.util.Collections; +import java.util.List; + +public class Example { + public static void main(String[] args) { + ApiClient defaultClient = ApiClient.getDefaultApiClient(); + LogsApi apiInstance = new LogsApi(defaultClient); + + List body = + Collections.singletonList( + new HTTPLogItem() + .ddsource("nginx") + .ddtags("env:staging,version:5.1") + .hostname("i-012345678") + .message("2019-11-19T14:37:58,995 INFO [process.name][20081] Hello World") + .service("payment")); + + try { + apiInstance.submitLog( + body, new SubmitLogOptionalParameters().contentEncoding(ContentEncoding.GZIP)); + } catch (ApiException e) { + System.err.println("Exception when calling LogsApi#submitLog"); + System.err.println("Status code: " + e.getCode()); + System.err.println("Reason: " + e.getResponseBody()); + System.err.println("Response headers: " + e.getResponseHeaders()); + e.printStackTrace(); + } + } +} diff --git a/src/test/java/com/datadog/api/RecorderSteps.java b/src/test/java/com/datadog/api/RecorderSteps.java index 620d3abc53b..97382974f68 100644 --- a/src/test/java/com/datadog/api/RecorderSteps.java +++ b/src/test/java/com/datadog/api/RecorderSteps.java @@ -6,6 +6,8 @@ import io.cucumber.java.BeforeAll; import io.cucumber.java.Scenario; import io.cucumber.java.Status; +import io.cucumber.java.en.Given; +import io.cucumber.java.en.Then; import io.cucumber.java.en.When; import java.io.File; import java.io.IOException; @@ -197,6 +199,16 @@ public String getCassetteName() { return world.getName() + ".json"; } + @Then("the request uses {string} compression") + public void theRequestUsesCompression(String compression) throws Exception { + TestRunner.assertLastRequestContentEncoding(world, compression); + } + + @Given("the client selects {string} compression") + public void theClientSelectsCompression(String compression) { + // The generated request plan passes the selected compression to the client call. + } + @When("the request is sent") public void theRequestIsSent() throws Exception { TestRunner.applyPlan(world, false); diff --git a/src/test/java/com/datadog/api/TestRunner.java b/src/test/java/com/datadog/api/TestRunner.java index b3213a0473f..244a52c74ad 100644 --- a/src/test/java/com/datadog/api/TestRunner.java +++ b/src/test/java/com/datadog/api/TestRunner.java @@ -138,6 +138,29 @@ public static void stopSession(World world) throws Exception { world.testServerSession = null; } + @SuppressWarnings("unchecked") + public static void assertLastRequestContentEncoding(World world, String expected) + throws Exception { + if (!serverEnabled()) { + return; + } + if (world.testServerSession == null) { + throw new IllegalStateException("Generated test-server session has not been started"); + } + Map result = + controlRequest("GET", "/sessions/" + world.testServerSession + "/last-request", null); + Map request = (Map) result.get("request"); + if (request == null) { + throw new AssertionError("Generated test server has not received a request"); + } + Map headers = (Map) request.get("headers"); + String actual = headers == null ? null : (String) headers.get("content-encoding"); + if (actual == null || !actual.equalsIgnoreCase(expected)) { + throw new AssertionError( + String.format("Expected Content-Encoding %s, got %s", expected, actual)); + } + } + public static void applyPlan(World world, boolean pagination) throws Exception { if (!runnerEnabled()) { return; @@ -167,6 +190,10 @@ public static void applyPlan(World world, boolean pagination) throws Exception { Object value = materialize(body.get("value"), world); world.addMaterializedRequestParameter("body", MAPPER.writeValueAsString(value)); } + if (request.get("selected_compression") != null) { + world.addMaterializedRequestParameter( + "Content-Encoding", MAPPER.writeValueAsString(request.get("selected_compression"))); + } for (Map parameter : parameters) { if (!"path".equals(parameter.get("in")) && !Boolean.TRUE.equals(parameter.get("required"))) { applyParameter(world, parameter); @@ -254,9 +281,14 @@ private static Map readMap(Path path) throws Exception { private static Map controlRequest(String endpoint, Map payload) throws Exception { + return controlRequest("POST", endpoint, payload); + } + + private static Map controlRequest( + String method, String endpoint, Map payload) throws Exception { URL url = new URL(serverUrl() + CONTROL_ROOT + endpoint); HttpURLConnection connection = (HttpURLConnection) url.openConnection(); - connection.setRequestMethod("POST"); + connection.setRequestMethod(method); connection.setRequestProperty("Connection", "close"); if (payload != null) { byte[] body = MAPPER.writeValueAsBytes(payload); @@ -283,7 +315,7 @@ private static Map controlRequest(String endpoint, Map= HttpURLConnection.HTTP_BAD_REQUEST) { throw new IllegalStateException( - String.format("Test server POST %s failed (%d): %s", endpoint, status, response)); + String.format("Test server %s %s failed (%d): %s", method, endpoint, status, response)); } return MAPPER.readValue(response.toString(), new TypeReference>() {}); } diff --git a/src/test/resources/com/datadog/api/client/v1/api/logs.feature b/src/test/resources/com/datadog/api/client/v1/api/logs.feature index 5ae909c3aae..1a4ba2de4c2 100644 --- a/src/test/resources/com/datadog/api/client/v1/api/logs.feature +++ b/src/test/resources/com/datadog/api/client/v1/api/logs.feature @@ -40,6 +40,15 @@ Feature: Logs When the request is sent Then the response status is 200 Response from server (always 200 empty JSON). + @skip-terraform-config @skip-validation @team:DataDog/event-platform-intake + Scenario: Send deflate logs with explicit compression returns "Response from server (always 200 empty JSON)." response + Given new "SubmitLog" request + And body with value [{"message": "{{ unique }}", "ddtags": "host:{{ unique_alnum }}"}] + And the client selects "deflate" compression + When the request is sent + Then the response status is 200 Response from server (always 200 empty JSON). + And the request uses "deflate" compression + @integration-only @skip-terraform-config @skip-validation @team:DataDog/event-platform-intake Scenario: Send gzip logs returns "Response from server (always 200 empty JSON)." response Given new "SubmitLog" request @@ -48,6 +57,15 @@ Feature: Logs When the request is sent Then the response status is 200 Response from server (always 200 empty JSON). + @skip-terraform-config @skip-validation @team:DataDog/event-platform-intake + Scenario: Send gzip logs with explicit compression returns "Response from server (always 200 empty JSON)." response + Given new "SubmitLog" request + And body with value [{"message": "{{ unique }}", "ddtags": "host:{{ unique_alnum }}"}] + And the client selects "gzip" compression + When the request is sent + Then the response status is 200 Response from server (always 200 empty JSON). + And the request uses "gzip" compression + @team:DataDog/event-platform-intake Scenario: Send logs returns "Response from server (always 200 empty JSON)." response Given new "SubmitLog" request diff --git a/src/test/resources/com/datadog/api/client/v2/api/logs.feature b/src/test/resources/com/datadog/api/client/v2/api/logs.feature index b7accd07dd9..1e1e000e3aa 100644 --- a/src/test/resources/com/datadog/api/client/v2/api/logs.feature +++ b/src/test/resources/com/datadog/api/client/v2/api/logs.feature @@ -136,6 +136,15 @@ Feature: Logs When the request is sent Then the response status is 202 Response from server (always 202 empty JSON). + @skip-terraform-config @skip-validation @team:DataDog/event-platform-intake @team:DataDog/logs-backend @team:DataDog/logs-ingestion + Scenario: Send deflate logs with explicit compression returns "Request accepted for processing (always 202 empty JSON)." response + Given new "SubmitLog" request + And body with value [{"ddsource": "nginx", "ddtags": "env:staging,version:5.1", "hostname": "i-012345678", "message": "2019-11-19T14:37:58,995 INFO [process.name][20081] Hello World", "service": "payment"}] + And the client selects "deflate" compression + When the request is sent + Then the response status is 202 Response from server (always 202 empty JSON). + And the request uses "deflate" compression + @integration-only @skip-terraform-config @skip-validation @team:DataDog/event-platform-intake @team:DataDog/logs-backend @team:DataDog/logs-ingestion Scenario: Send gzip logs returns "Request accepted for processing (always 202 empty JSON)." response Given new "SubmitLog" request @@ -144,6 +153,15 @@ Feature: Logs When the request is sent Then the response status is 202 Request accepted for processing (always 202 empty JSON). + @skip-terraform-config @skip-validation @team:DataDog/event-platform-intake @team:DataDog/logs-backend @team:DataDog/logs-ingestion + Scenario: Send gzip logs with explicit compression returns "Request accepted for processing (always 202 empty JSON)." response + Given new "SubmitLog" request + And body with value [{"ddsource": "nginx", "ddtags": "env:staging,version:5.1", "hostname": "i-012345678", "message": "2019-11-19T14:37:58,995 INFO [process.name][20081] Hello World", "service": "payment"}] + And the client selects "gzip" compression + When the request is sent + Then the response status is 202 Request accepted for processing (always 202 empty JSON). + And the request uses "gzip" compression + @generated @skip @team:DataDog/event-platform-intake @team:DataDog/logs-backend @team:DataDog/logs-ingestion Scenario: Send logs returns "Bad Request" response Given new "SubmitLog" request diff --git a/src/test/resources/generated-test/test-runner-data/manifest.json b/src/test/resources/generated-test/test-runner-data/manifest.json index c28289cd887..e17a377c310 100644 --- a/src/test/resources/generated-test/test-runner-data/manifest.json +++ b/src/test/resources/generated-test/test-runner-data/manifest.json @@ -1211,6 +1211,20 @@ "scenario": "Search test logs returns \"OK\" response", "version": "v1" }, + { + "feature": "Logs", + "feature_file": "../../com/datadog/api/client/v1/api/logs.feature", + "file": "v1/logs/send-deflate-logs-with-explicit-compression-returns-response-from-server-always-200-empty-json-response.json", + "scenario": "Send deflate logs with explicit compression returns \"Response from server (always 200 empty JSON).\" response", + "version": "v1" + }, + { + "feature": "Logs", + "feature_file": "../../com/datadog/api/client/v1/api/logs.feature", + "file": "v1/logs/send-gzip-logs-with-explicit-compression-returns-response-from-server-always-200-empty-json-response.json", + "scenario": "Send gzip logs with explicit compression returns \"Response from server (always 200 empty JSON).\" response", + "version": "v1" + }, { "feature": "Logs", "feature_file": "../../com/datadog/api/client/v1/api/logs.feature", @@ -6391,6 +6405,20 @@ "scenario": "Search logs returns \"OK\" response with pagination", "version": "v2" }, + { + "feature": "Logs", + "feature_file": "../../com/datadog/api/client/v2/api/logs.feature", + "file": "v2/logs/send-deflate-logs-with-explicit-compression-returns-request-accepted-for-processing-always-202-empty-json-response.json", + "scenario": "Send deflate logs with explicit compression returns \"Request accepted for processing (always 202 empty JSON).\" response", + "version": "v2" + }, + { + "feature": "Logs", + "feature_file": "../../com/datadog/api/client/v2/api/logs.feature", + "file": "v2/logs/send-gzip-logs-with-explicit-compression-returns-request-accepted-for-processing-always-202-empty-json-response.json", + "scenario": "Send gzip logs with explicit compression returns \"Request accepted for processing (always 202 empty JSON).\" response", + "version": "v2" + }, { "feature": "Logs", "feature_file": "../../com/datadog/api/client/v2/api/logs.feature", diff --git a/src/test/resources/generated-test/test-runner-data/v1/logs/send-deflate-logs-with-explicit-compression-returns-response-from-server-always-200-empty-json-response.json b/src/test/resources/generated-test/test-runner-data/v1/logs/send-deflate-logs-with-explicit-compression-returns-response-from-server-always-200-empty-json-response.json new file mode 100644 index 00000000000..feeb55fe462 --- /dev/null +++ b/src/test/resources/generated-test/test-runner-data/v1/logs/send-deflate-logs-with-explicit-compression-returns-response-from-server-always-200-empty-json-response.json @@ -0,0 +1,37 @@ +{ + "api": "Logs", + "expected_status": 200, + "feature": "Logs", + "id": "v1/Logs/Send deflate logs with explicit compression returns \"Response from server (always 200 empty JSON).\" response", + "operation_id": "SubmitLog", + "request": { + "body": { + "schema": { + "format": null, + "items": { + "format": null, + "ref": "HTTPLogItem", + "type": "object" + }, + "ref": "HTTPLog", + "type": "array" + }, + "source": "inline", + "value": [ + { + "ddtags": "host:{{ unique_alnum }}", + "message": "{{ unique }}" + } + ] + }, + "content_type": "application/json", + "method": "POST", + "pagination": false, + "parameters": [], + "path": "/v1/input", + "selected_compression": "deflate" + }, + "scenario": "Send deflate logs with explicit compression returns \"Response from server (always 200 empty JSON).\" response", + "schema_version": 1, + "version": "v1" +} diff --git a/src/test/resources/generated-test/test-runner-data/v1/logs/send-gzip-logs-with-explicit-compression-returns-response-from-server-always-200-empty-json-response.json b/src/test/resources/generated-test/test-runner-data/v1/logs/send-gzip-logs-with-explicit-compression-returns-response-from-server-always-200-empty-json-response.json new file mode 100644 index 00000000000..dcbdc13ade9 --- /dev/null +++ b/src/test/resources/generated-test/test-runner-data/v1/logs/send-gzip-logs-with-explicit-compression-returns-response-from-server-always-200-empty-json-response.json @@ -0,0 +1,37 @@ +{ + "api": "Logs", + "expected_status": 200, + "feature": "Logs", + "id": "v1/Logs/Send gzip logs with explicit compression returns \"Response from server (always 200 empty JSON).\" response", + "operation_id": "SubmitLog", + "request": { + "body": { + "schema": { + "format": null, + "items": { + "format": null, + "ref": "HTTPLogItem", + "type": "object" + }, + "ref": "HTTPLog", + "type": "array" + }, + "source": "inline", + "value": [ + { + "ddtags": "host:{{ unique_alnum }}", + "message": "{{ unique }}" + } + ] + }, + "content_type": "application/json", + "method": "POST", + "pagination": false, + "parameters": [], + "path": "/v1/input", + "selected_compression": "gzip" + }, + "scenario": "Send gzip logs with explicit compression returns \"Response from server (always 200 empty JSON).\" response", + "schema_version": 1, + "version": "v1" +} diff --git a/src/test/resources/generated-test/test-runner-data/v2/logs/send-deflate-logs-with-explicit-compression-returns-request-accepted-for-processing-always-202-empty-json-response.json b/src/test/resources/generated-test/test-runner-data/v2/logs/send-deflate-logs-with-explicit-compression-returns-request-accepted-for-processing-always-202-empty-json-response.json new file mode 100644 index 00000000000..8184ebb0f2b --- /dev/null +++ b/src/test/resources/generated-test/test-runner-data/v2/logs/send-deflate-logs-with-explicit-compression-returns-request-accepted-for-processing-always-202-empty-json-response.json @@ -0,0 +1,40 @@ +{ + "api": "Logs", + "expected_status": 202, + "feature": "Logs", + "id": "v2/Logs/Send deflate logs with explicit compression returns \"Request accepted for processing (always 202 empty JSON).\" response", + "operation_id": "SubmitLog", + "request": { + "body": { + "schema": { + "format": null, + "items": { + "format": null, + "ref": "HTTPLogItem", + "type": "object" + }, + "ref": "HTTPLog", + "type": "array" + }, + "source": "inline", + "value": [ + { + "ddsource": "nginx", + "ddtags": "env:staging,version:5.1", + "hostname": "i-012345678", + "message": "2019-11-19T14:37:58,995 INFO [process.name][20081] Hello World", + "service": "payment" + } + ] + }, + "content_type": "application/json", + "method": "POST", + "pagination": false, + "parameters": [], + "path": "/api/v2/logs", + "selected_compression": "deflate" + }, + "scenario": "Send deflate logs with explicit compression returns \"Request accepted for processing (always 202 empty JSON).\" response", + "schema_version": 1, + "version": "v2" +} diff --git a/src/test/resources/generated-test/test-runner-data/v2/logs/send-gzip-logs-with-explicit-compression-returns-request-accepted-for-processing-always-202-empty-json-response.json b/src/test/resources/generated-test/test-runner-data/v2/logs/send-gzip-logs-with-explicit-compression-returns-request-accepted-for-processing-always-202-empty-json-response.json new file mode 100644 index 00000000000..f3f391f0d8a --- /dev/null +++ b/src/test/resources/generated-test/test-runner-data/v2/logs/send-gzip-logs-with-explicit-compression-returns-request-accepted-for-processing-always-202-empty-json-response.json @@ -0,0 +1,40 @@ +{ + "api": "Logs", + "expected_status": 202, + "feature": "Logs", + "id": "v2/Logs/Send gzip logs with explicit compression returns \"Request accepted for processing (always 202 empty JSON).\" response", + "operation_id": "SubmitLog", + "request": { + "body": { + "schema": { + "format": null, + "items": { + "format": null, + "ref": "HTTPLogItem", + "type": "object" + }, + "ref": "HTTPLog", + "type": "array" + }, + "source": "inline", + "value": [ + { + "ddsource": "nginx", + "ddtags": "env:staging,version:5.1", + "hostname": "i-012345678", + "message": "2019-11-19T14:37:58,995 INFO [process.name][20081] Hello World", + "service": "payment" + } + ] + }, + "content_type": "application/json", + "method": "POST", + "pagination": false, + "parameters": [], + "path": "/api/v2/logs", + "selected_compression": "gzip" + }, + "scenario": "Send gzip logs with explicit compression returns \"Request accepted for processing (always 202 empty JSON).\" response", + "schema_version": 1, + "version": "v2" +} diff --git a/src/test/resources/generated-test/test-server b/src/test/resources/generated-test/test-server index 417fcc086cf..acd7c193121 100755 --- a/src/test/resources/generated-test/test-server +++ b/src/test/resources/generated-test/test-server @@ -10,12 +10,15 @@ from __future__ import annotations import argparse import base64 +import ctypes import json import os import re import tempfile import threading import uuid +import zlib +from ctypes.util import find_library from datetime import datetime, timezone from http import HTTPStatus from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer @@ -44,6 +47,14 @@ SAFE_BROWSER_HEADERS = { "content-security-policy": "default-src 'none'; sandbox", "x-content-type-options": "nosniff", } +ZSTD_CONTENTSIZE_UNKNOWN = (1 << 64) - 1 +ZSTD_CONTENTSIZE_ERROR = (1 << 64) - 2 +ZSTD_MAGIC = b"\x28\xb5\x2f\xfd" +ZSTD_SINGLE_SEGMENT_FLAG = 0x20 +ZSTD_DESCRIPTOR_LOW_BITS_MASK = 0x3F +ZSTD_MIN_RAW_FRAME_SIZE = 6 +ZSTD_FCS_TWO_BYTE_OFFSET = 0x100 +MAX_DECOMPRESSED_BODY_SIZE = 256 * 1024 * 1024 class RecordingDatabase: @@ -92,6 +103,7 @@ class RecordingDatabase: "scenario": scenario, "recording": recording, "cursor": 0, + "requests": [], "captures": [], "frozen_at": frozen_at, } @@ -109,6 +121,7 @@ class RecordingDatabase: if cursor >= len(interactions): raise LookupError(f"Recording has no interaction #{cursor + 1}") expected = interactions[cursor] + session["requests"].append(actual) if not _requests_match(expected["request"], actual): raise RequestMismatchError(expected["request"], actual, cursor) session["cursor"] += 1 @@ -139,6 +152,19 @@ class RecordingDatabase: request = interactions[cursor]["request"] if cursor < len(interactions) else None return {"request": request} + def last_request(self, session_id: str) -> dict[str, Any]: + """Return the most recent request received for a session.""" + with self.lock: + session = self.sessions.get(session_id) + if session is None: + raise LookupError(f"Unknown session: {session_id}") + requests = ( + session["requests"] + if session["mode"] == "replay" + else [capture["request"] for capture in session["captures"]] + ) + return {"request": requests[-1] if requests else None} + def capture( self, session_id: str | None, @@ -187,13 +213,22 @@ class RecordingDatabase: "scenario": session["scenario"], "version": key[0], "frozen_at": session["frozen_at"], - "interactions": session["captures"], + "interactions": [ + { + **capture, + "request": {key: value for key, value in capture["request"].items() if key != "headers"}, + } + for capture in session["captures"] + ], } shard["recordings"] = [item for item in shard["recordings"] if item["scenario"] != session["scenario"]] shard["recordings"].append(recording) shard["recordings"].sort(key=lambda item: item["scenario"]) _write_json_atomic(self.shard_paths[key], shard) - return {"interactions": len(session["captures"]), "file": str(self.shard_paths[key])} + return { + "interactions": len(session["captures"]), + "file": str(self.shard_paths[key]), + } def _add_manifest_entry(self, key: tuple[str, str], shard_path: Path) -> None: manifest_path = self.root / "manifest.json" @@ -216,7 +251,13 @@ class RequestMismatchError(Exception): class TestServer(ThreadingHTTPServer): daemon_threads = True - def __init__(self, address: tuple[str, int], database: RecordingDatabase, mode: str, upstream: str | None): + def __init__( + self, + address: tuple[str, int], + database: RecordingDatabase, + mode: str, + upstream: str | None, + ): super().__init__(address, TestRequestHandler) self.database = database self.mode = mode @@ -259,11 +300,22 @@ class TestRequestHandler(BaseHTTPRequestHandler): if self.path == f"{CONTROL_ROOT}/sessions" and self.command == "POST": self._start_session() return - next_match = re.fullmatch(rf"{re.escape(CONTROL_ROOT)}/sessions/([a-f0-9]+)/next-request", self.path) + next_match = re.fullmatch( + rf"{re.escape(CONTROL_ROOT)}/sessions/([a-f0-9]+)/next-request", + self.path, + ) if next_match and self.command == "GET": result = self.server.database.next_request(next_match.group(1)) self._send_json(HTTPStatus.OK, result) return + last_match = re.fullmatch( + rf"{re.escape(CONTROL_ROOT)}/sessions/([a-f0-9]+)/last-request", + self.path, + ) + if last_match and self.command == "GET": + result = self.server.database.last_request(last_match.group(1)) + self._send_json(HTTPStatus.OK, result) + return stop_match = re.fullmatch(rf"{re.escape(CONTROL_ROOT)}/sessions/([a-f0-9]+)/stop", self.path) if stop_match and self.command == "POST": result = self.server.database.stop(stop_match.group(1)) @@ -284,7 +336,11 @@ class TestRequestHandler(BaseHTTPRequestHandler): except (LookupError, ValueError) as error: self._send_json(HTTPStatus.NOT_FOUND, {"error": str(error)}, error="recording-not-found") except Exception as error: - self._send_json(HTTPStatus.INTERNAL_SERVER_ERROR, {"error": str(error)}, error="internal-error") + self._send_json( + HTTPStatus.INTERNAL_SERVER_ERROR, + {"error": str(error)}, + error="internal-error", + ) def _start_session(self) -> None: payload = json.loads(self._read_body().decode("utf-8")) @@ -298,7 +354,13 @@ class TestRequestHandler(BaseHTTPRequestHandler): def _handle_api_request(self) -> None: body = self._read_body() - actual = _normalise_request(self.command, self.path, self.headers.get("content-type", ""), body) + actual = _normalise_request( + self.command, + self.path, + self.headers.get("content-type", ""), + self.headers.get("content-encoding", ""), + body, + ) session_id = self.headers.get(SESSION_HEADER) if self.server.mode == "replay": response = self.server.database.replay(session_id, actual) @@ -332,7 +394,10 @@ class TestRequestHandler(BaseHTTPRequestHandler): "status": remote.status, "reason": remote.reason, "headers": response_headers, - "body": {"encoding": "base64", "value": base64.b64encode(response_body).decode("ascii")}, + "body": { + "encoding": "base64", + "value": base64.b64encode(response_body).decode("ascii"), + }, } finally: remote.close() @@ -348,8 +413,9 @@ class TestRequestHandler(BaseHTTPRequestHandler): else: body = body_data.get("value", "").encode("utf-8") status = response["status"] + response_headers = response.get("headers", {}) self.send_response(status, response.get("reason")) - for key, value in response.get("headers", {}).items(): + for key, value in response_headers.items(): if key.lower() not in HOP_BY_HOP_HEADERS | {"content-length"} | SAFE_BROWSER_HEADERS.keys(): self.send_header(key, value) for key, value in SAFE_BROWSER_HEADERS.items(): @@ -372,21 +438,48 @@ class TestRequestHandler(BaseHTTPRequestHandler): self.wfile.write(body) -def _normalise_request(method: str, raw_path: str, content_type: str, body: bytes) -> dict[str, Any]: +def _normalise_request( + method: str, + raw_path: str, + content_type: str, + content_encoding: str, + body: bytes, +) -> dict[str, Any]: parsed = urlsplit(raw_path) - normalised_body = _normalise_body(body, content_type) + compression = content_encoding.strip().casefold() + normalised_body = _normalise_body(body, content_type, compression) return { "method": method.upper(), "path": _normalise_path(parsed.path), "query": sorted([list(pair) for pair in parse_qsl(parsed.query, keep_blank_values=True)]), "content_type": _normalise_content_type(content_type, normalised_body), + "headers": {"content-encoding": compression} if compression else {}, "body": normalised_body, } -def _normalise_body(body: bytes, content_type: str) -> dict[str, Any]: +def _normalise_body(body: bytes, content_type: str, compression: str = "") -> dict[str, Any]: if not body: + if compression in {"gzip", "deflate", "zstd1"}: + return {"type": "invalid-compression", "value": ""} return {"type": "empty", "value": None} + if compression == "zstd1": + decompressed = _decompress_zstd(body) + if decompressed is None: + return { + "type": "invalid-compression", + "value": base64.b64encode(body).decode("ascii"), + } + body = decompressed + elif compression in {"gzip", "deflate"}: + wbits = zlib.MAX_WBITS | 16 if compression == "gzip" else zlib.MAX_WBITS + decompressed = _decompress_zlib(body, wbits=wbits) + if decompressed is None: + return { + "type": "invalid-compression", + "value": base64.b64encode(body).decode("ascii"), + } + body = decompressed media_type = _media_type(content_type) text = body.decode("utf-8", errors="surrogateescape") if media_type.endswith("json"): @@ -401,6 +494,150 @@ def _normalise_body(body: bytes, content_type: str) -> dict[str, Any]: return {"type": "text", "value": text} +def _decompress_zlib( + body: bytes, + *, + wbits: int, + max_size: int = MAX_DECOMPRESSED_BODY_SIZE, +) -> bytes | None: + """Decompress exactly one complete zlib stream with bounded output.""" + decompressor = zlib.decompressobj(wbits) + output = bytearray() + pending = body + try: + while pending: + remaining = max_size - len(output) + 1 + chunk = decompressor.decompress(pending, remaining) + output.extend(chunk) + if len(output) > max_size: + return None + unconsumed = decompressor.unconsumed_tail + if not unconsumed: + break + if len(unconsumed) == len(pending) and not chunk: + return None + pending = unconsumed + except zlib.error: + return None + if not decompressor.eof or decompressor.unused_data or decompressor.unconsumed_tail: + return None + return bytes(output) + + +def _load_zstd() -> Any | None: + candidates = ( + find_library("zstd"), + "libzstd.so.1", + "libzstd.dylib", + "/opt/homebrew/lib/libzstd.dylib", + "/usr/local/lib/libzstd.dylib", + "libzstd.dll", + "zstd.dll", + ) + for candidate in dict.fromkeys(item for item in candidates if item): + try: + return ctypes.CDLL(candidate) + except OSError: + continue + return None + + +def _decompress_zstd(body: bytes) -> bytes | None: + """Decompress one complete Zstandard frame with the platform libzstd.""" + library = _load_zstd() + if library is None: + return _decompress_zstd_raw_frame(body) + + try: + library.ZSTD_isError.argtypes = [ctypes.c_size_t] + library.ZSTD_isError.restype = ctypes.c_uint + library.ZSTD_findFrameCompressedSize.argtypes = [ + ctypes.c_void_p, + ctypes.c_size_t, + ] + library.ZSTD_findFrameCompressedSize.restype = ctypes.c_size_t + library.ZSTD_getFrameContentSize.argtypes = [ctypes.c_void_p, ctypes.c_size_t] + library.ZSTD_getFrameContentSize.restype = ctypes.c_ulonglong + library.ZSTD_decompressBound.argtypes = [ctypes.c_void_p, ctypes.c_size_t] + library.ZSTD_decompressBound.restype = ctypes.c_ulonglong + library.ZSTD_decompress.argtypes = [ + ctypes.c_void_p, + ctypes.c_size_t, + ctypes.c_void_p, + ctypes.c_size_t, + ] + library.ZSTD_decompress.restype = ctypes.c_size_t + + source = ctypes.create_string_buffer(body) + source_pointer = ctypes.cast(source, ctypes.c_void_p) + compressed_size = library.ZSTD_findFrameCompressedSize(source_pointer, len(body)) + if library.ZSTD_isError(compressed_size) or compressed_size != len(body): + return None + + capacity = library.ZSTD_getFrameContentSize(source_pointer, len(body)) + if capacity == ZSTD_CONTENTSIZE_ERROR: + return None + if capacity == ZSTD_CONTENTSIZE_UNKNOWN: + capacity = library.ZSTD_decompressBound(source_pointer, len(body)) + if capacity > MAX_DECOMPRESSED_BODY_SIZE: + return None + + destination = ctypes.create_string_buffer(max(1, capacity)) + decompressed_size = library.ZSTD_decompress(destination, capacity, source_pointer, len(body)) + if library.ZSTD_isError(decompressed_size) or decompressed_size > capacity: + return None + return destination.raw[:decompressed_size] + except (AttributeError, OSError, OverflowError, TypeError): + return None + + +def _decompress_zstd_raw_frame(body: bytes) -> bytes | None: + """Decode the raw-block Zstandard frames emitted by the dependency-free fallback.""" + if len(body) < ZSTD_MIN_RAW_FRAME_SIZE or body[:4] != ZSTD_MAGIC: + return None + descriptor = body[4] + if descriptor & ZSTD_DESCRIPTOR_LOW_BITS_MASK != ZSTD_SINGLE_SEGMENT_FLAG: + return None + size_flag = descriptor >> 6 + size_bytes = (1, 2, 4, 8)[size_flag] + cursor = 5 + if len(body) < cursor + size_bytes: + return None + expected_size = int.from_bytes(body[cursor : cursor + size_bytes], "little") + if size_flag == 1: + expected_size += ZSTD_FCS_TWO_BYTE_OFFSET + if expected_size > MAX_DECOMPRESSED_BODY_SIZE: + return None + cursor += size_bytes + output = bytearray() + last_block = False + while not last_block: + if len(body) < cursor + 3: + return None + header = int.from_bytes(body[cursor : cursor + 3], "little") + cursor += 3 + last_block = bool(header & 1) + block_type = (header >> 1) & 0x3 + block_size = header >> 3 + if block_type == 0: + if len(body) < cursor + block_size: + return None + output.extend(body[cursor : cursor + block_size]) + cursor += block_size + elif block_type == 1: + if cursor >= len(body): + return None + output.extend(body[cursor : cursor + 1] * block_size) + cursor += 1 + else: + return None + if len(output) > MAX_DECOMPRESSED_BODY_SIZE: + return None + if cursor != len(body) or len(output) != expected_size: + return None + return bytes(output) + + def _normalise_json(value: Any) -> Any: if isinstance(value, dict): return {key: _normalise_json(item) for key, item in value.items()} @@ -446,7 +683,10 @@ def _media_type(value: str) -> str: def _status_allows_message_content(status: int) -> bool: - return status >= HTTPStatus.OK and status not in {HTTPStatus.NO_CONTENT, HTTPStatus.NOT_MODIFIED} + return status >= HTTPStatus.OK and status not in { + HTTPStatus.NO_CONTENT, + HTTPStatus.NOT_MODIFIED, + } def _normalise_path(path: str) -> str: diff --git a/src/test/resources/generated-test/test-server-data/v1/logs.json b/src/test/resources/generated-test/test-server-data/v1/logs.json index 3ade4a2a924..1735b20978e 100644 --- a/src/test/resources/generated-test/test-server-data/v1/logs.json +++ b/src/test/resources/generated-test/test-server-data/v1/logs.json @@ -41,6 +41,78 @@ "scenario": "Search test logs returns \"OK\" response", "version": "v1" }, + { + "feature": "Logs", + "frozen_at": "2026-09-28T14:19:35.421Z", + "interactions": [ + { + "request": { + "body": { + "type": "json", + "value": [ + { + "ddtags": "host:TestSenddeflatelogswithexplicitcompressionreturnsResponsefromserveralways200emptyJSONrespo1790605175", + "message": "Test-Send_deflate_logs_with_explicit_compression_returns_Response_from_server_always_200_empty_JSON_respo-1790605175" + } + ] + }, + "content_type": "application/json", + "method": "POST", + "path": "/v1/input", + "query": [] + }, + "response": { + "body": { + "encoding": "text", + "value": "{}" + }, + "headers": { + "content-type": "application/json" + }, + "reason": "OK", + "status": 200 + } + } + ], + "scenario": "Send deflate logs with explicit compression returns \"Response from server (always 200 empty JSON).\" response", + "version": "v1" + }, + { + "feature": "Logs", + "frozen_at": "2026-09-28T14:19:35.658Z", + "interactions": [ + { + "request": { + "body": { + "type": "json", + "value": [ + { + "ddtags": "host:TestSendgziplogswithexplicitcompressionreturnsResponsefromserveralways200emptyJSONresponse1790605175", + "message": "Test-Send_gzip_logs_with_explicit_compression_returns_Response_from_server_always_200_empty_JSON_response-1790605175" + } + ] + }, + "content_type": "application/json", + "method": "POST", + "path": "/v1/input", + "query": [] + }, + "response": { + "body": { + "encoding": "text", + "value": "{}" + }, + "headers": { + "content-type": "application/json" + }, + "reason": "OK", + "status": 200 + } + } + ], + "scenario": "Send gzip logs with explicit compression returns \"Response from server (always 200 empty JSON).\" response", + "version": "v1" + }, { "feature": "Logs", "frozen_at": "2022-04-12T14:46:46.337Z", diff --git a/src/test/resources/generated-test/test-server-data/v2/logs.json b/src/test/resources/generated-test/test-server-data/v2/logs.json index d4b92521f23..7b3b1563788 100644 --- a/src/test/resources/generated-test/test-server-data/v2/logs.json +++ b/src/test/resources/generated-test/test-server-data/v2/logs.json @@ -467,6 +467,84 @@ "scenario": "Search logs returns \"OK\" response with pagination", "version": "v2" }, + { + "feature": "Logs", + "frozen_at": "2026-09-28T14:19:38.034Z", + "interactions": [ + { + "request": { + "body": { + "type": "json", + "value": [ + { + "ddsource": "nginx", + "ddtags": "env:staging,version:5.1", + "hostname": "i-012345678", + "message": "2019-11-19T14:37:58,995 INFO [process.name][20081] Hello World", + "service": "payment" + } + ] + }, + "content_type": "application/json", + "method": "POST", + "path": "/api/v2/logs", + "query": [] + }, + "response": { + "body": { + "encoding": "text", + "value": "{}" + }, + "headers": { + "content-type": "application/json" + }, + "reason": "Accepted", + "status": 202 + } + } + ], + "scenario": "Send deflate logs with explicit compression returns \"Request accepted for processing (always 202 empty JSON).\" response", + "version": "v2" + }, + { + "feature": "Logs", + "frozen_at": "2026-09-28T14:19:39.717Z", + "interactions": [ + { + "request": { + "body": { + "type": "json", + "value": [ + { + "ddsource": "nginx", + "ddtags": "env:staging,version:5.1", + "hostname": "i-012345678", + "message": "2019-11-19T14:37:58,995 INFO [process.name][20081] Hello World", + "service": "payment" + } + ] + }, + "content_type": "application/json", + "method": "POST", + "path": "/api/v2/logs", + "query": [] + }, + "response": { + "body": { + "encoding": "text", + "value": "{}" + }, + "headers": { + "content-type": "application/json" + }, + "reason": "Accepted", + "status": 202 + } + } + ], + "scenario": "Send gzip logs with explicit compression returns \"Request accepted for processing (always 202 empty JSON).\" response", + "version": "v2" + }, { "feature": "Logs", "frozen_at": "2024-10-01T15:36:43.563Z",