Coverage Report

Created: 2026-09-28 07:36

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
/src/fluent-bit/plugins/in_elasticsearch/in_elasticsearch.c
Line
Count
Source
1
/* -*- Mode: C; tab-width: 4; indent-tabs-mode: nil; c-basic-offset: 4 -*- */
2
3
/*  Fluent Bit
4
 *  ==========
5
 *  Copyright (C) 2015-2026 The Fluent Bit Authors
6
 *
7
 *  Licensed under the Apache License, Version 2.0 (the "License");
8
 *  you may not use this file except in compliance with the License.
9
 *  You may obtain a copy of the License at
10
 *
11
 *      http://www.apache.org/licenses/LICENSE-2.0
12
 *
13
 *  Unless required by applicable law or agreed to in writing, software
14
 *  distributed under the License is distributed on an "AS IS" BASIS,
15
 *  WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
16
 *  See the License for the specific language governing permissions and
17
 *  limitations under the License.
18
 */
19
20
21
#include <fluent-bit/flb_input_plugin.h>
22
#include <fluent-bit/flb_config.h>
23
#include <fluent-bit/flb_random.h>
24
25
#include "in_elasticsearch.h"
26
#include "in_elasticsearch_config.h"
27
#include "in_elasticsearch_bulk_prot.h"
28
29
0
static void bytes_to_groupname(unsigned char *data, char *buf, size_t len) {
30
0
    int index;
31
0
    char charset[] = "0123456789"
32
0
                     "abcdefghijklmnopqrstuvwxyz"
33
0
                     "ABCDEFGHIJKLMNOPQRSTUVWXYZ";
34
35
0
    while (len-- > 0) {
36
0
        index = (int) data[len];
37
0
        index = index % (sizeof(charset) - 1);
38
0
        buf[len] = charset[index];
39
0
    }
40
0
}
41
42
0
static void bytes_to_nodename(unsigned char *data, char *buf, size_t len) {
43
0
    int index;
44
0
    char charset[] = "0123456789"
45
0
                     "abcdefghijklmnopqrstuvwxyz";
46
47
0
    while (len-- > 0) {
48
0
        index = (int) data[len];
49
0
        index = index % (sizeof(charset) - 1);
50
0
        buf[len] = charset[index];
51
0
    }
52
0
}
53
54
static int in_elasticsearch_bulk_init(struct flb_input_instance *ins,
55
                                      struct flb_config *config, void *data)
56
0
{
57
0
    int ret;
58
0
    struct flb_in_elasticsearch *ctx;
59
0
    unsigned char rand[16];
60
0
    struct flb_http_server_options http_server_options;
61
62
0
    (void) config;
63
0
    (void) data;
64
65
    /* Create context and basic conf */
66
0
    ctx = in_elasticsearch_config_create(ins);
67
0
    if (!ctx) {
68
0
        return -1;
69
0
    }
70
71
    /* Populate context with config map defaults and incoming properties */
72
0
    ret = flb_input_config_map_set(ins, (void *) ctx);
73
0
    if (ret == -1) {
74
0
        flb_plg_error(ctx->ins, "configuration error");
75
0
        in_elasticsearch_config_destroy(ctx);
76
0
        return -1;
77
0
    }
78
79
    /* Set the context */
80
0
    flb_input_set_context(ins, ctx);
81
82
0
    if (flb_random_bytes(rand, 16)) {
83
0
        flb_plg_error(ctx->ins, "cannot generate cluster name");
84
0
        in_elasticsearch_config_destroy(ctx);
85
0
        return -1;
86
0
    }
87
88
0
    bytes_to_groupname(rand, ctx->cluster_name, 16);
89
90
0
    if (flb_random_bytes(rand, 12)) {
91
0
        flb_plg_error(ctx->ins, "cannot generate node name");
92
0
        in_elasticsearch_config_destroy(ctx);
93
0
        return -1;
94
0
    }
95
96
0
    bytes_to_nodename(rand, ctx->node_name, 12);
97
98
0
    ret = flb_input_http_server_options_init(
99
0
            &http_server_options,
100
0
            ins,
101
0
            (FLB_HTTP_SERVER_FLAG_KEEPALIVE | FLB_HTTP_SERVER_FLAG_AUTO_INFLATE),
102
0
            in_elasticsearch_bulk_prot_handle_ng,
103
0
            ctx);
104
0
    if (ret == 0) {
105
0
        if (http_server_options.workers > 1) {
106
0
            ret = flb_input_ingress_enable(ins);
107
0
        }
108
0
    }
109
0
    if (ret == 0) {
110
0
        ret = flb_http_server_init_with_options(&ctx->http_server,
111
0
                                                &http_server_options);
112
113
0
        if (ret == 0) {
114
0
            ret = flb_http_server_start(&ctx->http_server);
115
0
        }
116
117
0
        if (ret == 0 && ctx->http_server.downstream != NULL) {
118
0
            ret = flb_input_downstream_set(ctx->http_server.downstream, ins);
119
0
        }
120
0
    }
121
122
0
    if (ret != 0) {
123
0
        flb_plg_error(ctx->ins,
124
0
                      "could not initialize http server on %s:%u. Aborting",
125
0
                      ins->host.listen, ins->host.port);
126
127
0
        in_elasticsearch_config_destroy(ctx);
128
129
0
        return -1;
130
0
    }
131
132
0
    flb_plg_info(ctx->ins, "listening on %s:%u with %i worker%s",
133
0
                 ins->host.listen,
134
0
                 ins->host.port,
135
0
                 ctx->http_server.workers,
136
0
                 ctx->http_server.workers == 1 ? "" : "s");
137
138
0
    return 0;
139
0
}
140
141
static int in_elasticsearch_bulk_exit(void *data, struct flb_config *config)
142
0
{
143
0
    struct flb_in_elasticsearch *ctx;
144
145
0
    (void) config;
146
147
0
    ctx = data;
148
149
0
    if (ctx != NULL) {
150
0
        in_elasticsearch_config_destroy(ctx);
151
0
    }
152
153
0
    return 0;
154
0
}
155
156
static int in_elasticsearch_bulk_pause(void *data, struct flb_config *config)
157
0
{
158
0
    struct flb_in_elasticsearch *ctx;
159
160
0
    (void) config;
161
162
0
    ctx = data;
163
0
    if (flb_http_server_pause(&ctx->http_server) != 0) {
164
0
        flb_plg_error(ctx->ins, "could not pause HTTP server");
165
0
        return -1;
166
0
    }
167
168
0
    return 0;
169
0
}
170
171
static int in_elasticsearch_bulk_resume(void *data, struct flb_config *config)
172
0
{
173
0
    struct flb_in_elasticsearch *ctx;
174
175
0
    (void) config;
176
177
0
    ctx = data;
178
0
    if (flb_http_server_resume(&ctx->http_server) != 0) {
179
0
        flb_plg_error(ctx->ins, "could not resume HTTP server");
180
0
        return -1;
181
0
    }
182
183
0
    return 0;
184
0
}
185
186
/* Configuration properties map */
187
static struct flb_config_map config_map[] = {
188
    {
189
     FLB_CONFIG_MAP_STR, "tag_key", NULL,
190
     0, FLB_TRUE, offsetof(struct flb_in_elasticsearch, tag_key),
191
     "Specify a key name for extracting as a tag"
192
    },
193
194
    {
195
     FLB_CONFIG_MAP_STR, "meta_key", "@meta",
196
     0, FLB_TRUE, offsetof(struct flb_in_elasticsearch, meta_key),
197
     "Specify a key name for meta information"
198
    },
199
200
    {
201
     FLB_CONFIG_MAP_STR, "hostname", "localhost",
202
     0, FLB_TRUE, offsetof(struct flb_in_elasticsearch, hostname),
203
     "Specify hostname or FQDN. This parameter is effective for sniffering node information."
204
    },
205
206
    {
207
     FLB_CONFIG_MAP_STR, "version", "8.0.0",
208
     0, FLB_TRUE, offsetof(struct flb_in_elasticsearch, es_version),
209
     "Specify returning Elasticsearch server version."
210
    },
211
212
    /* EOF */
213
    {0}
214
};
215
216
/* Plugin reference */
217
struct flb_input_plugin in_elasticsearch_plugin = {
218
    .name         = "elasticsearch",
219
    .description  = "HTTP Endpoints for Elasticsearch (Bulk API)",
220
    .cb_init      = in_elasticsearch_bulk_init,
221
    .cb_pre_run   = NULL,
222
    .cb_collect   = NULL,
223
    .cb_flush_buf = NULL,
224
    .cb_pause_checked = in_elasticsearch_bulk_pause,
225
    .cb_resume_checked = in_elasticsearch_bulk_resume,
226
    .cb_exit      = in_elasticsearch_bulk_exit,
227
    .config_map   = config_map,
228
    .flags        = FLB_INPUT_NET_SERVER | FLB_INPUT_HTTP_SERVER | FLB_IO_OPT_TLS
229
};