From b6d737e959435c2f736b5c1c69a3c20c02a330d5 Mon Sep 17 00:00:00 2001 From: Ricken BAZOLO Date: Fri, 31 Jul 2026 12:15:53 +0200 Subject: [PATCH 1/7] Make Anthropic message mapping reusable Signed-off-by: Ricken BAZOLO --- .../ai/anthropic/AnthropicChatModel.java | 75 ++++++++++++------- 1 file changed, 46 insertions(+), 29 deletions(-) diff --git a/models/spring-ai-anthropic/src/main/java/org/springframework/ai/anthropic/AnthropicChatModel.java b/models/spring-ai-anthropic/src/main/java/org/springframework/ai/anthropic/AnthropicChatModel.java index 6b3ea00ec7..2dbea2d15f 100644 --- a/models/spring-ai-anthropic/src/main/java/org/springframework/ai/anthropic/AnthropicChatModel.java +++ b/models/spring-ai-anthropic/src/main/java/org/springframework/ai/anthropic/AnthropicChatModel.java @@ -591,7 +591,7 @@ public ChatResponse internalCall(Prompt prompt, @Nullable ChatResponse previousC return response; } - private static AnthropicChatOptions resolveAnthropicOptions(Prompt prompt) { + static AnthropicChatOptions resolveAnthropicOptions(Prompt prompt) { ChatOptions options = prompt.getOptions(); return options instanceof AnthropicChatOptions anthropicOptions ? anthropicOptions : AnthropicChatOptions.builder().build(); @@ -615,6 +615,23 @@ private static RequestOptions requestOptionsFor(Prompt prompt) { * @return the constructed request parameters */ MessageCreateParams createRequest(Prompt prompt, boolean stream) { + return createRequest(prompt, this.toolCallingManager, this.options.getSkillContainer()); + } + + /** + * Same mapping as {@link #createRequest(Prompt, boolean)}, with the collaborators + * passed in explicitly so that other Anthropic models in this module — notably + * {@link AnthropicBatchModel} — reuse the exact same {@link Prompt} to + * {@link MessageCreateParams} conversion instead of duplicating it. + * @param prompt the prompt with message history and options + * @param toolCallingManager resolves the tool definitions to advertise + * @param defaultSkillContainer fallback skill container when the prompt options carry + * none; may be {@code null} + * @return the constructed request parameters + * @since 2.0.0 + */ + static MessageCreateParams createRequest(Prompt prompt, ToolCallingManager toolCallingManager, + @Nullable AnthropicSkillContainer defaultSkillContainer) { MessageCreateParams.Builder builder = MessageCreateParams.builder(); @@ -856,9 +873,9 @@ else if (message.getMessageType() == MessageType.TOOL) { List allTools = new ArrayList<>(); // Add user-defined tool definitions - List toolDefinitions = this.toolCallingManager.resolveToolDefinitions(requestOptions); + List toolDefinitions = toolCallingManager.resolveToolDefinitions(requestOptions); if (!CollectionUtils.isEmpty(toolDefinitions)) { - List tools = toolDefinitions.stream().map(this::toAnthropicTool).toList(); + List tools = toolDefinitions.stream().map(AnthropicChatModel::toAnthropicTool).toList(); // Apply cache control to the last tool if caching strategy includes tools CacheControlEphemeral toolCacheControl = cacheResolver.resolveToolCacheControl(); @@ -907,8 +924,8 @@ else if (Boolean.TRUE.equals(requestOptions.getDisableParallelToolUse())) { // Skills support AnthropicSkillContainer skillContainer = requestOptions.getSkillContainer(); - if (skillContainer == null && this.options.getSkillContainer() != null) { - skillContainer = this.options.getSkillContainer(); + if (skillContainer == null && defaultSkillContainer != null) { + skillContainer = defaultSkillContainer; } if (skillContainer != null) { // Add container with skills config @@ -953,7 +970,7 @@ else if (Boolean.TRUE.equals(requestOptions.getDisableParallelToolUse())) { * @param lastUserIndex the index of the last user message (inclusive) * @return the combined text of eligible messages */ - private String combineEligibleMessagesText(List messages, + private static String combineEligibleMessagesText(List messages, int lastUserIndex) { StringBuilder combined = new StringBuilder(); for (int i = 0; i <= lastUserIndex && i < messages.size(); i++) { @@ -965,7 +982,7 @@ private String combineEligibleMessagesText(List responses) { + private static String combineToolResponsesText(List responses) { StringBuilder combined = new StringBuilder(); for (ToolResponseMessage.ToolResponse response : responses) { String data = response.responseData(); @@ -984,7 +1001,7 @@ private String combineToolResponsesText(List r * @param webSearchAccumulator collects web search results found in response * @return list of generations with text, tool calls, and/or thinking content */ - private List buildGenerations(Message message, List citationAccumulator, + static List buildGenerations(Message message, List citationAccumulator, List webSearchAccumulator) { List generations = new ArrayList<>(); @@ -1067,7 +1084,7 @@ else if (block.isContainerUpload() || block.isServerToolUse() || block.isBashCod * @param usage the usage information * @return the chat response metadata */ - private ChatResponseMetadata from(Message message, Usage usage, List citations, + static ChatResponseMetadata from(Message message, Usage usage, List citations, List webSearchResults, RateLimit rateLimit) { Assert.notNull(message, "Anthropic Message must not be null"); ChatResponseMetadata.Builder metadataBuilder = ChatResponseMetadata.builder() @@ -1090,7 +1107,7 @@ private ChatResponseMetadata from(Message message, Usage usage, List c * @param usage the Anthropic SDK usage * @return the Spring AI usage */ - private Usage getDefaultUsage(com.anthropic.models.messages.Usage usage) { + static Usage getDefaultUsage(com.anthropic.models.messages.Usage usage) { if (usage == null) { return new EmptyUsage(); } @@ -1103,7 +1120,7 @@ private Usage getDefaultUsage(com.anthropic.models.messages.Usage usage) { Integer.valueOf(Math.toIntExact(inputTokens + outputTokens)), usage, cacheRead, cacheWrite); } - private @Nullable Citation convertTextCitation(TextCitation textCitation) { + private static @Nullable Citation convertTextCitation(TextCitation textCitation) { if (textCitation.isCharLocation()) { return fromCharLocation(textCitation.asCharLocation()); } @@ -1119,7 +1136,7 @@ else if (textCitation.isWebSearchResultLocation()) { return null; } - private @Nullable Citation convertStreamingCitation(CitationsDelta.Citation citation) { + private static @Nullable Citation convertStreamingCitation(CitationsDelta.Citation citation) { if (citation.isCharLocation()) { return fromCharLocation(citation.asCharLocation()); } @@ -1135,22 +1152,22 @@ else if (citation.isWebSearchResultLocation()) { return null; } - private Citation fromCharLocation(CitationCharLocation loc) { + private static Citation fromCharLocation(CitationCharLocation loc) { return Citation.ofCharLocation(loc.citedText(), (int) loc.documentIndex(), loc.documentTitle().orElse(null), (int) loc.startCharIndex(), (int) loc.endCharIndex()); } - private Citation fromPageLocation(CitationPageLocation loc) { + private static Citation fromPageLocation(CitationPageLocation loc) { return Citation.ofPageLocation(loc.citedText(), (int) loc.documentIndex(), loc.documentTitle().orElse(null), (int) loc.startPageNumber(), (int) loc.endPageNumber()); } - private Citation fromContentBlockLocation(CitationContentBlockLocation loc) { + private static Citation fromContentBlockLocation(CitationContentBlockLocation loc) { return Citation.ofContentBlockLocation(loc.citedText(), (int) loc.documentIndex(), loc.documentTitle().orElse(null), (int) loc.startBlockIndex(), (int) loc.endBlockIndex()); } - private Citation fromWebSearchResultLocation(CitationsWebSearchResultLocation loc) { + private static Citation fromWebSearchResultLocation(CitationsWebSearchResultLocation loc) { return Citation.ofWebSearchResultLocation(loc.citedText(), loc.url(), loc.title().orElse(null)); } @@ -1162,7 +1179,7 @@ private Citation fromWebSearchResultLocation(CitationsWebSearchResultLocation lo * @return a valid JSON string * @throws RuntimeException if serialization fails */ - private String convertJsonValueToString(JsonValue jsonValue) { + private static String convertJsonValueToString(JsonValue jsonValue) { try { var jsonMapper = tools.jackson.databind.json.JsonMapper.builder().build(); // Convert to native Java objects first, then serialize with Jackson @@ -1180,7 +1197,7 @@ private String convertJsonValueToString(JsonValue jsonValue) { * @param jsonValue the SDK's JsonValue to convert * @return the equivalent native Java object, or null for JSON null */ - private @Nullable Object convertJsonValueToNative(JsonValue jsonValue) { + private static @Nullable Object convertJsonValueToNative(JsonValue jsonValue) { return jsonValue.accept(new JsonValue.Visitor<@Nullable Object>() { @Override public @Nullable Object visitNull() { @@ -1232,7 +1249,7 @@ public Object visitObject(java.util.Map values) { * @param argumentsJson the JSON string containing tool call arguments * @return a ToolUseBlockParam.Input with the parsed arguments */ - private ToolUseBlockParam.Input buildToolInput(String argumentsJson) { + private static ToolUseBlockParam.Input buildToolInput(String argumentsJson) { ToolUseBlockParam.Input.Builder inputBuilder = ToolUseBlockParam.Input.builder(); if (argumentsJson != null && !argumentsJson.isEmpty()) { try { @@ -1267,7 +1284,7 @@ private ToolUseBlockParam.Input buildToolInput(String argumentsJson) { * @throws RuntimeException if the JSON schema cannot be parsed */ @SuppressWarnings("unchecked") - private Tool toAnthropicTool(ToolDefinition toolDefinition) { + private static Tool toAnthropicTool(ToolDefinition toolDefinition) { try { // Parse the JSON schema string into a Map var jsonMapper = tools.jackson.databind.json.JsonMapper.builder().build(); @@ -1314,7 +1331,7 @@ private Tool toAnthropicTool(ToolDefinition toolDefinition) { * @param webSearchTool the web search configuration * @return the SDK web search tool */ - private WebSearchTool20260209 toSdkWebSearchTool(AnthropicWebSearchTool webSearchTool) { + private static WebSearchTool20260209 toSdkWebSearchTool(AnthropicWebSearchTool webSearchTool) { WebSearchTool20260209.Builder sdkBuilder = WebSearchTool20260209.builder(); if (webSearchTool.getAllowedDomains() != null) { @@ -1355,7 +1372,7 @@ private WebSearchTool20260209 toSdkWebSearchTool(AnthropicWebSearchTool webSearc * @return the appropriate ContentBlockParam (ImageBlockParam or DocumentBlockParam) * @throws IllegalArgumentException if the media type is unsupported */ - private ContentBlockParam getContentBlockParamByMedia(Media media) { + private static ContentBlockParam getContentBlockParamByMedia(Media media) { MimeType mimeType = media.getMimeType(); String data = fromMediaData(media.getData()); @@ -1374,7 +1391,7 @@ else if (isPdfMedia(mimeType)) { * @param mimeType the MIME type to check * @return true if the type is image/* */ - private boolean isImageMedia(MimeType mimeType) { + private static boolean isImageMedia(MimeType mimeType) { return "image".equals(mimeType.getType()); } @@ -1383,7 +1400,7 @@ private boolean isImageMedia(MimeType mimeType) { * @param mimeType the MIME type to check * @return true if the type is application/pdf */ - private boolean isPdfMedia(MimeType mimeType) { + private static boolean isPdfMedia(MimeType mimeType) { return "application".equals(mimeType.getType()) && "pdf".equals(mimeType.getSubtype()); } @@ -1394,7 +1411,7 @@ private boolean isPdfMedia(MimeType mimeType) { * @return base64-encoded string or URL string * @throws IllegalArgumentException if data type is unsupported */ - private String fromMediaData(Object mediaData) { + private static String fromMediaData(Object mediaData) { if (mediaData instanceof byte[] bytes) { return Base64.getEncoder().encodeToString(bytes); } @@ -1411,7 +1428,7 @@ else if (mediaData instanceof String text) { * @param data base64-encoded image data or HTTPS URL * @return the ImageBlockParam wrapped in ContentBlockParam */ - private ContentBlockParam createImageBlockParam(MimeType mimeType, String data) { + private static ContentBlockParam createImageBlockParam(MimeType mimeType, String data) { ImageBlockParam.Source source; if (data.startsWith("https://")) { source = ImageBlockParam.Source.ofUrl(UrlImageSource.builder().url(data).build()); @@ -1428,7 +1445,7 @@ private ContentBlockParam createImageBlockParam(MimeType mimeType, String data) * @param data base64-encoded PDF data or HTTPS URL * @return the DocumentBlockParam wrapped in ContentBlockParam */ - private ContentBlockParam createDocumentBlockParam(String data) { + private static ContentBlockParam createDocumentBlockParam(String data) { DocumentBlockParam.Source source; if (data.startsWith("https://")) { source = DocumentBlockParam.Source.ofUrl(UrlPdfSource.builder().url(data).build()); @@ -1445,7 +1462,7 @@ private ContentBlockParam createDocumentBlockParam(String data) { * @return the SDK media type enum value * @throws IllegalArgumentException if the image type is unsupported */ - private Base64ImageSource.MediaType toSdkImageMediaType(MimeType mimeType) { + private static Base64ImageSource.MediaType toSdkImageMediaType(MimeType mimeType) { String subtype = mimeType.getSubtype(); return switch (subtype) { case "png" -> Base64ImageSource.MediaType.IMAGE_PNG; @@ -1461,7 +1478,7 @@ private Base64ImageSource.MediaType toSdkImageMediaType(MimeType mimeType) { * Applies {@code disableParallelToolUse} to an existing {@link ToolChoice} by * rebuilding the appropriate subtype with the flag set to {@code true}. */ - private ToolChoice applyDisableParallelToolUse(ToolChoice toolChoice) { + private static ToolChoice applyDisableParallelToolUse(ToolChoice toolChoice) { if (toolChoice.isAuto()) { return ToolChoice.ofAuto(toolChoice.asAuto().toBuilder().disableParallelToolUse(true).build()); } From dc654d3686e3793428aad6e367ae7c24d569ab00 Mon Sep 17 00:00:00 2001 From: Ricken BAZOLO Date: Fri, 31 Jul 2026 12:16:15 +0200 Subject: [PATCH 2/7] Add Anthropic batch API types Signed-off-by: Ricken BAZOLO --- .../ai/anthropic/AnthropicBatch.java | 78 ++++++++++++ .../ai/anthropic/AnthropicBatchError.java | 115 ++++++++++++++++++ .../ai/anthropic/AnthropicBatchRequest.java | 93 ++++++++++++++ .../AnthropicBatchRequestCounts.java | 60 +++++++++ .../ai/anthropic/AnthropicBatchResult.java | 112 +++++++++++++++++ .../anthropic/AnthropicBatchResultStatus.java | 57 +++++++++ .../ai/anthropic/AnthropicBatchStatus.java | 83 +++++++++++++ .../anthropic/AnthropicBatchRequestTests.java | 78 ++++++++++++ 8 files changed, 676 insertions(+) create mode 100644 models/spring-ai-anthropic/src/main/java/org/springframework/ai/anthropic/AnthropicBatch.java create mode 100644 models/spring-ai-anthropic/src/main/java/org/springframework/ai/anthropic/AnthropicBatchError.java create mode 100644 models/spring-ai-anthropic/src/main/java/org/springframework/ai/anthropic/AnthropicBatchRequest.java create mode 100644 models/spring-ai-anthropic/src/main/java/org/springframework/ai/anthropic/AnthropicBatchRequestCounts.java create mode 100644 models/spring-ai-anthropic/src/main/java/org/springframework/ai/anthropic/AnthropicBatchResult.java create mode 100644 models/spring-ai-anthropic/src/main/java/org/springframework/ai/anthropic/AnthropicBatchResultStatus.java create mode 100644 models/spring-ai-anthropic/src/main/java/org/springframework/ai/anthropic/AnthropicBatchStatus.java create mode 100644 models/spring-ai-anthropic/src/test/java/org/springframework/ai/anthropic/AnthropicBatchRequestTests.java diff --git a/models/spring-ai-anthropic/src/main/java/org/springframework/ai/anthropic/AnthropicBatch.java b/models/spring-ai-anthropic/src/main/java/org/springframework/ai/anthropic/AnthropicBatch.java new file mode 100644 index 0000000000..789b4c81a2 --- /dev/null +++ b/models/spring-ai-anthropic/src/main/java/org/springframework/ai/anthropic/AnthropicBatch.java @@ -0,0 +1,78 @@ +/* + * Copyright 2023-present the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.ai.anthropic; + +import java.time.OffsetDateTime; + +import com.anthropic.models.messages.batches.MessageBatch; +import org.jspecify.annotations.Nullable; + +/** + * State of an Anthropic message batch, as returned by + * {@link AnthropicBatchModel#submit(java.util.List)}, + * {@link AnthropicBatchModel#retrieve(String)} and + * {@link AnthropicBatchModel#cancel(String)}. + * + *

+ * Batch processing is asynchronous and can take up to 24 hours. Applications are expected + * to persist {@link #id()} and poll {@link AnthropicBatchModel#retrieve(String)} on their + * own schedule — Spring AI performs no polling. + * + * @param id the batch identifier, to be persisted by the application for later polling + * and result retrieval + * @param status the current processing status + * @param requestCounts per-outcome request counters + * @param createdAt when the batch was created + * @param expiresAt when the batch expires; requests not completed by then are reported as + * {@link AnthropicBatchResultStatus#EXPIRED} + * @param endedAt when processing finished, or {@code null} while still in progress + * @param cancelInitiatedAt when cancellation was requested, or {@code null} + * @param archivedAt when the batch was archived, or {@code null} + * @param resultsUrl the URL of the JSONL results, or {@code null} until processing ends; + * prefer {@link AnthropicBatchModel#results(String)} over fetching it directly + * @author Ricken Bazolo + * @since 2.0.0 + */ +public record AnthropicBatch(String id, AnthropicBatchStatus status, AnthropicBatchRequestCounts requestCounts, + OffsetDateTime createdAt, OffsetDateTime expiresAt, @Nullable OffsetDateTime endedAt, + @Nullable OffsetDateTime cancelInitiatedAt, @Nullable OffsetDateTime archivedAt, @Nullable String resultsUrl) { + + /** + * Whether processing has finished and results can be read. + * @return {@code true} when the status is {@link AnthropicBatchStatus#ENDED} + */ + public boolean isEnded() { + return this.status == AnthropicBatchStatus.ENDED; + } + + /** + * Whether cancellation has been requested for this batch. + * @return {@code true} when the status is {@link AnthropicBatchStatus#CANCELING} + */ + public boolean isCanceling() { + return this.status == AnthropicBatchStatus.CANCELING; + } + + static AnthropicBatch from(MessageBatch messageBatch) { + return new AnthropicBatch(messageBatch.id(), AnthropicBatchStatus.from(messageBatch.processingStatus()), + AnthropicBatchRequestCounts.from(messageBatch.requestCounts()), messageBatch.createdAt(), + messageBatch.expiresAt(), messageBatch.endedAt().orElse(null), + messageBatch.cancelInitiatedAt().orElse(null), messageBatch.archivedAt().orElse(null), + messageBatch.resultsUrl().orElse(null)); + } + +} diff --git a/models/spring-ai-anthropic/src/main/java/org/springframework/ai/anthropic/AnthropicBatchError.java b/models/spring-ai-anthropic/src/main/java/org/springframework/ai/anthropic/AnthropicBatchError.java new file mode 100644 index 0000000000..ca217561b7 --- /dev/null +++ b/models/spring-ai-anthropic/src/main/java/org/springframework/ai/anthropic/AnthropicBatchError.java @@ -0,0 +1,115 @@ +/* + * Copyright 2023-present the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.ai.anthropic; + +import com.anthropic.core.JsonValue; +import com.anthropic.models.ApiErrorObject; +import com.anthropic.models.AuthenticationError; +import com.anthropic.models.BillingError; +import com.anthropic.models.ErrorObject; +import com.anthropic.models.ErrorResponse; +import com.anthropic.models.GatewayTimeoutError; +import com.anthropic.models.InvalidRequestError; +import com.anthropic.models.NotFoundError; +import com.anthropic.models.OverloadedError; +import com.anthropic.models.PermissionError; +import com.anthropic.models.RateLimitError; +import org.jspecify.annotations.Nullable; + +/** + * Failure of a single request inside an Anthropic message batch. + * + *

+ * Individual failures are surfaced per request rather than thrown, so that a single bad + * entry does not hide the results of the rest of the batch. + * + * @param type the Anthropic error type, for example {@code invalid_request_error}, + * {@code rate_limit_error} or {@code overloaded_error} + * @param message the human-readable error message + * @param requestId the Anthropic request identifier, when reported + * @author Ricken Bazolo + * @since 2.0.0 + * @see Anthropic error types + */ +public record AnthropicBatchError(String type, String message, @Nullable String requestId) { + + private static final String UNKNOWN_TYPE = "unknown_error"; + + static AnthropicBatchError from(ErrorResponse errorResponse) { + String requestId = errorResponse.requestId().orElse(null); + AnthropicBatchError error = errorResponse.error().accept(new ErrorObject.Visitor() { + @Override + public AnthropicBatchError visitInvalidRequestError(InvalidRequestError error) { + return of(error._type(), error.message(), requestId); + } + + @Override + public AnthropicBatchError visitAuthenticationError(AuthenticationError error) { + return of(error._type(), error.message(), requestId); + } + + @Override + public AnthropicBatchError visitBillingError(BillingError error) { + return of(error._type(), error.message(), requestId); + } + + @Override + public AnthropicBatchError visitPermissionError(PermissionError error) { + return of(error._type(), error.message(), requestId); + } + + @Override + public AnthropicBatchError visitNotFoundError(NotFoundError error) { + return of(error._type(), error.message(), requestId); + } + + @Override + public AnthropicBatchError visitRateLimitError(RateLimitError error) { + return of(error._type(), error.message(), requestId); + } + + @Override + public AnthropicBatchError visitTimeoutError(GatewayTimeoutError error) { + return of(error._type(), error.message(), requestId); + } + + @Override + public AnthropicBatchError visitApiError(ApiErrorObject error) { + return of(error._type(), error.message(), requestId); + } + + @Override + public AnthropicBatchError visitOverloadedError(OverloadedError error) { + return of(error._type(), error.message(), requestId); + } + + @Override + public AnthropicBatchError unknown(@Nullable JsonValue json) { + return new AnthropicBatchError(UNKNOWN_TYPE, String.valueOf(json), requestId); + } + }); + return error; + } + + private static AnthropicBatchError of(JsonValue type, String message, @Nullable String requestId) { + // JsonValue extends the raw JsonField type, so asString() erases to + // Optional. + Object typeValue = type.asString().orElse(null); + return new AnthropicBatchError(typeValue instanceof String text ? text : UNKNOWN_TYPE, message, requestId); + } + +} diff --git a/models/spring-ai-anthropic/src/main/java/org/springframework/ai/anthropic/AnthropicBatchRequest.java b/models/spring-ai-anthropic/src/main/java/org/springframework/ai/anthropic/AnthropicBatchRequest.java new file mode 100644 index 0000000000..e71559d046 --- /dev/null +++ b/models/spring-ai-anthropic/src/main/java/org/springframework/ai/anthropic/AnthropicBatchRequest.java @@ -0,0 +1,93 @@ +/* + * Copyright 2023-present the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.ai.anthropic; + +import java.util.regex.Pattern; + +import org.springframework.ai.chat.prompt.Prompt; +import org.springframework.util.Assert; + +/** + * A single entry of a batch submitted through {@link AnthropicBatchModel}. + * + *

+ * The {@code customId} is the only way to correlate a result with its request: the + * Anthropic API does not guarantee that results come back in submission order. Use + * an identifier your application can resolve back to its own domain entity, and keep it + * stored alongside the batch id so that correlation survives a restart. + * + *

+ * The {@link Prompt} is mapped exactly like it would be for + * {@link AnthropicChatModel#call(Prompt)} — including system messages, conversation + * history, images and PDF documents, prompt caching, thinking, structured output and tool + * definitions. Per-request model and options are taken from {@link Prompt#getOptions()} + * when it carries an {@link AnthropicChatOptions}; otherwise the batch model's default + * options apply. + * + * @param customId the caller-defined correlation identifier; 1 to 64 characters, limited + * to letters, digits, underscores and hyphens + * @param prompt the prompt to run + * @author Ricken Bazolo + * @since 2.0.0 + */ +public record AnthropicBatchRequest(String customId, Prompt prompt) { + + /** + * The identifier format accepted by the Anthropic Message Batches API. + */ + private static final Pattern CUSTOM_ID_PATTERN = Pattern.compile("^[a-zA-Z0-9_-]{1,64}$"); + + public AnthropicBatchRequest { + Assert.hasText(customId, "customId must not be empty"); + Assert.isTrue(CUSTOM_ID_PATTERN.matcher(customId).matches(), + () -> "customId must match " + CUSTOM_ID_PATTERN.pattern() + " but was: '" + customId + "'"); + Assert.notNull(prompt, "prompt must not be null"); + } + + /** + * Creates a batch request from a correlation identifier and a prompt. + * @param customId the caller-defined correlation identifier + * @param prompt the prompt to run + * @return the batch request + */ + public static AnthropicBatchRequest of(String customId, Prompt prompt) { + return new AnthropicBatchRequest(customId, prompt); + } + + /** + * Creates a batch request from a correlation identifier and a plain user message. + * @param customId the caller-defined correlation identifier + * @param userText the user message content + * @return the batch request, using the batch model's default options + */ + public static AnthropicBatchRequest of(String customId, String userText) { + return new AnthropicBatchRequest(customId, new Prompt(userText)); + } + + /** + * Creates a batch request from a correlation identifier, a plain user message and + * per-request options. + * @param customId the caller-defined correlation identifier + * @param userText the user message content + * @param options the Anthropic options for this entry + * @return the batch request + */ + public static AnthropicBatchRequest of(String customId, String userText, AnthropicChatOptions options) { + return new AnthropicBatchRequest(customId, new Prompt(userText, options)); + } + +} diff --git a/models/spring-ai-anthropic/src/main/java/org/springframework/ai/anthropic/AnthropicBatchRequestCounts.java b/models/spring-ai-anthropic/src/main/java/org/springframework/ai/anthropic/AnthropicBatchRequestCounts.java new file mode 100644 index 0000000000..c923e369b7 --- /dev/null +++ b/models/spring-ai-anthropic/src/main/java/org/springframework/ai/anthropic/AnthropicBatchRequestCounts.java @@ -0,0 +1,60 @@ +/* + * Copyright 2023-present the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.ai.anthropic; + +import com.anthropic.models.messages.batches.MessageBatchRequestCounts; + +/** + * Per-outcome request counters of an Anthropic message batch. + * + *

+ * While a batch is {@link AnthropicBatchStatus#IN_PROGRESS in progress} these counters + * let an application report progress without reading the (potentially large) result + * stream. + * + * @param processing number of requests still being processed + * @param succeeded number of requests that completed successfully + * @param errored number of requests that failed + * @param canceled number of requests canceled before completion + * @param expired number of requests that expired before completion + * @author Ricken Bazolo + * @since 2.0.0 + */ +public record AnthropicBatchRequestCounts(long processing, long succeeded, long errored, long canceled, long expired) { + + /** + * Returns the total number of requests tracked by this batch. + * @return the sum of every counter + */ + public long total() { + return this.processing + this.succeeded + this.errored + this.canceled + this.expired; + } + + /** + * Returns the number of requests that reached a terminal outcome. + * @return the sum of every counter except {@link #processing()} + */ + public long completed() { + return this.succeeded + this.errored + this.canceled + this.expired; + } + + static AnthropicBatchRequestCounts from(MessageBatchRequestCounts counts) { + return new AnthropicBatchRequestCounts(counts.processing(), counts.succeeded(), counts.errored(), + counts.canceled(), counts.expired()); + } + +} diff --git a/models/spring-ai-anthropic/src/main/java/org/springframework/ai/anthropic/AnthropicBatchResult.java b/models/spring-ai-anthropic/src/main/java/org/springframework/ai/anthropic/AnthropicBatchResult.java new file mode 100644 index 0000000000..509106e277 --- /dev/null +++ b/models/spring-ai-anthropic/src/main/java/org/springframework/ai/anthropic/AnthropicBatchResult.java @@ -0,0 +1,112 @@ +/* + * Copyright 2023-present the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.ai.anthropic; + +import org.jspecify.annotations.Nullable; + +import org.springframework.ai.chat.metadata.Usage; +import org.springframework.ai.chat.model.ChatResponse; + +/** + * Result of a single request inside an Anthropic message batch, correlated to its request + * through {@link #customId()}. + * + *

+ * Results are not returned in submission order. Key them by {@code customId} + * rather than by position in the stream. + * + *

+ * Exactly one of {@link #chatResponse()} / {@link #error()} is populated: a succeeded + * result carries the {@link ChatResponse}, an errored one carries the + * {@link AnthropicBatchError}, and canceled or expired results carry neither. + * + * @param customId the correlation identifier supplied on the matching + * {@link AnthropicBatchRequest} + * @param status the terminal outcome of this request + * @param chatResponse the response, converted exactly as + * {@link AnthropicChatModel#call(org.springframework.ai.chat.prompt.Prompt)} would; only + * present when the status is {@link AnthropicBatchResultStatus#SUCCEEDED} + * @param error the failure detail; only present when the status is + * {@link AnthropicBatchResultStatus#ERRORED} + * @author Ricken Bazolo + * @since 2.0.0 + */ +public record AnthropicBatchResult(String customId, AnthropicBatchResultStatus status, + @Nullable ChatResponse chatResponse, @Nullable AnthropicBatchError error) { + + /** + * Creates a succeeded result. + * @param customId the correlation identifier + * @param chatResponse the converted response + * @return the result + */ + static AnthropicBatchResult succeeded(String customId, ChatResponse chatResponse) { + return new AnthropicBatchResult(customId, AnthropicBatchResultStatus.SUCCEEDED, chatResponse, null); + } + + /** + * Creates an errored result. + * @param customId the correlation identifier + * @param error the failure detail + * @return the result + */ + static AnthropicBatchResult errored(String customId, AnthropicBatchError error) { + return new AnthropicBatchResult(customId, AnthropicBatchResultStatus.ERRORED, null, error); + } + + /** + * Creates a result with no payload, for canceled, expired or unrecognised outcomes. + * @param customId the correlation identifier + * @param status the terminal outcome + * @return the result + */ + static AnthropicBatchResult of(String customId, AnthropicBatchResultStatus status) { + return new AnthropicBatchResult(customId, status, null, null); + } + + /** + * Whether this request completed successfully. + * @return {@code true} when the status is + * {@link AnthropicBatchResultStatus#SUCCEEDED} + */ + public boolean isSucceeded() { + return this.status == AnthropicBatchResultStatus.SUCCEEDED; + } + + /** + * Returns the token usage reported for this request, for cost accounting. + * @return the usage, or {@code null} when this request produced no response + */ + public @Nullable Usage usage() { + if (this.chatResponse == null || this.chatResponse.getMetadata() == null) { + return null; + } + return this.chatResponse.getMetadata().getUsage(); + } + + /** + * Returns the aggregated text of the response. + * @return the response text, or {@code null} when this request produced no response + */ + public @Nullable String getText() { + if (this.chatResponse == null || this.chatResponse.getResult() == null) { + return null; + } + return this.chatResponse.getResult().getOutput().getText(); + } + +} diff --git a/models/spring-ai-anthropic/src/main/java/org/springframework/ai/anthropic/AnthropicBatchResultStatus.java b/models/spring-ai-anthropic/src/main/java/org/springframework/ai/anthropic/AnthropicBatchResultStatus.java new file mode 100644 index 0000000000..d1e328aff2 --- /dev/null +++ b/models/spring-ai-anthropic/src/main/java/org/springframework/ai/anthropic/AnthropicBatchResultStatus.java @@ -0,0 +1,57 @@ +/* + * Copyright 2023-present the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.ai.anthropic; + +/** + * Terminal outcome of a single request inside an Anthropic message batch. + * + *

+ * Outcomes are per request, not per batch: one errored entry does not prevent the other + * entries of the same batch from succeeding. + * + * @author Ricken Bazolo + * @since 2.0.0 + */ +public enum AnthropicBatchResultStatus { + + /** + * The request completed and a message is available. + */ + SUCCEEDED, + + /** + * The request failed; see {@link AnthropicBatchResult#error()} for the reason. + */ + ERRORED, + + /** + * The request was canceled before completion, following a + * {@link AnthropicBatchModel#cancel(String)} call. + */ + CANCELED, + + /** + * The request did not complete before the batch expired. + */ + EXPIRED, + + /** + * An outcome returned by the API that this version of Spring AI does not know about. + */ + UNKNOWN + +} diff --git a/models/spring-ai-anthropic/src/main/java/org/springframework/ai/anthropic/AnthropicBatchStatus.java b/models/spring-ai-anthropic/src/main/java/org/springframework/ai/anthropic/AnthropicBatchStatus.java new file mode 100644 index 0000000000..340079685c --- /dev/null +++ b/models/spring-ai-anthropic/src/main/java/org/springframework/ai/anthropic/AnthropicBatchStatus.java @@ -0,0 +1,83 @@ +/* + * Copyright 2023-present the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.ai.anthropic; + +import com.anthropic.models.messages.batches.MessageBatch; + +/** + * Processing status of an Anthropic message batch. + * + * @author Ricken Bazolo + * @since 2.0.0 + * @see Anthropic + * Message Batches API + */ +public enum AnthropicBatchStatus { + + /** + * The batch has been accepted and its requests are being processed. + */ + IN_PROGRESS("in_progress"), + + /** + * Cancellation has been requested; requests already completed keep their result while + * the remaining ones are canceled. + */ + CANCELING("canceling"), + + /** + * Processing has finished. Results are available for reading, and every request has a + * terminal outcome (succeeded, errored, canceled or expired). + */ + ENDED("ended"), + + /** + * A status returned by the API that this version of Spring AI does not know about. + */ + UNKNOWN("unknown"); + + private final String value; + + AnthropicBatchStatus(String value) { + this.value = value; + } + + /** + * Returns the wire value used by the Anthropic API. + * @return the wire value + */ + public String getValue() { + return this.value; + } + + /** + * Maps the SDK processing status onto this enum, returning {@link #UNKNOWN} for + * values added by the API after this release. + * @param processingStatus the SDK processing status + * @return the corresponding status + */ + static AnthropicBatchStatus from(MessageBatch.ProcessingStatus processingStatus) { + String value = processingStatus.asString(); + for (AnthropicBatchStatus status : values()) { + if (status.value.equals(value)) { + return status; + } + } + return UNKNOWN; + } + +} diff --git a/models/spring-ai-anthropic/src/test/java/org/springframework/ai/anthropic/AnthropicBatchRequestTests.java b/models/spring-ai-anthropic/src/test/java/org/springframework/ai/anthropic/AnthropicBatchRequestTests.java new file mode 100644 index 0000000000..2486f9f383 --- /dev/null +++ b/models/spring-ai-anthropic/src/test/java/org/springframework/ai/anthropic/AnthropicBatchRequestTests.java @@ -0,0 +1,78 @@ +/* + * Copyright 2023-present the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.ai.anthropic; + +import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.ValueSource; + +import org.springframework.ai.chat.prompt.Prompt; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatIllegalArgumentException; + +/** + * Tests for {@link AnthropicBatchRequest} correlation-identifier validation. A malformed + * {@code customId} is rejected client-side rather than costing a round trip, because it + * is the only handle the caller has to match a result back to its request. + * + * @author Ricken Bazolo + */ +class AnthropicBatchRequestTests { + + @Test + void acceptsValidCustomIds() { + assertThat(AnthropicBatchRequest.of("invoice_42-A", "hello").customId()).isEqualTo("invoice_42-A"); + assertThat(AnthropicBatchRequest.of("a".repeat(64), "hello").customId()).hasSize(64); + } + + @Test + void carriesThePromptAndItsOptions() { + AnthropicChatOptions options = AnthropicChatOptions.builder().model("claude-haiku-4-5").maxTokens(32).build(); + + AnthropicBatchRequest request = AnthropicBatchRequest.of("req-1", "hello", options); + + assertThat(request.prompt().getOptions()).isSameAs(options); + assertThat(request.prompt().getContents()).contains("hello"); + } + + @Test + void acceptsAnExplicitPrompt() { + Prompt prompt = new Prompt("hello"); + + assertThat(AnthropicBatchRequest.of("req-1", prompt).prompt()).isSameAs(prompt); + } + + @ParameterizedTest + @ValueSource(strings = { " ", "with space", "with/slash", "with:colon", "accentué" }) + void rejectsMalformedCustomIds(String customId) { + assertThatIllegalArgumentException().isThrownBy(() -> AnthropicBatchRequest.of(customId, "hello")); + } + + @Test + void rejectsAnEmptyCustomId() { + assertThatIllegalArgumentException().isThrownBy(() -> AnthropicBatchRequest.of("", "hello")) + .withMessageContaining("customId must not be empty"); + } + + @Test + void rejectsACustomIdLongerThanTheApiLimit() { + assertThatIllegalArgumentException().isThrownBy(() -> AnthropicBatchRequest.of("a".repeat(65), "hello")) + .withMessageContaining("customId must match"); + } + +} From 05b338ec2b4de3f4a3b7923782bccab7081dbad7 Mon Sep 17 00:00:00 2001 From: Ricken BAZOLO Date: Fri, 31 Jul 2026 12:16:32 +0200 Subject: [PATCH 3/7] Add Anthropic batch observability Signed-off-by: Ricken BAZOLO --- .../AnthropicBatchObservationContext.java | 206 ++++++++++++++++++ .../AnthropicBatchObservationConvention.java | 36 +++ ...nthropicBatchObservationDocumentation.java | 178 +++++++++++++++ ...ltAnthropicBatchObservationConvention.java | 105 +++++++++ 4 files changed, 525 insertions(+) create mode 100644 models/spring-ai-anthropic/src/main/java/org/springframework/ai/anthropic/AnthropicBatchObservationContext.java create mode 100644 models/spring-ai-anthropic/src/main/java/org/springframework/ai/anthropic/AnthropicBatchObservationConvention.java create mode 100644 models/spring-ai-anthropic/src/main/java/org/springframework/ai/anthropic/AnthropicBatchObservationDocumentation.java create mode 100644 models/spring-ai-anthropic/src/main/java/org/springframework/ai/anthropic/DefaultAnthropicBatchObservationConvention.java diff --git a/models/spring-ai-anthropic/src/main/java/org/springframework/ai/anthropic/AnthropicBatchObservationContext.java b/models/spring-ai-anthropic/src/main/java/org/springframework/ai/anthropic/AnthropicBatchObservationContext.java new file mode 100644 index 0000000000..18e801fa38 --- /dev/null +++ b/models/spring-ai-anthropic/src/main/java/org/springframework/ai/anthropic/AnthropicBatchObservationContext.java @@ -0,0 +1,206 @@ +/* + * Copyright 2023-present the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.ai.anthropic; + +import io.micrometer.observation.Observation; +import org.jspecify.annotations.Nullable; + +import org.springframework.util.Assert; + +/** + * Context used to store metadata for Anthropic Message Batches operations. + * + *

+ * Deliberately carries no batch identifier and no prompt or generated content: batch ids + * are unbounded in cardinality and prompts must not leak into metric tags. + * + * @author Ricken Bazolo + * @since 2.0.0 + */ +public class AnthropicBatchObservationContext extends Observation.Context { + + private final Operation operation; + + private final String provider; + + private final @Nullable String requestModel; + + private final @Nullable Integer requestCount; + + private @Nullable AnthropicBatch batch; + + private @Nullable AnthropicBatchRequestCounts requestCounts; + + AnthropicBatchObservationContext(Operation operation, String provider, @Nullable String requestModel, + @Nullable Integer requestCount) { + this.operation = operation; + this.provider = provider; + this.requestModel = requestModel; + this.requestCount = requestCount; + } + + public static Builder builder() { + return new Builder(); + } + + /** + * The batch operation being observed. + * @return the operation + */ + public Operation getOperation() { + return this.operation; + } + + /** + * The model provider as identified by the client instrumentation. + * @return the provider + */ + public String getProvider() { + return this.provider; + } + + /** + * The model requested for the batch entries, when a single model applies to all of + * them. + * @return the model, or {@code null} + */ + public @Nullable String getRequestModel() { + return this.requestModel; + } + + /** + * The number of requests submitted, for {@link Operation#CREATE}. + * @return the request count, or {@code null} + */ + public @Nullable Integer getRequestCount() { + return this.requestCount; + } + + /** + * The batch returned by the operation, when it returns one. + * @return the batch, or {@code null} + */ + public @Nullable AnthropicBatch getBatch() { + return this.batch; + } + + public void setBatch(@Nullable AnthropicBatch batch) { + this.batch = batch; + if (batch != null) { + this.requestCounts = batch.requestCounts(); + } + } + + /** + * The per-outcome counters observed for this operation. Set from the batch for + * control-plane operations, and accumulated while streaming for + * {@link Operation#RESULTS}. + * @return the counters, or {@code null} + */ + public @Nullable AnthropicBatchRequestCounts getRequestCounts() { + return this.requestCounts; + } + + public void setRequestCounts(@Nullable AnthropicBatchRequestCounts requestCounts) { + this.requestCounts = requestCounts; + } + + /** + * The Anthropic Message Batches operations that Spring AI observes. + */ + public enum Operation { + + /** + * Batch submission, {@code POST /v1/messages/batches}. + */ + CREATE("batch_create"), + + /** + * Status lookup, {@code GET /v1/messages/batches/{id}}. + */ + RETRIEVE("batch_retrieve"), + + /** + * Result streaming, {@code GET /v1/messages/batches/{id}/results}. + */ + RESULTS("batch_results"), + + /** + * Cancellation, {@code POST /v1/messages/batches/{id}/cancel}. + */ + CANCEL("batch_cancel"), + + /** + * Deletion, {@code DELETE /v1/messages/batches/{id}}. + */ + DELETE("batch_delete"); + + private final String value; + + Operation(String value) { + this.value = value; + } + + public String value() { + return this.value; + } + + } + + public static final class Builder { + + private @Nullable Operation operation; + + private @Nullable String provider; + + private @Nullable String requestModel; + + private @Nullable Integer requestCount; + + private Builder() { + } + + public Builder operation(Operation operation) { + this.operation = operation; + return this; + } + + public Builder provider(String provider) { + this.provider = provider; + return this; + } + + public Builder requestModel(@Nullable String requestModel) { + this.requestModel = requestModel; + return this; + } + + public Builder requestCount(@Nullable Integer requestCount) { + this.requestCount = requestCount; + return this; + } + + public AnthropicBatchObservationContext build() { + Assert.state(this.operation != null, "Operation must not be null"); + Assert.state(this.provider != null, "Provider must not be null"); + return new AnthropicBatchObservationContext(this.operation, this.provider, this.requestModel, + this.requestCount); + } + + } + +} diff --git a/models/spring-ai-anthropic/src/main/java/org/springframework/ai/anthropic/AnthropicBatchObservationConvention.java b/models/spring-ai-anthropic/src/main/java/org/springframework/ai/anthropic/AnthropicBatchObservationConvention.java new file mode 100644 index 0000000000..c352577bff --- /dev/null +++ b/models/spring-ai-anthropic/src/main/java/org/springframework/ai/anthropic/AnthropicBatchObservationConvention.java @@ -0,0 +1,36 @@ +/* + * Copyright 2023-present the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.ai.anthropic; + +import io.micrometer.observation.Observation; +import io.micrometer.observation.ObservationConvention; + +/** + * Interface for an {@link ObservationConvention} for Anthropic Message Batches + * operations. + * + * @author Ricken Bazolo + * @since 2.0.0 + */ +public interface AnthropicBatchObservationConvention extends ObservationConvention { + + @Override + default boolean supportsContext(Observation.Context context) { + return context instanceof AnthropicBatchObservationContext; + } + +} diff --git a/models/spring-ai-anthropic/src/main/java/org/springframework/ai/anthropic/AnthropicBatchObservationDocumentation.java b/models/spring-ai-anthropic/src/main/java/org/springframework/ai/anthropic/AnthropicBatchObservationDocumentation.java new file mode 100644 index 0000000000..81fe6f38fe --- /dev/null +++ b/models/spring-ai-anthropic/src/main/java/org/springframework/ai/anthropic/AnthropicBatchObservationDocumentation.java @@ -0,0 +1,178 @@ +/* + * Copyright 2023-present the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.ai.anthropic; + +import io.micrometer.common.docs.KeyName; +import io.micrometer.observation.Observation; +import io.micrometer.observation.ObservationConvention; +import io.micrometer.observation.docs.ObservationDocumentation; + +import org.springframework.ai.observation.conventions.AiObservationAttributes; + +/** + * Documented conventions for Anthropic Message Batches observations. + * + *

+ * Neither the batch identifier nor any prompt or generated content is exposed as a tag: + * batch ids are unbounded in cardinality and prompt content must not reach a metrics + * backend. + * + * @author Ricken Bazolo + * @since 2.0.0 + */ +public enum AnthropicBatchObservationDocumentation implements ObservationDocumentation { + + /** + * Observation emitted around each Anthropic Message Batches operation. + */ + BATCH_MODEL_OPERATION { + @Override + public Class> getDefaultConvention() { + return DefaultAnthropicBatchObservationConvention.class; + } + + @Override + public KeyName[] getLowCardinalityKeyNames() { + return LowCardinalityKeyNames.values(); + } + + @Override + public KeyName[] getHighCardinalityKeyNames() { + return HighCardinalityKeyNames.values(); + } + + }; + + /** + * Low-cardinality observation key names for batch operations. + */ + public enum LowCardinalityKeyNames implements KeyName { + + /** + * The batch operation being performed: {@code batch_create}, + * {@code batch_retrieve}, {@code batch_results}, {@code batch_cancel} or + * {@code batch_delete}. + */ + AI_OPERATION_TYPE { + @Override + public String asString() { + return AiObservationAttributes.AI_OPERATION_TYPE.value(); + } + }, + + /** + * The model provider as identified by the client instrumentation. + */ + AI_PROVIDER { + @Override + public String asString() { + return AiObservationAttributes.AI_PROVIDER.value(); + } + }, + + /** + * The name of the model the batch entries target, or {@code none} when the + * entries do not share a single model. + */ + REQUEST_MODEL { + @Override + public String asString() { + return AiObservationAttributes.REQUEST_MODEL.value(); + } + }, + + /** + * The processing status of the batch: {@code in_progress}, {@code canceling}, + * {@code ended}, or {@code none} when the operation returns no batch. + */ + BATCH_STATUS { + @Override + public String asString() { + return "spring.ai.anthropic.batch.status"; + } + } + + } + + /** + * High-cardinality observation key names for batch operations. + */ + public enum HighCardinalityKeyNames implements KeyName { + + /** + * The number of requests submitted with the batch. + */ + BATCH_REQUEST_COUNT { + @Override + public String asString() { + return "spring.ai.anthropic.batch.request.count"; + } + }, + + /** + * The number of requests still being processed. + */ + BATCH_PROCESSING_COUNT { + @Override + public String asString() { + return "spring.ai.anthropic.batch.counts.processing"; + } + }, + + /** + * The number of requests that completed successfully. + */ + BATCH_SUCCEEDED_COUNT { + @Override + public String asString() { + return "spring.ai.anthropic.batch.counts.succeeded"; + } + }, + + /** + * The number of requests that failed. + */ + BATCH_ERRORED_COUNT { + @Override + public String asString() { + return "spring.ai.anthropic.batch.counts.errored"; + } + }, + + /** + * The number of requests canceled before completion. + */ + BATCH_CANCELED_COUNT { + @Override + public String asString() { + return "spring.ai.anthropic.batch.counts.canceled"; + } + }, + + /** + * The number of requests that expired before completion. + */ + BATCH_EXPIRED_COUNT { + @Override + public String asString() { + return "spring.ai.anthropic.batch.counts.expired"; + } + } + + } + +} diff --git a/models/spring-ai-anthropic/src/main/java/org/springframework/ai/anthropic/DefaultAnthropicBatchObservationConvention.java b/models/spring-ai-anthropic/src/main/java/org/springframework/ai/anthropic/DefaultAnthropicBatchObservationConvention.java new file mode 100644 index 0000000000..6dde49a0f3 --- /dev/null +++ b/models/spring-ai-anthropic/src/main/java/org/springframework/ai/anthropic/DefaultAnthropicBatchObservationConvention.java @@ -0,0 +1,105 @@ +/* + * Copyright 2023-present the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.ai.anthropic; + +import io.micrometer.common.KeyValues; + +import org.springframework.ai.anthropic.AnthropicBatchObservationDocumentation.HighCardinalityKeyNames; +import org.springframework.ai.anthropic.AnthropicBatchObservationDocumentation.LowCardinalityKeyNames; + +/** + * Default conventions to populate observations for Anthropic Message Batches operations. + * + * @author Ricken Bazolo + * @since 2.0.0 + */ +public class DefaultAnthropicBatchObservationConvention implements AnthropicBatchObservationConvention { + + public static final String DEFAULT_NAME = "spring.ai.anthropic.batch.operation"; + + private static final String KEY_VALUE_NONE = "none"; + + @Override + public String getName() { + return DEFAULT_NAME; + } + + @Override + public String getContextualName(AnthropicBatchObservationContext context) { + return "%s %s".formatted(context.getOperation().value(), context.getProvider()); + } + + @Override + public KeyValues getLowCardinalityKeyValues(AnthropicBatchObservationContext context) { + KeyValues keyValues = KeyValues.empty(); + keyValues = aiOperationType(keyValues, context); + keyValues = aiProvider(keyValues, context); + keyValues = requestModel(keyValues, context); + keyValues = batchStatus(keyValues, context); + return keyValues; + } + + @Override + public KeyValues getHighCardinalityKeyValues(AnthropicBatchObservationContext context) { + KeyValues keyValues = KeyValues.empty(); + keyValues = requestCount(keyValues, context); + keyValues = requestCounts(keyValues, context); + return keyValues; + } + + private KeyValues aiOperationType(KeyValues keyValues, AnthropicBatchObservationContext context) { + return keyValues.and(LowCardinalityKeyNames.AI_OPERATION_TYPE.asString(), context.getOperation().value()); + } + + private KeyValues aiProvider(KeyValues keyValues, AnthropicBatchObservationContext context) { + return keyValues.and(LowCardinalityKeyNames.AI_PROVIDER.asString(), context.getProvider()); + } + + private KeyValues requestModel(KeyValues keyValues, AnthropicBatchObservationContext context) { + String requestModel = context.getRequestModel(); + return keyValues.and(LowCardinalityKeyNames.REQUEST_MODEL.asString(), + requestModel != null ? requestModel : KEY_VALUE_NONE); + } + + private KeyValues batchStatus(KeyValues keyValues, AnthropicBatchObservationContext context) { + AnthropicBatch batch = context.getBatch(); + return keyValues.and(LowCardinalityKeyNames.BATCH_STATUS.asString(), + batch != null ? batch.status().getValue() : KEY_VALUE_NONE); + } + + private KeyValues requestCount(KeyValues keyValues, AnthropicBatchObservationContext context) { + Integer requestCount = context.getRequestCount(); + if (requestCount == null) { + return keyValues; + } + return keyValues.and(HighCardinalityKeyNames.BATCH_REQUEST_COUNT.asString(), String.valueOf(requestCount)); + } + + private KeyValues requestCounts(KeyValues keyValues, AnthropicBatchObservationContext context) { + AnthropicBatchRequestCounts counts = context.getRequestCounts(); + if (counts == null) { + return keyValues; + } + return keyValues + .and(HighCardinalityKeyNames.BATCH_PROCESSING_COUNT.asString(), String.valueOf(counts.processing())) + .and(HighCardinalityKeyNames.BATCH_SUCCEEDED_COUNT.asString(), String.valueOf(counts.succeeded())) + .and(HighCardinalityKeyNames.BATCH_ERRORED_COUNT.asString(), String.valueOf(counts.errored())) + .and(HighCardinalityKeyNames.BATCH_CANCELED_COUNT.asString(), String.valueOf(counts.canceled())) + .and(HighCardinalityKeyNames.BATCH_EXPIRED_COUNT.asString(), String.valueOf(counts.expired())); + } + +} From 62b2282792009ba71fad713f7c2300d285e8f07c Mon Sep 17 00:00:00 2001 From: Ricken BAZOLO Date: Fri, 31 Jul 2026 12:16:51 +0200 Subject: [PATCH 4/7] Add Anthropic Message Batches model Add AnthropicBatchModel, exposing the five provider operations of the Message Batches API, and its default implementation over the official Anthropic Java SDK. Batch processing is asynchronous and can take up to 24 hours, so this is deliberately not a ChatModel: submit has no immediate ChatResponse to return. The control operations are synchronous; only results is reactive, backed by Flux.using so the SDK stream is closed on termination and on cancellation and a batch with a very large number of entries is never held in memory. Entries are mapped through the same Prompt conversion as AnthropicChatModel.call, so prompt caching, thinking, images and PDF documents, structured output and tool definitions behave identically on both paths, and a succeeded result is converted back into a ChatResponse with the same generations, metadata and usage. Tool calls are not executed: a batch entry cannot be continued mid-flight, so a tool_use response is returned as is. Spring AI provides the provider access, not the orchestration. There is no polling, no persistence and no scheduling: those stay with the application, which owns the batch id and the customId correlations. Signed-off-by: Ricken BAZOLO --- .../ai/anthropic/AnthropicBatchModel.java | 129 ++++ .../anthropic/DefaultAnthropicBatchModel.java | 623 ++++++++++++++++++ .../anthropic/AnthropicBatchModelTests.java | 414 ++++++++++++ 3 files changed, 1166 insertions(+) create mode 100644 models/spring-ai-anthropic/src/main/java/org/springframework/ai/anthropic/AnthropicBatchModel.java create mode 100644 models/spring-ai-anthropic/src/main/java/org/springframework/ai/anthropic/DefaultAnthropicBatchModel.java create mode 100644 models/spring-ai-anthropic/src/test/java/org/springframework/ai/anthropic/AnthropicBatchModelTests.java diff --git a/models/spring-ai-anthropic/src/main/java/org/springframework/ai/anthropic/AnthropicBatchModel.java b/models/spring-ai-anthropic/src/main/java/org/springframework/ai/anthropic/AnthropicBatchModel.java new file mode 100644 index 0000000000..255769fda9 --- /dev/null +++ b/models/spring-ai-anthropic/src/main/java/org/springframework/ai/anthropic/AnthropicBatchModel.java @@ -0,0 +1,129 @@ +/* + * Copyright 2023-present the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.ai.anthropic; + +import java.util.List; + +import reactor.core.publisher.Flux; + +import org.springframework.ai.chat.model.ChatModel; +import org.springframework.ai.chat.model.ChatResponse; +import org.springframework.ai.chat.prompt.Prompt; + +/** + * Anthropic Message + * Batches abstraction: submit many prompts at once, poll the batch, then read the + * per-request results. + * + *

+ * Batch processing is asynchronous and can take up to 24 hours, so this is + * deliberately not a {@link ChatModel}: {@link #submit(List)} has no immediate + * {@link ChatResponse} to return. The control operations are synchronous; only + * {@link #results(String)} is reactive, so a batch with a very large number of entries + * never has to be held in memory. + * + *

+ * Spring AI provides the provider access, not the orchestration. There is no + * automatic polling, no persistence, no scheduling and no business retry: how often to + * call {@link #retrieve(String)}, where to store the batch id and the {@code customId} + * correlations, when to notify and how to account for cost all stay with the application. + * + *

+ * Requests are mapped exactly as they would be for {@link ChatModel#call(Prompt)}, and a + * succeeded result is converted back into a {@link ChatResponse} with the same + * generations, metadata and usage — see {@link DefaultAnthropicBatchModel} for the + * details and for the tool-calling limitation. + * + *

+ * Results are unordered. Correlate them through + * {@link AnthropicBatchResult#customId()}, never by position. + * + * @author Ricken Bazolo + * @since 2.0.0 + * @see AnthropicBatchRequest + * @see AnthropicBatchResult + * @see DefaultAnthropicBatchModel + */ +public interface AnthropicBatchModel { + + /** + * Creates a builder for the default {@link AnthropicBatchModel} implementation. + * @return a new builder instance + */ + static DefaultAnthropicBatchModel.Builder builder() { + return DefaultAnthropicBatchModel.builder(); + } + + /** + * Submits a batch of prompts for asynchronous processing. + * + *

+ * Returns as soon as Anthropic accepts the batch; no request has been processed yet. + * Persist {@link AnthropicBatch#id()} together with the {@code customId} of every + * entry so that polling and correlation survive an application restart. + * @param requests the batch entries; must be non-empty and carry distinct + * {@code customId} values + * @return the accepted batch, in state {@link AnthropicBatchStatus#IN_PROGRESS} + * @throws IllegalArgumentException if the list is empty or a {@code customId} is + * duplicated + */ + AnthropicBatch submit(List requests); + + /** + * Retrieves the current state of a batch. + * + *

+ * Call this on the application's own schedule; Spring AI performs no polling. Results + * become readable once {@link AnthropicBatch#isEnded()} is {@code true}. + * @param batchId the batch identifier returned by {@link #submit(List)} + * @return the current batch state + */ + AnthropicBatch retrieve(String batchId); + + /** + * Streams the results of an ended batch. + * + *

+ * Results arrive in an unspecified order: key them by + * {@link AnthropicBatchResult#customId()}. Individual failures are emitted as + * {@link AnthropicBatchResultStatus#ERRORED} items rather than thrown, so one bad + * entry never hides the rest. + * @param batchId the batch identifier + * @return a lazily-populated flux of per-request results + */ + Flux results(String batchId); + + /** + * Requests cancellation of a batch. + * + *

+ * Cancellation is not immediate: the batch moves to + * {@link AnthropicBatchStatus#CANCELING} and requests that already completed keep + * their result, while the remaining ones end up as + * {@link AnthropicBatchResultStatus#CANCELED}. + * @param batchId the batch identifier + * @return the batch state after the cancellation request + */ + AnthropicBatch cancel(String batchId); + + /** + * Deletes a batch. Only batches whose processing has ended can be deleted. + * @param batchId the batch identifier + */ + void delete(String batchId); + +} diff --git a/models/spring-ai-anthropic/src/main/java/org/springframework/ai/anthropic/DefaultAnthropicBatchModel.java b/models/spring-ai-anthropic/src/main/java/org/springframework/ai/anthropic/DefaultAnthropicBatchModel.java new file mode 100644 index 0000000000..5d580798d0 --- /dev/null +++ b/models/spring-ai-anthropic/src/main/java/org/springframework/ai/anthropic/DefaultAnthropicBatchModel.java @@ -0,0 +1,623 @@ +/* + * Copyright 2023-present the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.ai.anthropic; + +import java.time.Duration; +import java.util.ArrayList; +import java.util.LinkedHashMap; +import java.util.LinkedHashSet; +import java.util.List; +import java.util.Map; +import java.util.Objects; +import java.util.Set; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.atomic.AtomicLong; + +import com.anthropic.client.AnthropicClient; +import com.anthropic.core.RequestOptions; +import com.anthropic.core.http.Headers; +import com.anthropic.core.http.StreamResponse; +import com.anthropic.models.messages.Message; +import com.anthropic.models.messages.MessageCreateParams; +import com.anthropic.models.messages.batches.BatchCreateParams; +import com.anthropic.models.messages.batches.MessageBatchIndividualResponse; +import com.anthropic.models.messages.batches.MessageBatchResult; +import io.micrometer.core.instrument.MeterRegistry; +import io.micrometer.observation.Observation; +import io.micrometer.observation.ObservationRegistry; +import org.apache.commons.logging.Log; +import org.apache.commons.logging.LogFactory; +import org.jspecify.annotations.Nullable; +import reactor.core.publisher.Flux; +import reactor.core.scheduler.Schedulers; + +import org.springframework.ai.anthropic.http.okhttp.AnthropicHttpClientBuilderCustomizer; +import org.springframework.ai.chat.metadata.EmptyRateLimit; +import org.springframework.ai.chat.metadata.Usage; +import org.springframework.ai.chat.model.ChatResponse; +import org.springframework.ai.chat.model.Generation; +import org.springframework.ai.chat.prompt.Prompt; +import org.springframework.ai.model.tool.ToolCallingManager; +import org.springframework.ai.observation.conventions.AiProvider; +import org.springframework.util.Assert; +import org.springframework.util.CollectionUtils; + +/** + * Default {@link AnthropicBatchModel} implementation, built on the official + * Anthropic Java SDK. + * + *

+ * Requests are mapped exactly like realtime ones. Each + * {@link AnthropicBatchRequest} carries a {@link Prompt} that goes through the very same + * {@code Prompt} to {@code MessageCreateParams} conversion used by + * {@link AnthropicChatModel#call(Prompt)}, so system messages, conversation history, + * images and PDF documents, prompt caching, thinking, structured output and tool + * definitions behave the same on both paths. Likewise, a succeeded result is converted + * into a {@link ChatResponse} with the same generations, metadata keys and usage as a + * realtime call. + * + *

+ * Tool calls are not executed. Tool definitions are sent, but a batch response + * containing {@code tool_use} blocks is returned as-is: there is no interactive + * tool-execution loop, because a batch entry cannot be continued mid-flight. Surface the + * tool calls from {@link org.springframework.ai.chat.model.Generation#getOutput() the + * assistant message} and submit a follow-up batch if you need a second turn. + * + *

+ * Results are unordered. Correlate them through + * {@link AnthropicBatchResult#customId()}, never by position. + * + *

+ * Typical usage: + * + *

{@code
+ * AnthropicBatch batch = batchModel.submit(List.of(
+ *         AnthropicBatchRequest.of("invoice-1", "Summarize invoice 1"),
+ *         AnthropicBatchRequest.of("invoice-2", "Summarize invoice 2")));
+ *
+ * // later, on the application's own schedule
+ * if (batchModel.retrieve(batch.id()).isEnded()) {
+ *     batchModel.results(batch.id())
+ *         .doOnNext(result -> store(result.customId(), result))
+ *         .blockLast();
+ * }
+ * }
+ * + * @author Ricken Bazolo + * @since 2.0.0 + * @see AnthropicBatchRequest + * @see AnthropicBatchResult + * @see AnthropicChatModel + */ +public final class DefaultAnthropicBatchModel implements AnthropicBatchModel { + + private static final Log logger = LogFactory.getLog(DefaultAnthropicBatchModel.class); + + private static final AnthropicBatchObservationConvention DEFAULT_OBSERVATION_CONVENTION = new DefaultAnthropicBatchObservationConvention(); + + private static final ToolCallingManager DEFAULT_TOOL_CALLING_MANAGER = ToolCallingManager.builder().build(); + + private final AnthropicClient anthropicClient; + + private final AnthropicChatOptions options; + + private final ToolCallingManager toolCallingManager; + + private final ObservationRegistry observationRegistry; + + private AnthropicBatchObservationConvention observationConvention = DEFAULT_OBSERVATION_CONVENTION; + + private DefaultAnthropicBatchModel(@Nullable AnthropicClient anthropicClient, + @Nullable AnthropicChatOptions options, @Nullable ToolCallingManager toolCallingManager, + @Nullable ObservationRegistry observationRegistry, @Nullable MeterRegistry meterRegistry, + @Nullable ExecutorService dispatcherExecutor, + List httpClientCustomizers) { + + this.options = options != null ? options : AnthropicChatOptions.builder().build(); + this.observationRegistry = Objects.requireNonNullElse(observationRegistry, ObservationRegistry.NOOP); + this.toolCallingManager = Objects.requireNonNullElse(toolCallingManager, DEFAULT_TOOL_CALLING_MANAGER); + + this.anthropicClient = Objects.requireNonNullElseGet(anthropicClient, + () -> AnthropicSetup.setupSyncClient(this.options.getBaseUrl(), this.options.getApiKey(), + this.options.getTimeout(), this.options.getMaxRetries(), this.options.getProxy(), + this.options.getCustomHeaders(), this.observationRegistry, meterRegistry, dispatcherExecutor, + httpClientCustomizers)); + } + + /** + * Creates a new builder for {@link AnthropicBatchModel}. + * @return a new builder instance + */ + public static Builder builder() { + return new Builder(); + } + + /** + * Gets the default options applied to batch entries whose prompt carries none. + * @return the default options + */ + public AnthropicChatOptions getOptions() { + return this.options; + } + + /** + * Returns the underlying synchronous Anthropic SDK client, for accessing SDK features + * this model does not expose (for example listing batches). + * @return the sync client + */ + public AnthropicClient getAnthropicClient() { + return this.anthropicClient; + } + + /** + * Submits a batch of prompts for asynchronous processing. + * + *

+ * Returns as soon as Anthropic accepts the batch; no request has been processed yet. + * Persist {@link AnthropicBatch#id()} together with the {@code customId} of every + * entry so that polling and correlation survive an application restart. + * @param requests the batch entries; must be non-empty and carry distinct + * {@code customId} values + * @return the accepted batch, in state {@link AnthropicBatchStatus#IN_PROGRESS} + * @throws IllegalArgumentException if the list is empty or a {@code customId} is + * duplicated + * @throws com.anthropic.errors.AnthropicServiceException if the API rejects the batch + */ + @Override + public AnthropicBatch submit(List requests) { + Assert.notEmpty(requests, "requests must not be empty"); + + Set customIds = new LinkedHashSet<>(); + List sdkRequests = new ArrayList<>(requests.size()); + Map> additionalHeaders = new LinkedHashMap<>(); + Set requestModels = new LinkedHashSet<>(); + + for (AnthropicBatchRequest request : requests) { + Assert.isTrue(customIds.add(request.customId()), () -> "Duplicate customId in batch: '" + request.customId() + + "'. Results are correlated by customId, so it must be unique within a batch."); + + MessageCreateParams params = AnthropicChatModel.createRequest(buildRequestPrompt(request.prompt()), + this.toolCallingManager, this.options.getSkillContainer()); + + requestModels.add(params.model().asString()); + collectAdditionalHeaders(params._additionalHeaders(), additionalHeaders); + + sdkRequests.add(BatchCreateParams.Request.builder() + .customId(request.customId()) + .params(toBatchRequestParams(params)) + .build()); + } + + BatchCreateParams.Builder builder = BatchCreateParams.builder().requests(sdkRequests); + // Per-request headers cannot be expressed per batch entry, so entry-level headers + // (for example the beta headers implied by a skill container) are merged onto the + // batch request itself. + additionalHeaders.forEach(builder::replaceAdditionalHeaders); + + AnthropicBatchObservationContext observationContext = AnthropicBatchObservationContext.builder() + .operation(AnthropicBatchObservationContext.Operation.CREATE) + .provider(AiProvider.ANTHROPIC.value()) + .requestModel(requestModels.size() == 1 ? requestModels.iterator().next() : null) + .requestCount(requests.size()) + .build(); + + return observe(observationContext, () -> { + AnthropicBatch batch = AnthropicBatch + .from(this.anthropicClient.messages().batches().create(builder.build(), requestOptions())); + observationContext.setBatch(batch); + if (logger.isDebugEnabled()) { + logger.debug("Submitted Anthropic message batch " + batch.id() + " with " + requests.size() + + " request(s)."); + } + return batch; + }); + } + + /** + * Retrieves the current state of a batch. + * + *

+ * Call this on the application's own schedule; Spring AI performs no polling. Results + * become readable once {@link AnthropicBatch#isEnded()} is {@code true}. + * @param batchId the batch identifier returned by {@link #submit(List)} + * @return the current batch state + * @throws com.anthropic.errors.NotFoundException if no such batch exists + */ + @Override + public AnthropicBatch retrieve(String batchId) { + Assert.hasText(batchId, "batchId must not be empty"); + + AnthropicBatchObservationContext observationContext = observationContext( + AnthropicBatchObservationContext.Operation.RETRIEVE); + + return observe(observationContext, () -> { + AnthropicBatch batch = AnthropicBatch + .from(this.anthropicClient.messages().batches().retrieve(batchId, requestOptions())); + observationContext.setBatch(batch); + return batch; + }); + } + + /** + * Streams the results of an ended batch. + * + *

+ * The underlying JSONL stream is consumed lazily and the SDK stream is closed when + * the returned {@link Flux} terminates or is cancelled, so a batch with a very large + * number of entries never has to be held in memory. Because the SDK exposes a + * blocking stream, items are emitted on {@link Schedulers#boundedElastic()}. + * + *

+ * Results arrive in an unspecified order: key them by + * {@link AnthropicBatchResult#customId()}. Individual failures are emitted as + * {@link AnthropicBatchResultStatus#ERRORED} items rather than thrown, so one bad + * entry never hides the rest. + * @param batchId the batch identifier + * @return a lazily-populated flux of per-request results + */ + @Override + public Flux results(String batchId) { + Assert.hasText(batchId, "batchId must not be empty"); + + return Flux.defer(() -> { + AnthropicBatchObservationContext observationContext = observationContext( + AnthropicBatchObservationContext.Operation.RESULTS); + Observation observation = AnthropicBatchObservationDocumentation.BATCH_MODEL_OPERATION.observation( + this.observationConvention, DEFAULT_OBSERVATION_CONVENTION, () -> observationContext, + this.observationRegistry); + observation.start(); + + ResultCounters counters = new ResultCounters(); + + return Flux + .using(() -> this.anthropicClient.messages().batches().resultsStreaming(batchId, requestOptions()), + streamResponse -> Flux.fromStream(streamResponse.stream()).map(response -> { + AnthropicBatchResult result = toBatchResult(response); + counters.record(result.status()); + return result; + }), StreamResponse::close) + .doOnError(observation::error) + .doFinally(signal -> { + observationContext.setRequestCounts(counters.snapshot()); + observation.stop(); + }) + .subscribeOn(Schedulers.boundedElastic()); + }); + } + + /** + * Requests cancellation of a batch. + * + *

+ * Cancellation is not immediate: the batch moves to + * {@link AnthropicBatchStatus#CANCELING} and requests that already completed keep + * their result, while the remaining ones end up as + * {@link AnthropicBatchResultStatus#CANCELED}. + * @param batchId the batch identifier + * @return the batch state after the cancellation request + */ + @Override + public AnthropicBatch cancel(String batchId) { + Assert.hasText(batchId, "batchId must not be empty"); + + AnthropicBatchObservationContext observationContext = observationContext( + AnthropicBatchObservationContext.Operation.CANCEL); + + return observe(observationContext, () -> { + AnthropicBatch batch = AnthropicBatch + .from(this.anthropicClient.messages().batches().cancel(batchId, requestOptions())); + observationContext.setBatch(batch); + return batch; + }); + } + + /** + * Deletes a batch. Only batches whose processing has ended can be deleted. + * @param batchId the batch identifier + */ + @Override + public void delete(String batchId) { + Assert.hasText(batchId, "batchId must not be empty"); + + AnthropicBatchObservationContext observationContext = observationContext( + AnthropicBatchObservationContext.Operation.DELETE); + + AnthropicBatchObservationDocumentation.BATCH_MODEL_OPERATION + .observation(this.observationConvention, DEFAULT_OBSERVATION_CONVENTION, () -> observationContext, + this.observationRegistry) + .observe(() -> this.anthropicClient.messages().batches().delete(batchId, requestOptions())); + } + + /** + * Use the provided convention for reporting observation data. + * @param observationConvention the provided convention + */ + public void setObservationConvention(AnthropicBatchObservationConvention observationConvention) { + Assert.notNull(observationConvention, "observationConvention cannot be null"); + this.observationConvention = observationConvention; + } + + private T observe(AnthropicBatchObservationContext observationContext, java.util.function.Supplier action) { + T result = AnthropicBatchObservationDocumentation.BATCH_MODEL_OPERATION + .observation(this.observationConvention, DEFAULT_OBSERVATION_CONVENTION, () -> observationContext, + this.observationRegistry) + .observe(action); + Assert.state(result != null, "Anthropic batch operation returned no result"); + return result; + } + + private AnthropicBatchObservationContext observationContext(AnthropicBatchObservationContext.Operation operation) { + return AnthropicBatchObservationContext.builder() + .operation(operation) + .provider(AiProvider.ANTHROPIC.value()) + .build(); + } + + /** + * Applies this model's default options when the prompt carries none, mirroring + * {@link AnthropicChatModel}'s behaviour. + */ + private Prompt buildRequestPrompt(Prompt prompt) { + return prompt.getOptions() == null ? prompt.mutate().chatOptions(this.options).build() : prompt; + } + + private RequestOptions requestOptions() { + Duration timeout = this.options.getTimeout(); + return timeout != null ? RequestOptions.builder().timeout(timeout).build() : RequestOptions.none(); + } + + private static void collectAdditionalHeaders(Headers headers, Map> target) { + for (String name : headers.names()) { + target.computeIfAbsent(name, key -> new LinkedHashSet<>()).addAll(headers.values(name)); + } + } + + /** + * Copies a realtime {@link MessageCreateParams} onto the batch entry parameter shape. + * The two SDK types are structurally identical but nominally distinct, so every field + * is transferred explicitly; {@code stream} is deliberately never set, as batch + * entries cannot stream. + */ + private static BatchCreateParams.Request.Params toBatchRequestParams(MessageCreateParams params) { + BatchCreateParams.Request.Params.Builder builder = BatchCreateParams.Request.Params.builder() + .maxTokens(params._maxTokens()) + .messages(params._messages()) + .model(params._model()); + + params.cacheControl().ifPresent(builder::cacheControl); + params.container().ifPresent(builder::container); + params.inferenceGeo().ifPresent(builder::inferenceGeo); + params.metadata().ifPresent(builder::metadata); + params.outputConfig().ifPresent(builder::outputConfig); + params.serviceTier() + .ifPresent(serviceTier -> builder + .serviceTier(BatchCreateParams.Request.Params.ServiceTier.of(serviceTier.asString()))); + params.stopSequences().ifPresent(builder::stopSequences); + params.system().ifPresent(system -> { + if (system.isString()) { + builder.system(BatchCreateParams.Request.Params.System.ofString(system.asString())); + } + else if (system.isTextBlockParams()) { + builder.system(BatchCreateParams.Request.Params.System.ofTextBlockParams(system.asTextBlockParams())); + } + }); + params.temperature().ifPresent(builder::temperature); + params.thinking().ifPresent(builder::thinking); + params.toolChoice().ifPresent(builder::toolChoice); + params.tools().ifPresent(builder::tools); + params.topK().ifPresent(builder::topK); + params.topP().ifPresent(builder::topP); + + // Carries anything set through putAdditionalBodyProperty(), notably the skills + // container. + params._additionalBodyProperties().forEach(builder::putAdditionalProperty); + + return builder.build(); + } + + /** + * Converts one JSONL entry into a Spring AI result, reusing the realtime response + * conversion so that a batched message yields the same {@link ChatResponse} shape as + * {@link AnthropicChatModel#call(Prompt)}. + */ + private static AnthropicBatchResult toBatchResult(MessageBatchIndividualResponse response) { + String customId = response.customId(); + MessageBatchResult result = response.result(); + + if (result.isSucceeded()) { + Message message = result.asSucceeded().message(); + List citations = new ArrayList<>(); + List webSearchResults = new ArrayList<>(); + List generations = AnthropicChatModel.buildGenerations(message, citations, webSearchResults); + Usage usage = AnthropicChatModel.getDefaultUsage(message.usage()); + // Batch results carry no per-request rate-limit headers. + ChatResponse chatResponse = new ChatResponse(generations, + AnthropicChatModel.from(message, usage, citations, webSearchResults, new EmptyRateLimit())); + return AnthropicBatchResult.succeeded(customId, chatResponse); + } + if (result.isErrored()) { + return AnthropicBatchResult.errored(customId, AnthropicBatchError.from(result.asErrored().error())); + } + if (result.isCanceled()) { + return AnthropicBatchResult.of(customId, AnthropicBatchResultStatus.CANCELED); + } + if (result.isExpired()) { + return AnthropicBatchResult.of(customId, AnthropicBatchResultStatus.EXPIRED); + } + if (logger.isWarnEnabled()) { + logger.warn("Unrecognised batch result type for customId '" + customId + "': " + result); + } + return AnthropicBatchResult.of(customId, AnthropicBatchResultStatus.UNKNOWN); + } + + /** + * Accumulates per-outcome counters while the result stream is consumed, so the + * {@link AnthropicBatchObservationContext.Operation#RESULTS} observation can report + * them without buffering the results themselves. + */ + private static final class ResultCounters { + + private final AtomicLong succeeded = new AtomicLong(); + + private final AtomicLong errored = new AtomicLong(); + + private final AtomicLong canceled = new AtomicLong(); + + private final AtomicLong expired = new AtomicLong(); + + void record(AnthropicBatchResultStatus status) { + switch (status) { + case SUCCEEDED -> this.succeeded.incrementAndGet(); + case ERRORED -> this.errored.incrementAndGet(); + case CANCELED -> this.canceled.incrementAndGet(); + case EXPIRED -> this.expired.incrementAndGet(); + case UNKNOWN -> { + } + } + } + + AnthropicBatchRequestCounts snapshot() { + return new AnthropicBatchRequestCounts(0, this.succeeded.get(), this.errored.get(), this.canceled.get(), + this.expired.get()); + } + + } + + /** + * Builder for {@link AnthropicBatchModel}. Accepts the same connection and + * observability configuration as {@link AnthropicChatModel.Builder}, so a batch model + * reuses the application's Anthropic credentials, base URL, timeout, retries, proxy, + * custom headers and HTTP client customizers. + */ + public static final class Builder { + + private @Nullable AnthropicClient anthropicClient; + + private @Nullable AnthropicChatOptions options; + + private @Nullable ToolCallingManager toolCallingManager; + + private @Nullable ObservationRegistry observationRegistry; + + private @Nullable MeterRegistry meterRegistry; + + private @Nullable ExecutorService dispatcherExecutor; + + private List httpClientCustomizers = new ArrayList<>(); + + private Builder() { + } + + /** + * Sets a pre-configured Anthropic SDK client. When supplied, the + * connection-related options are ignored. + * @param anthropicClient the client + * @return this builder + */ + public Builder anthropicClient(AnthropicClient anthropicClient) { + this.anthropicClient = anthropicClient; + return this; + } + + /** + * Sets the default options applied to batch entries whose prompt carries none, + * and the connection settings used when no client is supplied. + * @param options the options + * @return this builder + */ + public Builder options(AnthropicChatOptions options) { + this.options = options; + return this; + } + + /** + * Sets the tool calling manager used to resolve the tool definitions advertised + * to the model. Tool calls returned by a batch are never executed; see the + * class-level documentation. + * @param toolCallingManager the tool calling manager + * @return this builder + */ + public Builder toolCallingManager(ToolCallingManager toolCallingManager) { + this.toolCallingManager = toolCallingManager; + return this; + } + + /** + * Sets the observation registry batch operations report to. + * @param observationRegistry the observation registry + * @return this builder + */ + public Builder observationRegistry(ObservationRegistry observationRegistry) { + this.observationRegistry = observationRegistry; + return this; + } + + /** + * Sets the meter registry OkHttp connection-pool gauges are bound to. + * @param meterRegistry the meter registry, or {@code null} to disable the gauges + * @return this builder + */ + public Builder meterRegistry(@Nullable MeterRegistry meterRegistry) { + this.meterRegistry = meterRegistry; + return this; + } + + /** + * Sets the OkHttp dispatcher executor. The caller owns its lifecycle. + * @param dispatcherExecutor the executor, or {@code null} for the library default + * @return this builder + */ + public Builder dispatcherExecutor(@Nullable ExecutorService dispatcherExecutor) { + this.dispatcherExecutor = dispatcherExecutor; + return this; + } + + /** + * Adds a customizer applied to the underlying OkHttp client builder. + * @param customizer the customizer + * @return this builder + */ + public Builder httpClientBuilderCustomizer(AnthropicHttpClientBuilderCustomizer customizer) { + Assert.notNull(customizer, "customizer cannot be null"); + this.httpClientCustomizers.add(customizer); + return this; + } + + /** + * Replaces the customizers applied to the underlying OkHttp client builder. + * @param customizers the customizers + * @return this builder + */ + public Builder httpClientBuilderCustomizers(List customizers) { + Assert.notNull(customizers, "customizers cannot be null"); + this.httpClientCustomizers = CollectionUtils.isEmpty(customizers) ? new ArrayList<>() + : new ArrayList<>(customizers); + return this; + } + + /** + * Builds the batch model. + * @return a new {@link DefaultAnthropicBatchModel} + */ + public DefaultAnthropicBatchModel build() { + return new DefaultAnthropicBatchModel(this.anthropicClient, this.options, this.toolCallingManager, + this.observationRegistry, this.meterRegistry, this.dispatcherExecutor, this.httpClientCustomizers); + } + + } + +} diff --git a/models/spring-ai-anthropic/src/test/java/org/springframework/ai/anthropic/AnthropicBatchModelTests.java b/models/spring-ai-anthropic/src/test/java/org/springframework/ai/anthropic/AnthropicBatchModelTests.java new file mode 100644 index 0000000000..3c10a79575 --- /dev/null +++ b/models/spring-ai-anthropic/src/test/java/org/springframework/ai/anthropic/AnthropicBatchModelTests.java @@ -0,0 +1,414 @@ +/* + * Copyright 2023-present the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.ai.anthropic; + +import java.time.OffsetDateTime; +import java.util.List; +import java.util.Map; +import java.util.Optional; +import java.util.stream.Stream; + +import com.anthropic.client.AnthropicClient; +import com.anthropic.core.RequestOptions; +import com.anthropic.core.http.StreamResponse; +import com.anthropic.models.ErrorObject; +import com.anthropic.models.ErrorResponse; +import com.anthropic.models.InvalidRequestError; +import com.anthropic.models.messages.ContentBlock; +import com.anthropic.models.messages.Message; +import com.anthropic.models.messages.Model; +import com.anthropic.models.messages.StopReason; +import com.anthropic.models.messages.TextBlock; +import com.anthropic.models.messages.Usage; +import com.anthropic.models.messages.batches.BatchCreateParams; +import com.anthropic.models.messages.batches.DeletedMessageBatch; +import com.anthropic.models.messages.batches.MessageBatch; +import com.anthropic.models.messages.batches.MessageBatchErroredResult; +import com.anthropic.models.messages.batches.MessageBatchIndividualResponse; +import com.anthropic.models.messages.batches.MessageBatchRequestCounts; +import com.anthropic.models.messages.batches.MessageBatchResult; +import com.anthropic.models.messages.batches.MessageBatchSucceededResult; +import com.anthropic.services.blocking.MessageService; +import com.anthropic.services.blocking.messages.BatchService; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.ArgumentCaptor; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; +import org.mockito.junit.jupiter.MockitoSettings; +import org.mockito.quality.Strictness; + +import org.springframework.ai.chat.messages.SystemMessage; +import org.springframework.ai.chat.messages.UserMessage; +import org.springframework.ai.chat.prompt.Prompt; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatIllegalArgumentException; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.BDDMockito.given; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; + +/** + * Unit tests for {@link DefaultAnthropicBatchModel}. Exercises SDK parameter + * construction, batch and result mapping, out-of-order results, per-request errors and + * stream cleanup with a mocked SDK client — no API key and no network access required. + * + * @author Ricken Bazolo + */ +@ExtendWith(MockitoExtension.class) +@MockitoSettings(strictness = Strictness.LENIENT) +class AnthropicBatchModelTests { + + @Mock + private AnthropicClient anthropicClient; + + @Mock + private MessageService messageService; + + @Mock + private BatchService batchService; + + private AnthropicBatchModel batchModel; + + @BeforeEach + void setUp() { + given(this.anthropicClient.messages()).willReturn(this.messageService); + given(this.messageService.batches()).willReturn(this.batchService); + + this.batchModel = AnthropicBatchModel.builder() + .anthropicClient(this.anthropicClient) + .options(AnthropicChatOptions.builder().model("claude-haiku-4-5").maxTokens(256).build()) + .build(); + } + + @Test + void submitMapsEachPromptThroughTheRealtimeRequestConversion() { + MessageBatch accepted = mockMessageBatch("msgbatch_1", MessageBatch.ProcessingStatus.IN_PROGRESS); + given(this.batchService.create(any(BatchCreateParams.class), any(RequestOptions.class))).willReturn(accepted); + + AnthropicBatch batch = this.batchModel.submit(List.of( + AnthropicBatchRequest.of("req-1", + new Prompt(List.of(new SystemMessage("Be brief."), new UserMessage("Hello")))), + AnthropicBatchRequest.of("req-2", "World"))); + + ArgumentCaptor captor = ArgumentCaptor.forClass(BatchCreateParams.class); + verify(this.batchService).create(captor.capture(), any(RequestOptions.class)); + + List requests = captor.getValue().requests(); + assertThat(requests).hasSize(2); + assertThat(requests.stream().map(BatchCreateParams.Request::customId)).containsExactly("req-1", "req-2"); + + BatchCreateParams.Request.Params first = requests.get(0).params(); + assertThat(first.model().asString()).isEqualTo("claude-haiku-4-5"); + assertThat(first.maxTokens()).isEqualTo(256L); + assertThat(first.system().orElseThrow().asString()).isEqualTo("Be brief."); + assertThat(first.messages()).hasSize(1); + // Batch entries cannot stream. + assertThat(first._additionalProperties()).doesNotContainKey("stream"); + + assertThat(batch.id()).isEqualTo("msgbatch_1"); + assertThat(batch.status()).isEqualTo(AnthropicBatchStatus.IN_PROGRESS); + assertThat(batch.isEnded()).isFalse(); + } + + @Test + void submitHonoursPerRequestOptions() { + MessageBatch accepted = mockMessageBatch("msgbatch_2", MessageBatch.ProcessingStatus.IN_PROGRESS); + given(this.batchService.create(any(BatchCreateParams.class), any(RequestOptions.class))).willReturn(accepted); + + this.batchModel.submit(List.of(AnthropicBatchRequest.of("req-1", "Hello", + AnthropicChatOptions.builder().model("claude-opus-4-5").maxTokens(64).temperature(0.2).build()))); + + ArgumentCaptor captor = ArgumentCaptor.forClass(BatchCreateParams.class); + verify(this.batchService).create(captor.capture(), any(RequestOptions.class)); + + BatchCreateParams.Request.Params params = captor.getValue().requests().get(0).params(); + assertThat(params.model().asString()).isEqualTo("claude-opus-4-5"); + assertThat(params.maxTokens()).isEqualTo(64L); + assertThat(params.temperature()).contains(0.2); + } + + @Test + void submitRejectsAnEmptyBatch() { + assertThatIllegalArgumentException().isThrownBy(() -> this.batchModel.submit(List.of())) + .withMessageContaining("requests must not be empty"); + } + + @Test + void submitRejectsDuplicateCustomIds() { + assertThatIllegalArgumentException() + .isThrownBy(() -> this.batchModel + .submit(List.of(AnthropicBatchRequest.of("same", "a"), AnthropicBatchRequest.of("same", "b")))) + .withMessageContaining("Duplicate customId"); + } + + @Test + void retrieveMapsStatusCountersAndTimestamps() { + MessageBatch messageBatch = mockMessageBatch("msgbatch_3", MessageBatch.ProcessingStatus.ENDED); + given(this.batchService.retrieve(eq("msgbatch_3"), any(RequestOptions.class))).willReturn(messageBatch); + + AnthropicBatch batch = this.batchModel.retrieve("msgbatch_3"); + + assertThat(batch.status()).isEqualTo(AnthropicBatchStatus.ENDED); + assertThat(batch.isEnded()).isTrue(); + assertThat(batch.requestCounts()).isEqualTo(new AnthropicBatchRequestCounts(1, 2, 1, 0, 0)); + assertThat(batch.requestCounts().total()).isEqualTo(4); + assertThat(batch.requestCounts().completed()).isEqualTo(3); + assertThat(batch.resultsUrl()).isEqualTo("https://api.anthropic.com/v1/messages/batches/msgbatch_3/results"); + assertThat(batch.endedAt()).isNotNull(); + } + + @Test + void retrieveMapsUnknownStatusesWithoutFailing() { + MessageBatch messageBatch = mockMessageBatch("msgbatch_x", + MessageBatch.ProcessingStatus.of("brand_new_status")); + given(this.batchService.retrieve(eq("msgbatch_x"), any(RequestOptions.class))).willReturn(messageBatch); + + assertThat(this.batchModel.retrieve("msgbatch_x").status()).isEqualTo(AnthropicBatchStatus.UNKNOWN); + } + + @Test + void resultsCorrelateByCustomIdEvenWhenReturnedOutOfOrder() { + StreamResponse streamResponse = mockResults( + succeededResponse("req-2", "second"), succeededResponse("req-1", "first")); + given(this.batchService.resultsStreaming(eq("msgbatch_4"), any(RequestOptions.class))) + .willReturn(streamResponse); + + Map byCustomId = this.batchModel.results("msgbatch_4") + .collectMap(AnthropicBatchResult::customId) + .block(); + + assertThat(byCustomId).containsOnlyKeys("req-1", "req-2"); + assertThat(byCustomId.get("req-1").getText()).isEqualTo("first"); + assertThat(byCustomId.get("req-2").getText()).isEqualTo("second"); + assertThat(byCustomId.get("req-1").isSucceeded()).isTrue(); + verify(streamResponse).close(); + } + + @Test + void resultsConvertSucceededMessagesLikeARealtimeCall() { + StreamResponse streamResponse = mockResults( + succeededResponse("req-1", "Hello there")); + given(this.batchService.resultsStreaming(eq("msgbatch_5"), any(RequestOptions.class))) + .willReturn(streamResponse); + + AnthropicBatchResult result = this.batchModel.results("msgbatch_5").blockFirst(); + + assertThat(result).isNotNull(); + assertThat(result.status()).isEqualTo(AnthropicBatchResultStatus.SUCCEEDED); + assertThat(result.error()).isNull(); + assertThat(result.chatResponse()).isNotNull(); + assertThat(result.chatResponse().getResult().getOutput().getText()).isEqualTo("Hello there"); + assertThat(result.chatResponse().getResult().getMetadata().getFinishReason()) + .isEqualTo(StopReason.END_TURN.toString()); + assertThat(result.chatResponse().getMetadata().getId()).isEqualTo("msg_batch_req-1"); + assertThat(result.usage()).isNotNull(); + assertThat(result.usage().getPromptTokens()).isEqualTo(10); + assertThat(result.usage().getCompletionTokens()).isEqualTo(20); + assertThat(result.usage().getTotalTokens()).isEqualTo(30); + } + + @Test + void resultsSurfaceIndividualErrorsWithoutHidingTheOtherEntries() { + StreamResponse streamResponse = mockResults( + erroredResponse("req-1", "max_tokens must be positive"), succeededResponse("req-2", "fine"), + terminalResponse("req-3", TerminalKind.CANCELED), terminalResponse("req-4", TerminalKind.EXPIRED)); + given(this.batchService.resultsStreaming(eq("msgbatch_6"), any(RequestOptions.class))) + .willReturn(streamResponse); + + List results = this.batchModel.results("msgbatch_6").collectList().block(); + + assertThat(results).hasSize(4); + AnthropicBatchResult errored = results.get(0); + assertThat(errored.status()).isEqualTo(AnthropicBatchResultStatus.ERRORED); + assertThat(errored.chatResponse()).isNull(); + assertThat(errored.error()).isNotNull(); + assertThat(errored.error().type()).isEqualTo("invalid_request_error"); + assertThat(errored.error().message()).isEqualTo("max_tokens must be positive"); + assertThat(errored.error().requestId()).isEqualTo("req_abc"); + + assertThat(results.get(1).status()).isEqualTo(AnthropicBatchResultStatus.SUCCEEDED); + assertThat(results.get(2).status()).isEqualTo(AnthropicBatchResultStatus.CANCELED); + assertThat(results.get(3).status()).isEqualTo(AnthropicBatchResultStatus.EXPIRED); + assertThat(results.get(2).chatResponse()).isNull(); + assertThat(results.get(3).error()).isNull(); + } + + @Test + void resultsCloseTheSdkStreamWhenTheSubscriberCancels() { + StreamResponse streamResponse = mockResults(succeededResponse("req-1", "a"), + succeededResponse("req-2", "b")); + given(this.batchService.resultsStreaming(eq("msgbatch_7"), any(RequestOptions.class))) + .willReturn(streamResponse); + + AnthropicBatchResult first = this.batchModel.results("msgbatch_7").next().block(); + + assertThat(first).isNotNull(); + verify(streamResponse).close(); + } + + @Test + void resultsAreLazyAndDoNotCallTheApiUntilSubscribed() { + this.batchModel.results("msgbatch_8"); + + verify(this.batchService, never()).resultsStreaming(any(String.class), any(RequestOptions.class)); + } + + @Test + void cancelDelegatesToTheProvider() { + MessageBatch canceling = mockMessageBatch("msgbatch_9", MessageBatch.ProcessingStatus.CANCELING); + given(this.batchService.cancel(eq("msgbatch_9"), any(RequestOptions.class))).willReturn(canceling); + + AnthropicBatch batch = this.batchModel.cancel("msgbatch_9"); + + assertThat(batch.status()).isEqualTo(AnthropicBatchStatus.CANCELING); + assertThat(batch.isCanceling()).isTrue(); + verify(this.batchService).cancel(eq("msgbatch_9"), any(RequestOptions.class)); + } + + @Test + void deleteDelegatesToTheProvider() { + DeletedMessageBatch deleted = mock(DeletedMessageBatch.class); + given(this.batchService.delete(eq("msgbatch_10"), any(RequestOptions.class))).willReturn(deleted); + + this.batchModel.delete("msgbatch_10"); + + verify(this.batchService).delete(eq("msgbatch_10"), any(RequestOptions.class)); + } + + @Test + void controlOperationsRejectBlankBatchIds() { + assertThatIllegalArgumentException().isThrownBy(() -> this.batchModel.retrieve(" ")); + assertThatIllegalArgumentException().isThrownBy(() -> this.batchModel.cancel("")); + assertThatIllegalArgumentException().isThrownBy(() -> this.batchModel.delete("")); + assertThatIllegalArgumentException().isThrownBy(() -> this.batchModel.results("")); + } + + // --- fixtures --- + + private static MessageBatch mockMessageBatch(String id, MessageBatch.ProcessingStatus status) { + MessageBatchRequestCounts counts = mock(MessageBatchRequestCounts.class); + given(counts.processing()).willReturn(1L); + given(counts.succeeded()).willReturn(2L); + given(counts.errored()).willReturn(1L); + given(counts.canceled()).willReturn(0L); + given(counts.expired()).willReturn(0L); + + OffsetDateTime now = OffsetDateTime.parse("2026-07-30T10:15:30Z"); + + MessageBatch messageBatch = mock(MessageBatch.class); + given(messageBatch.id()).willReturn(id); + given(messageBatch.processingStatus()).willReturn(status); + given(messageBatch.requestCounts()).willReturn(counts); + given(messageBatch.createdAt()).willReturn(now); + given(messageBatch.expiresAt()).willReturn(now.plusDays(1)); + given(messageBatch.endedAt()).willReturn(Optional.of(now.plusHours(2))); + given(messageBatch.cancelInitiatedAt()).willReturn(Optional.empty()); + given(messageBatch.archivedAt()).willReturn(Optional.empty()); + given(messageBatch.resultsUrl()) + .willReturn(Optional.of("https://api.anthropic.com/v1/messages/batches/" + id + "/results")); + return messageBatch; + } + + @SuppressWarnings("unchecked") + private static StreamResponse mockResults( + MessageBatchIndividualResponse... responses) { + StreamResponse streamResponse = mock(StreamResponse.class); + given(streamResponse.stream()).willReturn(Stream.of(responses)); + return streamResponse; + } + + private static MessageBatchIndividualResponse succeededResponse(String customId, String text) { + Message message = mockMessage("msg_batch_" + customId, text); + MessageBatchSucceededResult succeeded = mock(MessageBatchSucceededResult.class); + given(succeeded.message()).willReturn(message); + + MessageBatchResult result = mock(MessageBatchResult.class); + given(result.isSucceeded()).willReturn(true); + given(result.asSucceeded()).willReturn(succeeded); + + return individualResponse(customId, result); + } + + private static MessageBatchIndividualResponse erroredResponse(String customId, String message) { + ErrorResponse errorResponse = ErrorResponse.builder() + .error(ErrorObject.ofInvalidRequestError(InvalidRequestError.builder().message(message).build())) + .requestId("req_abc") + .build(); + MessageBatchErroredResult errored = mock(MessageBatchErroredResult.class); + given(errored.error()).willReturn(errorResponse); + + MessageBatchResult result = mock(MessageBatchResult.class); + given(result.isSucceeded()).willReturn(false); + given(result.isErrored()).willReturn(true); + given(result.asErrored()).willReturn(errored); + + return individualResponse(customId, result); + } + + private enum TerminalKind { + + CANCELED, EXPIRED + + } + + private static MessageBatchIndividualResponse terminalResponse(String customId, TerminalKind kind) { + MessageBatchResult result = mock(MessageBatchResult.class); + given(result.isSucceeded()).willReturn(false); + given(result.isErrored()).willReturn(false); + given(result.isCanceled()).willReturn(kind == TerminalKind.CANCELED); + given(result.isExpired()).willReturn(kind == TerminalKind.EXPIRED); + + return individualResponse(customId, result); + } + + private static MessageBatchIndividualResponse individualResponse(String customId, MessageBatchResult result) { + MessageBatchIndividualResponse response = mock(MessageBatchIndividualResponse.class); + given(response.customId()).willReturn(customId); + given(response.result()).willReturn(result); + return response; + } + + private static Message mockMessage(String id, String text) { + TextBlock textBlock = mock(TextBlock.class); + given(textBlock.text()).willReturn(text); + given(textBlock.citations()).willReturn(Optional.empty()); + + ContentBlock contentBlock = mock(ContentBlock.class); + given(contentBlock.isText()).willReturn(true); + given(contentBlock.asText()).willReturn(textBlock); + + Usage usage = mock(Usage.class); + given(usage.inputTokens()).willReturn(10L); + given(usage.outputTokens()).willReturn(20L); + given(usage.cacheReadInputTokens()).willReturn(Optional.empty()); + given(usage.cacheCreationInputTokens()).willReturn(Optional.empty()); + + Message message = mock(Message.class); + given(message.id()).willReturn(id); + given(message.model()).willReturn(Model.CLAUDE_HAIKU_4_5); + given(message.content()).willReturn(List.of(contentBlock)); + given(message.stopReason()).willReturn(Optional.of(StopReason.END_TURN)); + given(message.usage()).willReturn(usage); + return message; + } + +} From d3287bc05851b4675591bc2ce1252be4da94c6f3 Mon Sep 17 00:00:00 2001 From: Ricken BAZOLO Date: Fri, 31 Jul 2026 12:17:08 +0200 Subject: [PATCH 5/7] Auto-configure the Anthropic batch model Expose the batch model as a Spring Boot bean so an application can inject it without building an SDK client itself. The bean is opt-in behind spring.ai.anthropic.batch.enabled, so applications that never submit batches do not pay for a second HTTP client. Connection settings and the chat model defaults are reused as-is, which means a batch shares the credentials, base URL, timeout, retries, proxy, custom headers and HTTP client customizers of realtime calls; spring.ai.anthropic.batch.model and .max-tokens override the model and output ceiling for batch entries only. Signed-off-by: Ricken BAZOLO --- .../AnthropicBatchAutoConfiguration.java | 113 +++++++++++++ .../AnthropicBatchProperties.java | 84 ++++++++++ ...ot.autoconfigure.AutoConfiguration.imports | 1 + .../AnthropicBatchAutoConfigurationTests.java | 148 ++++++++++++++++++ 4 files changed, 346 insertions(+) create mode 100644 auto-configurations/models/spring-ai-autoconfigure-model-anthropic/src/main/java/org/springframework/ai/model/anthropic/autoconfigure/AnthropicBatchAutoConfiguration.java create mode 100644 auto-configurations/models/spring-ai-autoconfigure-model-anthropic/src/main/java/org/springframework/ai/model/anthropic/autoconfigure/AnthropicBatchProperties.java create mode 100644 auto-configurations/models/spring-ai-autoconfigure-model-anthropic/src/test/java/org/springframework/ai/model/anthropic/autoconfigure/AnthropicBatchAutoConfigurationTests.java diff --git a/auto-configurations/models/spring-ai-autoconfigure-model-anthropic/src/main/java/org/springframework/ai/model/anthropic/autoconfigure/AnthropicBatchAutoConfiguration.java b/auto-configurations/models/spring-ai-autoconfigure-model-anthropic/src/main/java/org/springframework/ai/model/anthropic/autoconfigure/AnthropicBatchAutoConfiguration.java new file mode 100644 index 0000000000..562f38c54e --- /dev/null +++ b/auto-configurations/models/spring-ai-autoconfigure-model-anthropic/src/main/java/org/springframework/ai/model/anthropic/autoconfigure/AnthropicBatchAutoConfiguration.java @@ -0,0 +1,113 @@ +/* + * Copyright 2023-present the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.ai.model.anthropic.autoconfigure; + +import java.util.List; + +import com.anthropic.client.AnthropicClient; +import io.micrometer.core.instrument.MeterRegistry; +import io.micrometer.observation.ObservationRegistry; + +import org.springframework.ai.anthropic.AnthropicBatchModel; +import org.springframework.ai.anthropic.AnthropicBatchObservationConvention; +import org.springframework.ai.anthropic.AnthropicChatOptions; +import org.springframework.ai.anthropic.DefaultAnthropicBatchModel; +import org.springframework.ai.anthropic.http.okhttp.AnthropicHttpClientBuilderCustomizer; +import org.springframework.ai.model.tool.ToolCallingManager; +import org.springframework.beans.factory.ObjectProvider; +import org.springframework.boot.autoconfigure.AutoConfiguration; +import org.springframework.boot.autoconfigure.condition.ConditionalOnClass; +import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; +import org.springframework.boot.context.properties.EnableConfigurationProperties; +import org.springframework.context.annotation.Bean; + +/** + * {@link AutoConfiguration Auto-configuration} for the Anthropic Message Batches model. + * + *

+ * Opt-in: the bean is only created when {@code spring.ai.anthropic.batch.enabled=true}, + * so applications that do not submit batches do not pay for a second HTTP client. + * Connection settings ({@code spring.ai.anthropic.*}) and model defaults + * ({@code spring.ai.anthropic.chat.*}) are shared with + * {@link AnthropicChatAutoConfiguration}, so a batch reuses the same credentials, base + * URL, timeout, retries, proxy, custom headers and HTTP client customizers as realtime + * calls. + * + * @author Ricken Bazolo + * @since 2.0.0 + */ +@AutoConfiguration(after = AnthropicChatAutoConfiguration.class) +@EnableConfigurationProperties({ AnthropicConnectionProperties.class, AnthropicChatProperties.class, + AnthropicBatchProperties.class }) +@ConditionalOnClass(AnthropicClient.class) +@ConditionalOnProperty(prefix = AnthropicBatchProperties.CONFIG_PREFIX, name = "enabled", havingValue = "true") +public class AnthropicBatchAutoConfiguration { + + @Bean + @ConditionalOnMissingBean + public AnthropicBatchModel anthropicBatchModel(AnthropicConnectionProperties connectionProperties, + AnthropicChatProperties chatProperties, AnthropicBatchProperties batchProperties, + ToolCallingManager toolCallingManager, ObjectProvider observationRegistry, + ObjectProvider meterRegistry, + ObjectProvider observationConvention, + ObjectProvider httpClientBuilderCustomizers) { + + AnthropicChatOptions.Builder builder = chatProperties.toOptions().mutate(); + if (connectionProperties.getApiKey() != null) { + builder.apiKey(connectionProperties.getApiKey()); + } + if (connectionProperties.getBaseUrl() != null) { + builder.baseUrl(connectionProperties.getBaseUrl()); + } + if (connectionProperties.getTimeout() != null) { + builder.timeout(connectionProperties.getTimeout()); + } + if (connectionProperties.getMaxRetries() != null) { + builder.maxRetries(connectionProperties.getMaxRetries()); + } + if (connectionProperties.getProxy() != null) { + builder.proxy(connectionProperties.getProxy()); + } + if (!connectionProperties.getCustomHeaders().isEmpty()) { + builder.customHeaders(connectionProperties.getCustomHeaders()); + } + // Batch-specific overrides of the shared chat defaults. + if (batchProperties.getModel() != null) { + builder.model(batchProperties.getModel()); + } + if (batchProperties.getMaxTokens() != null) { + builder.maxTokens(batchProperties.getMaxTokens()); + } + AnthropicChatOptions options = builder.build(); + + List customizers = httpClientBuilderCustomizers.orderedStream().toList(); + + DefaultAnthropicBatchModel batchModel = DefaultAnthropicBatchModel.builder() + .options(options) + .toolCallingManager(toolCallingManager) + .observationRegistry(observationRegistry.getIfUnique(() -> ObservationRegistry.NOOP)) + .meterRegistry(chatProperties.isConnectionPoolMetricsEnabled() ? meterRegistry.getIfAvailable() : null) + .httpClientBuilderCustomizers(customizers) + .build(); + + observationConvention.ifAvailable(batchModel::setObservationConvention); + + return batchModel; + } + +} diff --git a/auto-configurations/models/spring-ai-autoconfigure-model-anthropic/src/main/java/org/springframework/ai/model/anthropic/autoconfigure/AnthropicBatchProperties.java b/auto-configurations/models/spring-ai-autoconfigure-model-anthropic/src/main/java/org/springframework/ai/model/anthropic/autoconfigure/AnthropicBatchProperties.java new file mode 100644 index 0000000000..637d317600 --- /dev/null +++ b/auto-configurations/models/spring-ai-autoconfigure-model-anthropic/src/main/java/org/springframework/ai/model/anthropic/autoconfigure/AnthropicBatchProperties.java @@ -0,0 +1,84 @@ +/* + * Copyright 2023-present the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.ai.model.anthropic.autoconfigure; + +import org.jspecify.annotations.Nullable; + +import org.springframework.boot.context.properties.ConfigurationProperties; + +/** + * Anthropic Message Batches autoconfiguration properties. + * + *

+ * The batch model is opt-in: set {@code spring.ai.anthropic.batch.enabled=true} to get an + * {@link org.springframework.ai.anthropic.AnthropicBatchModel} bean. Connection settings + * ({@code spring.ai.anthropic.*}) and the model defaults + * ({@code spring.ai.anthropic.chat.*}) are shared with the chat model; the two properties + * here only override the model and output-token ceiling for batch entries when batches + * need to differ from realtime calls. + * + * @author Ricken Bazolo + * @since 2.0.0 + */ +@ConfigurationProperties(AnthropicBatchProperties.CONFIG_PREFIX) +public class AnthropicBatchProperties { + + public static final String CONFIG_PREFIX = "spring.ai.anthropic.batch"; + + /** + * Whether to expose an Anthropic batch model bean. Disabled by default so that + * applications that do not use batches do not pay for a second HTTP client. + */ + private boolean enabled = false; + + /** + * Model to use for batch entries. Falls back to + * {@code spring.ai.anthropic.chat.model}. + */ + private @Nullable String model; + + /** + * Maximum number of tokens to generate per batch entry. Falls back to + * {@code spring.ai.anthropic.chat.max-tokens}. + */ + private @Nullable Integer maxTokens; + + public boolean isEnabled() { + return this.enabled; + } + + public void setEnabled(boolean enabled) { + this.enabled = enabled; + } + + public @Nullable String getModel() { + return this.model; + } + + public void setModel(@Nullable String model) { + this.model = model; + } + + public @Nullable Integer getMaxTokens() { + return this.maxTokens; + } + + public void setMaxTokens(@Nullable Integer maxTokens) { + this.maxTokens = maxTokens; + } + +} diff --git a/auto-configurations/models/spring-ai-autoconfigure-model-anthropic/src/main/resources/META-INF/spring/org.springframework.boot.autoconfigure.AutoConfiguration.imports b/auto-configurations/models/spring-ai-autoconfigure-model-anthropic/src/main/resources/META-INF/spring/org.springframework.boot.autoconfigure.AutoConfiguration.imports index 3242a6f7a4..eb7a3faf3d 100644 --- a/auto-configurations/models/spring-ai-autoconfigure-model-anthropic/src/main/resources/META-INF/spring/org.springframework.boot.autoconfigure.AutoConfiguration.imports +++ b/auto-configurations/models/spring-ai-autoconfigure-model-anthropic/src/main/resources/META-INF/spring/org.springframework.boot.autoconfigure.AutoConfiguration.imports @@ -1 +1,2 @@ org.springframework.ai.model.anthropic.autoconfigure.AnthropicChatAutoConfiguration +org.springframework.ai.model.anthropic.autoconfigure.AnthropicBatchAutoConfiguration diff --git a/auto-configurations/models/spring-ai-autoconfigure-model-anthropic/src/test/java/org/springframework/ai/model/anthropic/autoconfigure/AnthropicBatchAutoConfigurationTests.java b/auto-configurations/models/spring-ai-autoconfigure-model-anthropic/src/test/java/org/springframework/ai/model/anthropic/autoconfigure/AnthropicBatchAutoConfigurationTests.java new file mode 100644 index 0000000000..d4802291b1 --- /dev/null +++ b/auto-configurations/models/spring-ai-autoconfigure-model-anthropic/src/test/java/org/springframework/ai/model/anthropic/autoconfigure/AnthropicBatchAutoConfigurationTests.java @@ -0,0 +1,148 @@ +/* + * Copyright 2023-present the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.ai.model.anthropic.autoconfigure; + +import java.time.Duration; +import java.util.List; + +import org.junit.jupiter.api.Test; +import reactor.core.publisher.Flux; + +import org.springframework.ai.anthropic.AnthropicBatch; +import org.springframework.ai.anthropic.AnthropicBatchModel; +import org.springframework.ai.anthropic.AnthropicBatchRequest; +import org.springframework.ai.anthropic.AnthropicBatchResult; +import org.springframework.ai.anthropic.DefaultAnthropicBatchModel; +import org.springframework.ai.model.tool.autoconfigure.ToolCallingAutoConfiguration; +import org.springframework.boot.autoconfigure.AutoConfigurations; +import org.springframework.boot.test.context.runner.ApplicationContextRunner; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; + +import static org.assertj.core.api.Assertions.assertThat; + +/** + * Tests for {@link AnthropicBatchAutoConfiguration}: the batch model must be opt-in, must + * back off when the application defines its own bean, and must inherit the shared + * Anthropic connection and chat defaults. + * + * @author Ricken Bazolo + */ +class AnthropicBatchAutoConfigurationTests { + + private final ApplicationContextRunner contextRunner = new ApplicationContextRunner() + .withConfiguration(AutoConfigurations.of(ToolCallingAutoConfiguration.class, + AnthropicChatAutoConfiguration.class, AnthropicBatchAutoConfiguration.class)) + .withPropertyValues("spring.ai.anthropic.api-key=test-key"); + + @Test + void batchModelIsNotCreatedByDefault() { + this.contextRunner.run(context -> assertThat(context).doesNotHaveBean(AnthropicBatchModel.class)); + } + + @Test + void batchModelIsNotCreatedWhenExplicitlyDisabled() { + this.contextRunner.withPropertyValues("spring.ai.anthropic.batch.enabled=false") + .run(context -> assertThat(context).doesNotHaveBean(AnthropicBatchModel.class)); + } + + @Test + void batchModelIsCreatedWhenEnabled() { + this.contextRunner.withPropertyValues("spring.ai.anthropic.batch.enabled=true").run(context -> { + assertThat(context).hasSingleBean(AnthropicBatchModel.class); + assertThat(context.getBean(AnthropicBatchModel.class)).isInstanceOf(DefaultAnthropicBatchModel.class); + }); + } + + @Test + void batchModelInheritsConnectionAndChatDefaults() { + this.contextRunner + .withPropertyValues("spring.ai.anthropic.batch.enabled=true", "spring.ai.anthropic.base-url=https://proxy", + "spring.ai.anthropic.timeout=45s", "spring.ai.anthropic.max-retries=5", + "spring.ai.anthropic.chat.model=claude-sonnet-4-5", "spring.ai.anthropic.chat.max-tokens=1234") + .run(context -> { + var options = context.getBean(DefaultAnthropicBatchModel.class).getOptions(); + assertThat(options.getApiKey()).isEqualTo("test-key"); + assertThat(options.getBaseUrl()).isEqualTo("https://proxy"); + assertThat(options.getTimeout()).isEqualTo(Duration.ofSeconds(45)); + assertThat(options.getMaxRetries()).isEqualTo(5); + assertThat(options.getModel()).isEqualTo("claude-sonnet-4-5"); + assertThat(options.getMaxTokens()).isEqualTo(1234); + }); + } + + @Test + void batchPropertiesOverrideChatDefaults() { + this.contextRunner + .withPropertyValues("spring.ai.anthropic.batch.enabled=true", + "spring.ai.anthropic.chat.model=claude-sonnet-4-5", "spring.ai.anthropic.chat.max-tokens=1234", + "spring.ai.anthropic.batch.model=claude-haiku-4-5", "spring.ai.anthropic.batch.max-tokens=64") + .run(context -> { + var options = context.getBean(DefaultAnthropicBatchModel.class).getOptions(); + assertThat(options.getModel()).isEqualTo("claude-haiku-4-5"); + assertThat(options.getMaxTokens()).isEqualTo(64); + }); + } + + @Test + void applicationDefinedBatchModelWins() { + this.contextRunner.withPropertyValues("spring.ai.anthropic.batch.enabled=true") + .withUserConfiguration(CustomBatchModelConfiguration.class) + .run(context -> { + assertThat(context).hasSingleBean(AnthropicBatchModel.class); + assertThat(context.getBean(AnthropicBatchModel.class)) + .isNotInstanceOf(DefaultAnthropicBatchModel.class); + }); + } + + @Configuration(proxyBeanMethods = false) + static class CustomBatchModelConfiguration { + + @Bean + AnthropicBatchModel customBatchModel() { + return new AnthropicBatchModel() { + + @Override + public AnthropicBatch submit(List requests) { + throw new UnsupportedOperationException(); + } + + @Override + public AnthropicBatch retrieve(String batchId) { + throw new UnsupportedOperationException(); + } + + @Override + public Flux results(String batchId) { + return Flux.empty(); + } + + @Override + public AnthropicBatch cancel(String batchId) { + throw new UnsupportedOperationException(); + } + + @Override + public void delete(String batchId) { + } + + }; + } + + } + +} From 3dffaab109e8b5d8176c821a7d2d470f0f98632b Mon Sep 17 00:00:00 2001 From: Ricken BAZOLO Date: Fri, 31 Jul 2026 12:17:23 +0200 Subject: [PATCH 6/7] Add Anthropic Message Batches integration test Cover the submit, poll, read, cancel and error paths against the real API, which the mocked unit tests cannot validate: the JSONL result stream, the out-of-order correlation by customId and the usage reported per entry. Two switches guard it, because a batch consumes quota and its completion time is not bounded by the API contract: ANTHROPIC_API_KEY must be set, the repository-wide convention, and ANTHROPIC_BATCH_IT_DISABLED turns it off even when a key is present. A batch is asynchronous by design, so a silent wait would be indistinguishable from a hung test. Each poll logs the elapsed time, the status and the counters, and past the budget the test aborts rather than fails: a slow queue on Anthropic's side is not a regression here. Only one test waits, so the suite waits once. Signed-off-by: Ricken BAZOLO --- .../ai/anthropic/AnthropicBatchIT.java | 268 ++++++++++++++++++ 1 file changed, 268 insertions(+) create mode 100644 models/spring-ai-anthropic/src/test/java/org/springframework/ai/anthropic/AnthropicBatchIT.java diff --git a/models/spring-ai-anthropic/src/test/java/org/springframework/ai/anthropic/AnthropicBatchIT.java b/models/spring-ai-anthropic/src/test/java/org/springframework/ai/anthropic/AnthropicBatchIT.java new file mode 100644 index 0000000000..953704d4ac --- /dev/null +++ b/models/spring-ai-anthropic/src/test/java/org/springframework/ai/anthropic/AnthropicBatchIT.java @@ -0,0 +1,268 @@ +/* + * Copyright 2023-present the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.ai.anthropic; + +import java.time.Duration; +import java.util.ArrayList; +import java.util.List; +import java.util.Map; + +import com.anthropic.client.AnthropicClient; +import com.anthropic.errors.AnthropicServiceException; +import com.anthropic.models.messages.Model; +import org.apache.commons.logging.Log; +import org.apache.commons.logging.LogFactory; +import org.awaitility.Awaitility; +import org.awaitility.core.ConditionTimeoutException; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.Assumptions; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.condition.DisabledIfEnvironmentVariable; +import org.junit.jupiter.api.condition.EnabledIfEnvironmentVariable; + +import org.springframework.ai.chat.messages.SystemMessage; +import org.springframework.ai.chat.messages.UserMessage; +import org.springframework.ai.chat.prompt.Prompt; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.SpringBootConfiguration; +import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.context.annotation.Bean; +import org.springframework.util.StringUtils; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatExceptionOfType; + +/** + * Integration tests for the Anthropic + * Message Batches + * API support: full submit → poll → read lifecycle, cancellation, and error handling + * against the real API. + * + *

+ * Gating. Two switches guard this test, because a batch consumes tokens and its + * completion time is not bounded by the API contract: + *

    + *
  • {@code ANTHROPIC_API_KEY} must be set — the repository-wide convention for provider + * integration tests.
  • + *
  • {@code ANTHROPIC_BATCH_IT_DISABLED=true} turns it off even when a key is present, + * for runs that must not spend batch quota.
  • + *
+ * Anthropic ITs are also excluded from the {@code ci-fast-integration-tests} profile, so + * this only runs under {@code -Pintegration-tests}. + * + *

+ * Expect minutes, not seconds. Batches are asynchronous by design and the API + * allows up to 24 hours, though a small batch normally ends within a few minutes of + * queueing. {@link #submitPollAndReadResults()} is the only test that waits, and it logs + * every poll with the status and counters so a slow queue is visibly a slow queue rather + * than a hang. Past {@link #COMPLETION_TIMEOUT} — 5 minutes by default, override with + * {@code ANTHROPIC_BATCH_IT_TIMEOUT_MINUTES} — it {@link Assumptions#abort aborts} + * instead of failing, because a slow queue on Anthropic's side is not a Spring AI + * regression. + * + * @author Ricken Bazolo + * @since 2.0.0 + */ +@SpringBootTest(classes = AnthropicBatchIT.Config.class) +@EnabledIfEnvironmentVariable(named = "ANTHROPIC_API_KEY", matches = ".+") +@DisabledIfEnvironmentVariable(named = "ANTHROPIC_BATCH_IT_DISABLED", matches = "(?i)true") +class AnthropicBatchIT { + + private static final Log logger = LogFactory.getLog(AnthropicBatchIT.class); + + private static final Duration COMPLETION_TIMEOUT = completionTimeout(); + + private static final Duration POLL_INTERVAL = Duration.ofSeconds(20); + + private static final Duration RESULTS_TIMEOUT = Duration.ofMinutes(2); + + @Autowired + private AnthropicBatchModel batchModel; + + private final List createdBatchIds = new ArrayList<>(); + + private static Duration completionTimeout() { + String minutes = System.getenv("ANTHROPIC_BATCH_IT_TIMEOUT_MINUTES"); + return StringUtils.hasText(minutes) ? Duration.ofMinutes(Long.parseLong(minutes.trim())) + : Duration.ofMinutes(5); + } + + @AfterEach + void deleteCreatedBatches() { + for (String batchId : this.createdBatchIds) { + cancelQuietly(batchId); + deleteQuietly(batchId); + } + this.createdBatchIds.clear(); + } + + @Test + void submitPollAndReadResults() { + AnthropicBatch submitted = submit(List.of( + AnthropicBatchRequest.of("spring-ai-batch-it-1", "Reply with exactly one word: ONE"), + AnthropicBatchRequest.of("spring-ai-batch-it-2", + new Prompt(List.of(new SystemMessage("Answer with a single uppercase word and nothing else."), + new UserMessage("Reply with exactly one word: TWO")))), + AnthropicBatchRequest.of("spring-ai-batch-it-3", "Reply with exactly one word: THREE", + AnthropicChatOptions.builder().model(Model.CLAUDE_HAIKU_4_5).maxTokens(64).build()))); + + assertThat(submitted.id()).isNotBlank(); + assertThat(submitted.status()).isIn(AnthropicBatchStatus.IN_PROGRESS, AnthropicBatchStatus.ENDED); + assertThat(submitted.expiresAt()).isAfter(submitted.createdAt()); + + AnthropicBatch ended = awaitEnded(submitted.id()); + assertThat(ended.requestCounts().total()).isEqualTo(3); + assertThat(ended.requestCounts().processing()).isZero(); + assertThat(ended.resultsUrl()).isNotBlank(); + assertThat(ended.endedAt()).isNotNull(); + + Map byCustomId = this.batchModel.results(submitted.id()) + .collectMap(AnthropicBatchResult::customId) + .block(RESULTS_TIMEOUT); + + // Correlation is by customId only: the API does not preserve submission order. + assertThat(byCustomId).containsOnlyKeys("spring-ai-batch-it-1", "spring-ai-batch-it-2", "spring-ai-batch-it-3"); + + byCustomId.forEach((customId, result) -> { + assertThat(result.status()).as("%s outcome, error was %s", customId, result.error()) + .isEqualTo(AnthropicBatchResultStatus.SUCCEEDED); + assertThat(result.error()).isNull(); + assertThat(result.getText()).as("%s response text", customId).isNotBlank(); + // A batched message must carry the same metadata and usage as a realtime + // call. + assertThat(result.chatResponse()).isNotNull(); + assertThat(result.chatResponse().getMetadata().getId()).isNotBlank(); + assertThat(result.chatResponse().getMetadata().getModel()).isNotBlank(); + assertThat(result.chatResponse().getResult().getMetadata().getFinishReason()).isNotBlank(); + assertThat(result.usage()).isNotNull(); + assertThat(result.usage().getPromptTokens()).isPositive(); + assertThat(result.usage().getCompletionTokens()).isPositive(); + assertThat(result.usage().getTotalTokens()).isPositive(); + }); + + // Same batch, fresh subscription: taking one element must not drain the whole + // JSONL stream. Asserted here so the suite waits for a batch only once. + AnthropicBatchResult first = this.batchModel.results(submitted.id()).next().block(RESULTS_TIMEOUT); + assertThat(first).isNotNull(); + assertThat(first.customId()).startsWith("spring-ai-batch-it-"); + } + + @Test + void cancelIsAcknowledged() { + AnthropicBatch submitted = submit(List.of(AnthropicBatchRequest.of("spring-ai-batch-it-cancel", + "Write a detailed multi-paragraph essay about the history of gardening."))); + + AnthropicBatch canceling = this.batchModel.cancel(submitted.id()); + + assertThat(canceling.id()).isEqualTo(submitted.id()); + // A tiny batch can finish before the cancellation request lands. + assertThat(canceling.status()).isIn(AnthropicBatchStatus.CANCELING, AnthropicBatchStatus.ENDED); + if (canceling.isCanceling()) { + assertThat(canceling.cancelInitiatedAt()).isNotNull(); + } + } + + @Test + void retrievingAnUnknownBatchFailsWithAClientError() { + assertThatExceptionOfType(AnthropicServiceException.class) + .isThrownBy(() -> this.batchModel.retrieve("msgbatch_01SpringAiNoSuchBatch00")) + .satisfies(ex -> assertThat(ex.statusCode()).isBetween(400, 499)); + } + + private AnthropicBatch submit(List requests) { + AnthropicBatch batch = this.batchModel.submit(requests); + this.createdBatchIds.add(batch.id()); + return batch; + } + + /** + * Polls until the batch ends, logging every attempt. + *

+ * The logging is not decoration: a batch is asynchronous by design, so a silent wait + * is indistinguishable from a hung test. Each line reports the elapsed time, the + * processing status and the per-outcome counters, so a slow queue is visibly a slow + * queue. + */ + private AnthropicBatch awaitEnded(String batchId) { + long startedAt = System.nanoTime(); + logger.info("Waiting up to %s for batch %s to end (polling every %s)".formatted(COMPLETION_TIMEOUT, batchId, + POLL_INTERVAL)); + try { + Awaitility.await() + .atMost(COMPLETION_TIMEOUT) + .pollInterval(POLL_INTERVAL) + .pollDelay(Duration.ofSeconds(2)) + .until(() -> { + AnthropicBatch current = this.batchModel.retrieve(batchId); + logger.info(" [%3ds] batch %s status=%s counts=%s".formatted( + Duration.ofNanos(System.nanoTime() - startedAt).toSeconds(), batchId, + current.status().getValue(), current.requestCounts())); + return current.isEnded(); + }); + } + catch (ConditionTimeoutException ex) { + Assumptions.abort( + "Batch %s had not ended after %s. The API allows up to 24 hours, so a slow queue on Anthropic's side is not a Spring AI failure; raise ANTHROPIC_BATCH_IT_TIMEOUT_MINUTES to wait longer." + .formatted(batchId, COMPLETION_TIMEOUT)); + } + AnthropicBatch ended = this.batchModel.retrieve(batchId); + assertThat(ended.isEnded()).isTrue(); + return ended; + } + + private void cancelQuietly(String batchId) { + try { + this.batchModel.cancel(batchId); + } + catch (RuntimeException ex) { + // Best-effort cleanup: an already-ended batch cannot be canceled. + } + } + + private void deleteQuietly(String batchId) { + try { + this.batchModel.delete(batchId); + } + catch (RuntimeException ex) { + // Best-effort cleanup: a batch that has not ended yet cannot be deleted. + } + } + + @SpringBootConfiguration + public static class Config { + + @Bean + public AnthropicClient anthropicClient() { + String apiKey = System.getenv("ANTHROPIC_API_KEY"); + if (!StringUtils.hasText(apiKey)) { + throw new IllegalArgumentException( + "You must provide an API key. Put it in an environment variable under the name ANTHROPIC_API_KEY"); + } + return AnthropicSetup.setupSyncClient(null, apiKey, null, null, null, null); + } + + @Bean + public AnthropicBatchModel anthropicBatchModel(AnthropicClient client) { + return AnthropicBatchModel.builder() + .anthropicClient(client) + .options(AnthropicChatOptions.builder().model(Model.CLAUDE_HAIKU_4_5).maxTokens(256).build()) + .build(); + } + + } + +} From f67ec9a4e4815004fdb72b1b1d325052974a70ae Mon Sep 17 00:00:00 2001 From: Ricken BAZOLO Date: Fri, 31 Jul 2026 12:17:38 +0200 Subject: [PATCH 7/7] Document Anthropic Message Batches support Add a dedicated anthropic-batch.adoc page rather than growing anthropic-chat.adoc, which was already the largest page under api/chat at over 1700 lines. It is nested under the Anthropic entry in the navigation, the shape the embeddings section already uses for Bedrock and VertexAI, and anthropic-chat.adoc keeps a short pointer so the feature stays discoverable from the provider's entry page. The page covers submitting a batch, polling it, streaming the results, cancelling and deleting, the auto-configuration properties, and the two constraints callers have to design around: correlation by customId because results are unordered, and the absence of an interactive tool-execution loop. Signed-off-by: Ricken BAZOLO --- models/spring-ai-anthropic/README.md | 61 ++++++ .../src/main/antora/modules/ROOT/nav.adoc | 1 + .../ROOT/pages/api/chat/anthropic-batch.adoc | 177 ++++++++++++++++++ .../ROOT/pages/api/chat/anthropic-chat.adoc | 7 + 4 files changed, 246 insertions(+) create mode 100644 spring-ai-docs/src/main/antora/modules/ROOT/pages/api/chat/anthropic-batch.adoc diff --git a/models/spring-ai-anthropic/README.md b/models/spring-ai-anthropic/README.md index a04992fab6..ab37855db9 100644 --- a/models/spring-ai-anthropic/README.md +++ b/models/spring-ai-anthropic/README.md @@ -32,6 +32,7 @@ This module supports: - **Citations** - Document-grounded responses with source attribution - **Prompt Caching** - Reduce costs for repeated context with configurable strategies - **Structured Output** - JSON schema-constrained responses with effort control +- **Message Batches** - Asynchronous bulk processing at reduced cost, correlated by `custom_id` - **Per-Request HTTP Headers** - Custom headers per API call for tracking, beta features, and routing - **Observability** - Micrometer-based metrics and tracing @@ -266,6 +267,66 @@ var options = AnthropicChatOptions.builder() ChatResponse response = chatModel.call(new Prompt("Hello", options)); ``` +## Message Batches + +Submit many prompts at once for asynchronous processing. Batch results are not returned in +submission order, so every entry carries a `customId` used to correlate its result. + +Processing can take up to 24 hours. Spring AI deliberately performs **no polling, no +persistence and no scheduling**: it exposes the five provider operations, and the application +owns the orchestration. + +```java +AnthropicBatchModel batchModel = AnthropicBatchModel.builder() + .options(AnthropicChatOptions.builder().model("claude-haiku-4-5").maxTokens(1024).build()) + .build(); + +// 1. Submit — returns as soon as Anthropic accepts the batch +AnthropicBatch batch = batchModel.submit(List.of( + AnthropicBatchRequest.of("invoice-1", "Summarize invoice 1"), + AnthropicBatchRequest.of("invoice-2", new Prompt("Summarize invoice 2", perRequestOptions)))); + +// Persist batch.id() and the customIds so polling survives a restart +String batchId = batch.id(); + +// 2. Poll on your own schedule +if (batchModel.retrieve(batchId).isEnded()) { + + // 3. Stream the results — never buffered wholesale + batchModel.results(batchId) + .doOnNext(result -> { + switch (result.status()) { + case SUCCEEDED -> store(result.customId(), result.getText(), result.usage()); + case ERRORED -> logFailure(result.customId(), result.error()); + case CANCELED, EXPIRED, UNKNOWN -> requeue(result.customId()); + } + }) + .blockLast(); +} + +// Optional lifecycle control +batchModel.cancel(batchId); +batchModel.delete(batchId); +``` + +In Spring Boot the batch model is opt-in, so applications that never submit batches do not +pay for a second HTTP client: + +```properties +spring.ai.anthropic.batch.enabled=true +``` + +Connection settings (`spring.ai.anthropic.*`) and model defaults +(`spring.ai.anthropic.chat.*`) are shared with the chat model; `spring.ai.anthropic.batch.model` +and `spring.ai.anthropic.batch.max-tokens` override them for batch entries only. + +Requests go through the same `Prompt` conversion as `AnthropicChatModel.call(...)`, and a +succeeded result is converted back into the same `ChatResponse` shape. + +> **Tool calls are not executed.** Tool definitions are sent, but a batch response containing +> `tool_use` blocks is returned as-is — a batch entry cannot be continued mid-flight. Read the +> tool calls off the assistant message and submit a follow-up batch for the next turn. + ## Logging Enable SDK logging by setting the environment variable: diff --git a/spring-ai-docs/src/main/antora/modules/ROOT/nav.adoc b/spring-ai-docs/src/main/antora/modules/ROOT/nav.adoc index bfdb21c478..66057fcde6 100644 --- a/spring-ai-docs/src/main/antora/modules/ROOT/nav.adoc +++ b/spring-ai-docs/src/main/antora/modules/ROOT/nav.adoc @@ -19,6 +19,7 @@ **** xref:api/chat/comparison.adoc[Chat Models Comparison] **** xref:api/chat/bedrock-converse.adoc[Amazon Bedrock Converse] **** xref:api/chat/anthropic-chat.adoc[Anthropic] +***** xref:api/chat/anthropic-batch.adoc[Message Batches] **** xref:api/chat/azure-openai-chat.adoc[Azure OpenAI] **** xref:api/chat/deepseek-chat.adoc[DeepSeek] **** xref:api/chat/dmr-chat.adoc[Docker Model Runner] diff --git a/spring-ai-docs/src/main/antora/modules/ROOT/pages/api/chat/anthropic-batch.adoc b/spring-ai-docs/src/main/antora/modules/ROOT/pages/api/chat/anthropic-batch.adoc new file mode 100644 index 0000000000..85a9201a5d --- /dev/null +++ b/spring-ai-docs/src/main/antora/modules/ROOT/pages/api/chat/anthropic-batch.adoc @@ -0,0 +1,177 @@ +[[anthropic-batch]] += Anthropic Message Batches + +The link:https://platform.claude.com/docs/en/api/messages/batches[Message Batches API] processes a large number of prompts asynchronously at reduced cost. +Spring AI exposes it through `AnthropicBatchModel`, a dedicated abstraction rather than a `ChatModel`: `submit(...)` returns as soon as Anthropic accepts the batch, and there is no `ChatResponse` to return yet. + +This page covers the batch model only. +See xref:api/chat/anthropic-chat.adoc[Anthropic Chat] for the dependency, the Boot starter and the `spring.ai.anthropic.*` connection properties, which the batch model shares. + +[IMPORTANT] +==== +**Spring AI provides the provider access, not the orchestration.** + +There is no automatic polling, no persistence, no scheduling and no business retry. +Deciding how often to call `retrieve(...)`, where to store the batch id and the `customId` correlations, when to notify and how to account for cost all stay with your application. +==== + +== Dependency and auto-configuration + +The batch model ships in the same `spring-ai-starter-model-anthropic` starter as the chat model, but it is **opt-in** so that applications which never submit batches do not pay for a second HTTP client: + +[source,properties] +---- +spring.ai.anthropic.batch.enabled=true +---- + +Connection settings (`spring.ai.anthropic.*`) and model defaults (`spring.ai.anthropic.chat.*`) are shared with the chat model, so a batch reuses the same credentials, base URL, timeout, retries, proxy, custom headers and `AnthropicHttpClientBuilderCustomizer` beans as realtime calls. + +[cols="3,4"] +|==== +| Property | Description + +| `spring.ai.anthropic.batch.enabled` +| Whether to expose an `AnthropicBatchModel` bean. Default `false`. + +| `spring.ai.anthropic.batch.model` +| Model for batch entries. Falls back to `spring.ai.anthropic.chat.model`. + +| `spring.ai.anthropic.batch.max-tokens` +| Output-token ceiling per batch entry. Falls back to `spring.ai.anthropic.chat.max-tokens`. +|==== + +Outside Spring Boot, build one directly — the builder accepts the same connection and observability configuration as `AnthropicChatModel.builder()`: + +[source,java] +---- +AnthropicBatchModel batchModel = AnthropicBatchModel.builder() + .options(AnthropicChatOptions.builder() + .model("claude-haiku-4-5") + .maxTokens(1024) + .build()) + .observationRegistry(observationRegistry) + .build(); +---- + +== Submitting a batch + +Every entry pairs a `customId` with a `Prompt`. +The prompt is mapped through the *same* conversion used by `AnthropicChatModel.call(...)`, so system messages, conversation history, images and PDF documents, prompt caching, thinking, structured output and tool definitions all behave identically on both paths. +Per-entry options are taken from the prompt when it carries `AnthropicChatOptions`; otherwise the batch model's defaults apply. + +[source,java] +---- +AnthropicBatch batch = batchModel.submit(List.of( + AnthropicBatchRequest.of("invoice-1", "Summarize invoice 1"), + AnthropicBatchRequest.of("invoice-2", "Summarize invoice 2", + AnthropicChatOptions.builder().model("claude-sonnet-4-6").maxTokens(2048).build()), + AnthropicBatchRequest.of("invoice-3", new Prompt(List.of( + new SystemMessage("Answer in one sentence."), + new UserMessage("Summarize invoice 3")))))); + +// Persist these before returning: they are the only handle on the work in flight +String batchId = batch.id(); +---- + +`customId` must be 1–64 characters limited to letters, digits, underscores and hyphens, and must be unique within the batch. +Both rules are enforced client-side, so a malformed identifier fails fast instead of costing a round trip. + +== Checking status + +Batch processing is asynchronous and **can take up to 24 hours**. +Poll on your own schedule: + +[source,java] +---- +AnthropicBatch current = batchModel.retrieve(batchId); + +if (current.isEnded()) { + // results can now be read +} + +// Progress reporting without reading the result stream +AnthropicBatchRequestCounts counts = current.requestCounts(); +log.info("{}/{} done ({} errored)", counts.completed(), counts.total(), counts.errored()); +---- + +`AnthropicBatch` also exposes `createdAt()`, `expiresAt()`, `endedAt()`, `cancelInitiatedAt()`, `archivedAt()` and `resultsUrl()`. +`status()` is an `AnthropicBatchStatus`: `IN_PROGRESS`, `CANCELING`, `ENDED`, or `UNKNOWN` for a status added by the API after your Spring AI release. + +== Reading results + +`results(...)` returns a `Flux` that consumes the JSONL stream lazily and closes the underlying SDK stream when the flux terminates or is cancelled, so a batch with a very large number of entries is never held in memory. + +[IMPORTANT] +==== +**Results are not returned in submission order.** Key them by `customId`, never by position in the stream. +==== + +Individual failures are emitted as `ERRORED` items rather than thrown, so one bad entry never hides the rest of the batch: + +[source,java] +---- +batchModel.results(batchId) + .doOnNext(result -> { + switch (result.status()) { + case SUCCEEDED -> { + ChatResponse response = result.chatResponse(); + store(result.customId(), response.getResult().getOutput().getText()); + recordCost(result.customId(), result.usage()); + } + case ERRORED -> { + AnthropicBatchError error = result.error(); + log.warn("{} failed: {} — {}", result.customId(), error.type(), error.message()); + } + case CANCELED, EXPIRED, UNKNOWN -> requeue(result.customId()); + } + }) + .blockLast(); +---- + +A succeeded result carries a full `ChatResponse` — same generations, same metadata keys (including `anthropic-response`, citations and web-search results) and same `Usage` as a realtime call — so downstream code that already consumes `ChatResponse` needs no change. +`AnthropicBatchResult` adds `getText()` and `usage()` shortcuts for the common cases. + +== Cancelling and deleting + +[source,java] +---- +// Cancellation is not immediate: the batch moves to CANCELING, entries that already +// completed keep their result, and the rest end up CANCELED. +AnthropicBatch canceling = batchModel.cancel(batchId); + +// Only batches whose processing has ended can be deleted. +batchModel.delete(batchId); +---- + +== Correlating results with your domain + +`customId` is the only correlation handle the API offers, and results outlive the JVM that submitted them. +Persist the batch id together with the `customId` of every entry, mapped to your own identifiers, before `submit(...)` returns to the caller — otherwise a restart mid-flight leaves results that cannot be attributed. + +[source,java] +---- +// customId built from a domain identifier so correlation survives a restart +AnthropicBatchRequest.of("invoice-" + invoice.getId(), prompt(invoice)); +---- + +== Tool calling limitation + +Tool definitions resolved from the request options *are* sent with each entry, but a batch response containing `tool_use` blocks is returned as-is: there is **no interactive tool-execution loop**, because a batch entry cannot be continued mid-flight. +Read the tool calls off the assistant message and submit a follow-up batch if you need a second turn: + +[source,java] +---- +AssistantMessage output = result.chatResponse().getResult().getOutput(); +if (!output.getToolCalls().isEmpty()) { + // execute locally, then submit a follow-up batch with the tool results appended +} +---- + +== Observability + +Every batch operation emits a `spring.ai.anthropic.batch.operation` observation whose `gen_ai.operation.name` tag is `batch_create`, `batch_retrieve`, `batch_results`, `batch_cancel` or `batch_delete`, alongside `gen_ai.system`, `gen_ai.request.model` and `spring.ai.anthropic.batch.status`. +The submitted request count and the per-outcome counters are recorded as high-cardinality tags. + +Neither the batch id nor any prompt or generated content is exposed as a tag: batch ids are unbounded in cardinality, and prompt content must not reach a metrics backend. + +Supply your own `AnthropicBatchObservationConvention` bean to customise the tags. diff --git a/spring-ai-docs/src/main/antora/modules/ROOT/pages/api/chat/anthropic-chat.adoc b/spring-ai-docs/src/main/antora/modules/ROOT/pages/api/chat/anthropic-chat.adoc index 89b7f60285..6de1c45f19 100644 --- a/spring-ai-docs/src/main/antora/modules/ROOT/pages/api/chat/anthropic-chat.adoc +++ b/spring-ai-docs/src/main/antora/modules/ROOT/pages/api/chat/anthropic-chat.adoc @@ -1642,6 +1642,13 @@ spring.ai.anthropic.chat.web-search-tool.user-location.city=San Francisco spring.ai.anthropic.chat.web-search-tool.user-location.country=US ---- +== Message Batches + +Anthropic can process a large number of prompts asynchronously at reduced cost through the Message Batches API. +Spring AI exposes it as a dedicated `AnthropicBatchModel` rather than a `ChatModel`, because a submitted batch has no immediate response to return. + +See xref:api/chat/anthropic-batch.adoc[Anthropic Message Batches] for submitting a batch, polling it, reading results, and the auto-configuration properties. + == Observability Spring AI emits Micrometer observations at two layers for every Anthropic call: