From 7dff44a626ab66ed6ef01c85c57bb9597a5d5d86 Mon Sep 17 00:00:00 2001 From: Stas Date: Mon, 14 Sep 2026 22:10:48 +0200 Subject: [PATCH 1/3] build: register Manticore output plugin Signed-off-by: Stas --- CMakeLists.txt | 1 + cmake/plugins_options.cmake | 1 + cmake/windows-setup.cmake | 1 + plugins/CMakeLists.txt | 1 + 4 files changed, 4 insertions(+) diff --git a/CMakeLists.txt b/CMakeLists.txt index b5639277495..dbe900d44c4 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -377,6 +377,7 @@ if(FLB_ALL) set(FLB_OUT_GCS 1) set(FLB_OUT_GELF 1) set(FLB_OUT_HTTP 1) + set(FLB_OUT_MANTICORE 1) set(FLB_OUT_NATS 1) set(FLB_OUT_NULL 1) set(FLB_OUT_PLOT 1) diff --git a/cmake/plugins_options.cmake b/cmake/plugins_options.cmake index e6e12e1aee9..75f30ffd608 100644 --- a/cmake/plugins_options.cmake +++ b/cmake/plugins_options.cmake @@ -119,6 +119,7 @@ DEFINE_OPTION(FLB_OUT_CLOUDWATCH_LOGS "Enable AWS CloudWatch output plug DEFINE_OPTION(FLB_OUT_COUNTER "Enable Counter output plugin" ON) DEFINE_OPTION(FLB_OUT_DATADOG "Enable DataDog output plugin" ON) DEFINE_OPTION(FLB_OUT_ES "Enable Elasticsearch output plugin" ON) +DEFINE_OPTION(FLB_OUT_MANTICORE "Enable Manticore Search output plugin" ON) DEFINE_OPTION(FLB_OUT_EXIT "Enable Exit output plugin" ON) DEFINE_OPTION(FLB_OUT_FILE "Enable file output plugin" ON) DEFINE_OPTION(FLB_OUT_FLOWCOUNTER "Enable flowcount output plugin" ON) diff --git a/cmake/windows-setup.cmake b/cmake/windows-setup.cmake index b3306266e8f..2536c0a767c 100644 --- a/cmake/windows-setup.cmake +++ b/cmake/windows-setup.cmake @@ -95,6 +95,7 @@ if(FLB_WINDOWS_DEFAULTS) set(FLB_OUT_CHRONICLE Yes) set(FLB_OUT_DATADOG Yes) set(FLB_OUT_ES Yes) + set(FLB_OUT_MANTICORE Yes) set(FLB_OUT_EXIT Yes) set(FLB_OUT_FORWARD Yes) set(FLB_OUT_GELF Yes) diff --git a/plugins/CMakeLists.txt b/plugins/CMakeLists.txt index af43a8fc2e4..85e72ce3b4c 100644 --- a/plugins/CMakeLists.txt +++ b/plugins/CMakeLists.txt @@ -372,6 +372,7 @@ REGISTER_OUT_PLUGIN("out_file") REGISTER_OUT_PLUGIN("out_forward") REGISTER_OUT_PLUGIN("out_http") REGISTER_OUT_PLUGIN("out_influxdb") +REGISTER_OUT_PLUGIN("out_manticore") REGISTER_OUT_PLUGIN("out_logdna") REGISTER_OUT_PLUGIN("out_loki") From 991d7c3e17fe6ab1786e4a37f63b3e018d732b5a Mon Sep 17 00:00:00 2001 From: Stas Date: Mon, 14 Sep 2026 22:25:40 +0200 Subject: [PATCH 2/3] out_manticore: add Manticore bulk import output Signed-off-by: Stas --- conf/out_manticore.conf | 22 + plugins/out_manticore/CMakeLists.txt | 5 + plugins/out_manticore/README.md | 146 +++ plugins/out_manticore/manticore.c | 1818 ++++++++++++++++++++++++++ plugins/out_manticore/manticore.h | 59 + 5 files changed, 2050 insertions(+) create mode 100644 conf/out_manticore.conf create mode 100644 plugins/out_manticore/CMakeLists.txt create mode 100644 plugins/out_manticore/README.md create mode 100644 plugins/out_manticore/manticore.c create mode 100644 plugins/out_manticore/manticore.h diff --git a/conf/out_manticore.conf b/conf/out_manticore.conf new file mode 100644 index 00000000000..fd24014bbd2 --- /dev/null +++ b/conf/out_manticore.conf @@ -0,0 +1,22 @@ +[SERVICE] + Flush 1 + +[INPUT] + Name dummy + Tag manticore + Dummy {"id":1,"message":"hello from Fluent Bit","status":200} + +[OUTPUT] + Name manticore + Match manticore + Host 127.0.0.1 + Port 9308 + Table logs + Action insert + Id_Key id + Stream_Chunk_Size 64K + + # For finite imports that must produce exactly one Manticore disk chunk: + # Workers 1 + # Single_Chunk On + # Spool_Path /var/lib/fluent-bit/manticore-logs.ndjson diff --git a/plugins/out_manticore/CMakeLists.txt b/plugins/out_manticore/CMakeLists.txt new file mode 100644 index 00000000000..4c1533f1f1b --- /dev/null +++ b/plugins/out_manticore/CMakeLists.txt @@ -0,0 +1,5 @@ +set(src + manticore.c + ) + +FLB_PLUGIN(out_manticore "${src}" "mk_core") diff --git a/plugins/out_manticore/README.md b/plugins/out_manticore/README.md new file mode 100644 index 00000000000..969f87e4aee --- /dev/null +++ b/plugins/out_manticore/README.md @@ -0,0 +1,146 @@ +# Manticore Search output + +The `manticore` output sends log records to Manticore Search's direct-to-disk +bulk import endpoint, `/bulk?bulk_import=`, as newline-delimited JSON +over HTTP/1.1 chunked transfer encoding. + +The plugin formats each Fluent Bit record as one Manticore operation: + +```json +{"insert":{"table":"logs","id":42,"doc":{"message":"hello","status":200}}} +``` + +It converts records incrementally and buffers at most `stream_chunk_size` bytes +before writing an HTTP chunk. A record larger than that limit is sent as its own +chunk. The plugin never builds the complete NDJSON request body in memory. + +## Requirements + +The target table must exist before Fluent Bit sends data. Manticore's native +`/bulk` endpoint does not create tables automatically. + +Each request is imported directly into one disk chunk and published at request +EOF. The `table` option is sent both in `bulk_import=
` and in every +operation. Empty NDJSON lines are never generated. In the default mode, every +Fluent Bit flush is one request and therefore one disk chunk. The plugin closes +the HTTP connection after each request to release Manticore's bulk import +reservation and unblock ordinary writes. For this reason, `net.keepalive` is +disabled for this output. + +Every record must contain a stable, non-zero numeric ID in `id_key`. Decimal +strings are accepted and normalized to JSON numbers. Missing, zero, negative, +non-decimal, overflowing, or duplicate IDs are rejected during preflight before +the HTTP connection is opened. Preflight retains one 64-bit ID per record to +verify uniqueness; it does not buffer the encoded NDJSON body. + +## Configuration + +```ini +[OUTPUT] + Name manticore + Match * + Host manticore + Port 9308 + Table logs + Action insert + Id_Key id + Stream_Chunk_Size 64K +``` + +### One chunk for a finite import + +Fluent Bit normally splits a large input into several internal chunks. To +publish all of them as one Manticore disk chunk, enable `single_chunk` and use a +dedicated durable spool path: + +```ini +[INPUT] + Name tail + Path /data/import.ndjson + Read_From_Head On + Exit_On_Eof On + Parser json + +[OUTPUT] + Name manticore + Match * + Host manticore + Port 9308 + Table logs + Action insert + Workers 1 + Single_Chunk On + Spool_Path /var/lib/fluent-bit/manticore-logs.ndjson +``` + +In this mode, `FLB_OK` means that the callback is durably staged locally; it +does not mean Manticore has published it yet. Flush callbacks serialize records +to the spool, `fsync` the data, and then `fsync` the callback's committed byte +offset and checksum to a companion `.commit` journal. On recovery, +committed-data corruption stops startup without a network request; incomplete +data or journal tails are truncated to the last complete, checksummed commit +record. Retained serialized data is replayed using the current output +configuration. When Fluent Bit shuts down after input EOF, the plugin first +requires zero running tasks and zero pending storage chunks, then streams the +committed spool in one HTTP request. Both files are removed only after Manticore +acknowledges the request. A failed final upload keeps both files and is reported +in the log, but does not change Fluent Bit's process exit status: shutdown +success means that Fluent Bit stopped, not that Manticore published the import. +The next run replays the committed spool before accepting new input. + +Use this mode only for finite imports with an explicit process completion +boundary such as `Exit_On_Eof On`. A continuously running Fluent Bit process +does not publish the session until shutdown. `single_chunk` requires `insert`, +a writable `spool_path` dedicated to one output instance, and one output worker. +It is not currently supported on Windows. A permanent record error aborts the +whole session without contacting Manticore and removes the partial spool; rerun +the finite source after correcting the record. + +Session-wide duplicate detection is memory-resident and capped by +`max_session_ids` (default 1,048,576 IDs). Raise that explicit limit for larger +imports; the hash table uses up to roughly 16 bytes per allowed ID. + +Recovery is at-least-once. If the process loses power after Manticore commits the +request but before the local commit journal is cleared, the next run can replay +the same stable IDs. Rows remain correct because `insert` replaces those IDs, +but the replay can create one additional Manticore disk chunk. + +TLS uses the standard Fluent Bit output options: + +```ini + TLS On + TLS.Verify On +``` + +HTTP Basic authentication is available through `HTTP_User` and `HTTP_Passwd`. + +| Option | Description | Default | +|---|---|---| +| `table` | Existing target Manticore table. Required. | none | +| `action` | Direct-to-disk `/bulk` action: `insert` or `create`. | `insert` | +| `id_key` | Required top-level non-zero numeric ID, moved to `id` and removed from `doc`. | `id` | +| `single_chunk` | Stage all Fluent Bit flushes and publish one request at shutdown. | `false` | +| `spool_path` | Exclusive durable spool file required when `single_chunk` is enabled. | none | +| `max_session_ids` | Maximum IDs retained for session-wide duplicate detection. | `1M` | +| `stream_chunk_size` | Maximum NDJSON bytes buffered before an HTTP chunk is written. A single larger record is sent separately. | `64K` | +| `buffer_size` | Maximum buffer used to read the Manticore response. | `64K` | +| `http_user` | HTTP Basic authentication user. | none | +| `http_passwd` | HTTP Basic authentication password. | empty | + +## Response and retry behavior + +- HTTP `2xx` with `"errors": false`: the Fluent Bit chunk is acknowledged. +- A response with `"errors": true` and any item status `408`, `429`, or `5xx`: + Fluent Bit retries the whole chunk. +- HTTP `408` or `429`, transport failures, and HTTP `5xx` without a readable + item status: Fluent Bit retries the whole chunk. +- A response with only permanent item statuses is rejected permanently, even + when Manticore reports the request itself as HTTP `5xx`. Other HTTP `4xx` + responses are also permanent. + +A mixed `/bulk` response can contain successful and transiently failed batches. +Fluent Bit can only retry its original chunk, so successful batches are +replayed. Bulk import publication replaces rows with matching IDs already in +the table, making replay safe when every record has a stable ID. Within one +batch, duplicate numeric IDs are invalid; the plugin rejects them before +delivery. diff --git a/plugins/out_manticore/manticore.c b/plugins/out_manticore/manticore.c new file mode 100644 index 00000000000..5475039b4cc --- /dev/null +++ b/plugins/out_manticore/manticore.c @@ -0,0 +1,1818 @@ +/* -*- Mode: C; tab-width: 4; indent-tabs-mode: nil; c-basic-offset: 4 -*- */ + +/* Fluent Bit + * ========== + * Copyright (C) 2015-2026 The Fluent Bit Authors + * + * Licensed 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. + */ + +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include + +#include + +#include +#include +#include +#include +#include +#include + +#ifndef FLB_SYSTEM_WINDOWS +#include +#include +#include +#include +#endif + +#include "manticore.h" + +#define MANTICORE_STREAM_OK 0 +#define MANTICORE_STREAM_RETRY 1 +#define MANTICORE_STREAM_RECORD_ERROR 2 +struct manticore_id_list { + uint64_t *values; + size_t count; + size_t capacity; +}; + +static int append(flb_sds_t *buf, const char *data, size_t len) +{ + return flb_sds_cat_safe(buf, data, len); +} + +static int is_query_component_char(unsigned char c) +{ + return (c >= 'a' && c <= 'z') || + (c >= 'A' && c <= 'Z') || + (c >= '0' && c <= '9') || + c == '-' || c == '_' || c == '.' || c == '~'; +} + +static flb_sds_t build_bulk_uri(const char *table) +{ + int ret; + size_t i; + char encoded[4]; + flb_sds_t uri; + + uri = flb_sds_create(FLB_MANTICORE_BULK_URI); + if (uri == NULL) { + return NULL; + } + + for (i = 0; table[i] != '\0'; i++) { + if (is_query_component_char((unsigned char) table[i])) { + ret = append(&uri, &table[i], 1); + } + else { + snprintf(encoded, sizeof(encoded), "%%%02X", + (unsigned char) table[i]); + ret = append(&uri, encoded, 3); + } + + if (ret != 0) { + flb_sds_destroy(uri); + return NULL; + } + } + + return uri; +} + +static char *object_to_json(const msgpack_object *obj, int escape_unicode) +{ + return flb_msgpack_to_json_str(256, obj, escape_unicode); +} + +static int key_equals(const msgpack_object *key, const char *name) +{ + size_t len; + + if (key->type != MSGPACK_OBJECT_STR || name == NULL) { + return FLB_FALSE; + } + + len = strlen(name); + if (key->via.str.size != len) { + return FLB_FALSE; + } + + return memcmp(key->via.str.ptr, name, len) == 0; +} + +static int parse_id(const msgpack_object *value, uint64_t *id) +{ + size_t i; + uint64_t current; + unsigned int digit; + + if (value->type == MSGPACK_OBJECT_POSITIVE_INTEGER) { + if (value->via.u64 == 0) { + return -1; + } + *id = value->via.u64; + return 0; + } + + if (value->type != MSGPACK_OBJECT_STR || value->via.str.size == 0) { + return -1; + } + + current = 0; + for (i = 0; i < value->via.str.size; i++) { + if (value->via.str.ptr[i] < '0' || value->via.str.ptr[i] > '9') { + return -1; + } + digit = value->via.str.ptr[i] - '0'; + if (current > (UINT64_MAX - digit) / 10) { + return -1; + } + current = current * 10 + digit; + } + + if (current == 0) { + return -1; + } + *id = current; + return 0; +} + +static int append_id(struct manticore_id_list *ids, uint64_t id) +{ + size_t capacity; + uint64_t *values; + + if (ids->count == ids->capacity) { + capacity = ids->capacity == 0 ? 64 : ids->capacity * 2; + if (capacity < ids->capacity || + capacity > SIZE_MAX / sizeof(uint64_t)) { + return -1; + } + values = flb_realloc(ids->values, capacity * sizeof(uint64_t)); + if (values == NULL) { + return -1; + } + ids->values = values; + ids->capacity = capacity; + } + + ids->values[ids->count++] = id; + return 0; +} + +static int compare_ids(const void *left, const void *right) +{ + uint64_t a; + uint64_t b; + + a = *(const uint64_t *) left; + b = *(const uint64_t *) right; + return (a > b) - (a < b); +} + +static int validate_record(struct flb_out_manticore *ctx, + const msgpack_object *body, uint64_t *id) +{ + int i; + int id_found; + const msgpack_object_kv *entry; + + if (body == NULL || body->type != MSGPACK_OBJECT_MAP) { + flb_plg_error(ctx->ins, "log record body must be a map"); + return -1; + } + + entry = body->via.map.ptr; + id_found = FLB_FALSE; + for (i = 0; i < body->via.map.size; i++) { + if (entry[i].key.type != MSGPACK_OBJECT_STR) { + flb_plg_error(ctx->ins, "record keys must be strings"); + return -1; + } + + if (!key_equals(&entry[i].key, ctx->id_key)) { + continue; + } + + if (id_found || parse_id(&entry[i].val, id) != 0) { + flb_plg_error(ctx->ins, + "record key '%s' must be a unique, non-zero numeric ID", + ctx->id_key); + return -1; + } + id_found = FLB_TRUE; + } + + if (!id_found) { + flb_plg_error(ctx->ins, + "record key '%s' must contain a non-zero numeric ID", + ctx->id_key); + return -1; + } + + return 0; +} + +static int validate_events(struct flb_out_manticore *ctx, + const void *data, size_t bytes, + struct manticore_id_list *collected_ids) +{ + int ret; + size_t i; + uint64_t id; + struct flb_log_event event; + struct flb_log_event_decoder decoder; + struct manticore_id_list ids = {0}; + + ret = flb_log_event_decoder_init(&decoder, (char *) data, bytes); + if (ret != FLB_EVENT_DECODER_SUCCESS) { + flb_plg_error(ctx->ins, "could not initialize log event decoder: %s", + flb_log_event_decoder_get_error_description(ret)); + return MANTICORE_STREAM_RETRY; + } + + while (flb_log_event_decoder_next(&decoder, &event) == + FLB_EVENT_DECODER_SUCCESS) { + if (validate_record(ctx, event.body, &id) != 0) { + flb_free(ids.values); + flb_log_event_decoder_destroy(&decoder); + return MANTICORE_STREAM_RECORD_ERROR; + } + if (append_id(&ids, id) != 0) { + flb_plg_error(ctx->ins, "could not allocate document ID preflight"); + flb_free(ids.values); + flb_log_event_decoder_destroy(&decoder); + return MANTICORE_STREAM_RETRY; + } + } + + ret = flb_log_event_decoder_get_last_result(&decoder); + if (ret != FLB_EVENT_DECODER_SUCCESS) { + flb_plg_error(ctx->ins, "could not decode log event: %s", + flb_log_event_decoder_get_error_description(ret)); + } + else if (ids.count == 0) { + flb_plg_error(ctx->ins, "bulk_import requires at least one record"); + ret = MANTICORE_STREAM_RECORD_ERROR; + } + else if (ids.count > 1) { + qsort(ids.values, ids.count, sizeof(uint64_t), compare_ids); + for (i = 1; i < ids.count; i++) { + if (ids.values[i - 1] == ids.values[i]) { + flb_plg_error(ctx->ins, + "record key '%s' must be unique within a chunk", + ctx->id_key); + ret = MANTICORE_STREAM_RECORD_ERROR; + break; + } + } + } + + flb_log_event_decoder_destroy(&decoder); + if (ret != FLB_EVENT_DECODER_SUCCESS) { + flb_free(ids.values); + return MANTICORE_STREAM_RECORD_ERROR; + } + + if (collected_ids != NULL) { + *collected_ids = ids; + } + else { + flb_free(ids.values); + } + return MANTICORE_STREAM_OK; +} + +static flb_sds_t format_record(struct flb_out_manticore *ctx, + const msgpack_object *body) +{ + int ret; + int i; + int id_len; + int fields; + uint64_t numeric_id; + char id_json[32]; + char *key_json; + char *value_json; + flb_sds_t out; + const msgpack_object *id; + const msgpack_object_kv *entry; + + if (body == NULL || body->type != MSGPACK_OBJECT_MAP) { + return NULL; + } + + id = NULL; + fields = 0; + entry = body->via.map.ptr; + + for (i = 0; i < body->via.map.size; i++) { + if (key_equals(&entry[i].key, ctx->id_key)) { + id = &entry[i].val; + } + else { + fields++; + } + } + + if (id == NULL || parse_id(id, &numeric_id) != 0) { + return NULL; + } + id_len = snprintf(id_json, sizeof(id_json), "%" PRIu64, numeric_id); + if (id_len <= 0 || id_len >= sizeof(id_json)) { + return NULL; + } + + out = flb_sds_create_size(512); + if (out == NULL) { + return NULL; + } + + ret = append(&out, "{\"", sizeof("{\"") - 1); + ret |= append(&out, ctx->bulk_action, strlen(ctx->bulk_action)); + ret |= append(&out, "\":{\"table\":", sizeof("\":{\"table\":") - 1); + ret |= append(&out, ctx->table_json, flb_sds_len(ctx->table_json)); + ret |= append(&out, ",\"id\":", sizeof(",\"id\":") - 1); + ret |= append(&out, id_json, id_len); + ret |= append(&out, ",\"doc\":{", sizeof(",\"doc\":{") - 1); + + fields = 0; + for (i = 0; i < body->via.map.size; i++) { + if (key_equals(&entry[i].key, ctx->id_key)) { + continue; + } + + key_json = object_to_json(&entry[i].key, + ctx->config->json_escape_unicode); + value_json = object_to_json(&entry[i].val, + ctx->config->json_escape_unicode); + if (key_json == NULL || value_json == NULL) { + flb_free(key_json); + flb_free(value_json); + flb_sds_destroy(out); + return NULL; + } + + if (fields++ > 0) { + ret |= append(&out, ",", sizeof(",") - 1); + } + ret |= append(&out, key_json, strlen(key_json)); + ret |= append(&out, ":", sizeof(":") - 1); + ret |= append(&out, value_json, strlen(value_json)); + flb_free(key_json); + flb_free(value_json); + + if (ret != 0) { + flb_sds_destroy(out); + return NULL; + } + } + + ret |= append(&out, "}}}\n", sizeof("}}}\n") - 1); + if (ret != 0) { + flb_sds_destroy(out); + return NULL; + } + + return out; +} + +static int write_all(struct flb_connection *connection, + const void *data, size_t length) +{ + int ret; + size_t written; + + written = 0; + ret = flb_io_net_write(connection, data, length, &written); + if (ret == -1 || written != length) { + return -1; + } + + return 0; +} + +static int write_chunk(struct flb_connection *connection, + const void *data, size_t length) +{ + int len; + char header[32]; + + len = snprintf(header, sizeof(header), "%zx\r\n", length); + if (len <= 0 || len >= sizeof(header)) { + return -1; + } + + if (write_all(connection, header, len) != 0 || + write_all(connection, data, length) != 0 || + write_all(connection, "\r\n", 2) != 0) { + return -1; + } + + return 0; +} + +static void inspect_item_status(const msgpack_object *item, + int *has_status, int *retryable) +{ + int i; + msgpack_object action; + msgpack_object key; + msgpack_object value; + + if (item->type != MSGPACK_OBJECT_MAP || item->via.map.size != 1) { + return; + } + + action = item->via.map.ptr[0].val; + if (action.type != MSGPACK_OBJECT_MAP) { + return; + } + + for (i = 0; i < action.via.map.size; i++) { + key = action.via.map.ptr[i].key; + value = action.via.map.ptr[i].val; + if (key.type != MSGPACK_OBJECT_STR || key.via.str.size != 6 || + memcmp(key.via.str.ptr, "status", 6) != 0) { + continue; + } + + if (value.type == MSGPACK_OBJECT_POSITIVE_INTEGER) { + *has_status = FLB_TRUE; + if (value.via.u64 == 408 || value.via.u64 == 429 || + value.via.u64 >= 500) { + *retryable = FLB_TRUE; + } + return; + } + } +} + +static void inspect_response_items(const msgpack_object *items, + int *has_status, int *retryable) +{ + int i; + + if (items->type != MSGPACK_OBJECT_ARRAY) { + return; + } + + for (i = 0; i < items->via.array.size; i++) { + inspect_item_status(&items->via.array.ptr[i], has_status, retryable); + } +} + +static int response_ok(struct flb_out_manticore *ctx, + struct flb_http_client *client) +{ + int i; + int ret; + int root_type; + int errors; + int has_status; + int retryable; + char *packed; + size_t packed_size; + size_t offset; + msgpack_object root; + msgpack_object key; + msgpack_object value; + msgpack_unpacked result; + + packed = NULL; + packed_size = 0; + errors = -1; + has_status = FLB_FALSE; + retryable = FLB_FALSE; + if (client->resp.payload_size > 0) { + ret = flb_pack_json(client->resp.payload, client->resp.payload_size, + &packed, &packed_size, &root_type, NULL); + if (ret == 0) { + msgpack_unpacked_init(&result); + offset = 0; + ret = msgpack_unpack_next(&result, packed, packed_size, &offset); + if (ret == MSGPACK_UNPACK_SUCCESS) { + root = result.data; + if (root.type == MSGPACK_OBJECT_MAP) { + for (i = 0; i < root.via.map.size; i++) { + key = root.via.map.ptr[i].key; + value = root.via.map.ptr[i].val; + if (key.type == MSGPACK_OBJECT_STR && + key.via.str.size == 6 && + memcmp(key.via.str.ptr, "errors", 6) == 0 && + value.type == MSGPACK_OBJECT_BOOLEAN) { + errors = value.via.boolean; + } + else if (key.type == MSGPACK_OBJECT_STR && + key.via.str.size == 5 && + memcmp(key.via.str.ptr, "items", 5) == 0) { + inspect_response_items(&value, &has_status, &retryable); + } + } + } + } + msgpack_unpacked_destroy(&result); + } + flb_free(packed); + } + + if (client->resp.status >= 200 && client->resp.status < 300) { + if (errors == FLB_FALSE) { + return FLB_OK; + } + + if (errors == FLB_TRUE && retryable == FLB_TRUE) { + flb_plg_warn(ctx->ins, + "Manticore /bulk returned a retryable item error"); + return FLB_RETRY; + } + + flb_plg_error(ctx->ins, "invalid or failed Manticore /bulk response: %.*s", + (int) client->resp.payload_size, + client->resp.payload); + return FLB_ERROR; + } + + if (client->resp.payload_size > 0) { + flb_plg_error(ctx->ins, "Manticore /bulk returned HTTP %d: %.*s", + client->resp.status, + (int) client->resp.payload_size, + client->resp.payload); + } + else { + flb_plg_error(ctx->ins, "Manticore /bulk returned HTTP %d", + client->resp.status); + } + + if (retryable == FLB_TRUE || client->resp.status == 408 || + client->resp.status == 429) { + return FLB_RETRY; + } + + if (client->resp.status >= 500 && has_status == FLB_FALSE) { + return FLB_RETRY; + } + + return FLB_ERROR; +} + +static int stream_events(struct flb_out_manticore *ctx, + struct flb_connection *connection, + const void *data, size_t bytes) +{ + int ret; + flb_sds_t line; + flb_sds_t chunk; + struct flb_log_event event; + struct flb_log_event_decoder decoder; + + ret = flb_log_event_decoder_init(&decoder, (char *) data, bytes); + if (ret != FLB_EVENT_DECODER_SUCCESS) { + return MANTICORE_STREAM_RETRY; + } + + chunk = flb_sds_create_size(ctx->stream_chunk_size); + if (chunk == NULL) { + flb_log_event_decoder_destroy(&decoder); + return MANTICORE_STREAM_RETRY; + } + + while ((ret = flb_log_event_decoder_next(&decoder, &event)) == + FLB_EVENT_DECODER_SUCCESS) { + line = format_record(ctx, event.body); + if (line == NULL) { + ret = MANTICORE_STREAM_RETRY; + break; + } + + if (flb_sds_len(chunk) > 0 && + flb_sds_len(chunk) + flb_sds_len(line) > ctx->stream_chunk_size) { + if (write_chunk(connection, chunk, flb_sds_len(chunk)) != 0) { + flb_sds_destroy(line); + ret = MANTICORE_STREAM_RETRY; + break; + } + flb_sds_len_set(chunk, 0); + chunk[0] = '\0'; + } + + if (flb_sds_len(line) > ctx->stream_chunk_size) { + ret = write_chunk(connection, line, flb_sds_len(line)); + } + else { + ret = append(&chunk, line, flb_sds_len(line)); + } + flb_sds_destroy(line); + + if (ret != 0) { + ret = MANTICORE_STREAM_RETRY; + break; + } + } + + if (ret != MANTICORE_STREAM_RETRY && + ret != MANTICORE_STREAM_RECORD_ERROR) { + ret = flb_log_event_decoder_get_last_result(&decoder); + if (ret == FLB_EVENT_DECODER_SUCCESS) { + ret = MANTICORE_STREAM_OK; + } + else { + flb_plg_error(ctx->ins, "could not decode log event: %s", + flb_log_event_decoder_get_error_description(ret)); + ret = MANTICORE_STREAM_RECORD_ERROR; + } + } + + if (ret == MANTICORE_STREAM_OK && flb_sds_len(chunk) > 0) { + if (write_chunk(connection, chunk, flb_sds_len(chunk)) != 0) { + ret = MANTICORE_STREAM_RETRY; + } + } + + flb_sds_destroy(chunk); + flb_log_event_decoder_destroy(&decoder); + return ret; +} + +#ifndef FLB_SYSTEM_WINDOWS +static size_t session_id_slot(uint64_t id, size_t capacity) +{ + id ^= id >> 33; + id *= UINT64_C(0xff51afd7ed558ccd); + id ^= id >> 33; + id *= UINT64_C(0xc4ceb9fe1a85ec53); + id ^= id >> 33; + return (size_t) id & (capacity - 1); +} + +static int session_id_contains(struct flb_out_manticore *ctx, uint64_t id) +{ + size_t slot; + + if (ctx->session_id_capacity == 0) { + return FLB_FALSE; + } + slot = session_id_slot(id, ctx->session_id_capacity); + while (ctx->session_ids[slot] != 0) { + if (ctx->session_ids[slot] == id) { + return FLB_TRUE; + } + slot = (slot + 1) & (ctx->session_id_capacity - 1); + } + return FLB_FALSE; +} + +static void session_id_insert(struct flb_out_manticore *ctx, uint64_t id) +{ + size_t slot; + + slot = session_id_slot(id, ctx->session_id_capacity); + while (ctx->session_ids[slot] != 0) { + slot = (slot + 1) & (ctx->session_id_capacity - 1); + } + ctx->session_ids[slot] = id; + ctx->session_id_count++; +} + +static int prepare_session_ids(struct flb_out_manticore *ctx, + struct manticore_id_list *incoming) +{ + size_t i; + size_t required; + size_t capacity; + uint64_t *old_ids; + uint64_t *new_ids; + size_t old_capacity; + + if (incoming->count > SIZE_MAX - ctx->session_id_count) { + return -1; + } + required = incoming->count + ctx->session_id_count; + if (required > ctx->max_session_ids) { + flb_plg_error(ctx->ins, + "single-chunk session exceeds max_session_ids (%zu)", + ctx->max_session_ids); + return MANTICORE_STREAM_RECORD_ERROR; + } + capacity = ctx->session_id_capacity == 0 ? 128 : ctx->session_id_capacity; + while (required > capacity / 2) { + if (capacity > SIZE_MAX / 2) { + return -1; + } + capacity *= 2; + } + if (capacity > SIZE_MAX / sizeof(uint64_t)) { + return -1; + } + + for (i = 0; i < incoming->count; i++) { + if (session_id_contains(ctx, incoming->values[i])) { + flb_plg_error(ctx->ins, + "record key '%s' must be unique within a session", + ctx->id_key); + return MANTICORE_STREAM_RECORD_ERROR; + } + } + + if (capacity != ctx->session_id_capacity) { + new_ids = flb_calloc(capacity, sizeof(uint64_t)); + if (new_ids == NULL) { + return -1; + } + old_ids = ctx->session_ids; + old_capacity = ctx->session_id_capacity; + ctx->session_ids = new_ids; + ctx->session_id_capacity = capacity; + ctx->session_id_count = 0; + for (i = 0; i < old_capacity; i++) { + if (old_ids[i] != 0) { + session_id_insert(ctx, old_ids[i]); + } + } + flb_free(old_ids); + } + return MANTICORE_STREAM_OK; +} + +static int rollback_spool(struct flb_out_manticore *ctx, off_t offset) +{ + int failed; + + failed = FLB_FALSE; + clearerr(ctx->spool); + if (fflush(ctx->spool) != 0) { + failed = FLB_TRUE; + } + if (ftruncate(ctx->spool_fd, offset) != 0) { + failed = FLB_TRUE; + } + if (fseeko(ctx->spool, offset, SEEK_SET) != 0) { + failed = FLB_TRUE; + } + if (failed == FLB_TRUE) { + flb_plg_error(ctx->ins, "could not roll back spool '%s'", + ctx->spool_path); + return -1; + } + return 0; +} + +static int rollback_commit(struct flb_out_manticore *ctx, off_t offset) +{ + int failed; + + failed = FLB_FALSE; + clearerr(ctx->commit); + if (fflush(ctx->commit) != 0) { + failed = FLB_TRUE; + } + if (ftruncate(ctx->commit_fd, offset) != 0) { + failed = FLB_TRUE; + } + if (fseeko(ctx->commit, offset, SEEK_SET) != 0) { + failed = FLB_TRUE; + } + if (failed == FLB_TRUE) { + flb_plg_error(ctx->ins, "could not roll back commit journal '%s'", + ctx->commit_path); + return -1; + } + return 0; +} + +static uint64_t hash_bytes(uint64_t hash, const void *data, size_t length) +{ + size_t i; + const unsigned char *bytes; + + bytes = data; + for (i = 0; i < length; i++) { + hash ^= bytes[i]; + hash *= UINT64_C(1099511628211); + } + return hash; +} + +static int hash_file_range(int fd, off_t start, off_t end, uint64_t *result) +{ + ssize_t length; + off_t offset; + uint64_t hash; + char *buffer; + size_t requested; + + buffer = flb_malloc(65536); + if (buffer == NULL) { + return -1; + } + hash = UINT64_C(1469598103934665603); + offset = start; + while (offset < end) { + requested = (size_t) (end - offset); + if (requested > 65536) { + requested = 65536; + } + length = pread(fd, buffer, requested, offset); + if (length <= 0) { + flb_free(buffer); + return -1; + } + hash = hash_bytes(hash, buffer, (size_t) length); + offset += length; + } + flb_free(buffer); + *result = hash; + return 0; +} + +static int commit_spool_offset(struct flb_out_manticore *ctx, off_t start, + off_t offset, + off_t *journal_offset) +{ + uint64_t checksum; + uint64_t value; + + if (offset < 0) { + return -1; + } + *journal_offset = ftello(ctx->commit); + if (*journal_offset < 0) { + return -1; + } + value = (uint64_t) offset; + if (hash_file_range(ctx->spool_fd, start, offset, &checksum) != 0 || + fprintf(ctx->commit, + "%016" PRIx64 " %016" PRIx64 " %016" PRIx64 "\n", + value, ~value, checksum) != 51 || + fflush(ctx->commit) != 0 || fsync(ctx->commit_fd) != 0) { + if (rollback_commit(ctx, *journal_offset) != 0) { + return -2; + } + return -1; + } + return 0; +} + +static int spool_events(struct flb_out_manticore *ctx, + const void *data, size_t bytes) +{ + int ret; + off_t offset; + off_t journal_offset; + off_t committed_offset; + int commit_result; + size_t i; + + flb_sds_t line; + struct flb_log_event event; + struct flb_log_event_decoder decoder; + struct manticore_id_list ids = {0}; + + ret = validate_events(ctx, data, bytes, &ids); + if (ret != MANTICORE_STREAM_OK) { + return ret == MANTICORE_STREAM_RETRY ? FLB_RETRY : FLB_ERROR; + } + + ret = prepare_session_ids(ctx, &ids); + if (ret != MANTICORE_STREAM_OK) { + flb_free(ids.values); + return ret == MANTICORE_STREAM_RECORD_ERROR ? FLB_ERROR : FLB_RETRY; + } + + offset = ftello(ctx->spool); + if (offset < 0 || + flb_log_event_decoder_init(&decoder, (char *) data, bytes) != + FLB_EVENT_DECODER_SUCCESS) { + flb_free(ids.values); + return FLB_RETRY; + } + + ret = FLB_OK; + while (flb_log_event_decoder_next(&decoder, &event) == + FLB_EVENT_DECODER_SUCCESS) { + line = format_record(ctx, event.body); + if (line == NULL) { + ret = FLB_RETRY; + break; + } + if (fwrite(line, 1, flb_sds_len(line), ctx->spool) != + flb_sds_len(line)) { + flb_sds_destroy(line); + ret = FLB_RETRY; + break; + } + flb_sds_destroy(line); + } + flb_log_event_decoder_destroy(&decoder); + + if (ret == FLB_OK) { + committed_offset = ftello(ctx->spool); + commit_result = 0; + if (committed_offset < 0 || fflush(ctx->spool) != 0 || + fsync(ctx->spool_fd) != 0) { + ret = FLB_RETRY; + } + else { + commit_result = commit_spool_offset(ctx, offset, committed_offset, + &journal_offset); + if (commit_result != 0) { + ret = commit_result == -2 ? FLB_ERROR : FLB_RETRY; + } + } + } + if (ret != FLB_OK) { + if (rollback_spool(ctx, offset) != 0) { + ret = FLB_ERROR; + } + } + else { + for (i = 0; i < ids.count; i++) { + session_id_insert(ctx, ids.values[i]); + } + } + flb_free(ids.values); + return ret; +} + +static int stream_spool(struct flb_out_manticore *ctx, + struct flb_connection *connection) +{ + size_t length; + char *buffer; + + buffer = flb_malloc(ctx->stream_chunk_size); + if (buffer == NULL) { + return -1; + } + if (fflush(ctx->spool) != 0 || fseeko(ctx->spool, 0, SEEK_SET) != 0) { + flb_free(buffer); + return -1; + } + + while ((length = fread(buffer, 1, ctx->stream_chunk_size, + ctx->spool)) > 0) { + if (write_chunk(connection, buffer, length) != 0) { + flb_free(buffer); + return -1; + } + } + if (ferror(ctx->spool)) { + clearerr(ctx->spool); + flb_free(buffer); + return -1; + } + flb_free(buffer); + return 0; +} +#endif + +static int send_stream(struct flb_out_manticore *ctx, + const void *data, size_t bytes) +{ + int ret; + int result; + size_t sent; + struct flb_connection *connection; + struct flb_http_client *client; + + /* A permanent record error must not follow already transmitted records. */ + ret = validate_events(ctx, data, bytes, NULL); + if (ret != MANTICORE_STREAM_OK) { + return ret == MANTICORE_STREAM_RETRY ? FLB_RETRY : FLB_ERROR; + } + + connection = flb_upstream_conn_get(ctx->u); + if (connection == NULL) { + return FLB_RETRY; + } + + client = flb_http_client(connection, FLB_HTTP_POST, + ctx->bulk_uri, + NULL, 0, NULL, 0, NULL, 0); + if (client == NULL) { + flb_upstream_conn_release(connection); + return FLB_RETRY; + } + + flb_http_remove_header(client, "Content-Length", 14); + flb_http_remove_header(client, "Connection", 10); + client->body_len = -1; + flb_http_add_header(client, "Content-Type", 12, + "application/x-ndjson", 20); + flb_http_add_header(client, "Transfer-Encoding", 17, "chunked", 7); + flb_http_add_header(client, "Connection", 10, "close", 5); + flb_http_add_header(client, "User-Agent", 10, + "Fluent-Bit-Manticore", 20); + flb_http_buffer_size(client, ctx->buffer_size); + + if (ctx->http_user != NULL) { + flb_http_basic_auth(client, ctx->http_user, ctx->http_passwd); + } + + sent = 0; + ret = flb_http_do_request(client, &sent); + if (ret != FLB_HTTP_MORE) { + result = FLB_RETRY; + goto done; + } + + ret = stream_events(ctx, connection, data, bytes); + if (ret != MANTICORE_STREAM_OK) { + result = ret == MANTICORE_STREAM_RECORD_ERROR ? FLB_ERROR : FLB_RETRY; + goto done; + } + + if (write_all(connection, "0\r\n\r\n", 5) != 0) { + result = FLB_RETRY; + goto done; + } + + do { + ret = flb_http_get_response_data(client, 0); + } while (ret == FLB_HTTP_MORE || ret == FLB_HTTP_CHUNK_AVAILABLE); + + if (ret != FLB_HTTP_OK) { + result = FLB_RETRY; + goto done; + } + + result = response_ok(ctx, client); + +done: + /* Closing the session releases Manticore's bulk_import reservation. */ + flb_upstream_conn_recycle(connection, FLB_FALSE); + flb_http_client_destroy(client); + flb_upstream_conn_release(connection); + return result; +} + +#ifndef FLB_SYSTEM_WINDOWS +static int send_spool(struct flb_out_manticore *ctx) +{ + int ret; + int result; + size_t sent; + struct flb_connection *connection; + struct flb_http_client *client; + + connection = flb_upstream_conn_get(ctx->u); + if (connection == NULL) { + return FLB_RETRY; + } + + client = flb_http_client(connection, FLB_HTTP_POST, ctx->bulk_uri, + NULL, 0, NULL, 0, NULL, 0); + if (client == NULL) { + flb_upstream_conn_release(connection); + return FLB_RETRY; + } + + flb_http_remove_header(client, "Content-Length", 14); + flb_http_remove_header(client, "Connection", 10); + client->body_len = -1; + flb_http_add_header(client, "Content-Type", 12, + "application/x-ndjson", 20); + flb_http_add_header(client, "Transfer-Encoding", 17, "chunked", 7); + flb_http_add_header(client, "Connection", 10, "close", 5); + flb_http_add_header(client, "User-Agent", 10, + "Fluent-Bit-Manticore", 20); + flb_http_buffer_size(client, ctx->buffer_size); + if (ctx->http_user != NULL) { + flb_http_basic_auth(client, ctx->http_user, ctx->http_passwd); + } + + sent = 0; + ret = flb_http_do_request(client, &sent); + if (ret != FLB_HTTP_MORE || stream_spool(ctx, connection) != 0 || + write_all(connection, "0\r\n\r\n", 5) != 0) { + result = FLB_RETRY; + goto done; + } + + do { + ret = flb_http_get_response_data(client, 0); + } while (ret == FLB_HTTP_MORE || ret == FLB_HTTP_CHUNK_AVAILABLE); + result = ret == FLB_HTTP_OK ? response_ok(ctx, client) : FLB_RETRY; + +done: + flb_upstream_conn_recycle(connection, FLB_FALSE); + flb_http_client_destroy(client); + flb_upstream_conn_release(connection); + return result; +} + +static int recover_spool(struct flb_out_manticore *ctx) +{ + int consumed; + off_t data_size; + off_t good_end; + off_t current; + uint64_t checksum; + uint64_t expected_checksum; + uint64_t value; + uint64_t inverse; + char boundary; + char line[128]; + struct stat status; + + if (fstat(ctx->spool_fd, &status) != 0 || status.st_size < 0) { + return -1; + } + data_size = status.st_size; + current = 0; + good_end = 0; + clearerr(ctx->commit); + if (fseeko(ctx->commit, 0, SEEK_SET) != 0) { + return -1; + } + + if (fgets(line, sizeof(line), ctx->commit) == NULL) { + if (ferror(ctx->commit) || data_size != 0) { + return -1; + } + return 0; + } + + if (fseeko(ctx->commit, 0, SEEK_SET) != 0) { + return -1; + } + + /* A successful upload can crash after the spool was truncated but before + * its journal was cleared. The upload is already acknowledged by the + * server, so validate complete journal records and reset the empty pair. + */ + if (data_size == 0) { + while (fgets(line, sizeof(line), ctx->commit) != NULL) { + if (strchr(line, '\n') == NULL) { + if (!feof(ctx->commit)) { + flb_plg_error(ctx->ins, "commit journal '%s' is corrupt", + ctx->commit_path); + return -1; + } + break; + } + consumed = 0; + if (sscanf(line, + "%16" SCNx64 " %16" SCNx64 " %16" SCNx64 "%n", + &value, &inverse, &expected_checksum, &consumed) != 3 || + consumed != 50 || line[consumed] != '\n' || + line[consumed + 1] != '\0' || inverse != ~value || + value == 0) { + flb_plg_error(ctx->ins, "commit journal '%s' is corrupt", + ctx->commit_path); + return -1; + } + good_end = ftello(ctx->commit); + if (good_end < 0) { + return -1; + } + } + if (ferror(ctx->commit)) { + return -1; + } + clearerr(ctx->commit); + if (ftruncate(ctx->commit_fd, 0) != 0 || + fsync(ctx->commit_fd) != 0 || + fseeko(ctx->commit, 0, SEEK_END) != 0) { + return -1; + } + return 0; + } + + while (fgets(line, sizeof(line), ctx->commit) != NULL) { + if (strchr(line, '\n') == NULL) { + if (!feof(ctx->commit)) { + flb_plg_error(ctx->ins, "commit journal '%s' is corrupt", + ctx->commit_path); + return -1; + } + break; + } + consumed = 0; + if (sscanf(line, + "%16" SCNx64 " %16" SCNx64 " %16" SCNx64 "%n", + &value, &inverse, &expected_checksum, &consumed) != 3 || + consumed != 50 || line[consumed] != '\n' || + line[consumed + 1] != '\0' || inverse != ~value || + value <= (uint64_t) current || value > (uint64_t) data_size || + hash_file_range(ctx->spool_fd, current, (off_t) value, + &checksum) != 0 || checksum != expected_checksum || + pread(ctx->spool_fd, &boundary, 1, (off_t) value - 1) != 1 || + boundary != '\n') { + flb_plg_error(ctx->ins, "commit journal or spool '%s' is corrupt", + ctx->spool_path); + return -1; + } + current = (off_t) value; + good_end = ftello(ctx->commit); + if (good_end < 0) { + return -1; + } + } + if (ferror(ctx->commit)) { + return -1; + } + + clearerr(ctx->commit); + if (ftruncate(ctx->commit_fd, good_end) != 0 || + fsync(ctx->commit_fd) != 0 || + fseeko(ctx->commit, 0, SEEK_END) != 0) { + return -1; + } + if (ftruncate(ctx->spool_fd, current) != 0 || + fsync(ctx->spool_fd) != 0 || + fseeko(ctx->spool, current, SEEK_SET) != 0) { + return -1; + } + return 0; +} + +static int sync_parent_directory(const char *path) +{ + int fd; + int ret; + char *copy; + char *separator; + + copy = flb_strdup(path); + if (copy == NULL) { + return -1; + } + separator = strrchr(copy, '/'); + if (separator == NULL) { + flb_free(copy); + copy = flb_strdup("."); + if (copy == NULL) { + return -1; + } + } + else if (separator == copy) { + separator[1] = '\0'; + } + else { + *separator = '\0'; + } +#ifdef O_DIRECTORY + fd = open(copy, O_RDONLY | O_DIRECTORY); +#else + fd = open(copy, O_RDONLY); +#endif + flb_free(copy); + if (fd < 0) { + return -1; + } + ret = fsync(fd); + close(fd); + return ret; +} + +static int open_spool(struct flb_out_manticore *ctx) +{ + int flags; + struct stat status; + + flags = O_CREAT | O_RDWR | O_APPEND; +#ifdef O_CLOEXEC + flags |= O_CLOEXEC; +#endif +#ifdef O_NOFOLLOW + flags |= O_NOFOLLOW; +#endif + ctx->spool_fd = open(ctx->spool_path, flags, S_IRUSR | S_IWUSR); + if (ctx->spool_fd < 0 || flock(ctx->spool_fd, LOCK_EX | LOCK_NB) != 0) { + flb_plg_error(ctx->ins, "could not exclusively open spool '%s'", + ctx->spool_path); + if (ctx->spool_fd >= 0) { + close(ctx->spool_fd); + ctx->spool_fd = -1; + } + return -1; + } + if (fstat(ctx->spool_fd, &status) != 0 || !S_ISREG(status.st_mode)) { + flb_plg_error(ctx->ins, "spool '%s' must be a regular file", + ctx->spool_path); + close(ctx->spool_fd); + ctx->spool_fd = -1; + return -1; + } + ctx->spool = fdopen(ctx->spool_fd, "a+"); + if (ctx->spool == NULL) { + close(ctx->spool_fd); + ctx->spool_fd = -1; + return -1; + } + + if (strlen(ctx->spool_path) > SIZE_MAX - 8) { + fclose(ctx->spool); + ctx->spool = NULL; + ctx->spool_fd = -1; + return -1; + } + ctx->commit_path = flb_sds_create_size(strlen(ctx->spool_path) + 8); + if (ctx->commit_path == NULL || + flb_sds_printf(&ctx->commit_path, "%s.commit", ctx->spool_path) == NULL) { + flb_sds_destroy(ctx->commit_path); + ctx->commit_path = NULL; + fclose(ctx->spool); + ctx->spool = NULL; + ctx->spool_fd = -1; + return -1; + } + ctx->commit_fd = open(ctx->commit_path, flags, S_IRUSR | S_IWUSR); + if (ctx->commit_fd < 0 || fstat(ctx->commit_fd, &status) != 0 || + !S_ISREG(status.st_mode)) { + flb_plg_error(ctx->ins, "could not open commit journal '%s'", + ctx->commit_path); + if (ctx->commit_fd >= 0) { + close(ctx->commit_fd); + } + ctx->commit_fd = -1; + fclose(ctx->spool); + ctx->spool = NULL; + ctx->spool_fd = -1; + return -1; + } + ctx->commit = fdopen(ctx->commit_fd, "a+"); + if (ctx->commit == NULL) { + close(ctx->commit_fd); + ctx->commit_fd = -1; + fclose(ctx->spool); + ctx->spool = NULL; + ctx->spool_fd = -1; + return -1; + } + if (recover_spool(ctx) != 0 || + sync_parent_directory(ctx->spool_path) != 0) { + flb_plg_error(ctx->ins, "could not recover durable spool '%s'", + ctx->spool_path); + fclose(ctx->commit); + ctx->commit = NULL; + ctx->commit_fd = -1; + fclose(ctx->spool); + ctx->spool = NULL; + ctx->spool_fd = -1; + return -1; + } + return 0; +} + +static int clear_spool(struct flb_out_manticore *ctx) +{ + clearerr(ctx->spool); + if (fflush(ctx->spool) != 0 || + ftruncate(ctx->spool_fd, 0) != 0 || + fsync(ctx->spool_fd) != 0 || + fseeko(ctx->spool, 0, SEEK_SET) != 0) { + flb_plg_error(ctx->ins, "could not clear spool '%s'", + ctx->spool_path); + return -1; + } + + /* Retire commit records only after the empty spool is durable. */ + clearerr(ctx->commit); + if (fflush(ctx->commit) != 0 || + ftruncate(ctx->commit_fd, 0) != 0 || + fsync(ctx->commit_fd) != 0 || + fseeko(ctx->commit, 0, SEEK_END) != 0) { + flb_plg_error(ctx->ins, "could not clear commit journal '%s'", + ctx->commit_path); + return -1; + } + return 0; +} + +static void close_spool(struct flb_out_manticore *ctx) +{ + if (ctx->commit != NULL) { + fclose(ctx->commit); + ctx->commit = NULL; + ctx->commit_fd = -1; + } + if (ctx->spool != NULL) { + fclose(ctx->spool); + ctx->spool = NULL; + ctx->spool_fd = -1; + } +} + +static int remove_owned_path(struct flb_out_manticore *ctx, + const char *path, int fd, const char *kind) +{ + struct stat open_status; + struct stat path_status; + + if (fstat(fd, &open_status) != 0 || lstat(path, &path_status) != 0 || + open_status.st_dev != path_status.st_dev || + open_status.st_ino != path_status.st_ino) { + flb_plg_warn(ctx->ins, "refusing to remove replaced %s '%s'", kind, path); + return -1; + } + if (unlink(path) != 0) { + flb_plg_warn(ctx->ins, "could not remove %s '%s'", kind, path); + return -1; + } + return 0; +} + +static void remove_spool(struct flb_out_manticore *ctx, const char *reason) +{ + int failed; + + failed = FLB_FALSE; + if (remove_owned_path(ctx, ctx->spool_path, ctx->spool_fd, + reason) != 0) { + failed = FLB_TRUE; + } + if (ctx->commit_path != NULL && + remove_owned_path(ctx, ctx->commit_path, ctx->commit_fd, + "commit journal") != 0) { + failed = FLB_TRUE; + } + if (failed == FLB_FALSE && sync_parent_directory(ctx->spool_path) != 0) { + flb_plg_warn(ctx->ins, "could not sync spool directory for '%s'", + ctx->spool_path); + } +} +#endif + +static int cb_manticore_init(struct flb_output_instance *ins, + struct flb_config *config, void *data) +{ + int io_flags; + int ret; + char *table_json; + msgpack_object table; + struct flb_out_manticore *ctx; +#ifndef FLB_SYSTEM_WINDOWS + off_t spool_size; +#endif + + (void) data; + + ctx = flb_calloc(1, sizeof(struct flb_out_manticore)); + if (ctx == NULL) { + return -1; + } + + ctx->ins = ins; + ctx->config = config; + ctx->spool_fd = -1; + ctx->commit_fd = -1; + flb_output_net_default("127.0.0.1", FLB_MANTICORE_DEFAULT_PORT, ins); + + ret = flb_output_config_map_set(ins, ctx); + if (ret == -1 || ctx->table == NULL || ctx->table[0] == '\0') { + flb_plg_error(ins, "table is required"); + flb_free(ctx); + return -1; + } + + if (strcasecmp(ctx->action, "insert") != 0 && + strcasecmp(ctx->action, "create") != 0) { + flb_plg_error(ins, "action must be 'insert' or 'create'"); + flb_free(ctx); + return -1; + } + + ctx->bulk_action = strcasecmp(ctx->action, "create") == 0 ? + "create" : "insert"; + + if (ctx->single_chunk == FLB_TRUE) { +#ifdef FLB_SYSTEM_WINDOWS + flb_plg_error(ins, "single_chunk is not supported on Windows"); + flb_free(ctx); + return -1; +#else + if (ins->tp_workers != 1) { + flb_plg_error(ins, "single_chunk requires exactly one output worker"); + flb_free(ctx); + return -1; + } + if (ctx->spool_path == NULL || ctx->spool_path[0] == '\0') { + flb_plg_error(ins, "spool_path is required when single_chunk is enabled"); + flb_free(ctx); + return -1; + } + if (strcasecmp(ctx->action, "insert") != 0) { + flb_plg_error(ins, "single_chunk requires action 'insert' for replay safety"); + flb_free(ctx); + return -1; + } +#endif + } + + if (ctx->stream_chunk_size == 0 || ctx->max_session_ids == 0) { + flb_plg_error(ins, + "stream_chunk_size and max_session_ids must be greater than zero"); + flb_free(ctx); + return -1; + } + + table.type = MSGPACK_OBJECT_STR; + table.via.str.ptr = ctx->table; + table.via.str.size = strlen(ctx->table); + table_json = object_to_json(&table, config->json_escape_unicode); + if (table_json == NULL) { + flb_free(ctx); + return -1; + } + ctx->table_json = flb_sds_create(table_json); + flb_free(table_json); + if (ctx->table_json == NULL) { + flb_free(ctx); + return -1; + } + + ctx->bulk_uri = build_bulk_uri(ctx->table); + if (ctx->bulk_uri == NULL) { + flb_sds_destroy(ctx->table_json); + flb_free(ctx); + return -1; + } + + /* A persistent session would keep the bulk_import reservation active. */ + ins->net_setup.keepalive = FLB_FALSE; + + io_flags = ins->use_tls == FLB_TRUE ? FLB_IO_TLS : FLB_IO_TCP; + if (ins->host.ipv6 == FLB_TRUE) { + io_flags |= FLB_IO_IPV6; + } + + ctx->u = flb_upstream_create(config, ins->host.name, ins->host.port, + io_flags, ins->tls); + if (ctx->u == NULL) { + flb_sds_destroy(ctx->bulk_uri); + flb_sds_destroy(ctx->table_json); + flb_free(ctx); + return -1; + } + flb_output_upstream_set(ctx->u, ins); + +#ifndef FLB_SYSTEM_WINDOWS + if (ctx->single_chunk == FLB_TRUE) { + if (open_spool(ctx) != 0) { + flb_upstream_destroy(ctx->u); + flb_sds_destroy(ctx->bulk_uri); + flb_sds_destroy(ctx->table_json); + flb_sds_destroy(ctx->commit_path); + flb_free(ctx); + return -1; + } + spool_size = ftello(ctx->spool); + if (spool_size < 0) { + close_spool(ctx); + flb_upstream_destroy(ctx->u); + flb_sds_destroy(ctx->bulk_uri); + flb_sds_destroy(ctx->table_json); + flb_sds_destroy(ctx->commit_path); + flb_free(ctx); + return -1; + } + if (spool_size > 0) { + flb_plg_info(ins, "replaying pending single-chunk spool '%s'", + ctx->spool_path); + if (send_spool(ctx) != FLB_OK || clear_spool(ctx) != 0) { + flb_plg_error(ins, "could not replay pending spool '%s'", + ctx->spool_path); + close_spool(ctx); + flb_upstream_destroy(ctx->u); + flb_sds_destroy(ctx->bulk_uri); + flb_sds_destroy(ctx->table_json); + flb_sds_destroy(ctx->commit_path); + flb_free(ctx); + return -1; + } + } + } +#endif + flb_output_set_context(ins, ctx); + flb_output_set_http_debug_callbacks(ins); + return 0; +} + +static void cb_manticore_flush(struct flb_event_chunk *event_chunk, + struct flb_output_flush *out_flush, + struct flb_input_instance *ins, + void *out_context, + struct flb_config *config) +{ + int ret; + struct flb_out_manticore *ctx; + + (void) ins; + (void) config; + + ctx = out_context; + if (ctx->single_chunk == FLB_TRUE) { +#ifndef FLB_SYSTEM_WINDOWS + if (ctx->session_failed == FLB_TRUE) { + ret = FLB_ERROR; + } + else { + ret = spool_events(ctx, event_chunk->data, event_chunk->size); + if (ret == FLB_ERROR) { + ctx->session_failed = FLB_TRUE; + if (clear_spool(ctx) != 0) { + flb_plg_error(ctx->ins, + "could not durably abort single-chunk session"); + } + } + } +#else + ret = FLB_ERROR; +#endif + } + else { + ret = send_stream(ctx, event_chunk->data, event_chunk->size); + } + FLB_OUTPUT_RETURN(ret); +} + +static int cb_manticore_exit(void *data, struct flb_config *config) +{ + int fs_chunks; + int mem_chunks; + int tasks; + int ret; + struct flb_out_manticore *ctx; +#ifndef FLB_SYSTEM_WINDOWS + off_t spool_size; +#endif + + ctx = data; + if (ctx == NULL) { + return 0; + } + +#ifndef FLB_SYSTEM_WINDOWS + if (ctx->single_chunk == FLB_TRUE && ctx->spool != NULL) { + if (ctx->session_failed == FLB_TRUE) { + flb_plg_error(ctx->ins, + "single-chunk session aborted after a permanent record error"); + config->exit_status_code = 1; + ret = FLB_ERROR; + } + else if ((tasks = flb_task_running_count(config)) > 0) { + flb_plg_error(ctx->ins, + "single-chunk session is incomplete (%d running tasks)", + tasks); + ret = FLB_RETRY; + } + else { + flb_storage_chunk_count(config, &mem_chunks, &fs_chunks); + if (mem_chunks + fs_chunks > 0) { + flb_plg_error(ctx->ins, + "single-chunk session is incomplete (%d pending chunks)", + mem_chunks + fs_chunks); + ret = FLB_RETRY; + } + else if (fseeko(ctx->spool, 0, SEEK_END) != 0 || + (spool_size = ftello(ctx->spool)) < 0) { + ret = FLB_RETRY; + } + else if (spool_size > 0) { + ret = send_spool(ctx); + } + else { + ret = FLB_OK; + } + } + + if (ret == FLB_OK && clear_spool(ctx) != 0) { + ret = FLB_RETRY; + } + if (ctx->session_failed == FLB_TRUE) { + if (clear_spool(ctx) != 0) { + config->exit_status_code = 1; + } + remove_spool(ctx, "aborted"); + close_spool(ctx); + } + else if (ret == FLB_OK) { + remove_spool(ctx, "empty"); + close_spool(ctx); + } + else { + if (ret == FLB_ERROR) { + flb_plg_error(ctx->ins, + "single-chunk upload was rejected permanently"); + } + flb_plg_error(ctx->ins, + "single-chunk upload failed; preserving spool '%s'", + ctx->spool_path); + /* + * A remote final-publication failure during cb_exit is best + * effort. Preserve the durable spool for the next startup; + * Fluent Bit's shutdown status is not a Manticore publication + * certificate. + */ + close_spool(ctx); + } + } +#endif + + if (ctx->u != NULL) { + flb_upstream_destroy(ctx->u); + } + if (ctx->table_json != NULL) { + flb_sds_destroy(ctx->table_json); + } + if (ctx->bulk_uri != NULL) { + flb_sds_destroy(ctx->bulk_uri); + } + if (ctx->commit_path != NULL) { + flb_sds_destroy(ctx->commit_path); + } + flb_free(ctx->session_ids); + flb_free(ctx); + return 0; +} + +static struct flb_config_map config_map[] = { + { + FLB_CONFIG_MAP_STR, "table", NULL, + 0, FLB_TRUE, offsetof(struct flb_out_manticore, table), + "Target Manticore table (must already exist)" + }, + { + FLB_CONFIG_MAP_STR, "action", "insert", + 0, FLB_TRUE, offsetof(struct flb_out_manticore, action), + "Manticore bulk import action: insert or create" + }, + { + FLB_CONFIG_MAP_STR, "id_key", "id", + 0, FLB_TRUE, offsetof(struct flb_out_manticore, id_key), + "Required top-level non-zero numeric document ID, removed from doc" + }, + { + FLB_CONFIG_MAP_BOOL, "single_chunk", "false", + 0, FLB_TRUE, offsetof(struct flb_out_manticore, single_chunk), + "Stage all flushes locally and publish one bulk_import request on shutdown" + }, + { + FLB_CONFIG_MAP_STR, "spool_path", NULL, + 0, FLB_TRUE, offsetof(struct flb_out_manticore, spool_path), + "Durable spool file required by single_chunk" + }, + { + FLB_CONFIG_MAP_SIZE, "max_session_ids", "1M", + 0, FLB_TRUE, offsetof(struct flb_out_manticore, max_session_ids), + "Maximum IDs tracked for session-wide duplicate detection" + }, + { + FLB_CONFIG_MAP_SIZE, "stream_chunk_size", "64K", + 0, FLB_TRUE, offsetof(struct flb_out_manticore, stream_chunk_size), + "Maximum uncompressed NDJSON bytes buffered per HTTP chunk" + }, + { + FLB_CONFIG_MAP_SIZE, "buffer_size", "64K", + 0, FLB_TRUE, offsetof(struct flb_out_manticore, buffer_size), + "Maximum response buffer size" + }, + { + FLB_CONFIG_MAP_STR, "http_user", NULL, + 0, FLB_TRUE, offsetof(struct flb_out_manticore, http_user), + "HTTP Basic authentication user" + }, + { + FLB_CONFIG_MAP_STR, "http_passwd", "", + 0, FLB_TRUE, offsetof(struct flb_out_manticore, http_passwd), + "HTTP Basic authentication password" + }, + {0} +}; + +struct flb_output_plugin out_manticore_plugin = { + .name = "manticore", + .description = "Manticore Search native streaming output", + .cb_init = cb_manticore_init, + .cb_pre_run = NULL, + .cb_flush = cb_manticore_flush, + .cb_exit = cb_manticore_exit, + .workers = 2, + .config_map = config_map, + .event_type = FLB_OUTPUT_LOGS, + .flags = FLB_OUTPUT_NET | FLB_IO_OPT_TLS +}; diff --git a/plugins/out_manticore/manticore.h b/plugins/out_manticore/manticore.h new file mode 100644 index 00000000000..48a435e96a8 --- /dev/null +++ b/plugins/out_manticore/manticore.h @@ -0,0 +1,59 @@ +/* -*- Mode: C; tab-width: 4; indent-tabs-mode: nil; c-basic-offset: 4 -*- */ + +/* Fluent Bit + * ========== + * Copyright (C) 2015-2026 The Fluent Bit Authors + * + * Licensed 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. + */ + +#ifndef FLB_OUT_MANTICORE_H +#define FLB_OUT_MANTICORE_H + +#include + +#include +#include + +#define FLB_MANTICORE_DEFAULT_PORT 9308 +#define FLB_MANTICORE_BULK_URI "/bulk?bulk_import=" + +struct flb_out_manticore { + char *table; + char *action; + const char *bulk_action; + char *id_key; + char *http_user; + char *http_passwd; + int single_chunk; + int session_failed; + char *spool_path; + flb_sds_t commit_path; + FILE *spool; + int spool_fd; + FILE *commit; + int commit_fd; + uint64_t *session_ids; + size_t session_id_count; + size_t session_id_capacity; + size_t max_session_ids; + flb_sds_t table_json; + flb_sds_t bulk_uri; + size_t stream_chunk_size; + size_t buffer_size; + struct flb_upstream *u; + struct flb_output_instance *ins; + struct flb_config *config; +}; + +#endif From ab8722ab0bd4142cd35ec3e2b60f2fce25fc2acb Mon Sep 17 00:00:00 2001 From: Stas Date: Mon, 14 Sep 2026 22:25:47 +0200 Subject: [PATCH 3/3] tests: integration: add Manticore output coverage Signed-off-by: Stas --- .../config/out_manticore_permanent.yaml | 19 + .../config/out_manticore_recovery.yaml | 26 + .../config/out_manticore_retry.yaml | 20 + .../config/out_manticore_session.yaml | 28 + .../config/out_manticore_spool.yaml | 26 + .../config/out_manticore_wire.yaml | 22 + .../tests/test_out_manticore_001.py | 633 ++++++++++++++++++ .../tests/test_out_manticore_spool_001.py | 258 +++++++ 8 files changed, 1032 insertions(+) create mode 100644 tests/integration/scenarios/out_manticore/config/out_manticore_permanent.yaml create mode 100644 tests/integration/scenarios/out_manticore/config/out_manticore_recovery.yaml create mode 100644 tests/integration/scenarios/out_manticore/config/out_manticore_retry.yaml create mode 100644 tests/integration/scenarios/out_manticore/config/out_manticore_session.yaml create mode 100644 tests/integration/scenarios/out_manticore/config/out_manticore_spool.yaml create mode 100644 tests/integration/scenarios/out_manticore/config/out_manticore_wire.yaml create mode 100644 tests/integration/scenarios/out_manticore/tests/test_out_manticore_001.py create mode 100644 tests/integration/scenarios/out_manticore/tests/test_out_manticore_spool_001.py diff --git a/tests/integration/scenarios/out_manticore/config/out_manticore_permanent.yaml b/tests/integration/scenarios/out_manticore/config/out_manticore_permanent.yaml new file mode 100644 index 00000000000..ef5130d2d2b --- /dev/null +++ b/tests/integration/scenarios/out_manticore/config/out_manticore_permanent.yaml @@ -0,0 +1,19 @@ +service: + flush: 0.2 + log_level: info + http_server: on + http_port: ${FLUENT_BIT_HTTP_MONITORING_PORT} + +pipeline: + inputs: + - name: dummy + tag: manticore_permanent + dummy: '{"id":44,"message":"permanent-item"}' + samples: 1 + + outputs: + - name: manticore + match: '*' + host: 127.0.0.1 + port: ${MANTICORE_PORT} + table: wire_logs diff --git a/tests/integration/scenarios/out_manticore/config/out_manticore_recovery.yaml b/tests/integration/scenarios/out_manticore/config/out_manticore_recovery.yaml new file mode 100644 index 00000000000..472d2525f5f --- /dev/null +++ b/tests/integration/scenarios/out_manticore/config/out_manticore_recovery.yaml @@ -0,0 +1,26 @@ +service: + flush: 0.2 + grace: 3 + log_level: info + http_server: on + http_port: ${FLUENT_BIT_HTTP_MONITORING_PORT} + parsers_file: ${MANTICORE_PARSERS_FILE} + +pipeline: + inputs: + - name: tail + tag: manticore_recovery + path: ${MANTICORE_INPUT_PATH} + read_from_head: true + exit_on_eof: ${MANTICORE_EXIT_ON_EOF} + parser: recovery_json + + outputs: + - name: manticore + match: '*' + host: 127.0.0.1 + port: ${MANTICORE_PORT} + table: ${MANTICORE_TABLE} + workers: 1 + single_chunk: true + spool_path: ${MANTICORE_SPOOL_PATH} diff --git a/tests/integration/scenarios/out_manticore/config/out_manticore_retry.yaml b/tests/integration/scenarios/out_manticore/config/out_manticore_retry.yaml new file mode 100644 index 00000000000..d84c2cba975 --- /dev/null +++ b/tests/integration/scenarios/out_manticore/config/out_manticore_retry.yaml @@ -0,0 +1,20 @@ +service: + flush: 0.2 + log_level: info + http_server: on + http_port: ${FLUENT_BIT_HTTP_MONITORING_PORT} + +pipeline: + inputs: + - name: dummy + tag: manticore_retry + dummy: '{"id":"43","message":"retry-item"}' + samples: 1 + + outputs: + - name: manticore + match: '*' + host: 127.0.0.1 + port: ${MANTICORE_PORT} + table: wire_logs + action: create diff --git a/tests/integration/scenarios/out_manticore/config/out_manticore_session.yaml b/tests/integration/scenarios/out_manticore/config/out_manticore_session.yaml new file mode 100644 index 00000000000..e408d953d71 --- /dev/null +++ b/tests/integration/scenarios/out_manticore/config/out_manticore_session.yaml @@ -0,0 +1,28 @@ +service: + flush: 0.2 + grace: 3 + log_level: info + http_server: on + http_port: ${FLUENT_BIT_HTTP_MONITORING_PORT} + parsers_file: ${MANTICORE_PARSERS_FILE} + +pipeline: + inputs: + - name: tail + tag: manticore_session + path: ${MANTICORE_INPUT_PATH} + refresh_interval: 1 + read_from_head: true + exit_on_eof: ${MANTICORE_EXIT_ON_EOF} + parser: session_json + + outputs: + - name: manticore + match: '*' + host: 127.0.0.1 + port: ${MANTICORE_PORT} + table: session_logs + workers: 1 + single_chunk: true + spool_path: ${MANTICORE_SPOOL_PATH} + stream_chunk_size: 64K diff --git a/tests/integration/scenarios/out_manticore/config/out_manticore_spool.yaml b/tests/integration/scenarios/out_manticore/config/out_manticore_spool.yaml new file mode 100644 index 00000000000..0408f441099 --- /dev/null +++ b/tests/integration/scenarios/out_manticore/config/out_manticore_spool.yaml @@ -0,0 +1,26 @@ +service: + flush: 0.1 + grace: 3 + log_level: info + http_server: on + http_port: ${FLUENT_BIT_HTTP_MONITORING_PORT} + parsers_file: ${MANTICORE_PARSERS_FILE} + +pipeline: + inputs: + - name: tail + tag: manticore_spool + path: ${MANTICORE_INPUT_PATH} + read_from_head: true + exit_on_eof: ${MANTICORE_EXIT_ON_EOF} + parser: spool_json + + outputs: + - name: manticore + match: '*' + host: 127.0.0.1 + port: ${MANTICORE_PORT} + table: spool_logs + workers: 1 + single_chunk: true + spool_path: ${MANTICORE_SPOOL_PATH} diff --git a/tests/integration/scenarios/out_manticore/config/out_manticore_wire.yaml b/tests/integration/scenarios/out_manticore/config/out_manticore_wire.yaml new file mode 100644 index 00000000000..92a473092f5 --- /dev/null +++ b/tests/integration/scenarios/out_manticore/config/out_manticore_wire.yaml @@ -0,0 +1,22 @@ +service: + flush: 0.2 + log_level: info + http_server: on + http_port: ${FLUENT_BIT_HTTP_MONITORING_PORT} + parsers_file: ${MANTICORE_PARSERS_FILE} + +pipeline: + inputs: + - name: tail + tag: manticore_wire + path: ${MANTICORE_INPUT_PATH} + read_from_head: true + parser: manticore_json + + outputs: + - name: manticore + match: '*' + host: 127.0.0.1 + port: ${MANTICORE_PORT} + table: wire logs + stream_chunk_size: 128 diff --git a/tests/integration/scenarios/out_manticore/tests/test_out_manticore_001.py b/tests/integration/scenarios/out_manticore/tests/test_out_manticore_001.py new file mode 100644 index 00000000000..04ae197ad02 --- /dev/null +++ b/tests/integration/scenarios/out_manticore/tests/test_out_manticore_001.py @@ -0,0 +1,633 @@ +#!/usr/bin/env python3 + +import json +import os +import socket +import tempfile +import threading +import time + +import pytest + +from utils.fluent_bit_manager import FluentBitStartupError +from utils.test_service import FluentBitTestService + + +pytestmark = pytest.mark.skipif( + os.name == "nt", reason="Manticore single_chunk is not supported on Windows" +) + + +CONFIG_DIR = os.path.abspath(os.path.join(os.path.dirname(__file__), "../config")) + + +def config_path(config_file): + return os.path.join(CONFIG_DIR, config_file) + + +def config_environment(**values): + return {key: str(value) for key, value in values.items()} + + +def managed_service(config_file, environment): + return FluentBitTestService(config_path(config_file), extra_env=environment) + + +def read_service_log(service): + with open(service.flb.log_file, "r", encoding="utf-8", errors="replace") as stream: + return stream.read() + + +def stop_managed_service(service): + process = service.flb.process + log_file = service.flb.log_file + service.stop() + with open(log_file, "r", encoding="utf-8", errors="replace") as stream: + output = stream.read() + return process.returncode, output + + +def write_rejection_config(directory, payload, *, copies=None, output_options=None): + output_options = output_options or {} + input_path = os.path.join(directory, "records.json") + parsers_path = os.path.join(directory, "parsers.conf") + records = [payload] * (copies if copies is not None else 1) + with open(input_path, "w", encoding="utf-8") as stream: + for record in records: + stream.write(json.dumps(record, separators=(",", ":")) + "\n") + with open(parsers_path, "w", encoding="utf-8") as stream: + stream.write("[PARSER]\n Name rejection_json\n Format json\n") + + lines = [ + "service:", + " flush: 0.2", + " grace: 3", + " log_level: info", + " http_server: on", + " http_port: ${FLUENT_BIT_HTTP_MONITORING_PORT}", + " parsers_file: ${MANTICORE_PARSERS_FILE}", + "", + "pipeline:", + " inputs:", + " - name: tail", + " tag: manticore_rejection", + " path: ${MANTICORE_INPUT_PATH}", + " read_from_head: true", + " exit_on_eof: true", + " parser: rejection_json", + ] + lines.extend([ + "", + " outputs:", + " - name: manticore", + " match: '*'", + " host: 127.0.0.1", + " port: ${MANTICORE_PORT}", + " table: wire_logs", + ]) + for key, value in output_options.items(): + lines.append(" {}: {}".format(key, value)) + + path = os.path.join(directory, "rejection.yaml") + with open(path, "w", encoding="utf-8") as stream: + stream.write("\n".join(lines) + "\n") + return path, input_path, parsers_path + + +def read_request(connection): + stream = connection.makefile("rb") + request_line = stream.readline().decode().strip() + headers = {} + + while True: + line = stream.readline() + if line in (b"\r\n", b"\n", b""): + break + key, value = line.decode().split(":", 1) + headers[key.lower()] = value.strip() + + chunks = [] + while True: + size_line = stream.readline().strip() + size = int(size_line.split(b";", 1)[0], 16) + if size == 0: + stream.readline() + break + chunks.append(stream.read(size)) + if stream.read(2) != b"\r\n": + raise AssertionError("invalid HTTP chunk terminator") + + return request_line, headers, chunks + + +def send_response(connection, body, status="200 OK"): + response = ( + "HTTP/1.1 {}\r\n".format(status).encode() + + + b"Content-Type: application/json\r\n" + b"Content-Length: " + str(len(body)).encode() + b"\r\n" + b"Connection: close\r\n\r\n" + body + ) + connection.sendall(response) + connection.close() + + +def capture_request(listener, result): + connection, _ = listener.accept() + request_line, headers, chunks = read_request(connection) + result.update( + request_line=request_line, + headers=headers, + chunks=chunks, + ) + + body = ( + b'{"items":[{"bulk":{"created":20,"status":201}}],' + b'"current_line":20,"skipped_lines":0,"errors": false,"error":""}' + ) + send_response(connection, body) + listener.close() + + +def capture_retry(listener, result): + responses = [ + b'{"items":[{"bulk":{"status":503}}],"errors":true}', + b'{"items":[{"bulk":{"status":201}}],"errors":false}', + ] + result["requests"] = [] + + for body in responses: + connection, _ = listener.accept() + result["requests"].append(read_request(connection)) + send_response(connection, body) + + listener.close() + + +def capture_permanent_server_error(listener, result): + connection, _ = listener.accept() + result["request"] = read_request(connection) + body = b'{"items":[{"insert":{"status":409}}],"errors":true}' + send_response(connection, body, status="500 Internal Server Error") + listener.close() + + +def wait_for_service_log(service, expected, timeout=30): + deadline = time.monotonic() + timeout + while time.monotonic() < deadline: + output = read_service_log(service) + if expected in output: + return output + time.sleep(0.1) + raise AssertionError("Fluent Bit log did not contain {!r}\n{}".format( + expected, read_service_log(service))) + + +def wait_for_path(path, timeout=30): + deadline = time.monotonic() + timeout + while time.monotonic() < deadline: + if os.path.exists(path): + return + time.sleep(0.1) + raise AssertionError("Path did not appear before timeout: {}".format(path)) + + +def assert_rejected_before_connect(payload, expected, copies=None, + output_options=None): + listener = socket.socket() + listener.bind(("127.0.0.1", 0)) + listener.listen(1) + listener.settimeout(0.5) + port = listener.getsockname()[1] + + with tempfile.TemporaryDirectory(prefix="manticore-rejected-") as directory: + config_file, input_path, parsers_path = write_rejection_config( + directory, payload, copies=copies, output_options=output_options) + environment = config_environment( + MANTICORE_INPUT_PATH=input_path, + MANTICORE_PARSERS_FILE=parsers_path, + MANTICORE_PORT=port, + ) + service = FluentBitTestService(config_file, extra_env=environment) + with pytest.raises(FluentBitStartupError): + service.start() + output = read_service_log(service) + + connected = False + try: + connection, _ = listener.accept() + connection.close() + connected = True + except socket.timeout: + pass + finally: + listener.close() + + assert not connected + assert expected in output + assert "retry in" not in output + + +def test_out_manticore_chunked_and_recovery(): + result = {} + listener = socket.socket() + listener.bind(("127.0.0.1", 0)) + listener.listen(1) + port = listener.getsockname()[1] + + server = threading.Thread( + target=capture_request, + args=(listener, result), + daemon=True, + ) + server.start() + + fixture_dir = tempfile.TemporaryDirectory() + input_path = os.path.join(fixture_dir.name, "records.json") + parsers_path = os.path.join(fixture_dir.name, "parsers.conf") + with open(input_path, "w") as stream: + for document_id in range(1, 21): + stream.write(json.dumps({ + "id": document_id, + "message": "wire-test", + "status": 200, + }) + "\n") + with open(parsers_path, "w") as stream: + stream.write("[PARSER]\n Name manticore_json\n Format json\n") + + environment = config_environment( + MANTICORE_INPUT_PATH=input_path, + MANTICORE_PARSERS_FILE=parsers_path, + MANTICORE_PORT=port, + ) + service = managed_service("out_manticore_wire.yaml", environment) + service.start() + + server.join(20) + returncode, output = stop_managed_service(service) + + if server.is_alive(): + raise AssertionError("no request received\n{}".format(output)) + if returncode != 0: + raise AssertionError("Fluent Bit exited {}\n{}".format( + returncode, output)) + + headers = result["headers"] + chunks = result["chunks"] + records = [json.loads(line) for line in b"".join(chunks).splitlines()] + + assert result["request_line"] == ( + "POST /bulk?bulk_import=wire%20logs HTTP/1.1" + ) + assert headers.get("transfer-encoding") == "chunked" + assert headers.get("connection") == "close" + assert "content-length" not in headers + assert len(chunks) > 1 + assert len(records) == 20 + assert records == [ + { + "insert": { + "table": "wire logs", + "id": document_id, + "doc": {"message": "wire-test", "status": 200}, + } + } + for document_id in range(1, 21) + ] + fixture_dir.cleanup() + + print("captured {} records in {} HTTP chunks".format( + len(records), len(chunks))) + + session_result = {} + session_listener = socket.socket() + session_listener.bind(("127.0.0.1", 0)) + session_listener.listen(1) + session_port = session_listener.getsockname()[1] + session_server = threading.Thread( + target=capture_request, + args=(session_listener, session_result), + daemon=True, + ) + session_server.start() + session_dir = tempfile.TemporaryDirectory() + session_pattern = os.path.join(session_dir.name, "session-*.json") + session_input = os.path.join(session_dir.name, "session-records.json") + session_source = os.path.join(session_dir.name, "source.json") + session_parser = os.path.join(session_dir.name, "parsers.conf") + session_spool = os.path.join(session_dir.name, "session.ndjson") + session_record_count = 3000 if os.environ.get("VALGRIND") else 30000 + with open(session_source, "w") as stream: + for document_id in range(1, session_record_count + 1): + stream.write(json.dumps({ + "id": document_id, + "message": "session-test", + "payload": "x" * 128, + }, separators=(",", ":")) + "\n") + with open(session_parser, "w") as stream: + stream.write("[PARSER]\n Name session_json\n Format json\n") + with open(session_source, "rb") as stream: + session_data = stream.read() + session_environment = config_environment( + MANTICORE_INPUT_PATH=session_pattern, + MANTICORE_PARSERS_FILE=session_parser, + MANTICORE_PORT=session_port, + MANTICORE_SPOOL_PATH=session_spool, + MANTICORE_EXIT_ON_EOF="true", + ) + session_service = managed_service( + "out_manticore_session.yaml", session_environment) + session_service.start() + os.replace(session_source, session_input) + session_timeout = 180 if os.environ.get("VALGRIND") else 60 + session_server.join(timeout=session_timeout) + session_returncode, session_output = stop_managed_service(session_service) + assert session_returncode == 0, session_output + assert not session_server.is_alive() + session_records = b"".join(session_result["chunks"]).splitlines() + assert len(session_records) == session_record_count + assert session_result["request_line"] == ( + "POST /bulk?bulk_import=session_logs HTTP/1.1") + assert not os.path.exists(session_spool) + assert not os.path.exists(session_spool + ".commit") + assert "service has stopped (0 pending tasks)" in session_output + + rejected_listener = socket.socket() + rejected_listener.bind(("127.0.0.1", 0)) + rejected_listener.listen(1) + rejected_listener.settimeout(1) + rejected_port = rejected_listener.getsockname()[1] + rejected_pattern = os.path.join(session_dir.name, "rejected-*.json") + rejected_input = os.path.join(session_dir.name, "rejected-records.json") + rejected_data = session_data + json.dumps({ + "id": 1, + "message": "duplicate-late", + "payload": "x" * 128, + }, separators=(",", ":")).encode() + b"\n" + rejected_environment = config_environment( + MANTICORE_INPUT_PATH=rejected_pattern, + MANTICORE_PARSERS_FILE=session_parser, + MANTICORE_PORT=rejected_port, + MANTICORE_SPOOL_PATH=session_spool, + MANTICORE_EXIT_ON_EOF="true", + ) + rejected_service = managed_service( + "out_manticore_session.yaml", rejected_environment) + rejected_service.start() + with open(rejected_input, "wb") as stream: + stream.write(rejected_data) + wait_for_service_log(rejected_service, "must be unique within a session") + _, rejected_output = stop_managed_service(rejected_service) + connected = False + try: + rejected_connection, _ = rejected_listener.accept() + rejected_connection.close() + connected = True + except socket.timeout: + pass + rejected_listener.close() + assert not connected + assert "must be unique within a session" in rejected_output + assert "single-chunk session aborted" in rejected_output + assert not os.path.exists(session_spool) + assert not os.path.exists(session_spool + ".commit") + session_dir.cleanup() + print("single_chunk combined 30000 records and aborted a late duplicate") + + recovery_dir = tempfile.TemporaryDirectory() + recovery_input = os.path.join(recovery_dir.name, "recovery.json") + recovery_empty = os.path.join(recovery_dir.name, "empty.json") + recovery_parser = os.path.join(recovery_dir.name, "parsers.conf") + recovery_spool = os.path.join(recovery_dir.name, "recovery.ndjson") + with open(recovery_input, "w") as stream: + stream.write('{"id":70001,"message":"recover-me"}\n') + open(recovery_empty, "w").close() + with open(recovery_input, "rb") as stream: + recovery_data = stream.read() + with open(recovery_parser, "w") as stream: + stream.write("[PARSER]\n Name recovery_json\n Format json\n") + unavailable = socket.socket() + unavailable.bind(("127.0.0.1", 0)) + recovery_port = unavailable.getsockname()[1] + unavailable.close() + + def recovery_environment(path, table="recovery_logs", port=recovery_port, + exit_on_eof="false"): + return config_environment( + MANTICORE_INPUT_PATH=path, + MANTICORE_PARSERS_FILE=recovery_parser, + MANTICORE_PORT=port, + MANTICORE_SPOOL_PATH=recovery_spool, + MANTICORE_TABLE=table, + MANTICORE_EXIT_ON_EOF=exit_on_eof, + ) + + failed_service = managed_service( + "out_manticore_recovery.yaml", recovery_environment(recovery_input)) + failed_service.start() + wait_for_path(recovery_spool) + failed_returncode, failed_output = stop_managed_service(failed_service) + assert failed_returncode == 0 + assert os.path.getsize(recovery_spool) > 0 + assert os.path.getsize(recovery_spool + ".commit") > 0 + assert "preserving spool" in failed_output + + with open(recovery_spool, "ab") as stream: + stream.write(b'{"insert":{"table":"recovery_logs","id":999') + with open(recovery_spool + ".commit", "ab") as stream: + stream.write(b"000000000000") + + recovery_result = {} + recovery_listener = socket.socket() + recovery_listener.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) + recovery_listener.bind(("127.0.0.1", recovery_port)) + recovery_listener.listen(1) + recovery_server = threading.Thread( + target=capture_request, + args=(recovery_listener, recovery_result), + daemon=True, + ) + recovery_server.start() + open(recovery_empty, "w").close() + recovered_service = managed_service( + "out_manticore_recovery.yaml", + recovery_environment(recovery_empty, exit_on_eof="false")) + recovered_service.start() + recovered_returncode, recovered_output = stop_managed_service( + recovered_service) + recovery_server.join(timeout=30) + assert recovered_returncode == 0, recovered_output + assert not recovery_server.is_alive() + recovered_records = b"".join(recovery_result["chunks"]).splitlines() + assert len(recovered_records) == 1 + assert json.loads(recovered_records[0])["insert"]["id"] == 70001 + assert "replaying pending single-chunk spool" in recovered_output + assert not os.path.exists(recovery_spool) + assert not os.path.exists(recovery_spool + ".commit") + + with open(recovery_input, "wb") as stream: + stream.write(recovery_data) + second_failed_service = managed_service( + "out_manticore_recovery.yaml", recovery_environment(recovery_input)) + second_failed_service.start() + wait_for_path(recovery_spool) + second_failed_returncode, _ = stop_managed_service(second_failed_service) + assert second_failed_returncode == 0 + with open(recovery_spool, "r+b") as stream: + data = stream.read() + marker = data.index(b"recover-me") + stream.seek(marker) + stream.write(b"Recover-me") + corrupt_service = managed_service( + "out_manticore_recovery.yaml", + recovery_environment(recovery_empty, exit_on_eof="false")) + open(recovery_empty, "w").close() + with pytest.raises(FluentBitStartupError): + corrupt_service.start() + corrupt_output = read_service_log(corrupt_service) + assert "commit journal or spool" in corrupt_output + assert "is corrupt" in corrupt_output + os.remove(recovery_spool) + os.remove(recovery_spool + ".commit") + + with open(recovery_input, "wb") as stream: + stream.write(recovery_data) + third_failed_service = managed_service( + "out_manticore_recovery.yaml", recovery_environment(recovery_input)) + third_failed_service.start() + wait_for_path(recovery_spool) + third_failed_returncode, _ = stop_managed_service(third_failed_service) + assert third_failed_returncode == 0 + + changed_config_result = {} + changed_config_listener = socket.socket() + changed_config_listener.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) + changed_config_listener.bind(("127.0.0.1", recovery_port)) + changed_config_listener.listen(1) + changed_config_server = threading.Thread( + target=capture_request, + args=(changed_config_listener, changed_config_result), + daemon=True, + ) + changed_config_server.start() + changed_config_environment = recovery_environment( + recovery_empty, table="other_logs", exit_on_eof="false") + changed_config_service = managed_service( + "out_manticore_recovery.yaml", changed_config_environment) + open(recovery_empty, "w").close() + changed_config_service.start() + changed_config_returncode, changed_config_output = stop_managed_service( + changed_config_service) + changed_config_server.join(timeout=30) + assert changed_config_returncode == 0, changed_config_output + assert not changed_config_server.is_alive() + assert changed_config_result["request_line"] == \ + "POST /bulk?bulk_import=other_logs HTTP/1.1" + recovery_dir.cleanup() + print("recovery discarded torn tails and replayed under changed configuration") + + rejected_records = [ + ({"id": {"invalid": True}, "message": "poison"}, None, + "must be a unique, non-zero numeric ID"), + ({"message": "missing ID"}, None, + "must contain a non-zero numeric ID"), + ({"id": 0, "message": "zero ID"}, None, + "must be a unique, non-zero numeric ID"), + ({"id": 44, "message": "duplicate ID"}, 2, + "must be unique within a chunk"), + ] + for payload, copies, expected_error in rejected_records: + assert_rejected_before_connect(payload, expected_error, copies=copies) + print("permanent ID errors were rejected before delivery") + + retry_result = {} + retry_listener = socket.socket() + retry_listener.bind(("127.0.0.1", 0)) + retry_listener.listen(2) + retry_port = retry_listener.getsockname()[1] + retry_server = threading.Thread( + target=capture_retry, + args=(retry_listener, retry_result), + daemon=True, + ) + retry_server.start() + + retry_environment = config_environment(MANTICORE_PORT=retry_port) + retry_service = managed_service("out_manticore_retry.yaml", retry_environment) + retry_service.start() + retry_server.join(20) + retry_returncode, retry_output = stop_managed_service(retry_service) + + if retry_server.is_alive(): + raise AssertionError("transient item was not retried\n{}".format( + retry_output)) + + assert retry_returncode == 0, retry_output + assert len(retry_result["requests"]) == 2 + first_body = b"".join(retry_result["requests"][0][2]) + second_body = b"".join(retry_result["requests"][1][2]) + assert first_body == second_body + assert json.loads(first_body) == { + "create": { + "table": "wire_logs", + "id": 43, + "doc": {"message": "retry-item"}, + } + } + assert "retryable item error" in retry_output + print("transient item error retried the original chunk") + + permanent_result = {} + permanent_listener = socket.socket() + permanent_listener.bind(("127.0.0.1", 0)) + permanent_listener.listen(1) + permanent_port = permanent_listener.getsockname()[1] + permanent_server = threading.Thread( + target=capture_permanent_server_error, + args=(permanent_listener, permanent_result), + daemon=True, + ) + permanent_server.start() + permanent_environment = config_environment(MANTICORE_PORT=permanent_port) + permanent_service = managed_service( + "out_manticore_permanent.yaml", permanent_environment) + permanent_service.start() + permanent_server.join(timeout=20) + permanent_returncode, permanent_output = stop_managed_service( + permanent_service) + assert permanent_returncode == 0, permanent_output + assert not permanent_server.is_alive() + assert permanent_result["request"][0] == ( + "POST /bulk?bulk_import=wire_logs HTTP/1.1") + assert "returned HTTP 500" in permanent_output + assert "retry in" not in permanent_output + print("permanent item status inside HTTP 500 was not retried") + + assert_rejected_before_connect( + {"id": 44, "message": "invalid action"}, + "action must be 'insert' or 'create'", + output_options={"action": "replace"}, + ) + print("replace action rejected before delivery") + + for extra, expected in [ + ({"workers": "2"}, + "single_chunk requires exactly one output worker"), + ({"action": "create"}, + "single_chunk requires action 'insert' for replay safety"), + ({"max_session_ids": "0"}, + "stream_chunk_size and max_session_ids must be greater than zero"), + ]: + with tempfile.TemporaryDirectory(prefix="manticore-single-config-") as directory: + output_options = { + "workers": "1", + "single_chunk": "true", + "spool_path": os.path.join(directory, "spool.ndjson"), + } + output_options.update(extra) + assert_rejected_before_connect( + {"id": 44, "message": "invalid single_chunk option"}, + expected, + output_options=output_options, + ) + print("invalid single_chunk worker/action/limit configurations rejected") diff --git a/tests/integration/scenarios/out_manticore/tests/test_out_manticore_spool_001.py b/tests/integration/scenarios/out_manticore/tests/test_out_manticore_spool_001.py new file mode 100644 index 00000000000..cf096b4d529 --- /dev/null +++ b/tests/integration/scenarios/out_manticore/tests/test_out_manticore_spool_001.py @@ -0,0 +1,258 @@ +#!/usr/bin/env python3 + +"""Single-chunk retention and replay regressions using real Fluent Bit.""" + +import json +import os +from pathlib import Path +import socket +import tempfile +import time +import threading +import unittest + +import pytest + +from test_out_manticore_001 import ( + read_request, + read_service_log, + send_response, +) +from utils.fluent_bit_manager import FluentBitStartupError +from utils.test_service import FluentBitTestService + + +pytestmark = pytest.mark.skipif( + os.name == "nt", reason="Manticore single_chunk is not supported on Windows" +) + + +CONFIG_FILE = Path(__file__).resolve().parent.parent / "config" / "out_manticore_spool.yaml" + + +SUCCESS = ("200 OK", {"errors": False}) +RETRY = ("503 Service Unavailable", {"errors": True}) +UNAUTHORIZED = ("401 Unauthorized", {"error": "authentication required"}) +PERMANENT_ITEM = {"errors": True, "items": [{"insert": {"status": 409}}]} + + +class ServiceResult: + def __init__(self, returncode, stdout): + self.returncode = returncode + self.stdout = stdout + + +class Receiver: + def __init__(self, replies): + self.replies = replies + self.requests = [] + self.errors = [] + self.stopped = threading.Event() + self.listener = socket.socket() + self.listener.bind(("127.0.0.1", 0)) + self.listener.listen(1) + self.listener.settimeout(0.1) + self.port = self.listener.getsockname()[1] + self.thread = threading.Thread(target=self.serve, daemon=True) + self.thread.start() + + def serve(self): + while not self.stopped.is_set(): + try: + connection, _ = self.listener.accept() + except socket.timeout: + continue + except OSError: + return + try: + with connection: + connection.settimeout(30) + request, headers, chunks = read_request(connection) + number = len(self.requests) + self.requests.append((request, headers, b"".join(chunks))) + status, body = self.replies[number] + send_response(connection, json.dumps(body).encode(), status) + except Exception as error: + self.errors.append(repr(error)) + + def close(self): + self.stopped.set() + self.listener.close() + self.thread.join(timeout=35) + + +class SpoolRecoveryTests(unittest.TestCase): + def setUp(self): + evidence = os.environ.get("MANTICORE_TEST_RESULTS") + if evidence: + self.root = Path(evidence) / self._testMethodName + self.root.mkdir(parents=True, exist_ok=False) + else: + temporary = tempfile.TemporaryDirectory(prefix="manticore-spool-") + self.addCleanup(temporary.cleanup) + self.root = Path(temporary.name) + self.spool = self.root / "spool.ndjson" + self.journal = self.root / "spool.ndjson.commit" + self.parser = self.root / "parsers.conf" + self.parser.write_text("[PARSER]\n Name spool_json\n Format json\n") + self.receiver = None + + def stop_service(self, service): + process = service.flb.process + service.stop() + return process.returncode, read_service_log(service) + + def tearDown(self): + if self.receiver: + self.receiver.close() + requests = [] + for number, (request, headers, body) in enumerate(self.receiver.requests): + (self.root / "request-{}.ndjson".format(number)).write_bytes(body) + requests.append({"request": request, "headers": headers, + "bytes": len(body)}) + (self.root / "requests.json").write_text(json.dumps(requests, indent=2)) + self.assertFalse(self.receiver.thread.is_alive()) + self.assertEqual(self.receiver.errors, []) + + def start_receiver(self, replies): + self.receiver = Receiver(replies) + + def run_session(self, label, records, env=None, signal_after_spool=False): + source = self.root / (label + ".json") + source_text = "".join(json.dumps(record) + "\n" for record in records) + if records and not signal_after_spool: + source.write_text(source_text) + else: + source.write_text("") + session_env = { + "MANTICORE_INPUT_PATH": str(source), + "MANTICORE_PARSERS_FILE": str(self.parser), + "MANTICORE_PORT": str(self.receiver.port), + "MANTICORE_SPOOL_PATH": str(self.spool), + "MANTICORE_EXIT_ON_EOF": "false", + } + if env: + session_env.update(env) + (self.root / (label + "-config.txt")).write_text(str(CONFIG_FILE)) + expected_request_count = ( + len(self.receiver.requests) + 1 if not records else None + ) + service = FluentBitTestService(str(CONFIG_FILE), extra_env=session_env) + try: + service.start() + except FluentBitStartupError: + output = read_service_log(service) + (self.root / (label + ".log")).write_text(output) + (self.root / (label + "-exit.txt")).write_text("1") + return ServiceResult(1, output) + + process = service.flb.process + if signal_after_spool: + source.write_text(source_text) + + if signal_after_spool: + deadline = time.monotonic() + 30 + while not self.spool.exists() and process.poll() is None: + if time.monotonic() >= deadline: + service.stop() + self.fail("Manticore spool was not populated\n{}".format( + read_service_log(service))) + time.sleep(0.1) + if process.poll() is not None: + output = read_service_log(service) + service.stop() + self.fail("Fluent Bit exited before SIGTERM\n{}".format(output)) + returncode, output = self.stop_service(service) + if not records: + deadline = time.monotonic() + 30 + while len(self.receiver.requests) < expected_request_count: + if time.monotonic() >= deadline: + self.fail("Manticore request was not received\n{}".format( + output)) + time.sleep(0.1) + (self.root / (label + ".log")).write_text(output) + (self.root / (label + "-exit.txt")).write_text(str(returncode)) + for name, path in [("spool", self.spool), ("journal", self.journal)]: + if path.exists(): + (self.root / (label + "." + name)).write_bytes(path.read_bytes()) + return ServiceResult(returncode, output) + + def assert_exit(self, process, expected): + self.assertEqual(process.returncode, expected, process.stdout) + + def assert_removed(self): + self.assertFalse(self.spool.exists()) + self.assertFalse(self.journal.exists()) + + def assert_ids(self, expected): + actual = [] + for request, headers, body in self.receiver.requests: + self.assertEqual(request, "POST /bulk?bulk_import=spool_logs HTTP/1.1") + self.assertEqual(headers["transfer-encoding"].lower(), "chunked") + actual.append([json.loads(line)["insert"]["id"] + for line in body.splitlines()]) + self.assertEqual(actual, expected) + + def check_rejection(self, rejection): + self.start_receiver([rejection, SUCCESS]) + failed = self.run_session("rejected", [{"id": 101, "message": "retained"}]) + self.assert_exit(failed, 0) + self.assertIn("preserving spool", failed.stdout) + self.assertTrue(self.spool.exists()) + self.assertEqual(self.spool.read_bytes(), self.receiver.requests[0][2]) + recovered = self.run_session("recovered", []) + self.assert_exit(recovered, 0) + self.assert_ids([[101], [101]]) + self.assertEqual(self.receiver.requests[0][2], self.receiver.requests[1][2]) + self.assert_removed() + + def test_replay_then_new_input(self): + self.start_receiver([RETRY, SUCCESS, RETRY, SUCCESS]) + first = self.run_session("first", [{"id": 101, "message": "old"}]) + self.assert_exit(first, 0) + second = self.run_session("second", [{"id": 102, "message": "new"}]) + self.assert_exit(second, 0) + third = self.run_session("third", []) + self.assert_exit(third, 0) + self.assert_ids([[101], [101], [102], [102]]) + self.assertEqual(self.receiver.requests[2][2], self.receiver.requests[3][2]) + self.assert_removed() + + def test_auth_rejection_retains_accepted_data(self): + self.check_rejection(UNAUTHORIZED) + + def test_permanent_item_in_http_500_retains_accepted_data(self): + self.check_rejection(("500 Internal Server Error", PERMANENT_ITEM)) + + def test_permanent_item_in_http_200_retains_accepted_data(self): + self.check_rejection(("200 OK", PERMANENT_ITEM)) + + def test_transient_rejection_retains_accepted_data(self): + self.check_rejection(RETRY) + + def test_signal_final_rejection_is_best_effort(self): + self.start_receiver([RETRY, SUCCESS]) + failed = self.run_session("signal-failed", [{"id": 101, "message": "retained"}], + signal_after_spool=True) + self.assert_exit(failed, 0) + self.assertIn("preserving spool", failed.stdout) + self.assertTrue(self.spool.exists()) + self.assertTrue(self.journal.exists()) + recovered = self.run_session("signal-recovered", []) + self.assert_exit(recovered, 0) + self.assert_ids([[101], [101]]) + self.assert_removed() + + def test_rejected_startup_replay_preserves_data(self): + self.start_receiver([RETRY, UNAUTHORIZED, SUCCESS]) + first = self.run_session("first", [{"id": 101, "message": "pending"}]) + self.assert_exit(first, 0) + spool, journal = self.spool.read_bytes(), self.journal.read_bytes() + rejected = self.run_session("rejected-replay", []) + self.assertNotEqual(rejected.returncode, 0, rejected.stdout) + self.assertEqual(self.spool.read_bytes(), spool) + self.assertEqual(self.journal.read_bytes(), journal) + recovered = self.run_session("recovered", []) + self.assert_exit(recovered, 0) + self.assert_ids([[101], [101], [101]]) + self.assert_removed()