Coverage Report

Created: 2026-09-28 07:36

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
/src/fluent-bit/src/flb_notification.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
#include <monkey/mk_core.h>
21
#include <fluent-bit/flb_info.h>
22
#include <fluent-bit/flb_plugin.h>
23
#include <fluent-bit/flb_input.h>
24
#include <fluent-bit/flb_filter.h>
25
#include <fluent-bit/flb_output.h>
26
#include <fluent-bit/flb_engine.h>
27
#include <fluent-bit/flb_notification.h>
28
29
union generic_plugin_instance {
30
    struct flb_input_instance     *input;
31
    struct flb_output_instance    *output;
32
    struct flb_filter_instance    *filter;
33
    struct flb_processor_instance *processor;
34
    void                          *generic;
35
};
36
37
static struct flb_input_instance *find_input_instance(
38
                                    char *name,
39
                                    struct flb_config *config)
40
0
{
41
0
    struct mk_list              *iterator;
42
0
    struct flb_input_instance   *instance;
43
44
0
    mk_list_foreach(iterator, &config->inputs) {
45
0
        instance = mk_list_entry(iterator,
46
0
                                 struct flb_input_instance,
47
0
                                 _head);
48
49
0
        if (strcmp(flb_input_name(instance), name) == 0) {
50
0
            return instance;
51
0
        }
52
0
    }
53
54
0
    return NULL;
55
0
}
56
57
static struct flb_output_instance *find_output_instance(
58
                                    char *name,
59
                                    struct flb_config *config)
60
0
{
61
0
    struct mk_list              *iterator;
62
0
    struct flb_output_instance  *instance;
63
64
0
    mk_list_foreach(iterator, &config->outputs) {
65
0
        instance = mk_list_entry(iterator,
66
0
                                 struct flb_output_instance,
67
0
                                 _head);
68
69
0
        if (strcmp(flb_output_name(instance), name) == 0) {
70
0
            return instance;
71
0
        }
72
0
    }
73
74
0
    return NULL;
75
0
}
76
77
static struct flb_filter_instance *find_filter_instance(
78
                                    char *name,
79
                                    struct flb_config *config)
80
0
{
81
0
    struct mk_list              *iterator;
82
0
    struct flb_filter_instance  *instance;
83
84
0
    mk_list_foreach(iterator, &config->filters) {
85
0
        instance = mk_list_entry(iterator,
86
0
                                 struct flb_filter_instance,
87
0
                                 _head);
88
89
0
        if (strcmp(flb_filter_name(instance), name) == 0) {
90
0
            return instance;
91
0
        }
92
0
    }
93
94
0
    return NULL;
95
0
}
96
97
static void *find_processor_instance_internal_unit_level(
98
                char *name,
99
                int  *plugin_type,
100
                struct mk_list *processor_unit_list)
101
0
{
102
0
    struct mk_list                *iterator;
103
0
    struct flb_processor_unit     *processor_unit;
104
0
    struct flb_filter_instance    *filter_instance;
105
0
    struct flb_processor_instance *processor_instance;
106
107
0
    mk_list_foreach(iterator, processor_unit_list) {
108
0
        processor_unit = mk_list_entry(iterator,
109
0
                                        struct flb_processor_unit,
110
0
                                        _head);
111
112
0
        if (processor_unit->unit_type == FLB_PROCESSOR_UNIT_FILTER) {
113
0
            filter_instance = (struct flb_filter_instance *) \
114
0
                                processor_unit->ctx;
115
116
0
            if (strcmp(flb_filter_name(filter_instance), name) == 0) {
117
0
                *plugin_type = FLB_PLUGIN_FILTER;
118
119
0
                return (void *) filter_instance;
120
0
            }
121
0
        }
122
0
        else if (processor_unit->unit_type == FLB_PROCESSOR_UNIT_NATIVE) {
123
0
            processor_instance = (struct flb_processor_instance *) \
124
0
                                    processor_unit->ctx;
125
126
0
            if (strcmp(flb_processor_instance_get_name(processor_instance),
127
0
                        name) == 0) {
128
0
                *plugin_type = FLB_PLUGIN_PROCESSOR;
129
130
0
                return (void *) processor_instance;
131
0
            }
132
0
        }
133
0
    }
134
135
0
    return NULL;
136
0
}
137
138
139
static void *find_processor_instance_internal_processor_level(
140
                char *name,
141
                int  *plugin_type,
142
                struct flb_processor *processor)
143
0
{
144
0
    void *result;
145
146
0
    result = find_processor_instance_internal_unit_level(
147
0
                name,
148
0
                plugin_type,
149
0
                &processor->logs);
150
151
0
    if (result == NULL) {
152
0
        result = find_processor_instance_internal_unit_level(
153
0
                    name,
154
0
                    plugin_type,
155
0
                    &processor->metrics);
156
0
    }
157
158
0
    if (result == NULL) {
159
0
        result = find_processor_instance_internal_unit_level(
160
0
                    name,
161
0
                    plugin_type,
162
0
                    &processor->traces);
163
0
    }
164
165
0
    return result;
166
0
}
167
168
static void *find_processor_instance(
169
                char *name,
170
                int  *plugin_type,
171
                struct flb_config *config)
172
0
{
173
0
    struct flb_output_instance    *output_instance;
174
0
    struct flb_input_instance     *input_instance;
175
0
    struct mk_list                *iterator;
176
0
    void                          *result;
177
178
0
    mk_list_foreach(iterator, &config->inputs) {
179
0
        input_instance = mk_list_entry(iterator,
180
0
                                       struct flb_input_instance,
181
0
                                       _head);
182
183
0
        result = find_processor_instance_internal_processor_level(
184
0
                    name,
185
0
                    plugin_type,
186
0
                    input_instance->processor);
187
188
0
        if (result != NULL) {
189
0
            return result;
190
0
        }
191
0
    }
192
193
0
    mk_list_foreach(iterator, &config->outputs) {
194
0
        output_instance = mk_list_entry(iterator,
195
0
                                       struct flb_output_instance,
196
0
                                       _head);
197
198
0
        result = find_processor_instance_internal_processor_level(
199
0
                    name,
200
0
                    plugin_type,
201
0
                    output_instance->processor);
202
203
0
        if (result != NULL) {
204
0
            return result;
205
0
        }
206
0
    }
207
208
0
    return NULL;
209
0
}
210
211
int flb_notification_enqueue(int plugin_type,
212
                             char *instance_name,
213
                             struct flb_notification *notification,
214
                             struct flb_config *config)
215
0
{
216
0
    flb_pipefd_t                  notification_channel;
217
0
    union generic_plugin_instance plugin_instance;
218
0
    int                           result;
219
220
0
    plugin_instance.generic = NULL;
221
222
0
    if (plugin_instance.generic == NULL &&
223
0
        (plugin_type == FLB_PLUGIN_INPUT ||
224
0
         plugin_type == -1)) {
225
0
        plugin_instance.input = find_input_instance(instance_name, config);
226
0
        notification_channel = plugin_instance.input->notification_channel;
227
0
        plugin_type = FLB_PLUGIN_INPUT;
228
0
    }
229
230
0
    if (plugin_instance.generic == NULL &&
231
0
        (plugin_type == FLB_PLUGIN_OUTPUT ||
232
0
         plugin_type == -1)) {
233
0
        plugin_instance.output = find_output_instance(instance_name, config);
234
0
        notification_channel = plugin_instance.output->notification_channel;
235
0
        plugin_type = FLB_PLUGIN_OUTPUT;
236
0
    }
237
238
0
    if (plugin_instance.generic == NULL &&
239
0
        (plugin_type == FLB_PLUGIN_FILTER ||
240
0
         plugin_type == -1)) {
241
0
        plugin_instance.filter = find_filter_instance(instance_name, config);
242
0
        notification_channel = plugin_instance.filter->notification_channel;
243
0
        plugin_type = FLB_PLUGIN_FILTER;
244
0
    }
245
246
0
    if (plugin_instance.generic == NULL &&
247
0
        (plugin_type == FLB_PLUGIN_FILTER ||
248
0
         plugin_type == -1)) {
249
0
        plugin_instance.generic = find_processor_instance(instance_name,
250
0
                                                          &plugin_type,
251
0
                                                          config);
252
253
0
        if (plugin_instance.generic != NULL) {
254
0
            if (plugin_type == FLB_PLUGIN_FILTER) {
255
0
                notification_channel = plugin_instance.filter->notification_channel;
256
0
            }
257
0
            else if (plugin_type == FLB_PLUGIN_PROCESSOR) {
258
0
                notification_channel = plugin_instance.processor->notification_channel;
259
0
            }
260
0
        }
261
0
    }
262
263
0
    if (plugin_instance.generic == NULL) {
264
0
        flb_error("cannot enqueue notification for plugin \"%s\" with type %d",
265
0
                  instance_name, plugin_type);
266
267
0
        return -1;
268
0
    }
269
270
0
    notification->plugin_type = plugin_type;
271
0
    notification->plugin_instance = plugin_instance.generic;
272
273
0
    result = flb_pipe_w(notification_channel,
274
0
                        &notification,
275
0
                        sizeof(void *));
276
277
0
    if (result == -1) {
278
0
        flb_pipe_error();
279
280
0
        return -1;
281
0
    }
282
283
0
    return 0;
284
0
}
285
286
int flb_notification_receive(flb_pipefd_t channel,
287
                             struct flb_notification **notification)
288
0
{
289
0
    int result;
290
291
0
    result = flb_pipe_r(channel, notification, sizeof(struct flb_notification *));
292
293
0
    if (result <= 0) {
294
0
        flb_pipe_error();
295
0
        return -1;;
296
0
    }
297
298
0
    return 0;
299
0
}
300
301
int flb_notification_deliver(struct flb_notification *notification)
302
0
{
303
0
    int result;
304
0
    union generic_plugin_instance instance;
305
306
0
    if (notification == NULL) {
307
0
        flb_error("cannot deliver NULL notification instance");
308
309
0
        return -1;
310
0
    }
311
312
0
    instance.generic = notification->plugin_instance;
313
0
    result = -2;
314
315
0
    switch(notification->plugin_type) {
316
0
    case FLB_PLUGIN_INPUT:
317
0
        if (instance.input->p->cb_notification != NULL) {
318
0
            result = instance.input->p->cb_notification(
319
0
                                            instance.input->context,
320
0
                                            instance.input->config,
321
0
                                            (void *) notification);
322
0
        }
323
0
        else {
324
0
            result = -3;
325
0
        }
326
327
0
        break;
328
329
0
    case FLB_PLUGIN_OUTPUT:
330
0
        if (instance.output->p->cb_notification != NULL) {
331
0
            result = instance.output->p->cb_notification(
332
0
                                            instance.output->context,
333
0
                                            instance.output->config,
334
0
                                            (void *) notification);
335
0
        }
336
0
        else {
337
0
            result = -3;
338
0
        }
339
340
0
        break;
341
342
0
    case FLB_PLUGIN_FILTER:
343
0
        if (instance.filter->p->cb_notification != NULL) {
344
0
            result = instance.filter->p->cb_notification(
345
0
                                            instance.filter->context,
346
0
                                            instance.filter->config,
347
0
                                            (void *) notification);
348
0
        }
349
0
        else {
350
0
            result = -3;
351
0
        }
352
353
0
        break;
354
355
0
    case FLB_PLUGIN_PROCESSOR:
356
0
        if (instance.processor->p->cb_notification != NULL) {
357
0
            result = instance.processor->p->cb_notification(
358
0
                                                instance.processor->context,
359
0
                                                instance.processor->config,
360
0
                                                (void *) notification);
361
0
        }
362
0
        else {
363
0
            result = -3;
364
0
        }
365
366
0
        break;
367
0
    }
368
369
0
    return result;
370
0
}
371
372
void flb_notification_cleanup(struct flb_notification *notification)
373
0
{
374
0
    if (notification->destructor != NULL) {
375
0
        notification->destructor((void *) notification);
376
0
    }
377
378
0
    if (notification->dynamically_allocated == FLB_TRUE) {
379
0
        flb_free(notification);
380
0
    }
381
0
}