This is an automated email from the ASF dual-hosted git repository.

shreemaan-abhishek pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/apisix.git


The following commit(s) were added to refs/heads/master by this push:
     new 19e56c40b fix(ai-content-moderation): stop re-counting converted 
chunks in realtime batches (#13765)
19e56c40b is described below

commit 19e56c40b4d4f3c53f7a1aac56ab3694efdb966e
Author: Shreemaan Abhishek <[email protected]>
AuthorDate: Mon Aug 3 18:38:25 2026 +0800

    fix(ai-content-moderation): stop re-counting converted chunks in realtime 
batches (#13765)
    
    The realtime batch accumulator read ctx.llm_response_contents_in_chunk, 
which
    base.lua refills once per upstream chunk. With a protocol converter active 
one
    upstream chunk is dispatched as several converted chunks, so 
lua_body_filter ran
    once per converted chunk and appended the same texts every time: the 
moderation
    batch held the text N times over, stream_check_cache_size tripped N times
    sooner, and request volume scaled with the converter's fan-out.
    
    base.lua now bumps a per-upstream-chunk counter where it resets the texts, 
and
    each plugin takes them only on its first run for that counter.
---
 apisix/plugins/ai-aliyun-content-moderation.lua |  12 +-
 apisix/plugins/ai-aws-content-moderation.lua    |  12 +-
 apisix/plugins/ai-providers/base.lua            |   5 +
 t/plugin/ai-aliyun-content-moderation.t         | 155 ++++++++++++++++++++++++
 t/plugin/ai-aws-content-moderation.t            | 153 +++++++++++++++++++++++
 5 files changed, 333 insertions(+), 4 deletions(-)

diff --git a/apisix/plugins/ai-aliyun-content-moderation.lua 
b/apisix/plugins/ai-aliyun-content-moderation.lua
index 6ba2676ef..8778b754d 100644
--- a/apisix/plugins/ai-aliyun-content-moderation.lua
+++ b/apisix/plugins/ai-aliyun-content-moderation.lua
@@ -539,8 +539,16 @@ function _M.lua_body_filter(conf, ctx, headers, body)
     if conf.stream_check_mode == "realtime" then
         ctx.content_moderation_cache = ctx.content_moderation_cache or ""
         ctx.llm_response_contents_in_chunk = 
ctx.llm_response_contents_in_chunk or {}
-        local content = table.concat(ctx.llm_response_contents_in_chunk, "")
-        ctx.content_moderation_cache = ctx.content_moderation_cache .. content
+        -- With a protocol converter a single upstream chunk is dispatched as 
several
+        -- downstream chunks, so this filter runs once per converted chunk 
while
+        -- llm_response_contents_in_chunk is filled once per upstream chunk. 
Take the
+        -- texts on the first run only, otherwise the batch holds them N times 
over.
+        local chunk_seq = ctx.llm_response_chunk_seq
+        if not chunk_seq or chunk_seq ~= ctx.aliyun_cm_chunk_seq then
+            ctx.aliyun_cm_chunk_seq = chunk_seq
+            local content = table.concat(ctx.llm_response_contents_in_chunk, 
"")
+            ctx.content_moderation_cache = ctx.content_moderation_cache .. 
content
+        end
         local now_time = ngx.now()
         ctx.last_moderate_time = ctx.last_moderate_time or now_time
         if #ctx.content_moderation_cache < conf.stream_check_cache_size
diff --git a/apisix/plugins/ai-aws-content-moderation.lua 
b/apisix/plugins/ai-aws-content-moderation.lua
index c40fe84b4..01fb73ddc 100644
--- a/apisix/plugins/ai-aws-content-moderation.lua
+++ b/apisix/plugins/ai-aws-content-moderation.lua
@@ -401,8 +401,16 @@ function _M.lua_body_filter(conf, ctx, headers, body)
 
     -- realtime: moderate batches as they arrive so a hit can cut the stream 
off
     ctx.aws_cm_cache = ctx.aws_cm_cache or ""
-    ctx.aws_cm_cache = ctx.aws_cm_cache
-                       .. table.concat(ctx.llm_response_contents_in_chunk or 
{}, "")
+    -- With a protocol converter a single upstream chunk is dispatched as 
several
+    -- downstream chunks, so this filter runs once per converted chunk while
+    -- llm_response_contents_in_chunk is filled once per upstream chunk. Take 
the
+    -- texts on the first run only, otherwise the batch holds them N times 
over.
+    local chunk_seq = ctx.llm_response_chunk_seq
+    if not chunk_seq or chunk_seq ~= ctx.aws_cm_chunk_seq then
+        ctx.aws_cm_chunk_seq = chunk_seq
+        ctx.aws_cm_cache = ctx.aws_cm_cache
+                           .. table.concat(ctx.llm_response_contents_in_chunk 
or {}, "")
+    end
     local now = ngx.now()
     ctx.aws_cm_last_check = ctx.aws_cm_last_check or now
     if #ctx.aws_cm_cache < conf.stream_check_cache_size
diff --git a/apisix/plugins/ai-providers/base.lua 
b/apisix/plugins/ai-providers/base.lua
index 6ed57b84e..61b1c9529 100644
--- a/apisix/plugins/ai-providers/base.lua
+++ b/apisix/plugins/ai-providers/base.lua
@@ -632,6 +632,11 @@ function _M.parse_streaming_response(self, ctx, res, 
target_proto, converter, co
         end
         sse_parts = {remainder}
         ctx.llm_response_contents_in_chunk = {}
+        -- Bumped with every reset so body_filter hooks can tell which upstream
+        -- chunk the texts above belong to. A converter dispatches one upstream
+        -- chunk as several downstream chunks, running those hooks more than
+        -- once for the same texts.
+        ctx.llm_response_chunk_seq = (ctx.llm_response_chunk_seq or 0) + 1
         local converted_chunks = {}
 
         for _, event in ipairs(events) do
diff --git a/t/plugin/ai-aliyun-content-moderation.t 
b/t/plugin/ai-aliyun-content-moderation.t
index ddfd91183..e12ff01ca 100644
--- a/t/plugin/ai-aliyun-content-moderation.t
+++ b/t/plugin/ai-aliyun-content-moderation.t
@@ -2307,3 +2307,158 @@ bed_blk:bs1,bs2
 199: rejected
 600: rejected
 403: accepted
+
+
+
+=== TEST 74: realtime takes an upstream chunk's text once per chunk, not per 
converted chunk
+--- config
+    location /t {
+        content_by_lua_block {
+            local plugin = 
require("apisix.plugins.ai-aliyun-content-moderation")
+            local ctx = {
+                picked_ai_instance = { provider = "openai" },
+                ai_client_protocol = "anthropic-messages",
+                var = { request_type = "ai_stream" },
+                llm_response_contents_in_chunk = { "hello" },
+                llm_response_chunk_seq = 1,
+            }
+            local conf = {
+                endpoint = "http://localhost:6724";,
+                region_id = "cn-shanghai",
+                access_key_id = "fake-key-id",
+                access_key_secret = "fake-key-secret",
+                check_response = true,
+                stream_check_mode = "realtime",
+                stream_check_cache_size = 4096,
+                stream_check_interval = 60,
+            }
+            plugin.check_schema(conf)
+            -- a converter dispatches one upstream chunk as several downstream
+            -- chunks, so the filter runs once per converted chunk
+            for _ = 1, 3 do
+                plugin.lua_body_filter(conf, ctx, {}, "data: {}\n\n")
+            end
+            ngx.say("chunk 1: ", ctx.content_moderation_cache)
+
+            ctx.llm_response_contents_in_chunk = { " world" }
+            ctx.llm_response_chunk_seq = 2
+            for _ = 1, 2 do
+                plugin.lua_body_filter(conf, ctx, {}, "data: {}\n\n")
+            end
+            ngx.say("chunk 2: ", ctx.content_moderation_cache)
+        }
+    }
+--- response_body
+chunk 1: hello
+chunk 2: hello world
+--- no_error_log
+[error]
+
+
+
+=== TEST 75: set route serving an Anthropic client from an OpenAI upstream 
(converter active)
+--- config
+    location /t {
+        content_by_lua_block {
+            local t = require("lib.test_admin").test
+            local code, body = t('/apisix/admin/routes/1',
+                ngx.HTTP_PUT,
+                [[{
+                    "uri": "/v1/messages",
+                    "plugins": {
+                      "ai-proxy": {
+                          "provider": "openai",
+                          "auth": {
+                              "header": {
+                                  "Authorization": "Bearer token"
+                              }
+                          },
+                          "options": {
+                              "model": "claude-3-5-sonnet-20241022",
+                              "stream": true
+                          },
+                          "override": {
+                              "endpoint": 
"http://127.0.0.1:1980/v1/chat/completions";
+                          }
+                      },
+                      "ai-aliyun-content-moderation": {
+                        "endpoint": "http://localhost:6724";,
+                        "region_id": "cn-shanghai",
+                        "access_key_id": "fake-key-id",
+                        "access_key_secret": "fake-key-secret",
+                        "risk_level_bar": "high",
+                        "check_request": false,
+                        "check_response": true,
+                        "stream_check_mode": "realtime",
+                        "stream_check_cache_size": 4096,
+                        "stream_check_interval": 60
+                      }
+                    }
+                }]]
+            )
+
+            if code >= 300 then
+                ngx.status = code
+            end
+            ngx.say(body)
+        }
+    }
+--- response_body
+passed
+
+
+
+=== TEST 76: converter fan-out moderates the response text once, not once per 
converted chunk
+--- config
+    location /t {
+        content_by_lua_block {
+            local http = require("resty.http")
+            local httpc = http.new()
+
+            local ok, err = httpc:connect({
+                scheme = "http",
+                host = "localhost",
+                port = ngx.var.server_port,
+            })
+            if not ok then
+                ngx.status = 500
+                ngx.say(err)
+                return
+            end
+
+            local res, err = httpc:request({
+                method = "POST",
+                path = "/v1/messages",
+                headers = {
+                    ["Content-Type"] = "application/json",
+                    ["Connection"] = "close",
+                    ["X-AI-Fixture"] = 
"protocol-conversion/openai-to-anthropic-stream.sse",
+                },
+                body = [[{
+                    "model": "claude-3-5-sonnet-20241022",
+                    "messages": [{"role": "user", "content": "Hi"}],
+                    "stream": true
+                }]],
+            })
+            if not res then
+                ngx.status = 500
+                ngx.say(err)
+                return
+            end
+
+            local results = {}
+            while true do
+                local chunk = res.body_reader()
+                if not chunk then break end
+                table.insert(results, chunk)
+            end
+            ngx.print(table.concat(results, ""))
+        }
+    }
+--- error_code: 200
+--- response_body_like eval
+qr/event: message_stop/
+--- grep_error_log eval
+qr/execute content moderation/
+--- grep_error_log_out
+execute content moderation
diff --git a/t/plugin/ai-aws-content-moderation.t 
b/t/plugin/ai-aws-content-moderation.t
index 01872c510..12abbc2db 100644
--- a/t/plugin/ai-aws-content-moderation.t
+++ b/t/plugin/ai-aws-content-moderation.t
@@ -81,6 +81,7 @@ _EOC_
                         ngx.say("[INTERNAL FAILURE]: failed to decoded request 
body: ", err)
                     end
                     local result = body.TextSegments[1].Text
+                    ngx.log(ngx.WARN, "comprehend text: ", result)
                     local final_response = responses[result]
 
                     -- Response-side text is free-form LLM output, not a 
fixture
@@ -871,3 +872,155 @@ qr/code:nil\nrisk_level:none\n.*"risk_level":"none"/s
 qr/code:nil\n.*data: 12345.*"risk_level":"none"/s
 --- no_error_log
 [error]
+
+
+
+=== TEST 33: realtime takes an upstream chunk's text once per chunk, not per 
converted chunk
+--- config
+    location /t {
+        content_by_lua_block {
+            local plugin = require("apisix.plugins.ai-aws-content-moderation")
+            local ctx = {
+                picked_ai_instance = { provider = "openai" },
+                ai_client_protocol = "anthropic-messages",
+                var = { request_type = "ai_stream" },
+                llm_response_contents_in_chunk = { "hello" },
+                llm_response_chunk_seq = 1,
+            }
+            local conf = {
+                comprehend = {
+                    access_key_id = "access",
+                    secret_access_key = "secret",
+                    region = "us-east-1",
+                    endpoint = "http://localhost:2668";
+                },
+                check_response = true,
+                stream_check_mode = "realtime",
+                stream_check_cache_size = 4096,
+                stream_check_interval = 60,
+            }
+            plugin.check_schema(conf)
+            -- a converter dispatches one upstream chunk as several downstream
+            -- chunks, so the filter runs once per converted chunk
+            for _ = 1, 3 do
+                plugin.lua_body_filter(conf, ctx, {}, "data: {}\n\n")
+            end
+            ngx.say("chunk 1: ", ctx.aws_cm_cache)
+
+            ctx.llm_response_contents_in_chunk = { " world" }
+            ctx.llm_response_chunk_seq = 2
+            for _ = 1, 2 do
+                plugin.lua_body_filter(conf, ctx, {}, "data: {}\n\n")
+            end
+            ngx.say("chunk 2: ", ctx.aws_cm_cache)
+        }
+    }
+--- response_body
+chunk 1: hello
+chunk 2: hello world
+--- no_error_log
+[error]
+
+
+
+=== TEST 34: set route serving an Anthropic client from an OpenAI upstream 
(converter active)
+--- config
+    location /t {
+        content_by_lua_block {
+            local t = require("lib.test_admin").test
+            local code, body = t('/apisix/admin/routes/1',
+                ngx.HTTP_PUT,
+                [[{
+                    "uri": "/anthropic/v1/messages",
+                    "plugins": {
+                        "ai-proxy": {
+                            "provider": "openai",
+                            "auth": { "header": { "Authorization": "Bearer 
token" } },
+                            "options": {
+                                "model": "claude-3-5-sonnet-20241022",
+                                "stream": true
+                            },
+                            "override": { "endpoint": 
"http://127.0.0.1:1980/v1/chat/completions"; }
+                        },
+                        "ai-aws-content-moderation": {
+                            "comprehend": {
+                                "access_key_id": "access",
+                                "secret_access_key": "ea+secret",
+                                "region": "us-east-1",
+                                "endpoint": "http://localhost:2668";
+                            },
+                            "check_request": false,
+                            "check_response": true,
+                            "stream_check_mode": "realtime",
+                            "stream_check_cache_size": 4096,
+                            "stream_check_interval": 60
+                        }
+                    }
+                }]]
+            )
+
+            if code >= 300 then
+                ngx.status = code
+            end
+            ngx.say(body)
+        }
+    }
+--- response_body
+passed
+
+
+
+=== TEST 35: converter fan-out moderates the response text once, not once per 
converted chunk
+--- config
+    location /t {
+        content_by_lua_block {
+            local http = require("resty.http")
+            local httpc = http.new()
+
+            local ok, err = httpc:connect({
+                scheme = "http",
+                host = "localhost",
+                port = ngx.var.server_port,
+            })
+            if not ok then
+                ngx.status = 500
+                ngx.say(err)
+                return
+            end
+
+            local res, err = httpc:request({
+                method = "POST",
+                path = "/anthropic/v1/messages",
+                headers = {
+                    ["Content-Type"] = "application/json",
+                    ["Connection"] = "close",
+                    ["X-AI-Fixture"] = 
"protocol-conversion/openai-to-anthropic-stream.sse",
+                },
+                body = [[{
+                    "model": "claude-3-5-sonnet-20241022",
+                    "messages": [{"role": "user", "content": "Hi"}],
+                    "stream": true
+                }]],
+            })
+            if not res then
+                ngx.status = 500
+                ngx.say(err)
+                return
+            end
+
+            local results = {}
+            while true do
+                local chunk = res.body_reader()
+                if not chunk then break end
+                table.insert(results, chunk)
+            end
+            ngx.print(table.concat(results, ""))
+        }
+    }
+--- error_code: 200
+--- response_body_like eval
+qr/event: message_stop/
+--- grep_error_log eval
+qr/comprehend text: [^,]+/
+--- grep_error_log_out
+comprehend text: Hello world

Reply via email to