| # |
| # Licensed to the Apache Software Foundation (ASF) under one or more |
| # contributor license agreements. See the NOTICE file distributed with |
| # this work for additional information regarding copyright ownership. |
| # The ASF licenses this file to You 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 |
| # |
| # http://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. |
| # |
| use t::APISIX 'no_plan'; |
| |
| repeat_each(1); |
| no_long_string(); |
| no_root_location(); |
| |
| add_block_preprocessor(sub { |
| my ($block) = @_; |
| |
| if (!$block->request) { |
| $block->set_value("request", "GET /t"); |
| } |
| }); |
| |
| add_block_preprocessor(sub { |
| my ($block) = @_; |
| |
| # The plugin no longer logs the payload; reproduce the observability the |
| # tests rely on by logging each batch entry from a test-only hook. |
| my $extra_init_by_lua = <<_EOC_; |
| local bp_manager = require("apisix.utils.batch-processor-manager") |
| local core = require("apisix.core") |
| local function log_send_data(entry) |
| local data = type(entry) == "table" and core.json.encode(entry) or entry |
| core.log.info("send data to kafka: ", data) |
| end |
| local old_add = bp_manager.add_entry |
| bp_manager.add_entry = function(self, conf, entry) |
| local ok = old_add(self, conf, entry) |
| if ok then |
| log_send_data(entry) |
| end |
| return ok |
| end |
| local old_new = bp_manager.add_entry_to_new_processor |
| bp_manager.add_entry_to_new_processor = function(self, conf, entry, ctx, func) |
| local ok = old_new(self, conf, entry, ctx, func) |
| if ok then |
| log_send_data(entry) |
| end |
| return ok |
| end |
| _EOC_ |
| |
| if (!defined $block->extra_init_by_lua) { |
| $block->set_value("extra_init_by_lua", $extra_init_by_lua); |
| } |
| }); |
| |
| run_tests; |
| |
| __DATA__ |
| |
| === TEST 1: sanity |
| --- config |
| location /t { |
| content_by_lua_block { |
| local plugin = require("apisix.plugins.kafka-logger") |
| local ok, err = plugin.check_schema({ |
| kafka_topic = "test", |
| key = "key1", |
| broker_list = { |
| ["127.0.0.1"] = 3 |
| } |
| }) |
| if not ok then |
| ngx.say(err) |
| end |
| ngx.say("done") |
| } |
| } |
| --- response_body |
| done |
| |
| |
| |
| === TEST 2: missing broker list |
| --- config |
| location /t { |
| content_by_lua_block { |
| local plugin = require("apisix.plugins.kafka-logger") |
| local ok, err = plugin.check_schema({kafka_topic = "test", key= "key1"}) |
| if not ok then |
| ngx.say(err) |
| end |
| ngx.say("done") |
| } |
| } |
| --- response_body |
| value should match only one schema, but matches none |
| done |
| |
| |
| |
| === TEST 3: wrong type of string |
| --- config |
| location /t { |
| content_by_lua_block { |
| local plugin = require("apisix.plugins.kafka-logger") |
| local ok, err = plugin.check_schema({ |
| broker_list = { |
| ["127.0.0.1"] = 3000 |
| }, |
| timeout = "10", |
| kafka_topic ="test", |
| key= "key1" |
| }) |
| if not ok then |
| ngx.say(err) |
| end |
| ngx.say("done") |
| } |
| } |
| --- response_body |
| property "timeout" validation failed: wrong type: expected integer, got string |
| done |
| |
| |
| |
| === TEST 4: set route(id: 1) |
| --- 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, |
| [[{ |
| "plugins": { |
| "kafka-logger": { |
| "broker_list" : |
| { |
| "127.0.0.1":9092 |
| }, |
| "kafka_topic" : "test2", |
| "key" : "key1", |
| "timeout" : 1, |
| "batch_max_size": 1 |
| } |
| }, |
| "upstream": { |
| "nodes": { |
| "127.0.0.1:1980": 1 |
| }, |
| "type": "roundrobin" |
| }, |
| "uri": "/hello" |
| }]] |
| ) |
| if code >= 300 then |
| ngx.status = code |
| end |
| ngx.say(body) |
| } |
| } |
| --- response_body |
| passed |
| |
| |
| |
| === TEST 5: access |
| --- request |
| GET /hello |
| --- response_body |
| hello world |
| --- wait: 2 |
| --- ignore_error_log |
| |
| |
| |
| === TEST 6: error log |
| --- 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, |
| [[{ |
| "plugins": { |
| "kafka-logger": { |
| "broker_list" : |
| { |
| "127.0.0.1":9092, |
| "127.0.0.1":9093 |
| }, |
| "kafka_topic" : "test2", |
| "producer_type": "sync", |
| "key" : "key1", |
| "batch_max_size": 1 |
| } |
| }, |
| "upstream": { |
| "nodes": { |
| "127.0.0.1:1980": 1 |
| }, |
| "type": "roundrobin" |
| }, |
| "uri": "/hello" |
| }]] |
| ) |
| if code >= 300 then |
| ngx.status = code |
| end |
| ngx.say(body) |
| local http = require "resty.http" |
| local httpc = http.new() |
| local uri = "http://127.0.0.1:" .. ngx.var.server_port .. "/hello" |
| local res, err = httpc:request_uri(uri, {method = "GET"}) |
| } |
| } |
| --- error_log |
| failed to send data to Kafka topic |
| [error] |
| --- wait: 1 |
| |
| |
| |
| === TEST 7: set route(meta_format = origin, include_req_body = true) |
| --- 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, |
| [[{ |
| "plugins": { |
| "kafka-logger": { |
| "broker_list" : { |
| "127.0.0.1":9092 |
| }, |
| "kafka_topic" : "test2", |
| "key" : "key1", |
| "timeout" : 1, |
| "batch_max_size": 1, |
| "include_req_body": true, |
| "meta_format": "origin" |
| } |
| }, |
| "upstream": { |
| "nodes": { |
| "127.0.0.1:1980": 1 |
| }, |
| "type": "roundrobin" |
| }, |
| "uri": "/hello" |
| }]] |
| ) |
| if code >= 300 then |
| ngx.status = code |
| end |
| ngx.say(body) |
| } |
| } |
| --- response_body |
| passed |
| |
| |
| |
| === TEST 8: hit route, report log to kafka |
| --- request |
| GET /hello?ab=cd |
| abcdef |
| --- response_body |
| hello world |
| --- error_log |
| send data to kafka: GET /hello?ab=cd HTTP/1.1 |
| host: localhost |
| content-length: 6 |
| connection: close |
| |
| abcdef |
| --- wait: 2 |
| |
| |
| |
| === TEST 9: set route(meta_format = origin, include_req_body = false) |
| --- 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, |
| [[{ |
| "plugins": { |
| "kafka-logger": { |
| "broker_list" : { |
| "127.0.0.1":9092 |
| }, |
| "kafka_topic" : "test2", |
| "key" : "key1", |
| "timeout" : 1, |
| "batch_max_size": 1, |
| "include_req_body": false, |
| "meta_format": "origin" |
| } |
| }, |
| "upstream": { |
| "nodes": { |
| "127.0.0.1:1980": 1 |
| }, |
| "type": "roundrobin" |
| }, |
| "uri": "/hello" |
| }]] |
| ) |
| if code >= 300 then |
| ngx.status = code |
| end |
| ngx.say(body) |
| } |
| } |
| --- response_body |
| passed |
| |
| |
| |
| === TEST 10: hit route, report log to kafka |
| --- request |
| GET /hello?ab=cd |
| abcdef |
| --- response_body |
| hello world |
| --- error_log |
| send data to kafka: GET /hello?ab=cd HTTP/1.1 |
| host: localhost |
| content-length: 6 |
| connection: close |
| --- wait: 2 |
| |
| |
| |
| === TEST 11: set route(meta_format = 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, |
| [[{ |
| "plugins": { |
| "kafka-logger": { |
| "broker_list" : { |
| "127.0.0.1":9092 |
| }, |
| "kafka_topic" : "test2", |
| "key" : "key1", |
| "timeout" : 1, |
| "batch_max_size": 1, |
| "include_req_body": false |
| } |
| }, |
| "upstream": { |
| "nodes": { |
| "127.0.0.1:1980": 1 |
| }, |
| "type": "roundrobin" |
| }, |
| "uri": "/hello" |
| }]] |
| ) |
| if code >= 300 then |
| ngx.status = code |
| end |
| ngx.say(body) |
| } |
| } |
| --- response_body |
| passed |
| |
| |
| |
| === TEST 12: hit route, report log to kafka |
| --- request |
| GET /hello?ab=cd |
| abcdef |
| --- response_body |
| hello world |
| --- error_log_like eval |
| qr/send data to kafka: \{.*"upstream":"127.0.0.1:1980"/ |
| --- wait: 2 |
| |
| |
| |
| === TEST 13: set route(id: 1), missing key field |
| --- 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, |
| [[{ |
| "plugins": { |
| "kafka-logger": { |
| "broker_list" : |
| { |
| "127.0.0.1":9092 |
| }, |
| "kafka_topic" : "test2", |
| "timeout" : 1, |
| "batch_max_size": 1 |
| } |
| }, |
| "upstream": { |
| "nodes": { |
| "127.0.0.1:1980": 1 |
| }, |
| "type": "roundrobin" |
| }, |
| "uri": "/hello" |
| }]] |
| ) |
| if code >= 300 then |
| ngx.status = code |
| end |
| ngx.say(body) |
| } |
| } |
| --- response_body |
| passed |
| |
| |
| |
| === TEST 14: access, test key field is optional |
| --- request |
| GET /hello |
| --- response_body |
| hello world |
| --- wait: 2 |
| |
| |
| |
| === TEST 15: set route(meta_format = default), missing key field |
| --- 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, |
| [[{ |
| "plugins": { |
| "kafka-logger": { |
| "broker_list" : { |
| "127.0.0.1":9092 |
| }, |
| "kafka_topic" : "test2", |
| "timeout" : 1, |
| "batch_max_size": 1, |
| "include_req_body": false |
| } |
| }, |
| "upstream": { |
| "nodes": { |
| "127.0.0.1:1980": 1 |
| }, |
| "type": "roundrobin" |
| }, |
| "uri": "/hello" |
| }]] |
| ) |
| if code >= 300 then |
| ngx.status = code |
| end |
| ngx.say(body) |
| } |
| } |
| --- response_body |
| passed |
| |
| |
| |
| === TEST 16: hit route, report log to kafka |
| --- request |
| GET /hello?ab=cd |
| abcdef |
| --- response_body |
| hello world |
| --- error_log_like eval |
| qr/send data to kafka: \{.*"upstream":"127.0.0.1:1980"/ |
| --- wait: 2 |
| |
| |
| |
| === TEST 17: use the topic with 3 partitions |
| --- 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, |
| [[{ |
| "plugins": { |
| "kafka-logger": { |
| "broker_list" : { |
| "127.0.0.1": 9092 |
| }, |
| "kafka_topic" : "test3", |
| "timeout" : 1, |
| "batch_max_size": 1, |
| "include_req_body": false |
| } |
| }, |
| "upstream": { |
| "nodes": { |
| "127.0.0.1:1980": 1 |
| }, |
| "type": "roundrobin" |
| }, |
| "uri": "/hello" |
| }]] |
| ) |
| if code >= 300 then |
| ngx.status = code |
| end |
| ngx.say(body) |
| } |
| } |
| --- response_body |
| passed |
| |
| |
| |
| === TEST 18: report log to kafka by different partitions |
| --- 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, |
| [[{ |
| "plugins": { |
| "kafka-logger": { |
| "broker_list" : { |
| "127.0.0.1": 9092 |
| }, |
| "kafka_topic" : "test3", |
| "producer_type": "sync", |
| "timeout" : 1, |
| "batch_max_size": 1, |
| "include_req_body": false |
| } |
| }, |
| "upstream": { |
| "nodes": { |
| "127.0.0.1:1980": 1 |
| }, |
| "type": "roundrobin" |
| }, |
| "uri": "/hello" |
| }]] |
| ) |
| |
| t('/hello',ngx.HTTP_GET) |
| ngx.sleep(0.5) |
| t('/hello',ngx.HTTP_GET) |
| ngx.sleep(0.5) |
| t('/hello',ngx.HTTP_GET) |
| ngx.sleep(0.5) |
| } |
| } |
| --- timeout: 5s |
| --- ignore_response |
| --- error_log eval |
| [qr/partition_id: 1/, |
| qr/partition_id: 0/, |
| qr/partition_id: 2/] |
| |
| |
| |
| === TEST 19: report log to kafka by different partitions in async mode |
| --- 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, |
| [[{ |
| "plugins": { |
| "kafka-logger": { |
| "broker_list" : { |
| "127.0.0.1": 9092 |
| }, |
| "kafka_topic" : "test3", |
| "producer_type": "async", |
| "timeout" : 1, |
| "batch_max_size": 1, |
| "include_req_body": false |
| } |
| }, |
| "upstream": { |
| "nodes": { |
| "127.0.0.1:1980": 1 |
| }, |
| "type": "roundrobin" |
| }, |
| "uri": "/hello" |
| }]] |
| ) |
| t('/hello',ngx.HTTP_GET) |
| ngx.sleep(0.5) |
| t('/hello',ngx.HTTP_GET) |
| ngx.sleep(0.5) |
| t('/hello',ngx.HTTP_GET) |
| ngx.sleep(0.5) |
| } |
| } |
| --- timeout: 5s |
| --- ignore_response |
| --- error_log eval |
| [qr/partition_id: 1/, |
| qr/partition_id: 0/, |
| qr/partition_id: 2/] |
| |
| |
| |
| === TEST 20: set route with incorrect sasl_config |
| --- 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, |
| [[{ |
| "plugins":{ |
| "kafka-logger":{ |
| "brokers":[ |
| { |
| "host":"127.0.0.1", |
| "port":19094, |
| "sasl_config":{ |
| "mechanism":"PLAIN", |
| "user":"admin", |
| "password":"admin-secret2233" |
| } |
| }], |
| "kafka_topic":"test2", |
| "key":"key1", |
| "timeout":1, |
| "batch_max_size":1 |
| } |
| }, |
| "upstream":{ |
| "nodes":{ |
| "127.0.0.1:1980":1 |
| }, |
| "type":"roundrobin" |
| }, |
| "uri":"/hello" |
| }]] |
| ) |
| if code >= 300 then |
| ngx.status = code |
| end |
| ngx.say(body) |
| } |
| } |
| --- response_body |
| passed |
| |
| |
| |
| === TEST 21: hit route, failed to send data to kafka |
| --- request |
| GET /hello |
| --- response_body |
| hello world |
| --- error_log |
| failed to do PLAIN auth with 127.0.0.1:19094: Authentication failed: Invalid username or password |
| --- wait: 2 |
| |
| |
| |
| === TEST 22: set route with correct sasl_config |
| --- 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, |
| [[{ |
| "plugins":{ |
| "kafka-logger":{ |
| "brokers":[ |
| { |
| "host":"127.0.0.1", |
| "port":19094, |
| "sasl_config":{ |
| "mechanism":"PLAIN", |
| "user":"admin", |
| "password":"admin-secret" |
| } |
| }], |
| "kafka_topic":"test4", |
| "producer_type":"sync", |
| "key":"key1", |
| "timeout":1, |
| "batch_max_size":1, |
| "include_req_body": true |
| } |
| }, |
| "upstream":{ |
| "nodes":{ |
| "127.0.0.1:1980":1 |
| }, |
| "type":"roundrobin" |
| }, |
| "uri":"/hello" |
| }]] |
| ) |
| if code >= 300 then |
| ngx.status = code |
| end |
| ngx.say(body) |
| } |
| } |
| --- response_body |
| passed |
| |
| |
| |
| === TEST 23: hit route, send data to kafka successfully |
| --- request |
| POST /hello?name=qwerty |
| abcdef |
| --- response_body |
| hello world |
| --- error_log eval |
| qr/send data to kafka: \{.*"body":"abcdef"/ |
| --- no_error_log |
| [error] |
| --- wait: 2 |
| |
| |
| |
| === TEST 24: set route(batch_max_size = 2), check if prometheus is initialized properly |
| --- 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, |
| [[{ |
| "plugins": { |
| "kafka-logger": { |
| "broker_list" : |
| { |
| "127.0.0.1":9092 |
| }, |
| "kafka_topic" : "test2", |
| "key" : "key1", |
| "timeout" : 1, |
| "batch_max_size": 2 |
| } |
| }, |
| "upstream": { |
| "nodes": { |
| "127.0.0.1:1980": 1 |
| }, |
| "type": "roundrobin" |
| }, |
| "uri": "/hello" |
| }]] |
| ) |
| if code >= 300 then |
| ngx.status = code |
| end |
| ngx.say(body) |
| } |
| } |
| --- response_body |
| passed |
| |
| |
| |
| === TEST 25: access |
| --- extra_yaml_config |
| plugins: |
| - kafka-logger |
| --- request |
| GET /hello |
| --- response_body |
| hello world |
| --- wait: 2 |
| |
| |
| |
| === TEST 26: create a service with kafka-logger and three routes bound to it |
| --- config |
| location /t { |
| content_by_lua_block { |
| local t = require("lib.test_admin").test |
| local code, body = t('/apisix/admin/services/1', |
| ngx.HTTP_PUT, |
| [[{ |
| "plugins": { |
| "kafka-logger": { |
| "broker_list" : { |
| "127.0.0.1":9092 |
| }, |
| "kafka_topic" : "test2", |
| "key" : "key1", |
| "timeout" : 1, |
| "batch_max_size": 1, |
| "include_req_body": true, |
| "meta_format": "origin" |
| } |
| }, |
| "upstream": { |
| "nodes": { |
| "127.0.0.1:1980": 1 |
| }, |
| "type": "roundrobin" |
| } |
| }]] |
| ) |
| if code >= 300 then |
| ngx.say("create service failed") |
| return |
| end |
| for i = 1, 3 do |
| local code, body = t('/apisix/admin/routes/' .. i, |
| ngx.HTTP_PUT, |
| string.format([[{ |
| "uri": "/hello%d", |
| "service_id": "1" |
| }]], i) |
| ) |
| if code >= 300 then |
| ngx.say("create route failed") |
| return |
| end |
| end |
| ngx.say("passed") |
| } |
| } |
| --- response_body |
| passed |
| |
| |
| |
| === TEST 27: hit three routes, should create batch processor only once |
| --- log_level: debug |
| --- config |
| location /t { |
| content_by_lua_block { |
| local http = require "resty.http" |
| local httpc = http.new() |
| for i = 1, 3 do |
| local resp = httpc:request_uri("http://127.0.0.1:" .. ngx.var.server_port .. "/hello" .. i) |
| if not resp then |
| ngx.say("failed to request test server") |
| return |
| end |
| end |
| ngx.say("passed") |
| } |
| } |
| --- response_body |
| passed |
| --- grep_error_log eval |
| qr/creating new batch processor with config.*/ |
| --- grep_error_log_out eval |
| qr/creating new batch processor with config.*/ |
| |
| |
| |
| === TEST 28: check api_version schema: 2 is accepted, 3 is rejected |
| --- config |
| location /t { |
| content_by_lua_block { |
| local plugin = require("apisix.plugins.kafka-logger") |
| local ok, err = plugin.check_schema({ |
| broker_list = { |
| ["127.0.0.1"] = 9092 |
| }, |
| kafka_topic = "test", |
| api_version = 2 |
| }) |
| if not ok then |
| ngx.say(err) |
| end |
| |
| local ok, err = plugin.check_schema({ |
| broker_list = { |
| ["127.0.0.1"] = 9092 |
| }, |
| kafka_topic = "test", |
| api_version = 3 |
| }) |
| if not ok then |
| ngx.say(err) |
| end |
| ngx.say("done") |
| } |
| } |
| --- response_body |
| property "api_version" validation failed: matches none of the enum values |
| done |
| |
| |
| |
| === TEST 29: report log to kafka with api_version = 2, the broker should store the message timestamp |
| --- 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, |
| [[{ |
| "plugins": { |
| "kafka-logger": { |
| "broker_list" : { |
| "127.0.0.1":9092 |
| }, |
| "kafka_topic" : "test2", |
| "producer_type": "sync", |
| "timeout" : 1, |
| "batch_max_size": 1, |
| "api_version": 2 |
| } |
| }, |
| "upstream": { |
| "nodes": { |
| "127.0.0.1:1980": 1 |
| }, |
| "type": "roundrobin" |
| }, |
| "uri": "/hello" |
| }]] |
| ) |
| if code >= 300 then |
| ngx.status = code |
| ngx.say(body) |
| return |
| end |
| ngx.sleep(0.5) |
| |
| -- record the current end offset of the topic before sending the log |
| local bconsumer = require("resty.kafka.basic-consumer") |
| local pconsumer = require("resty.kafka.protocol.consumer") |
| local broker_list = {{host = "127.0.0.1", port = 9092}} |
| local consumer = bconsumer:new(broker_list, {}) |
| local offset, err = consumer:list_offset("test2", 0, |
| pconsumer.LIST_OFFSET_TIMESTAMP_LAST) |
| if not offset then |
| ngx.say("failed to list offset: ", err) |
| return |
| end |
| offset = tonumber(tostring(offset):match("^%-?%d+")) |
| |
| -- hit the route to send the log to kafka |
| t('/hello', ngx.HTTP_GET) |
| ngx.sleep(2) |
| |
| local data, err = consumer:fetch("test2", 0, offset) |
| if not data then |
| ngx.say("failed to fetch message: ", err) |
| return |
| end |
| local message = data.records[1] |
| if not message then |
| ngx.say("no message fetched") |
| return |
| end |
| if tonumber(message.timestamp) > 0 then |
| ngx.say("message timestamp is stored") |
| else |
| ngx.say("invalid message timestamp: ", tostring(message.timestamp)) |
| end |
| } |
| } |
| --- timeout: 10 |
| --- response_body |
| message timestamp is stored |