All Classes and Interfaces
Class
Description
A thread-safe state that ensures exactly-once increment and decrement operations.
An invoker for AddOffsetsToTxnRequestFilter.
An invoker for AddOffsetsToTxnResponseFilter.
An invoker for AddPartitionsToTxnRequestFilter.
An invoker for AddPartitionsToTxnResponseFilter.
An invoker for AddRaftVoterRequestFilter.
An invoker for AddRaftVoterResponseFilter.
An invoker for AllocateProducerIdsRequestFilter.
An invoker for AllocateProducerIdsResponseFilter.
An invoker for AlterClientQuotasRequestFilter.
An invoker for AlterClientQuotasResponseFilter.
An invoker for AlterConfigsRequestFilter.
An invoker for AlterConfigsResponseFilter.
An invoker for AlterPartitionReassignmentsRequestFilter.
An invoker for AlterPartitionReassignmentsResponseFilter.
An invoker for AlterPartitionRequestFilter.
An invoker for AlterPartitionResponseFilter.
An invoker for AlterReplicaLogDirsRequestFilter.
An invoker for AlterReplicaLogDirsResponseFilter.
An invoker for AlterShareGroupOffsetsRequestFilter.
An invoker for AlterShareGroupOffsetsResponseFilter.
An invoker for AlterUserScramCredentialsRequestFilter.
An invoker for AlterUserScramCredentialsResponseFilter.
Changes an API_VERSIONS response so that a client sees the intersection of supported version ranges for each
API key.
An invoker for ApiVersionsRequestFilter.
An invoker for ApiVersionsResponseFilter.
An invoker for AssignReplicasToDirsRequestFilter.
An invoker for AssignReplicasToDirsResponseFilter.
Thrown when TLS credentials are invalid — for example, when a private key does not match
the certificate, or when the certificate chain is structurally incorrect.
Ancient versions of Kafka implemented SASL/GSSAPI by sending the token
on the wire as length prefixed bytes (no Kafka protocol header).
Ancient versions of Kafka implemented SASL/GSSAPI by sending the response
on the wire as length prefixed bytes (no Kafka protocol header).
An invoker for BeginQuorumEpochRequestFilter.
An invoker for BeginQuorumEpochResponseFilter.
Bijective node ID mapping:
V = id + S × t, where
V is the virtual ID, t is the target node ID,
S is the total number of routes, and id is
the route's configured identifier.Decodes Kafka Readable into an ApiMessage
A bootstrap binding.
Strategy for selecting an upstream target from a given list of upstream targets for bootstrapping.
An internal filter that rewrites broker addresses in all relevant responses to the corresponding proxy address.
A broker specific endpoint binding.
An invoker for BrokerHeartbeatRequestFilter.
An invoker for BrokerHeartbeatResponseFilter.
An invoker for BrokerRegistrationRequestFilter.
An invoker for BrokerRegistrationResponseFilter.
Provides read and write access to byte buffer for serializing frames.
An implementation of Kafka's Readable and Writable abstraction in terms of
a Netty ByteBuf.
This class has been introduced as a work-around to allow using pooled
ByteBuf instances
that are allowed to grow on demand while used on MemoryRecordsBuilder
to create records (using MemoryRecordsHelper factory methods).Cache Configuration
The state machine for a single client's proxy session.
Enumeration of disconnect causes for tracking client to proxy disconnections.
A named cluster definition, referenced by routes and virtual clusters.
Exceptional-completion cause for
KafkaProxy.reconfigure(Configuration) when a
second reconfigure call arrives while another is in progress.The root of the proxy configuration.
Internal class that owns the
KafkaProxy.reconfigure(Configuration) pipeline.An invoker for ConsumerGroupDescribeRequestFilter.
An invoker for ConsumerGroupDescribeResponseFilter.
An invoker for ConsumerGroupHeartbeatRequestFilter.
An invoker for ConsumerGroupHeartbeatResponseFilter.
An invoker for ControllerRegistrationRequestFilter.
An invoker for ControllerRegistrationResponseFilter.
Allocates integer correlation IDs from a bounded circular range
[minInc, maxExc).Manages correlation ids for a single connection (across the proxy) between a single client
and a single broker.
A record for which responses should be decoded, together with their
API key and version.
An invoker for CreateAclsRequestFilter.
An invoker for CreateAclsResponseFilter.
An invoker for CreateDelegationTokenRequestFilter.
An invoker for CreateDelegationTokenResponseFilter.
An invoker for CreatePartitionsRequestFilter.
An invoker for CreatePartitionsResponseFilter.
An invoker for CreateTopicsRequestFilter.
An invoker for CreateTopicsResponseFilter.
A frame that has been decoded (as opposed to an
OpaqueFrame).A decoded request frame.
A decoded response frame.
Encapsulates decisions about whether requests and responses should be
fully deserialized into POJOs, or passed through as byte buffers with
minimal deserialization.
An invoker for DeleteAclsRequestFilter.
An invoker for DeleteAclsResponseFilter.
An invoker for DeleteGroupsRequestFilter.
An invoker for DeleteGroupsResponseFilter.
An invoker for DeleteRecordsRequestFilter.
An invoker for DeleteRecordsResponseFilter.
An invoker for DeleteShareGroupOffsetsRequestFilter.
An invoker for DeleteShareGroupOffsetsResponseFilter.
An invoker for DeleteShareGroupStateRequestFilter.
An invoker for DeleteShareGroupStateResponseFilter.
An invoker for DeleteTopicsRequestFilter.
An invoker for DeleteTopicsResponseFilter.
This class has been built to filter cipher suites based on requirements for Kroxylicious
If no ciphers are declared then the default platform ciphers will be used
If ciphers are declared then they are used instead of the default platform ciphers
Both lists would then have anything removed that wasn't supported or in the deniedCiphers
An invoker for DescribeAclsRequestFilter.
An invoker for DescribeAclsResponseFilter.
An invoker for DescribeClientQuotasRequestFilter.
An invoker for DescribeClientQuotasResponseFilter.
An invoker for DescribeClusterRequestFilter.
An invoker for DescribeClusterResponseFilter.
An invoker for DescribeConfigsRequestFilter.
An invoker for DescribeConfigsResponseFilter.
An invoker for DescribeDelegationTokenRequestFilter.
An invoker for DescribeDelegationTokenResponseFilter.
An invoker for DescribeGroupsRequestFilter.
An invoker for DescribeGroupsResponseFilter.
An invoker for DescribeLogDirsRequestFilter.
An invoker for DescribeLogDirsResponseFilter.
An invoker for DescribeProducersRequestFilter.
An invoker for DescribeProducersResponseFilter.
An invoker for DescribeQuorumRequestFilter.
An invoker for DescribeQuorumResponseFilter.
An invoker for DescribeShareGroupOffsetsRequestFilter.
An invoker for DescribeShareGroupOffsetsResponseFilter.
An invoker for DescribeTopicPartitionsRequestFilter.
An invoker for DescribeTopicPartitionsResponseFilter.
An invoker for DescribeTransactionsRequestFilter.
An invoker for DescribeTransactionsResponseFilter.
An invoker for DescribeUserScramCredentialsRequestFilter.
An invoker for DescribeUserScramCredentialsResponseFilter.
Routing model for a virtual cluster that forwards directly to a single, statically-configured
upstream Kafka cluster.
Thrown if there is some arithmetic exception constructing or operating on a
DurationRouting model for a virtual cluster that forwards to one or more upstream clusters via a named
router plugin.
An internal filter that causes the system to eagerly learn the cluster's topology by spontaneously emitting
an out-of-band Metadata request at the earliest legal point in the Kafka conversation.
An invoker for ElectLeadersRequestFilter.
An invoker for ElectLeadersResponseFilter.
Represents a network endpoint.
An endpoint binding.
Signals that an endpoint could not be bound due to an error condition.
Used by the
KafkaProxyInitializer to resolve incoming channel
metadata into a EndpointBinding.This class is the general class of exceptions produced by failed endpoint operations.
A gateway to an endpoint.
The endpoint registry is responsible for associating network endpoints with broker/bootstrap addresses of virtual clusters.
Signals that an endpoint could not be resolved into a known virtual cluster binding.
An invoker for EndQuorumEpochRequestFilter.
An invoker for EndQuorumEpochResponseFilter.
An invoker for EndTxnRequestFilter.
An invoker for EndTxnResponseFilter.
An invoker for EnvelopeRequestFilter.
An invoker for EnvelopeResponseFilter.
An invoker for ExpireDelegationTokenRequestFilter.
An invoker for ExpireDelegationTokenResponseFilter.
Configuration Features related to how we load configuration.
Represents the entire set of proxy features.
An invoker for FetchRequestFilter.
An invoker for FetchResponseFilter.
An invoker for FetchSnapshotRequestFilter.
An invoker for FetchSnapshotResponseFilter.
A Filter and it's respective invoker
Builds per-connection filter instances for a single virtual cluster's filter chain.
A
ChannelInboundHandler (for handling requests from downstream)
that applies a single Filter.The FilterInvoker connects Kroxylicious with the concrete implementation of a Filter.
Factory for FilterInvokers.
An invoker for FindCoordinatorRequestFilter.
An invoker for FindCoordinatorResponseFilter.
A frame in the Kafka protocol, which may or may not be fully decoded.
An invoker for GetTelemetrySubscriptionsRequestFilter.
An invoker for GetTelemetrySubscriptionsResponseFilter.
Immutable snapshot of a decoded PROXY protocol header.
A channel handler that intercepts
HAProxyMessage objects emitted by
Netty's HAProxyMessageDecoder and
stores the extracted context in the KafkaSession.Uses
HAProxyMessageDecoder.detectProtocol(ByteBuf) on the first bytes
of a connection to decide whether to install the PROXY protocol decoder.An invoker for HeartbeatRequestFilter.
An invoker for HeartbeatResponseFilter.
Represents a host port pair.
Identity node ID mapping for single-route configurations.
Signals that the configuration is syntactically/semantically correct but illegal
by some policy.
An invoker for IncrementalAlterConfigsRequestFilter.
An invoker for IncrementalAlterConfigsResponseFilter.
An invoker for InitializeShareGroupStateRequestFilter.
An invoker for InitializeShareGroupStateResponseFilter.
An invoker for InitProducerIdRequestFilter.
An invoker for InitProducerIdResponseFilter.
An invoker for JoinGroupRequestFilter.
An invoker for JoinGroupResponseFilter.
Callback invoked by the
KafkaMessageEncoder and KafkaMessageDecoder
on encode or decode of each message.In the operation of the proxy there are various exceptions which are "anticipated" but not necessarily handled directly.
Describes the possible states of a Kafka session as viewed from the proxy
Generally we expect a Kafka Session to transition through one of the following sequences.
An invoker for LeaveGroupRequestFilter.
An invoker for LeaveGroupResponseFilter.
An exception thrown during a lifecycle transition such as startup or shutdown
of the proxy.
An invoker for ListConfigResourcesRequestFilter.
An invoker for ListConfigResourcesResponseFilter.
An invoker for ListGroupsRequestFilter.
An invoker for ListGroupsResponseFilter.
An invoker for ListOffsetsRequestFilter.
An invoker for ListOffsetsResponseFilter.
An invoker for ListPartitionReassignmentsRequestFilter.
An invoker for ListPartitionReassignmentsResponseFilter.
An invoker for ListTransactionsRequestFilter.
An invoker for ListTransactionsResponseFilter.
Management configuration.
This introduces additional factory builder methods for
MemoryRecords that
accepts ByteBufOutputStreamA bootstrap binding which can only be used for metadata discovery.
An invoker for MetadataRequestFilter.
An invoker for MetadataResponseFilter.
A kafka message listener that emits message count and message size
metrics.
Service to build a MicrometerConfigurationHook for a configuration.
A named filter definition
Represents the inclusive set of integers between two integer endpoints.
Abstract encapsulation of a network binding operation.
Processes
NetworkBindingOperations on a suitable ServerBootstrap.Request for a network endpoint to be bound.
Request for a network endpoint to be unbound.
This is the Strategy for how we expose a virtual kafka cluster on the network.
Maps between target-cluster node IDs and the virtual node IDs
presented to clients.
A route name and target-cluster node ID pair.
An invoker for OffsetCommitRequestFilter.
An invoker for OffsetCommitResponseFilter.
An invoker for OffsetDeleteRequestFilter.
An invoker for OffsetDeleteResponseFilter.
An invoker for OffsetFetchRequestFilter.
An invoker for OffsetFetchResponseFilter.
An invoker for OffsetForLeaderEpochRequestFilter.
An invoker for OffsetForLeaderEpochResponseFilter.
A frame in the Kafka protocol which has not been decoded.
Used to represent Kafka requests that the proxy does not need to decode.
Used to represent Kafka responses that the proxy does not need to decode.
Problems with plugin interfaces implementations.
A PluginFactory is able to resolve references to a plugin implementation (i.e. a name) to that implementation.
Detects potential for port conflicts arising between virtual cluster configurations.
A NodeIdentificationStrategy implementation that uses a separate port per broker endpoint and that is aware of
distinct ranges of nodeIds present in the target cluster.
Configuration for a principal adder, which is responsible for contributing zero or more principals to the subject.
An invoker for ProduceRequestFilter.
An invoker for ProduceResponseFilter.
Configuration for the HaProxy PROXY protocol.
Controls how Kroxylicious handles HaProxy PROXY protocol headers on incoming connections.
An invoker for PushTelemetryRequestFilter.
An invoker for PushTelemetryResponseFilter.
BootstrapSelectionStrategy which selects a random server from the given list of servers as the bootstrap server.Represents the set of integers between two integer endpoints.
An invoker for ReadShareGroupStateRequestFilter.
An invoker for ReadShareGroupStateResponseFilter.
An invoker for ReadShareGroupStateSummaryRequestFilter.
An invoker for ReadShareGroupStateSummaryResponseFilter.
One per-component failure encountered during a
KafkaProxy.reconfigure(Configuration)
call.Outcome of a
KafkaProxy.reconfigure(Configuration) call.An invoker for RemoveRaftVoterRequestFilter.
An invoker for RemoveRaftVoterResponseFilter.
An invoker for RenewDelegationTokenRequestFilter.
An invoker for RenewDelegationTokenResponseFilter.
Utility for tagging message headers with internal tags.
While processing a request from a Client, we want to enable custom Protocol
Filters to decide to send a response toward the Client instead of forwarding
the message on towards the upstream broker.
BootstrapSelectionStrategy that selects a server from the given list of servers as the bootstrap server in a round-robin fashion.A route within a router definition.
Runtime representation of a resolved route within a router.
Abstracts the creation of router instances, hiding the configuration
required for instantiation at the point at which instances are created.
A named router definition referencing a
RouterFactory plugin type.Sits at the end of the VC-level filter chain (instead of
FilterChainCompletionHandler) when a
virtual cluster uses a router.A target for a route: exactly one of
cluster or router must be specified.Sealed hierarchy representing how a virtual cluster reaches its upstream Kafka cluster(s).
Wraps a delegate invoker so that onRequest and onResponse can be safely called even if this
Invoker does not want to handle this message, in this case the message will be forwarded without
the delegate doing anything with it.
An invoker for SaslAuthenticateRequestFilter.
An invoker for SaslAuthenticateResponseFilter.
An invoker for SaslHandshakeRequestFilter.
An invoker for SaslHandshakeResponseFilter.
Short circuit respond with an error code if we encounter a SASL v0 handshake, logging a meaningful warning and closing
the connection.
Runtime implementation of ServerTlsCredentialSupplierContext.
A
PluginFactoryRegistry that is implemented using ServiceLoader discovery.An invoker for ShareAcknowledgeRequestFilter.
An invoker for ShareAcknowledgeResponseFilter.
An invoker for ShareFetchRequestFilter.
An invoker for ShareFetchResponseFilter.
An invoker for ShareGroupDescribeRequestFilter.
An invoker for ShareGroupDescribeResponseFilter.
An invoker for ShareGroupHeartbeatRequestFilter.
An invoker for ShareGroupHeartbeatResponseFilter.
A NodeIdentificationStrategy implementation that binds to a single, shared, port for bootstrap and
all brokers.
Invoker for Filters that implement any number of Specific Message interfaces (for
example
AlterConfigsResponseFilter.Exceptional-completion cause for
KafkaProxy.reconfigure(Configuration) when the
submitted configuration differs from the running configuration in any static
section — i.e. a section the proxy does not reconcile at runtime.An invoker for StreamsGroupDescribeRequestFilter.
An invoker for StreamsGroupDescribeResponseFilter.
An invoker for StreamsGroupHeartbeatRequestFilter.
An invoker for StreamsGroupHeartbeatResponseFilter.
An invoker for SyncGroupRequestFilter.
An invoker for SyncGroupResponseFilter.
Represents the target (upstream) kafka cluster.
Runtime implementation of TlsCredentials containing the actual private key and certificate chain.
Manages the lifecycle of TLS credential supplier plugin instances.
Utility class for TLS credential validation.
A Filter that learns and caches all topic names, it is responsible for short circuit responding to internal
topic name retrievals.
An invoker for TxnOffsetCommitRequestFilter.
An invoker for TxnOffsetCommitResponseFilter.
An invoker for UnregisterBrokerRequestFilter.
An invoker for UnregisterBrokerResponseFilter.
An invoker for UpdateFeaturesRequestFilter.
An invoker for UpdateFeaturesResponseFilter.
An invoker for UpdateRaftVoterRequestFilter.
An invoker for UpdateRaftVoterResponseFilter.
Runtime representation of an upstream Kafka cluster, bundling its connection target with the
TLS resources needed to reach it.
A virtual cluster.
A virtual cluster listener.
Manages the lifecycle state of a single virtual cluster.
Lifecycle states for a virtual cluster.
New connections are rejected.
Configuration was not viable.
The cluster is being set up.
The proxy has completed setup and is serving traffic for this cluster.
Terminal state.
Runtime representation of a virtual cluster: its name, target Kafka cluster, gateways,
TLS configuration, and the components whose lifecycle is bound to this VC.
Represents a node in a virtual cluster, identified by its cluster name and an optional node ID.
Owns the virtual cluster configuration tree and lifecycle state.
An invoker for VoteRequestFilter.
An invoker for VoteResponseFilter.
An invoker for WriteShareGroupStateRequestFilter.
An invoker for WriteShareGroupStateResponseFilter.
An invoker for WriteTxnMarkersRequestFilter.
An invoker for WriteTxnMarkersResponseFilter.