Skip to content
Open
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
4 changes: 3 additions & 1 deletion apisix/plugins/tencent-cloud-cls.lua
Original file line number Diff line number Diff line change
Expand Up @@ -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 = {
Expand Down Expand Up @@ -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
Expand Down
28 changes: 25 additions & 3 deletions apisix/plugins/tencent-cloud-cls/cls-sdk.lua
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -196,20 +201,25 @@ 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,
topic = topic,
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
Expand Down Expand Up @@ -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,
}
Expand Down
120 changes: 120 additions & 0 deletions apisix/utils/zstd.lua
Original file line number Diff line number Diff line change
@@ -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
1 change: 1 addition & 0 deletions docs/en/latest/plugins/tencent-cloud-cls.md
Original file line number Diff line number Diff line change
Expand Up @@ -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. |
Expand Down
1 change: 1 addition & 0 deletions docs/zh/latest/plugins/tencent-cloud-cls.md
Original file line number Diff line number Diff line change
Expand Up @@ -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` 时,将对所有请求进行采样。 |
Expand Down
104 changes: 104 additions & 0 deletions t/plugin/tencent-cloud-cls.t
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Loading
Loading