Skip to content
Closed
Show file tree
Hide file tree
Changes from 1 commit
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions source/common/common/logger.h
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,7 @@ namespace Logger {
FUNCTION(http) \
FUNCTION(http2) \
FUNCTION(hystrix) \
FUNCTION(kafka) \
FUNCTION(lua) \
FUNCTION(main) \
FUNCTION(misc) \
Expand Down
1 change: 1 addition & 0 deletions source/extensions/extensions_build_config.bzl
Original file line number Diff line number Diff line change
Expand Up @@ -65,6 +65,7 @@ EXTENSIONS = {
"envoy.filters.network.echo": "//source/extensions/filters/network/echo:config",
"envoy.filters.network.ext_authz": "//source/extensions/filters/network/ext_authz:config",
"envoy.filters.network.http_connection_manager": "//source/extensions/filters/network/http_connection_manager:config",
"envoy.filters.network.kafka": "//source/extensions/filters/network/kafka:config",
"envoy.filters.network.mongo_proxy": "//source/extensions/filters/network/mongo_proxy:config",
"envoy.filters.network.ratelimit": "//source/extensions/filters/network/ratelimit:config",
"envoy.filters.network.rbac": "//source/extensions/filters/network/rbac:config",
Expand Down
114 changes: 114 additions & 0 deletions source/extensions/filters/network/kafka/BUILD
Original file line number Diff line number Diff line change
@@ -0,0 +1,114 @@
licenses(["notice"]) # Apache 2

# Kafka network filter.
# Public docs: docs/root/configuration/network_filters/kafka_filter.rst

load(
"//bazel:envoy_build_system.bzl",
"envoy_cc_library",
"envoy_package",
)

envoy_package()

envoy_cc_library(
name = "config",
srcs = ["config.cc"],
hdrs = ["metrics_holder.h"],
deps = [
":kafka_filter_lib",
":metrics_holder_lib",
"//include/envoy/registry",
"//include/envoy/server:filter_config_interface",
"//source/extensions/filters/network:well_known_names",
],
)

envoy_cc_library(
name = "kafka_filter_lib",
srcs = ["kafka_filter.cc"],
hdrs = ["kafka_filter.h"],
deps = [
":kafka_codec_lib",
":kafka_request_lib",
":kafka_response_lib",
":metrics_holder_lib",
"//include/envoy/buffer:buffer_interface",
"//include/envoy/network:connection_interface",
"//include/envoy/network:filter_interface",
"//source/common/common:assert_lib",
"//source/common/common:minimal_logger_lib",
],
)

envoy_cc_library(
name = "kafka_codec_lib",
srcs = ["codec.cc"],
hdrs = ["codec.h"],
deps = [
":kafka_request_lib",
":kafka_response_lib",
"//include/envoy/buffer:buffer_interface",
],
)

envoy_cc_library(
name = "kafka_request_lib",
srcs = ["kafka_request.cc"],
hdrs = ["kafka_request.h"],
deps = [
":parser_lib",
":serialization_lib",
"//source/common/common:minimal_logger_lib",
"//source/common/common:assert_lib",
],
)

envoy_cc_library(
name = "kafka_response_lib",
srcs = ["kafka_response.cc"],
hdrs = ["kafka_response.h"],
deps = [
":parser_lib",
":serialization_lib",
"//source/common/common:minimal_logger_lib",
],
)

envoy_cc_library(
name = "metrics_holder_lib",
srcs = ["metrics_holder.cc"],
hdrs = ["metrics_holder.h"],
deps = [
":kafka_protocol_lib",
"//include/envoy/stats:stats_interface",
"//source/common/common:macros",
"//source/common/common:to_lower_table_lib",
],
)

envoy_cc_library(
name = "parser_lib",
hdrs = ["parser.h"],
deps = [
":kafka_protocol_lib",
"//source/common/common:minimal_logger_lib",
],
)

envoy_cc_library(
name = "serialization_lib",
hdrs = ["serialization.h"],
deps = [
":kafka_protocol_lib",
],
)

envoy_cc_library(
name = "kafka_protocol_lib",
hdrs = ["kafka_types.h", "kafka_protocol.h"],
external_deps = ["abseil_optional"],
deps = [
"//source/common/common:macros",
],
)
52 changes: 52 additions & 0 deletions source/extensions/filters/network/kafka/codec.cc
Original file line number Diff line number Diff line change
@@ -0,0 +1,52 @@
#include "extensions/filters/network/kafka/codec.h"

#include "extensions/filters/network/kafka/kafka_protocol.h"

namespace Envoy {
namespace Extensions {
namespace NetworkFilters {
namespace Kafka {

void RequestDecoder::onData(Buffer::Instance& data) {
uint64_t num_slices = data.getRawSlices(nullptr, 0);
Buffer::RawSlice slices[num_slices];
data.getRawSlices(slices, num_slices);
for (const Buffer::RawSlice& slice : slices) {
doParse(request_parser_, slice);
}
}

void RequestDecoder::doParse(ParserSharedPtr& parser, const Buffer::RawSlice& slice) {
const char* buffer = reinterpret_cast<const char*>(slice.mem_);
uint64_t remaining = slice.len_;
while (remaining) {
ParseResponse result = parser->parse(buffer, remaining);
// this loop guarantees that parsers consuming 0 bytes also get processed
while (result.hasData()) {
if (!result.next_parser_) {

// next parser is not present, so we have finished parsing a message
MessageSharedPtr message = result.message_;
ENVOY_LOG(trace, "parsed message: {}", *message);
for (auto& listener : listeners_) {
listener->onMessage(result.message_);
}

// we finished parsing this request, start anew
parser = std::make_shared<RequestStartParser>(parser_resolver_);
} else {
parser = result.next_parser_;
}
result = parser->parse(buffer, remaining);
}
}
}

void ResponseDecoder::onWrite(Buffer::Instance&) {
/* not implemented yet */
}

} // namespace Kafka
} // namespace NetworkFilters
} // namespace Extensions
} // namespace Envoy
49 changes: 49 additions & 0 deletions source/extensions/filters/network/kafka/codec.h
Original file line number Diff line number Diff line change
@@ -0,0 +1,49 @@
#pragma once

#include "extensions/filters/network/kafka/parser.h"
#include "extensions/filters/network/kafka/kafka_request.h"

#include "envoy/buffer/buffer.h"
#include "envoy/common/pure.h"

namespace Envoy {
namespace Extensions {
namespace NetworkFilters {
namespace Kafka {

class MessageListener {
public:
virtual ~MessageListener() {};

virtual void onMessage(MessageSharedPtr) PURE;
};

typedef std::shared_ptr<MessageListener> MessageListenerPtr;

class RequestDecoder : public Logger::Loggable<Logger::Id::kafka> {
public:
RequestDecoder(const RequestParserResolver parserResolver, const std::vector<MessageListenerPtr> listeners):
parser_resolver_{parserResolver},
listeners_{listeners},
request_parser_{new RequestStartParser(parser_resolver_)}
{};

void onData(Buffer::Instance& data);
private:
void doParse(ParserSharedPtr& parser, const Buffer::RawSlice& slice);

const RequestParserResolver parser_resolver_;
const std::vector<MessageListenerPtr> listeners_;

ParserSharedPtr request_parser_;
};

class ResponseDecoder {
public:
void onWrite(Buffer::Instance& data);
};

} // namespace Kafka
} // namespace NetworkFilters
} // namespace Extensions
} // namespace Envoy
53 changes: 53 additions & 0 deletions source/extensions/filters/network/kafka/config.cc
Original file line number Diff line number Diff line change
@@ -0,0 +1,53 @@
#include "extensions/filters/network/kafka/kafka_filter.h"
#include "extensions/filters/network/kafka/metrics_holder.h"
#include "extensions/filters/network/well_known_names.h"

#include "envoy/registry/registry.h"
#include "envoy/server/filter_config.h"
#include "envoy/stats/scope.h"

namespace Envoy {
namespace Extensions {
namespace NetworkFilters {
namespace Kafka {

/**
* Config registration for the Kafka filter. @see NamedNetworkFilterConfigFactory.
*/
class KafkaConfigFactory : public Server::Configuration::NamedNetworkFilterConfigFactory {
public:
// NamedNetworkFilterConfigFactory
Network::FilterFactoryCb
createFilterFactory(const Json::Object&, Server::Configuration::FactoryContext& context) override {
return createInternal(context.scope());
}

Network::FilterFactoryCb
createFilterFactoryFromProto(const Protobuf::Message&, Server::Configuration::FactoryContext& context) override {
return createInternal(context.scope());
}

ProtobufTypes::MessagePtr createEmptyConfigProto() override {
return ProtobufTypes::MessagePtr{new Envoy::ProtobufWkt::Empty()};
}

std::string name() override { return NetworkFilterNames::get().Kafka; }

private:
Network::FilterFactoryCb createInternal(Stats::Scope& scope) {
std::shared_ptr<MetricsHolder> metrics_holder = std::make_shared<MetricsHolder>(scope);
return [&scope, metrics_holder](Network::FilterManager& filter_manager) -> void {
filter_manager.addFilter(std::make_shared<KafkaFilter>(scope, metrics_holder));
};
}
};

/**
* Static registration for the Kafka filter. @see RegisterFactory.
*/
static Registry::RegisterFactory<KafkaConfigFactory, Server::Configuration::NamedNetworkFilterConfigFactory> registered_;

} // namespace Kafka
} // namespace NetworkFilters
} // namespace Extensions
} // namespace Envoy
56 changes: 56 additions & 0 deletions source/extensions/filters/network/kafka/kafka_filter.cc
Original file line number Diff line number Diff line change
@@ -0,0 +1,56 @@
#include "extensions/filters/network/kafka/kafka_filter.h"

#include "extensions/filters/network/kafka/kafka_request.h"
#include "extensions/filters/network/kafka/kafka_response.h"

#include "envoy/buffer/buffer.h"
#include "envoy/network/connection.h"
#include "common/common/assert.h"

#include <sstream>

namespace Envoy {
namespace Extensions {
namespace NetworkFilters {
namespace Kafka {

class LoggingMessageListener : public MessageListener, public Logger::Loggable<Logger::Id::kafka> {
public:
void onMessage(MessageSharedPtr arg) override {
ENVOY_LOG(info, "received: {}", *arg);
}
};

KafkaFilter::KafkaFilter(Stats::Scope& scope, std::shared_ptr<MetricsHolder> metrics_holder):
stats_{KAFKA_STATS(POOL_COUNTER_PREFIX(scope, "kafka."))},
metrics_holder_{metrics_holder},
request_decoder_{
RequestParserResolver::KAFKA_0_11,
{ std::make_unique<LoggingMessageListener>() }
}
{
stats_.filters_created_.inc();
};

Network::FilterStatus KafkaFilter::onNewConnection() {
return Network::FilterStatus::Continue;
}

void KafkaFilter::initializeReadFilterCallbacks(Network::ReadFilterCallbacks& callbacks) {
read_callbacks_ = &callbacks;
}

Network::FilterStatus KafkaFilter::onData(Buffer::Instance& data, bool) {
request_decoder_.onData(data);
return Network::FilterStatus::Continue;
}

Network::FilterStatus KafkaFilter::onWrite(Buffer::Instance& data, bool) {
response_decoder_.onWrite(data);
return Network::FilterStatus::Continue;
}

} // namespace Kafka
} // namespace NetworkFilters
} // namespace Extensions
} // namespace Envoy
Loading