Skip to content
Merged
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
36 changes: 33 additions & 3 deletions src/flb_router_config.c
Original file line number Diff line number Diff line change
Expand Up @@ -1080,9 +1080,36 @@ static int parse_routes_block(struct cfl_variant *variant,
return -1;
}

static struct flb_input_instance *find_input_instance_by_index(struct flb_config *config,
size_t input_index)
{
size_t index;
struct mk_list *head;
struct flb_input_instance *ins;

if (!config) {
return NULL;
}

index = 0;

mk_list_foreach(head, &config->inputs) {
ins = mk_list_entry(head, struct flb_input_instance, _head);

if (index == input_index) {
return ins;
}

index++;
}

return NULL;
}
Comment thread
coderabbitai[bot] marked this conversation as resolved.

static int parse_input_section(struct flb_cf_section *section,
struct cfl_list *input_routes,
struct flb_config *config)
struct flb_config *config,
size_t input_index)
{
uint32_t mask;
size_t before_count;
Expand Down Expand Up @@ -1136,7 +1163,7 @@ static int parse_input_section(struct flb_cf_section *section,
cfl_list_init(&input->processors);
cfl_list_init(&input->routes);
input->has_alias = FLB_FALSE;
input->instance = NULL;
input->instance = find_input_instance_by_index(config, input_index);

input->plugin_name = copy_from_cfl_sds(name_var->data.as_string);
if (!input->plugin_name) {
Expand Down Expand Up @@ -1209,6 +1236,7 @@ int flb_router_config_parse(struct flb_cf *cf,
{
struct mk_list *head;
struct flb_cf_section *section;
size_t input_index;
int routes_found = FLB_FALSE;
int ret;

Expand All @@ -1218,9 +1246,10 @@ int flb_router_config_parse(struct flb_cf *cf,

cfl_list_init(input_routes);

input_index = 0;
mk_list_foreach(head, &cf->inputs) {
section = mk_list_entry(head, struct flb_cf_section, _head_section);
ret = parse_input_section(section, input_routes, config);
ret = parse_input_section(section, input_routes, config, input_index);
if (ret == -1) {
flb_router_routes_destroy(input_routes);
cfl_list_init(input_routes);
Expand All @@ -1229,6 +1258,7 @@ int flb_router_config_parse(struct flb_cf *cf,
else if (ret == 1) {
routes_found = FLB_TRUE;
}
input_index++;
}

if (cfl_list_is_empty(input_routes) == 1) {
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,47 @@
service:
flush: 1
grace: 1
log_level: trace
http_server: on
http_port: ${FLUENT_BIT_HTTP_MONITORING_PORT}

pipeline:
inputs:
- name: dummy
alias: dummy1
tag: firstdummy
dummy: '{"message":"custom dummy one"}'
samples: 1

- name: dummy
alias: dummy2
tag: seconddummy
dummy: '{"message":"custom dummy two"}'
samples: 1

- name: dummy
tag: somethingdiff
dummy: '{"message":"custom dummy"}'
samples: 1
processors:
logs:
- name: content_modifier
action: insert
key: topic
value: topic1
routes:
logs:
- name: topic1
condition:
op: and
rules:
- field: $topic
op: eq
value: topic1
to:
outputs:
- stdout1

outputs:
- name: stdout
alias: stdout1
Original file line number Diff line number Diff line change
@@ -0,0 +1,48 @@
service:
flush: 1
grace: 1
log_level: trace
http_server: on
http_port: ${FLUENT_BIT_HTTP_MONITORING_PORT}

pipeline:
inputs:
- name: dummy
alias: dummy1
tag: firstdummy
dummy: '{"message":"custom dummy one"}'
samples: 1

- name: dummy
alias: dummy2
tag: seconddummy
dummy: '{"message":"custom dummy two"}'
samples: 1

- name: dummy
alias: routed_dummy
tag: somethingdiff
dummy: '{"message":"custom dummy"}'
samples: 1
processors:
logs:
- name: content_modifier
action: insert
key: topic
value: topic1
routes:
logs:
- name: topic1
condition:
op: and
rules:
- field: $topic
op: eq
value: topic1
to:
outputs:
- stdout1

outputs:
- name: stdout
alias: stdout1
Original file line number Diff line number Diff line change
Expand Up @@ -206,3 +206,31 @@ def test_out_stdout_traces_accepts_otlp_json_ingestion():

assert "checkout-span" in log_text
assert "trace-scope" in log_text


def test_out_stdout_routes_without_alias_bind_to_third_input():
service = Service("out_stdout_routing_no_alias.yaml")
service.start()
service.wait_for_log_contains("[0] topic1:", timeout=10)
service.wait_for_log_contains("no matching route for input chunk", timeout=10)
log_text = service.wait_for_log_contains("tag 'firstdummy'", timeout=10)
service.stop()

assert "connected input 'dummy.2' route 'topic1' to output 'stdout1'" in log_text
assert "[0] topic1:" in log_text
assert "\"topic\"=>\"topic1\"" in log_text
assert "\"message\"=>\"custom dummy\"" in log_text
assert "no matching route for input chunk" in log_text
assert "tag 'firstdummy'" in log_text


def test_out_stdout_routes_with_alias_bind_to_third_input():
service = Service("out_stdout_routing_with_alias.yaml")
service.start()
log_text = service.wait_for_log_contains("[0] topic1:", timeout=10)
service.stop()

assert "connected input 'routed_dummy' route 'topic1' to output 'stdout1'" in log_text
assert "[0] topic1:" in log_text
assert "\"topic\"=>\"topic1\"" in log_text
assert "\"message\"=>\"custom dummy\"" in log_text
Loading
Loading