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 906d883e6 feat(ai-aws-content-moderation): moderate LLM responses, 
including streams (#13735)
906d883e6 is described below

commit 906d883e6ef865f238d1911efca71ad5e8931d75
Author: Shreemaan Abhishek <[email protected]>
AuthorDate: Thu Jul 30 19:19:19 2026 +0800

    feat(ai-aws-content-moderation): moderate LLM responses, including streams 
(#13735)
---
 apisix/plugins/ai-aws-content-moderation.lua       | 172 +++++++-
 .../en/latest/plugins/ai-aws-content-moderation.md |  14 +-
 t/fixtures/aws/chat-safe.json                      |  14 +
 t/fixtures/aws/chat-streaming-harmful.sse          |  12 +
 t/fixtures/aws/chat-with-harmful.json              |  14 +
 t/plugin/ai-aws-content-moderation.t               | 450 ++++++++++++++++++++-
 6 files changed, 665 insertions(+), 11 deletions(-)

diff --git a/apisix/plugins/ai-aws-content-moderation.lua 
b/apisix/plugins/ai-aws-content-moderation.lua
index d6efd830b..c40fe84b4 100644
--- a/apisix/plugins/ai-aws-content-moderation.lua
+++ b/apisix/plugins/ai-aws-content-moderation.lua
@@ -19,12 +19,14 @@ require("resty.aws.config") -- to read env vars before 
initing aws module
 local core      = require("apisix.core")
 local protocols = require("apisix.plugins.ai-protocols")
 local binding   = require("apisix.plugins.ai-protocols.binding")
+local sse       = require("apisix.plugins.ai-transport.sse")
 local aws       = require("resty.aws")
 local aws_instance
 
 local http = require("resty.http")
 
 local ngx     = ngx
+local ngx_ok  = ngx.OK
 local pairs   = pairs
 local unpack  = unpack
 local type    = type
@@ -73,6 +75,27 @@ local schema = {
             default = 0.5
         },
         check_request = { type = "boolean", default = true },
+        check_response = { type = "boolean", default = false },
+        stream_check_mode = {
+            type = "string",
+            enum = { "realtime", "final_packet" },
+            default = "final_packet",
+            description = "realtime: moderate batches while the response 
streams, replacing " ..
+                          "the rest of the stream on a hit | final_packet: 
moderate the " ..
+                          "assembled response and annotate the last chunk with 
risk_level.",
+        },
+        stream_check_cache_size = {
+            type = "integer",
+            minimum = 1,
+            default = 128,
+            description = "max characters per moderation batch in realtime 
mode",
+        },
+        stream_check_interval = {
+            type = "number",
+            minimum = 0.1,
+            default = 3,
+            description = "seconds between batch checks in realtime mode",
+        },
         deny_code = {
             type = "integer",
             minimum = 200,
@@ -103,10 +126,23 @@ function _M.check_schema(conf)
 end
 
 
+-- Comprehend scores content, it doesn't grade it, so the verdict is binary
+-- against the configured thresholds. Report it through the same ctx var the
+-- aliyun plugin uses so logging and downstream consumers stay uniform.
+local function set_risk_level(ctx, flagged)
+    if flagged then
+        ctx.var.llm_content_risk_level = "high"
+    elseif ctx.var.llm_content_risk_level ~= "high" then
+        ctx.var.llm_content_risk_level = "none"
+    end
+end
+
+
 -- Score content with AWS Comprehend detectToxicContent.
+-- `subject` names the moderated text in the deny reason ("request"/"response" 
body).
 -- Returns (reason, nil) when a category/toxicity threshold is exceeded,
 -- (nil, err) on a service error, and (nil, nil) when the content is clean.
-local function detect_toxic(conf, content)
+local function detect_toxic(conf, ctx, content, subject)
     local comprehend = conf.comprehend
 
     if not aws_instance then
@@ -151,15 +187,19 @@ local function detect_toxic(conf, content)
             for _, item in pairs(result.Labels) do
                 local threshold = conf.moderation_categories[item.Name]
                 if threshold and item.Score > threshold then
-                    return "request body exceeds " .. item.Name .. " threshold"
+                    set_risk_level(ctx, true)
+                    return subject .. " exceeds " .. item.Name .. " threshold"
                 end
             end
         end
 
         if result.Toxicity > conf.moderation_threshold then
-            return "request body exceeds toxicity threshold"
+            set_risk_level(ctx, true)
+            return subject .. " exceeds toxicity threshold"
         end
     end
+
+    set_risk_level(ctx, false)
 end
 
 
@@ -183,6 +223,66 @@ local function build_deny_message(ctx, conf, reason)
 end
 
 
+-- Moderate a piece of LLM response text.
+-- Returns (deny_code, deny_body) on a hit, (nil, nil, err) on a Comprehend
+-- failure, and nothing when the content is clean. The buffered caller fails
+-- closed on err (bytes not sent yet, same as the request side); streaming
+-- callers can't and let the content through.
+local function moderate_response(ctx, conf, content)
+    if not content or content == "" then
+        -- nothing to score, but keep the risk_level contract satisfied so the
+        -- streamed annotation still reports a verdict
+        set_risk_level(ctx, false)
+        return
+    end
+
+    local reason, err = detect_toxic(conf, ctx, content, "response body")
+    if err then
+        core.log.error(err)
+        return nil, nil, err
+    end
+    if reason then
+        return conf.deny_code, build_deny_message(ctx, conf, reason)
+    end
+end
+
+
+-- Annotate a streamed chunk with the verdict from the assembled response.
+-- The content already reached the client, so all we can do is tag it.
+local function annotate_stream(ctx, body)
+    local proto = protocols.get(ctx.ai_client_protocol)
+    if not proto or not proto.is_data_event then
+        return
+    end
+
+    local events = sse.decode(body)
+    local raw_events = {}
+    local contains_done_event = false
+    for _, event in ipairs(events) do
+        if proto.is_data_event(event) then
+            local data, err = core.json.decode(event.data)
+            -- a scalar or cjson.null decode can't be indexed; only annotate
+            -- well-formed JSON object frames and leave anything else untouched
+            if type(data) == "table" then
+                data.risk_level = ctx.var.llm_content_risk_level
+                event.data = core.json.encode(data)
+            else
+                core.log.warn("failed to decode SSE data as object: ", err or 
event.data)
+            end
+        end
+        if proto.is_done_event and proto.is_done_event(event) then
+            contains_done_event = true
+        end
+        table.insert(raw_events, sse.encode(event))
+    end
+
+    if not contains_done_event and proto.build_done_event and 
ctx.var.llm_request_done then
+        table.insert(raw_events, proto.build_done_event())
+    end
+    return table.concat(raw_events)
+end
+
+
 function _M.access(conf, ctx)
     if not ctx.picked_ai_instance then
         local handled, code, body = binding.on_unsupported(
@@ -241,7 +341,7 @@ function _M.access(conf, ctx)
         return
     end
 
-    local reason, err = detect_toxic(conf, content)
+    local reason, err = detect_toxic(conf, ctx, content, "request body")
     if err then
         core.log.error(err)
         return HTTP_INTERNAL_SERVER_ERROR, err
@@ -257,4 +357,68 @@ function _M.access(conf, ctx)
     end
 end
 
+
+function _M.lua_body_filter(conf, ctx, headers, body)
+    if not conf.check_response then
+        core.log.info("skip response check for this request")
+        return
+    end
+
+    if ngx.status >= 400 then
+        core.log.info("skip response check because upstream returned error 
status: ", ngx.status)
+        return
+    end
+
+    local request_type = ctx.var.request_type
+
+    -- ai-proxy hands us the fully assembled completion, so one check covers 
it.
+    if request_type == "ai_chat" then
+        local code, body, err = moderate_response(ctx, conf, 
ctx.var.llm_response_text)
+        if err then
+            -- the buffered body has not reached the client yet, so fail closed
+            -- like the request side instead of shipping unmoderated content
+            return HTTP_INTERNAL_SERVER_ERROR, err
+        end
+        return code, body
+    end
+
+    if request_type ~= "ai_stream" then
+        return
+    end
+
+    if conf.stream_check_mode == "final_packet" then
+        -- llm_response_text only appears once the stream is assembled, so
+        -- earlier chunks pass through untouched.
+        if not ctx.var.llm_response_text then
+            return
+        end
+        if not ctx.aws_cm_response_moderated then
+            ctx.aws_cm_response_moderated = true
+            moderate_response(ctx, conf, ctx.var.llm_response_text)
+        end
+        return nil, annotate_stream(ctx, body)
+    end
+
+    -- 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 
{}, "")
+    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
+            and now - ctx.aws_cm_last_check < conf.stream_check_interval
+            and not ctx.var.llm_request_done then
+        return
+    end
+
+    ctx.aws_cm_last_check = now
+    -- headers are already sent, so the deny body replaces the rest of the
+    -- stream rather than changing the status code
+    local _, message = moderate_response(ctx, conf, ctx.aws_cm_cache)
+    if message then
+        return ngx_ok, message
+    end
+    ctx.aws_cm_cache = ""
+end
+
 return _M
diff --git a/docs/en/latest/plugins/ai-aws-content-moderation.md 
b/docs/en/latest/plugins/ai-aws-content-moderation.md
index 063911ddd..dcddf60ff 100644
--- a/docs/en/latest/plugins/ai-aws-content-moderation.md
+++ b/docs/en/latest/plugins/ai-aws-content-moderation.md
@@ -36,10 +36,14 @@ import TabItem from '@theme/TabItem';
 
 ## Description
 
-The `ai-aws-content-moderation` Plugin integrates with [AWS 
Comprehend](https://aws.amazon.com/comprehend/) to check request content for 
toxicity when proxying to LLMs, such as profanity, hate speech, insult, 
harassment, violence, and more, rejecting requests if the evaluated outcome 
exceeds the configured threshold.
+The `ai-aws-content-moderation` Plugin integrates with [AWS 
Comprehend](https://aws.amazon.com/comprehend/) to check content for toxicity 
when proxying to LLMs, such as profanity, hate speech, insult, harassment, 
violence, and more, rejecting requests if the evaluated outcome exceeds the 
configured threshold.
 
 The Plugin is protocol-aware: it extracts the prompt content from the LLM 
request (for example `messages[].content`) and moderates only that decoded 
text, rather than the raw request body.
 
+Both directions can be moderated. Set `check_response` to moderate the LLM 
response as well. For streaming responses, `stream_check_mode` selects between 
`realtime`, which moderates batches as they arrive and replaces the remainder 
of the stream once a batch is flagged, and `final_packet`, which moderates the 
assembled response and annotates the last chunk with `risk_level`. The verdict 
is also exposed on the request context as `$llm_content_risk_level` (`high` or 
`none`) for logging.
+
+If AWS Comprehend cannot be reached, the request and the buffered 
(non-streaming) response both fail closed with a `500`, so unmoderated content 
is never proxied. Streaming response moderation is best-effort: once the first 
bytes have been sent to the client the response cannot be blocked, so a 
Comprehend failure there lets the remaining stream through.
+
 The `ai-aws-content-moderation` Plugin should be used with either 
[`ai-proxy`](./ai-proxy.md) or [`ai-proxy-multi`](./ai-proxy-multi.md) Plugin 
for proxying LLM requests.
 
 ## Plugin Attributes
@@ -55,8 +59,12 @@ The `ai-aws-content-moderation` Plugin should be used with 
either [`ai-proxy`](.
 | `moderation_categories` | object | False | | | Key-value pairs of moderation 
category and their corresponding threshold. In each pair, the key should be one 
of `PROFANITY`, `HATE_SPEECH`, `INSULT`, `HARASSMENT_OR_ABUSE`, `SEXUAL`, or 
`VIOLENCE_OR_THREAT`; and the threshold value should be between 0 and 1 
(inclusive). |
 | `moderation_threshold` | number | False | 0.5 | 0 - 1 | Overall toxicity 
threshold. A higher value means more toxic content allowed. This option differs 
from the individual category thresholds in `moderation_categories`. For 
example, if `moderation_categories` is set with a `PROFANITY` threshold of 
`0.5`, and a request has a `PROFANITY` score of `0.1`, the request will not 
exceed the category threshold. However, if the request has other categories 
like `SEXUAL` or `VIOLENCE_OR_THREAT`  [...]
 | `check_request` | boolean | False | `true` | | If `true`, moderate the 
request content. |
-| `deny_code` | integer | False | `200` | [200, 599] | HTTP status code 
returned when a request is rejected. Defaults to `200` so the 
provider-compatible refusal parses as a normal completion in client SDKs; set a 
4xx to surface denies as HTTP errors instead. |
-| `deny_message` | string | False | | | Message returned when a request is 
rejected. If unset, the moderation reason (for example `request body exceeds 
toxicity threshold`) is returned. |
+| `check_response` | boolean | False | `false` | | If `true`, moderate the LLM 
response content. |
+| `stream_check_mode` | string | False | `final_packet` | `realtime`, 
`final_packet` | Streaming moderation mode, used when `check_response` is 
`true`. `realtime`: moderate batches while the response streams, replacing the 
rest of the stream once a batch is flagged. `final_packet`: moderate the 
assembled response and annotate the last chunk with `risk_level`. |
+| `stream_check_cache_size` | integer | False | `128` | >= 1 | Maximum 
characters per moderation batch in `realtime` mode. |
+| `stream_check_interval` | number | False | `3` | >= 0.1 | Seconds between 
batch checks in `realtime` mode. |
+| `deny_code` | integer | False | `200` | [200, 599] | HTTP status code 
returned when a request is rejected. Defaults to `200` so the 
provider-compatible refusal parses as a normal completion in client SDKs; set a 
4xx to surface denies as HTTP errors instead. Streaming responses denied 
mid-stream keep the status already sent to the client. |
+| `deny_message` | string | False | | | Message returned when a request or 
response is rejected. If unset, the moderation reason (for example `request 
body exceeds toxicity threshold`) is returned. |
 | `fail_mode` | string | False | `skip` | `skip`, `warn`, `error` | Behavior 
when the request did not pass through `ai-proxy`/`ai-proxy-multi` and therefore 
cannot be moderated as an AI request. `skip`: let the request pass through 
unchecked; `warn`: pass through and log a warning; `error`: reject the request. 
|
 
 ## Examples
diff --git a/t/fixtures/aws/chat-safe.json b/t/fixtures/aws/chat-safe.json
new file mode 100644
index 000000000..a972b931b
--- /dev/null
+++ b/t/fixtures/aws/chat-safe.json
@@ -0,0 +1,14 @@
+{
+  "choices": [
+    {
+      "finish_reason": "stop",
+      "index": 0,
+      "message": { "content": "Hello there! How can I assist you today?", 
"role": "assistant" }
+    }
+  ],
+  "created": 1723780938,
+  "id": "chatcmpl-9wiSIg5LYrrpxwsr2PubSQnbtod1Q",
+  "model": "gpt-3.5-turbo",
+  "object": "chat.completion",
+  "usage": { "completion_tokens": 9, "prompt_tokens": 8, "total_tokens": 17 }
+}
diff --git a/t/fixtures/aws/chat-streaming-harmful.sse 
b/t/fixtures/aws/chat-streaming-harmful.sse
new file mode 100644
index 000000000..118645570
--- /dev/null
+++ b/t/fixtures/aws/chat-streaming-harmful.sse
@@ -0,0 +1,12 @@
+data: 
{"id":"chatcmpl-aws1","object":"chat.completion.chunk","created":1700000000,"model":"gpt-4o-mini","choices":[{"index":0,"delta":{"role":"assistant","content":""},"finish_reason":null}]}
+
+data: 
{"id":"chatcmpl-aws1","object":"chat.completion.chunk","created":1700000000,"model":"gpt-4o-mini","choices":[{"index":0,"delta":{"content":"I
 want to "},"finish_reason":null}]}
+
+data: 
{"id":"chatcmpl-aws1","object":"chat.completion.chunk","created":1700000000,"model":"gpt-4o-mini","choices":[{"index":0,"delta":{"content":"kill
 you "},"finish_reason":null}]}
+
+data: 
{"id":"chatcmpl-aws1","object":"chat.completion.chunk","created":1700000000,"model":"gpt-4o-mini","choices":[{"index":0,"delta":{"content":"right
 now!"},"finish_reason":null}]}
+
+data: 
{"id":"chatcmpl-aws1","object":"chat.completion.chunk","created":1700000000,"model":"gpt-4o-mini","choices":[{"index":0,"delta":{},"finish_reason":"stop"}],"usage":{"prompt_tokens":10,"completion_tokens":8,"total_tokens":18}}
+
+data: [DONE]
+
diff --git a/t/fixtures/aws/chat-with-harmful.json 
b/t/fixtures/aws/chat-with-harmful.json
new file mode 100644
index 000000000..13540cd7c
--- /dev/null
+++ b/t/fixtures/aws/chat-with-harmful.json
@@ -0,0 +1,14 @@
+{
+  "choices": [
+    {
+      "finish_reason": "stop",
+      "index": 0,
+      "message": { "content": "I will kill you.", "role": "assistant" }
+    }
+  ],
+  "created": 1723780938,
+  "id": "chatcmpl-9wiSIg5LYrrpxwsr2PubSQnbtod1P",
+  "model": "gpt-3.5-turbo",
+  "object": "chat.completion",
+  "usage": { "completion_tokens": 5, "prompt_tokens": 8, "total_tokens": 10 }
+}
diff --git a/t/plugin/ai-aws-content-moderation.t 
b/t/plugin/ai-aws-content-moderation.t
index f9561356a..01872c510 100644
--- a/t/plugin/ai-aws-content-moderation.t
+++ b/t/plugin/ai-aws-content-moderation.t
@@ -81,11 +81,18 @@ _EOC_
                         ngx.say("[INTERNAL FAILURE]: failed to decoded request 
body: ", err)
                     end
                     local result = body.TextSegments[1].Text
-                    local final_response = responses[result] or "invalid"
-
-                    if final_response == "invalid" then
-                        ngx.status = 500
+                    local final_response = responses[result]
+
+                    -- Response-side text is free-form LLM output, not a 
fixture
+                    -- key, so fall back to flagging anything violent.
+                    if not final_response then
+                        if result:find("kill", 1, true) then
+                            final_response = responses["toxic"]
+                        else
+                            final_response = responses["good_request"]
+                        end
                     end
+
                     ngx.say(json.encode(final_response))
                 }
             }
@@ -429,3 +436,438 @@ no ai instance picked, ai-aws-content-moderation plugin 
must be used with ai-pro
 199: rejected
 600: rejected
 403: accepted
+
+
+
+=== TEST 18: set route with check_response enabled (non-streaming)
+--- 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": "/chat",
+                    "plugins": {
+                        "ai-proxy": {
+                            "provider": "openai",
+                            "auth": { "header": { "Authorization": "Bearer 
token" } },
+                            "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,
+                            "deny_code": 400
+                        }
+                    }
+                }]]
+            )
+
+            if code >= 300 then
+                ngx.status = code
+            end
+            ngx.say(body)
+        }
+    }
+--- response_body
+passed
+
+
+
+=== TEST 19: toxic LLM response is denied
+--- request
+POST /chat
+{ "messages": [ { "role": "user", "content": "good_request" } ] }
+--- more_headers
+X-AI-Fixture: aws/chat-with-harmful.json
+--- error_code: 400
+--- response_body_like eval
+qr/response body exceeds toxicity threshold/
+
+
+
+=== TEST 20: clean LLM response passes through
+--- request
+POST /chat
+{ "messages": [ { "role": "user", "content": "good_request" } ] }
+--- more_headers
+X-AI-Fixture: aws/chat-safe.json
+--- error_code: 200
+--- response_body_like eval
+qr/How can I assist you today/
+
+
+
+=== TEST 21: set route with default deny_code and a custom deny_message
+--- 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": "/chat",
+                    "plugins": {
+                        "ai-proxy": {
+                            "provider": "openai",
+                            "auth": { "header": { "Authorization": "Bearer 
token" } },
+                            "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,
+                            "deny_message": "the response was withheld"
+                        }
+                    }
+                }]]
+            )
+
+            if code >= 300 then
+                ngx.status = code
+            end
+            ngx.say(body)
+        }
+    }
+--- response_body
+passed
+
+
+
+=== TEST 22: response deny is a provider-compatible completion the client can 
parse
+--- request
+POST /chat
+{ "messages": [ { "role": "user", "content": "good_request" } ] }
+--- more_headers
+X-AI-Fixture: aws/chat-with-harmful.json
+--- error_code: 200
+--- response_body_like eval
+qr/(?=.*"object":"chat\.completion")(?=.*the response was withheld)/s
+
+
+
+=== TEST 23: set route with check_response disabled (default)
+--- 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": "/chat",
+                    "plugins": {
+                        "ai-proxy": {
+                            "provider": "openai",
+                            "auth": { "header": { "Authorization": "Bearer 
token" } },
+                            "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
+                        }
+                    }
+                }]]
+            )
+
+            if code >= 300 then
+                ngx.status = code
+            end
+            ngx.say(body)
+        }
+    }
+--- response_body
+passed
+
+
+
+=== TEST 24: toxic LLM response passes when check_response is off
+--- request
+POST /chat
+{ "messages": [ { "role": "user", "content": "good_request" } ] }
+--- more_headers
+X-AI-Fixture: aws/chat-with-harmful.json
+--- error_code: 200
+--- response_body_like eval
+qr/I will kill you/
+
+
+
+=== TEST 25: set route with stream = true (SSE) and stream_check_mode = 
final_packet
+--- 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": "/chat",
+                    "plugins": {
+                        "ai-proxy": {
+                            "provider": "openai",
+                            "auth": { "header": { "Authorization": "Bearer 
token" } },
+                            "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": "final_packet"
+                        }
+                    }
+                }]]
+            )
+
+            if code >= 300 then
+                ngx.status = code
+            end
+            ngx.say(body)
+        }
+    }
+--- response_body
+passed
+
+
+
+=== TEST 26: final_packet annotates the streamed response with the verdict
+--- request
+POST /chat
+{ "messages": [ { "role": "user", "content": "good_request" } ], "stream": 
true }
+--- more_headers
+X-AI-Fixture: aws/chat-streaming-harmful.sse
+X-AI-Fixture-Flush-Events: true
+--- error_code: 200
+--- response_body_like eval
+qr/"risk_level":"high"/
+
+
+
+=== TEST 27: set route with stream_check_mode = realtime and a small batch size
+--- 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": "/chat",
+                    "plugins": {
+                        "ai-proxy": {
+                            "provider": "openai",
+                            "auth": { "header": { "Authorization": "Bearer 
token" } },
+                            "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,
+                            "deny_message": "the response was withheld",
+                            "stream_check_mode": "realtime",
+                            "stream_check_cache_size": 5
+                        }
+                    }
+                }]]
+            )
+
+            if code >= 300 then
+                ngx.status = code
+            end
+            ngx.say(body)
+        }
+    }
+--- response_body
+passed
+
+
+
+=== TEST 28: realtime cuts the stream off as soon as a batch is toxic
+--- request
+POST /chat
+{ "messages": [ { "role": "user", "content": "good_request" } ], "stream": 
true }
+--- more_headers
+X-AI-Fixture: aws/chat-streaming-harmful.sse
+X-AI-Fixture-Flush-Events: true
+--- error_code: 200
+--- response_body_like eval
+qr/the response was withheld/
+--- response_body_unlike eval
+qr/right now!/
+
+
+
+=== TEST 29: response moderation is skipped when upstream returns an error 
status
+--- config
+    location /t {
+        content_by_lua_block {
+            local plugin = require("apisix.plugins.ai-aws-content-moderation")
+            local ctx = {
+                picked_ai_instance = { provider = "openai" },
+                var = { request_type = "ai_stream" },
+            }
+            local conf = {
+                comprehend = {
+                    access_key_id = "access",
+                    secret_access_key = "secret",
+                    region = "us-east-1"
+                },
+                check_response = true,
+                stream_check_mode = "realtime",
+                stream_check_cache_size = 128,
+                stream_check_interval = 3,
+            }
+            ngx.status = 400
+            local code, msg = plugin.lua_body_filter(conf, ctx, {}, "body")
+            ngx.status = 200
+            ngx.say("code:", code or "nil", ", msg:", msg or "nil")
+        }
+    }
+--- response_body
+code:nil, msg:nil
+--- error_log
+skip response check because upstream returned error status: 400
+
+
+
+=== TEST 30: schema check: streaming knobs are validated
+--- config
+    location /t {
+        content_by_lua_block {
+            local plugin = require("apisix.plugins.ai-aws-content-moderation")
+            local function conf(extra)
+                local c = {
+                    comprehend = {
+                        access_key_id = "a",
+                        secret_access_key = "s",
+                        region = "us-east-1"
+                    }
+                }
+                for k, v in pairs(extra or {}) do
+                    c[k] = v
+                end
+                return c
+            end
+
+            ngx.say("bad mode: ",
+                plugin.check_schema(conf({stream_check_mode = "eventually"}))
+                    and "accepted" or "rejected")
+            ngx.say("zero cache size: ",
+                plugin.check_schema(conf({stream_check_cache_size = 0}))
+                    and "accepted" or "rejected")
+            ngx.say("tiny interval: ",
+                plugin.check_schema(conf({stream_check_interval = 0.01}))
+                    and "accepted" or "rejected")
+
+            local c = conf()
+            plugin.check_schema(c)
+            ngx.say("defaults: ", c.check_response and "on" or "off", " ",
+                    c.stream_check_mode, " ", c.stream_check_cache_size, " ",
+                    c.stream_check_interval)
+        }
+    }
+--- response_body
+bad mode: rejected
+zero cache size: rejected
+tiny interval: rejected
+defaults: off final_packet 128 3
+
+
+
+=== TEST 31: final_packet annotates an empty streamed response with risk_level 
none
+--- 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 = "openai-chat",
+                var = {
+                    request_type = "ai_stream",
+                    llm_response_text = "",
+                    llm_request_done = true,
+                },
+            }
+            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 = "final_packet",
+            }
+            plugin.check_schema(conf)
+            local body = 'data: {"choices":[{"delta":{"content":""}}]}\n\n'
+            local code, new_body = plugin.lua_body_filter(conf, ctx, {}, body)
+            ngx.say("code:", code or "nil")
+            ngx.say("risk_level:", ctx.var.llm_content_risk_level or "nil")
+            ngx.say(new_body)
+        }
+    }
+--- response_body_like eval
+qr/code:nil\nrisk_level:none\n.*"risk_level":"none"/s
+
+
+
+=== TEST 32: final_packet leaves a non-object SSE data frame untouched (no 
crash)
+--- 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 = "openai-chat",
+                var = {
+                    request_type = "ai_stream",
+                    llm_response_text = "good_request",
+                    llm_request_done = true,
+                },
+            }
+            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 = "final_packet",
+            }
+            plugin.check_schema(conf)
+            -- first frame is a bare scalar (not indexable), second is a valid 
object
+            local body = 'data: 12345\n\n'
+                      .. 'data: {"choices":[{"delta":{"content":"hi"}}]}\n\n'
+            local code, new_body = plugin.lua_body_filter(conf, ctx, {}, body)
+            ngx.say("code:", code or "nil")
+            ngx.say(new_body)
+        }
+    }
+--- response_body_like eval
+qr/code:nil\n.*data: 12345.*"risk_level":"none"/s
+--- no_error_log
+[error]

Reply via email to