diff --git a/apisix/plugins/tencent-cloud-cls.lua b/apisix/plugins/tencent-cloud-cls.lua index 18f0cd434f71..3e90faab06ff 100644 --- a/apisix/plugins/tencent-cloud-cls.lua +++ b/apisix/plugins/tencent-cloud-cls.lua @@ -32,6 +32,7 @@ local schema = { cls_topic = { type = "string" }, scheme = { type = "string", enum = {"http", "https"}, default = "https" }, ssl_verify = { type = "boolean", default = true }, + compress_type = { type = "string", enum = {"none", "zstd"}, default = "none" }, secret_id = { type = "string" }, secret_key = { type = "string" }, sample_ratio = { @@ -145,7 +146,8 @@ function _M.log(conf, ctx) local sdk, err = cls_sdk.new( conf.scheme, conf.cls_host, conf.cls_topic, conf.secret_id, - conf.secret_key, conf.ssl_verify) + conf.secret_key, conf.ssl_verify, + conf.compress_type) if err then core.log.error("init sdk failed err:", err) return false, err diff --git a/apisix/plugins/tencent-cloud-cls/cls-sdk.lua b/apisix/plugins/tencent-cloud-cls/cls-sdk.lua index 3e0373f371f0..2eb7ca9a3083 100644 --- a/apisix/plugins/tencent-cloud-cls/cls-sdk.lua +++ b/apisix/plugins/tencent-cloud-cls/cls-sdk.lua @@ -21,6 +21,7 @@ local http = require("resty.http") local socket = require("socket") local str_util = require("resty.string") local core = require("apisix.core") +local zstd = require("apisix.utils.zstd") local core_gethostname = require("apisix.core.utils").gethostname local json = core.json local json_encode = json.encode @@ -48,6 +49,10 @@ local MAX_SINGLE_VALUE_SIZE = 1 * 1024 * 1024 local MAX_LOG_GROUP_VALUE_SIZE = 5 * 1024 * 1024 -- 5MB local cls_api_path = "/structuredlog" +-- compress type used when uploading logs, see +-- https://www.tencentcloud.com/document/product/614/16873 +local COMPRESS_TYPE_NONE = "none" +local COMPRESS_TYPE_ZSTD = "zstd" local auth_expire_time = 60 local cls_conn_timeout = 1000 local cls_read_timeout = 10000 @@ -196,13 +201,17 @@ message LogGroupList end -function _M.new(scheme, host, topic, secret_id, secret_key, ssl_verify) +function _M.new(scheme, host, topic, secret_id, secret_key, ssl_verify, compress_type) if not pb_state then local err = init_pb_state() if err then return nil, err end end + if compress_type ~= nil and compress_type ~= COMPRESS_TYPE_NONE + and compress_type ~= COMPRESS_TYPE_ZSTD then + return nil, "unsupported compress type: " .. compress_type + end local self = { scheme = scheme, host = host, @@ -210,6 +219,7 @@ function _M.new(scheme, host, topic, secret_id, secret_key, ssl_verify) secret_id = secret_id, secret_key = secret_key, ssl_verify = ssl_verify, + compress_type = compress_type or COMPRESS_TYPE_NONE, } return setmetatable(self, mt) end @@ -240,10 +250,22 @@ function _M.send_cls_request(self, pb_obj) ["Authorization"] = sign(self.secret_id, self.secret_key, cls_api_path), } - -- TODO: support lz4/zstd compress + local body = pb_data + if self.compress_type == COMPRESS_TYPE_ZSTD then + local compressed, err = zstd.compress(pb_data) + if compressed then + body = compressed + headers["x-cls-compress-type"] = COMPRESS_TYPE_ZSTD + else + -- compression is only an optimization, never drop the logs because of it + core.log.error("failed to compress the log data with zstd, " + .. "upload it uncompressed, err: ", err) + end + end + local params = { method = "POST", - body = pb_data, + body = body, headers = headers, ssl_verify = self.ssl_verify, } diff --git a/apisix/utils/zstd.lua b/apisix/utils/zstd.lua new file mode 100644 index 000000000000..3d692878943f --- /dev/null +++ b/apisix/utils/zstd.lua @@ -0,0 +1,120 @@ +-- +-- 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. +-- + +-- zstd compression implemented by FFI, which requires libzstd to be +-- installed on the host (e.g. libzstd1 on Debian/Ubuntu, libzstd on CentOS). +local ffi = require("ffi") +local pcall = pcall +local ffi_load = ffi.load +local ffi_new = ffi.new +local ffi_copy = ffi.copy +local ffi_string = ffi.string +local type = type +local tonumber = tonumber + + +local DEFAULT_COMPRESS_LEVEL = 3 -- ZSTD_CLEVEL_DEFAULT +-- libzstd may only ship the versioned soname, so try all the common names +local LIB_NAMES = { "libzstd.so.1", "zstd", "libzstd" } + +local _M = {} + + +ffi.cdef[[ + size_t ZSTD_compressBound(size_t srcSize); + size_t ZSTD_compress(void *dst, size_t dstCapacity, + const void *src, size_t srcSize, int compressionLevel); + unsigned ZSTD_isError(size_t code); + const char *ZSTD_getErrorName(size_t code); +]] + + +local libzstd +local libzstd_err +local libzstd_loaded + + +-- load libzstd lazily and only once, the result (including the failure) +-- is cached so that we never call dlopen again. +local function load_libzstd() + if libzstd_loaded then + return libzstd, libzstd_err + end + + libzstd_loaded = true + for i = 1, #LIB_NAMES do + local ok, lib = pcall(ffi_load, LIB_NAMES[i]) + if ok and lib then + -- make sure the loaded library exports the symbols we need + local ok_sym = pcall(function() + return lib.ZSTD_compressBound(1) + end) + if ok_sym then + libzstd = lib + return libzstd + end + end + end + + libzstd_err = "failed to load libzstd, please make sure it is installed" + return nil, libzstd_err +end + + +-- tells whether zstd compression is usable on the current host +function _M.available() + local lib = load_libzstd() + return lib ~= nil +end + + +-- compress the given data into a zstd frame +-- returns the compressed data, or nil and an error message when the data +-- can not be compressed (e.g. libzstd is not installed) +function _M.compress(data, level) + if type(data) ~= "string" then + return nil, "invalid data type: " .. type(data) + end + + if level == nil then + level = DEFAULT_COMPRESS_LEVEL + elseif type(level) ~= "number" then + return nil, "invalid compression level" + end + + local lib, err = load_libzstd() + if not lib then + return nil, err + end + + local src_size = #data + local capacity = tonumber(lib.ZSTD_compressBound(src_size)) + local src = ffi_new("char[?]", src_size) + local dst = ffi_new("char[?]", capacity) + ffi_copy(src, data, src_size) + + local compressed_size = tonumber(lib.ZSTD_compress(dst, capacity, src, src_size, level)) + if lib.ZSTD_isError(compressed_size) ~= 0 then + return nil, "failed to compress the data: " + .. ffi_string(lib.ZSTD_getErrorName(compressed_size)) + end + + return ffi_string(dst, compressed_size) +end + + +return _M diff --git a/docs/en/latest/plugins/tencent-cloud-cls.md b/docs/en/latest/plugins/tencent-cloud-cls.md index 6b02a4cf7cc6..02ec80e57bfc 100644 --- a/docs/en/latest/plugins/tencent-cloud-cls.md +++ b/docs/en/latest/plugins/tencent-cloud-cls.md @@ -39,6 +39,7 @@ The `tencent-cloud-cls` Plugin uses [TencentCloud CLS](https://cloud.tencent.com | cls_host | string | Yes | | | CLS API host,please refer [Uploading Structured Logs](https://www.tencentcloud.com/document/api/614/16873). | | cls_topic | string | Yes | | | topic id of CLS. | | scheme | string | No | https | ["http", "https"] | The protocol scheme to use when connecting to CLS. Defaults to `https` for secure connections. | +| compress_type | string | No | none | ["none", "zstd"] | Compression algorithm used when uploading logs. Set it to `zstd` to compress the log data with zstd before uploading it, which reduces the log write traffic. `zstd` requires libzstd to be installed in the runtime environment, otherwise the logs are uploaded uncompressed. | | secret_id | string | Yes | | | SecretId of your API key. | | secret_key | string | Yes | | | SecretKey of your API key. | | sample_ratio | number | No | 1 | [0.00001, 1] | How often to sample the requests. Setting to `1` will sample all requests. | diff --git a/docs/zh/latest/plugins/tencent-cloud-cls.md b/docs/zh/latest/plugins/tencent-cloud-cls.md index 1c8486b4a788..fcade4a529e4 100644 --- a/docs/zh/latest/plugins/tencent-cloud-cls.md +++ b/docs/zh/latest/plugins/tencent-cloud-cls.md @@ -39,6 +39,7 @@ description: API 网关 Apache APISIX tencent-cloud-cls 插件可用于将日志 | cls_host | string | 是 | | | CLS API 域名,参考[使用 API 上传日志](https://cloud.tencent.com/document/api/614/16873)。| | cls_topic | string | 是 | | | CLS 日志主题 id。 | | scheme | string | 否 | https | ["http", "https"] | 连接 CLS 时使用的协议方案。默认为 `https` 以实现安全连接。 | +| compress_type | string | 否 | none | ["none", "zstd"] | 上传日志时使用的压缩方式。设置为 `zstd` 时,日志会先使用 zstd 压缩再上传,以减少日志写入流量;`none` 表示不压缩。使用 `zstd` 需要运行环境中已安装 libzstd 库,否则将退化为不压缩上传。 | | secret_id | string | 是 | | | 云 API 密钥的 id。 | | secret_key | string | 是 | | | 云 API 密钥的 key。 | | sample_ratio | number | 否 | 1 | [0.00001, 1] | 采样的比例。设置为 `1` 时,将对所有请求进行采样。 | diff --git a/t/plugin/tencent-cloud-cls.t b/t/plugin/tencent-cloud-cls.t index 9a33fbed618c..1d78419be4a1 100644 --- a/t/plugin/tencent-cloud-cls.t +++ b/t/plugin/tencent-cloud-cls.t @@ -37,6 +37,12 @@ add_block_preprocessor(sub { local data = ngx.req.get_body_data() local headers = ngx.req.get_headers() ngx.log(ngx.WARN, "tencent-cloud-cls body: ", data) + if data and #data >= 4 then + -- the first 4 bytes of the body, used to check the payload format + ngx.log(ngx.WARN, "tencent-cloud-cls body head: ", + string.format("%02x%02x%02x%02x", data:byte(1), + data:byte(2), data:byte(3), data:byte(4))) + end for k, v in pairs(headers) do ngx.log(ngx.WARN, "tencent-cloud-cls headers: " .. k .. ":" .. v) end @@ -769,3 +775,101 @@ opentracing --- error_log Batch Processor[tencent-cloud-cls] successfully processed the entries --- wait: 0.5 + + + +=== TEST 23: schema check, unsupported compress type +--- config + location /t { + content_by_lua_block { + local plugin = require("apisix.plugins.tencent-cloud-cls") + local ok, err = plugin.check_schema({ + cls_host = "ap-guangzhou.cls.tencentyun.com", + cls_topic = "143b5d70-139b-4aec-b54e-bb97756916de", + secret_id = "secret_id", + secret_key = "secret_key", + compress_type = "lz4", + }) + if not ok then + ngx.say(err) + end + + ngx.say("done") + } + } +--- response_body +property "compress_type" validation failed: matches none of the enum values +done + + + +=== TEST 24: add plugin with zstd compress +--- 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": { + "tencent-cloud-cls": { + "scheme": "http", + "cls_host": "127.0.0.1:10420", + "cls_topic": "143b5d70-139b-4aec-b54e-bb97756916de", + "secret_id": "secret_id", + "secret_key": "secret_key", + "compress_type": "zstd", + "batch_max_size": 1, + "max_retry_count": 1, + "retry_delay": 2, + "buffer_duration": 2, + "inactive_timeout": 2 + } + }, + "upstream": { + "nodes": { + "127.0.0.1:1982": 1 + }, + "type": "roundrobin" + }, + "uri": "/opentracing" + }]] + ) + + if code >= 300 then + ngx.status = code + end + ngx.say(body) + } + } +--- response_body +passed + + + +=== TEST 25: upload log with zstd compress +--- request +GET /opentracing +--- response_body +opentracing +--- error_log eval +qr/tencent-cloud-cls body head: 28b52ffd[\s\S]*tencent-cloud-cls headers: x-cls-compress-type:zstd/ +--- wait: 0.5 +--- skip_eval +3: system("ldconfig -p 2>/dev/null | grep -q libzstd") + + + +=== TEST 26: fall back to uncompressed upload when compress fails +--- extra_init_by_lua + local zstd = require("apisix.utils.zstd") + zstd.compress = function(data) + return nil, "mock compress error" + end +--- request +GET /opentracing +--- response_body +opentracing +--- error_log eval +qr/failed to compress the log data with zstd, upload it uncompressed, err: mock compress error[\s\S]*Batch Processor\[tencent-cloud-cls\] successfully processed the entries/ +--- wait: 0.5 diff --git a/t/utils/zstd.t b/t/utils/zstd.t new file mode 100644 index 000000000000..05c56efc70e2 --- /dev/null +++ b/t/utils/zstd.t @@ -0,0 +1,117 @@ +# +# 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 (!defined $block->request) { + $block->set_value("request", "GET /t"); + } +}); + +run_tests(); + +__DATA__ + +=== TEST 1: available() reports whether libzstd is usable +--- config + location /t { + content_by_lua_block { + local zstd = require("apisix.utils.zstd") + ngx.say("type: ", type(zstd.available())) + } + } +--- response_body +type: boolean + + + +=== TEST 2: reject invalid arguments +--- config + location /t { + content_by_lua_block { + local zstd = require("apisix.utils.zstd") + + local _, err = zstd.compress(123) + ngx.say("data: ", err) + + _, err = zstd.compress("hello", "fastest") + ngx.say("level: ", err) + } + } +--- response_body +data: invalid data type: number +level: invalid compression level + + + +=== TEST 3: compress data into a zstd frame +--- config + location /t { + content_by_lua_block { + local zstd = require("apisix.utils.zstd") + + local data = string.rep("hello apisix tencent-cloud-cls ", 1000) + local compressed, err = zstd.compress(data) + if not compressed then + ngx.say("failed to compress: ", err) + return + end + + -- the zstd frame magic number is 0xFD2FB528, stored in little endian + ngx.say("magic: ", string.format("%02x%02x%02x%02x", + compressed:byte(1), compressed:byte(2), + compressed:byte(3), compressed:byte(4))) + ngx.say("smaller: ", #compressed < #data) + } + } +--- response_body +magic: 28b52ffd +smaller: true +--- skip_eval +3: system("ldconfig -p 2>/dev/null | grep -q libzstd") + + + +=== TEST 4: compress with the given level +--- config + location /t { + content_by_lua_block { + local zstd = require("apisix.utils.zstd") + + local data = string.rep("{\"key\":\"value\"},", 2000) + for _, level in ipairs({1, 3, 9}) do + local compressed, err = zstd.compress(data, level) + if not compressed then + ngx.say("failed to compress with level ", level, ": ", err) + return + end + ngx.say("level ", level, ": ", #compressed < #data) + end + } + } +--- response_body +level 1: true +level 3: true +level 9: true +--- skip_eval +3: system("ldconfig -p 2>/dev/null | grep -q libzstd") diff --git a/utils/install-dependencies.sh b/utils/install-dependencies.sh index cc901e4f6e14..6ba8443e990b 100755 --- a/utils/install-dependencies.sh +++ b/utils/install-dependencies.sh @@ -56,7 +56,7 @@ function install_dependencies_with_yum() { sudo yum install -y \ gcc gcc-c++ curl wget unzip xz gnupg perl-ExtUtils-Embed cpanminus patch libyaml-devel \ perl perl-devel pcre pcre-devel pcre2 pcre2-devel openldap-devel \ - openresty-zlib-devel openresty-pcre-devel libxml2-devel libxslt-devel zlib-devel + openresty-zlib-devel openresty-pcre-devel libxml2-devel libxslt-devel zlib-devel libzstd-devel } # Install dependencies on ubuntu and debian @@ -81,7 +81,7 @@ function install_dependencies_with_apt() { sudo apt-get update # install some compilation tools - sudo apt-get install -y curl make gcc g++ cpanminus libpcre3 libpcre3-dev libpcre2-dev libyaml-dev unzip openresty-zlib-dev openresty-pcre-dev libxml2-dev libxslt-dev zlib1g-dev + sudo apt-get install -y curl make gcc g++ cpanminus libpcre3 libpcre3-dev libpcre2-dev libyaml-dev unzip openresty-zlib-dev openresty-pcre-dev libxml2-dev libxslt-dev zlib1g-dev libzstd-dev } # Identify the different distributions and call the corresponding function