Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
106 changes: 98 additions & 8 deletions crates/adaptive/src/response_cache/anthropic.rs
Original file line number Diff line number Diff line change
Expand Up @@ -12,26 +12,32 @@ use serde_json::{Value as Json, json};
use crate::response_cache::mark::CacheReason;

const THRESHOLD_BETA: &str = "compact-2026-01-12";
const ON_DEMAND_BETA: &str = "compact-2026-09-04";
const ANTHROPIC_API_VERSION: &str = "2023-06-01";
const MIN_THRESHOLD_TOKENS: u64 = 50_000;
const MAX_INSTRUCTIONS_CHARS: usize = 16_384;

#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
pub(crate) enum AnthropicProtocol {
ThresholdV1,
OnDemandV1,
}

#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
pub(crate) enum AnthropicOperation {
PausedThreshold,
ThresholdContinuation,
ThresholdRecompact,
OnDemandSummarize,
OnDemandContinuation,
OnDemandRecompact,
}

#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum AnthropicResponseKind {
Ordinary,
ThresholdCompaction,
OnDemandCompaction,
}

#[derive(Debug, Clone, PartialEq, Eq)]
Expand Down Expand Up @@ -96,6 +102,7 @@ pub(crate) fn cache_context(
.collect();
let protocol = match compact_tokens.as_slice() {
[token] if *token == THRESHOLD_BETA && !on_demand_related => AnthropicProtocol::ThresholdV1,
[token] if *token == ON_DEMAND_BETA && !threshold_related => AnthropicProtocol::OnDemandV1,
_ => return Err(CacheReason::AnthropicCompaction),
};

Expand All @@ -106,7 +113,14 @@ pub(crate) fn cache_context(
{
return Err(CacheReason::AnthropicCompaction);
}
if continuation && !has_supported_continuation(body) {
if on_demand_related
&& !body
.get("compaction")
.is_some_and(is_supported_on_demand_request)
{
return Err(CacheReason::AnthropicCompaction);
}
if continuation && !has_supported_continuation(body, protocol) {
return Err(CacheReason::AnthropicCompaction);
}

Expand All @@ -118,6 +132,13 @@ pub(crate) fn cache_context(
(AnthropicProtocol::ThresholdV1, true, false, true) => {
AnthropicOperation::ThresholdRecompact
}
(AnthropicProtocol::OnDemandV1, false, true, false) => {
AnthropicOperation::OnDemandSummarize
}
(AnthropicProtocol::OnDemandV1, false, false, true) => {
AnthropicOperation::OnDemandContinuation
}
(AnthropicProtocol::OnDemandV1, false, true, true) => AnthropicOperation::OnDemandRecompact,
_ => return Err(CacheReason::AnthropicCompaction),
};

Expand Down Expand Up @@ -221,6 +242,16 @@ fn context_management_is_compaction_related(value: &Json) -> bool {
})
}

fn is_supported_on_demand_request(value: &Json) -> bool {
value.as_object().is_some_and(|object| {
object
.keys()
.all(|key| matches!(key.as_str(), "type" | "instructions"))
&& object.get("type").and_then(Json::as_str) == Some("summarize")
&& object.get("instructions").is_none_or(valid_instructions)
})
}

fn compaction_blocks(body: &serde_json::Map<String, Json>) -> Vec<&Json> {
body.get("messages")
.and_then(Json::as_array)
Expand All @@ -236,12 +267,25 @@ fn contains_compaction_block(body: &serde_json::Map<String, Json>) -> bool {
!compaction_blocks(body).is_empty()
}

fn has_supported_continuation(body: &serde_json::Map<String, Json>) -> bool {
fn has_supported_continuation(
body: &serde_json::Map<String, Json>,
protocol: AnthropicProtocol,
) -> bool {
let blocks = compaction_blocks(body);
!blocks.is_empty()
&& blocks
.iter()
.all(|block| threshold_request_block_is_valid(block))
match protocol {
AnthropicProtocol::ThresholdV1 => {
!blocks.is_empty()
&& blocks
.iter()
.all(|block| threshold_request_block_is_valid(block))
}
AnthropicProtocol::OnDemandV1 => {
let [block] = blocks.as_slice() else {
return false;
};
on_demand_request_block_is_valid(block)
}
}
}

fn cache_control_is_valid(value: &Json) -> bool {
Expand All @@ -263,6 +307,18 @@ fn threshold_request_block_is_valid(block: &Json) -> bool {
.is_none_or(cache_control_is_valid)
}

fn on_demand_request_block_is_valid(block: &Json) -> bool {
has_only_fields(block, &["type", "content", "signature", "cache_control"])
&& block.get("content").is_some_and(Json::is_string)
&& block
.get("signature")
.and_then(Json::as_str)
.is_some_and(|value| !value.is_empty())
&& block
.get("cache_control")
.is_none_or(cache_control_is_valid)
}

pub(crate) fn classify_aggregate(
response: &Json,
context: &AnthropicCacheContext,
Expand Down Expand Up @@ -335,6 +391,12 @@ pub(crate) fn classify_aggregate(
{
Some(AnthropicResponseKind::ThresholdCompaction)
}
AnthropicProtocol::OnDemandV1
if operation_allows_on_demand_compaction(context.operation)
&& on_demand_response_block_is_valid(block) =>
{
Some(AnthropicResponseKind::OnDemandCompaction)
}
_ => None,
}
}
Expand All @@ -345,6 +407,7 @@ fn operation_allows_ordinary(operation: AnthropicOperation) -> bool {
AnthropicOperation::PausedThreshold
| AnthropicOperation::ThresholdContinuation
| AnthropicOperation::ThresholdRecompact
| AnthropicOperation::OnDemandContinuation
)
}

Expand All @@ -355,6 +418,13 @@ fn operation_allows_threshold_compaction(operation: AnthropicOperation) -> bool
)
}

fn operation_allows_on_demand_compaction(operation: AnthropicOperation) -> bool {
matches!(
operation,
AnthropicOperation::OnDemandSummarize | AnthropicOperation::OnDemandRecompact
)
}

fn valid_compaction_iterations(value: &Json) -> bool {
value.as_array().is_some_and(|iterations| {
!iterations.is_empty()
Expand All @@ -379,6 +449,15 @@ fn threshold_response_block_is_valid(block: &Json) -> bool {
.is_none_or(|value| value.is_null() || value.is_string())
}

fn on_demand_response_block_is_valid(block: &Json) -> bool {
has_only_fields(block, &["type", "content", "signature"])
&& block.get("content").is_some_and(Json::is_string)
&& block
.get("signature")
.and_then(Json::as_str)
.is_some_and(|value| !value.is_empty())
}

fn has_only_fields(value: &Json, allowed: &[&str]) -> bool {
value
.as_object()
Expand Down Expand Up @@ -471,6 +550,12 @@ impl AnthropicStreamValidator {
&& compaction.len() == 1
&& compaction[0].compaction_delta_count == 1
}
AnthropicResponseKind::OnDemandCompaction => {
self.protocol == AnthropicProtocol::OnDemandV1
&& self.blocks.len() == 1
&& compaction.len() == 1
&& compaction[0].compaction_delta_count == 0
}
}
}

Expand Down Expand Up @@ -542,8 +627,13 @@ impl AnthropicStreamValidator {
return false;
};
if kind == "compaction" {
let valid = has_only_fields(block, &["type", "content"])
&& block.get("content").is_some_and(Json::is_null);
let valid = match self.protocol {
AnthropicProtocol::ThresholdV1 => {
has_only_fields(block, &["type", "content"])
&& block.get("content").is_some_and(Json::is_null)
}
AnthropicProtocol::OnDemandV1 => on_demand_response_block_is_valid(block),
};
if !valid {
return false;
}
Expand Down
Loading
Loading