Interface FilterInvoker

All Known Implementing Classes:
AddOffsetsToTxnRequestFilterInvoker, AddOffsetsToTxnResponseFilterInvoker, AddPartitionsToTxnRequestFilterInvoker, AddPartitionsToTxnResponseFilterInvoker, AddRaftVoterRequestFilterInvoker, AddRaftVoterResponseFilterInvoker, AllocateProducerIdsRequestFilterInvoker, AllocateProducerIdsResponseFilterInvoker, AlterClientQuotasRequestFilterInvoker, AlterClientQuotasResponseFilterInvoker, AlterConfigsRequestFilterInvoker, AlterConfigsResponseFilterInvoker, AlterPartitionReassignmentsRequestFilterInvoker, AlterPartitionReassignmentsResponseFilterInvoker, AlterPartitionRequestFilterInvoker, AlterPartitionResponseFilterInvoker, AlterReplicaLogDirsRequestFilterInvoker, AlterReplicaLogDirsResponseFilterInvoker, AlterShareGroupOffsetsRequestFilterInvoker, AlterShareGroupOffsetsResponseFilterInvoker, AlterUserScramCredentialsRequestFilterInvoker, AlterUserScramCredentialsResponseFilterInvoker, ApiVersionsRequestFilterInvoker, ApiVersionsResponseFilterInvoker, AssignReplicasToDirsRequestFilterInvoker, AssignReplicasToDirsResponseFilterInvoker, BeginQuorumEpochRequestFilterInvoker, BeginQuorumEpochResponseFilterInvoker, BrokerHeartbeatRequestFilterInvoker, BrokerHeartbeatResponseFilterInvoker, BrokerRegistrationRequestFilterInvoker, BrokerRegistrationResponseFilterInvoker, ConsumerGroupDescribeRequestFilterInvoker, ConsumerGroupDescribeResponseFilterInvoker, ConsumerGroupHeartbeatRequestFilterInvoker, ConsumerGroupHeartbeatResponseFilterInvoker, ControllerRegistrationRequestFilterInvoker, ControllerRegistrationResponseFilterInvoker, CreateAclsRequestFilterInvoker, CreateAclsResponseFilterInvoker, CreateDelegationTokenRequestFilterInvoker, CreateDelegationTokenResponseFilterInvoker, CreatePartitionsRequestFilterInvoker, CreatePartitionsResponseFilterInvoker, CreateTopicsRequestFilterInvoker, CreateTopicsResponseFilterInvoker, DeleteAclsRequestFilterInvoker, DeleteAclsResponseFilterInvoker, DeleteGroupsRequestFilterInvoker, DeleteGroupsResponseFilterInvoker, DeleteRecordsRequestFilterInvoker, DeleteRecordsResponseFilterInvoker, DeleteShareGroupOffsetsRequestFilterInvoker, DeleteShareGroupOffsetsResponseFilterInvoker, DeleteShareGroupStateRequestFilterInvoker, DeleteShareGroupStateResponseFilterInvoker, DeleteTopicsRequestFilterInvoker, DeleteTopicsResponseFilterInvoker, DescribeAclsRequestFilterInvoker, DescribeAclsResponseFilterInvoker, DescribeClientQuotasRequestFilterInvoker, DescribeClientQuotasResponseFilterInvoker, DescribeClusterRequestFilterInvoker, DescribeClusterResponseFilterInvoker, DescribeConfigsRequestFilterInvoker, DescribeConfigsResponseFilterInvoker, DescribeDelegationTokenRequestFilterInvoker, DescribeDelegationTokenResponseFilterInvoker, DescribeGroupsRequestFilterInvoker, DescribeGroupsResponseFilterInvoker, DescribeLogDirsRequestFilterInvoker, DescribeLogDirsResponseFilterInvoker, DescribeProducersRequestFilterInvoker, DescribeProducersResponseFilterInvoker, DescribeQuorumRequestFilterInvoker, DescribeQuorumResponseFilterInvoker, DescribeShareGroupOffsetsRequestFilterInvoker, DescribeShareGroupOffsetsResponseFilterInvoker, DescribeTopicPartitionsRequestFilterInvoker, DescribeTopicPartitionsResponseFilterInvoker, DescribeTransactionsRequestFilterInvoker, DescribeTransactionsResponseFilterInvoker, DescribeUserScramCredentialsRequestFilterInvoker, DescribeUserScramCredentialsResponseFilterInvoker, ElectLeadersRequestFilterInvoker, ElectLeadersResponseFilterInvoker, EndQuorumEpochRequestFilterInvoker, EndQuorumEpochResponseFilterInvoker, EndTxnRequestFilterInvoker, EndTxnResponseFilterInvoker, EnvelopeRequestFilterInvoker, EnvelopeResponseFilterInvoker, ExpireDelegationTokenRequestFilterInvoker, ExpireDelegationTokenResponseFilterInvoker, FetchRequestFilterInvoker, FetchResponseFilterInvoker, FetchSnapshotRequestFilterInvoker, FetchSnapshotResponseFilterInvoker, FindCoordinatorRequestFilterInvoker, FindCoordinatorResponseFilterInvoker, GetTelemetrySubscriptionsRequestFilterInvoker, GetTelemetrySubscriptionsResponseFilterInvoker, HandleNothingFilterInvoker, HeartbeatRequestFilterInvoker, HeartbeatResponseFilterInvoker, IncrementalAlterConfigsRequestFilterInvoker, IncrementalAlterConfigsResponseFilterInvoker, InitializeShareGroupStateRequestFilterInvoker, InitializeShareGroupStateResponseFilterInvoker, InitProducerIdRequestFilterInvoker, InitProducerIdResponseFilterInvoker, JoinGroupRequestFilterInvoker, JoinGroupResponseFilterInvoker, LeaveGroupRequestFilterInvoker, LeaveGroupResponseFilterInvoker, ListConfigResourcesRequestFilterInvoker, ListConfigResourcesResponseFilterInvoker, ListGroupsRequestFilterInvoker, ListGroupsResponseFilterInvoker, ListOffsetsRequestFilterInvoker, ListOffsetsResponseFilterInvoker, ListPartitionReassignmentsRequestFilterInvoker, ListPartitionReassignmentsResponseFilterInvoker, ListTransactionsRequestFilterInvoker, ListTransactionsResponseFilterInvoker, MetadataRequestFilterInvoker, MetadataResponseFilterInvoker, OffsetCommitRequestFilterInvoker, OffsetCommitResponseFilterInvoker, OffsetDeleteRequestFilterInvoker, OffsetDeleteResponseFilterInvoker, OffsetFetchRequestFilterInvoker, OffsetFetchResponseFilterInvoker, OffsetForLeaderEpochRequestFilterInvoker, OffsetForLeaderEpochResponseFilterInvoker, ProduceRequestFilterInvoker, ProduceResponseFilterInvoker, PushTelemetryRequestFilterInvoker, PushTelemetryResponseFilterInvoker, ReadShareGroupStateRequestFilterInvoker, ReadShareGroupStateResponseFilterInvoker, ReadShareGroupStateSummaryRequestFilterInvoker, ReadShareGroupStateSummaryResponseFilterInvoker, RemoveRaftVoterRequestFilterInvoker, RemoveRaftVoterResponseFilterInvoker, RenewDelegationTokenRequestFilterInvoker, RenewDelegationTokenResponseFilterInvoker, RequestFilterInvoker, RequestResponseInvoker, ResponseFilterInvoker, SafeInvoker, SaslAuthenticateRequestFilterInvoker, SaslAuthenticateResponseFilterInvoker, SaslHandshakeRequestFilterInvoker, SaslHandshakeResponseFilterInvoker, ShareAcknowledgeRequestFilterInvoker, ShareAcknowledgeResponseFilterInvoker, ShareFetchRequestFilterInvoker, ShareFetchResponseFilterInvoker, ShareGroupDescribeRequestFilterInvoker, ShareGroupDescribeResponseFilterInvoker, ShareGroupHeartbeatRequestFilterInvoker, ShareGroupHeartbeatResponseFilterInvoker, SpecificFilterArrayInvoker, StreamsGroupDescribeRequestFilterInvoker, StreamsGroupDescribeResponseFilterInvoker, StreamsGroupHeartbeatRequestFilterInvoker, StreamsGroupHeartbeatResponseFilterInvoker, SyncGroupRequestFilterInvoker, SyncGroupResponseFilterInvoker, TxnOffsetCommitRequestFilterInvoker, TxnOffsetCommitResponseFilterInvoker, UnregisterBrokerRequestFilterInvoker, UnregisterBrokerResponseFilterInvoker, UpdateFeaturesRequestFilterInvoker, UpdateFeaturesResponseFilterInvoker, UpdateRaftVoterRequestFilterInvoker, UpdateRaftVoterResponseFilterInvoker, VoteRequestFilterInvoker, VoteResponseFilterInvoker, WriteShareGroupStateRequestFilterInvoker, WriteShareGroupStateResponseFilterInvoker, WriteTxnMarkersRequestFilterInvoker, WriteTxnMarkersResponseFilterInvoker

public interface FilterInvoker
The FilterInvoker connects Kroxylicious with the concrete implementation of a Filter.

When handling a message, we want to avoid the penalty of deserializing the bytes into an ApiMessage. When Kroxylicious receives a message, all the Filters in the filter chain will be consulted (via their invoker) to see if any want to handle that message. If any filter wants to handle it, then the message will be deserialized. Then onRequest|onResponse will be eligible to be called in the filter chain for that message (i.e. if filter A wants to handle request Y but filter B doesn't, only the onRequest of A would be invoked)

Guarantees

Implementors of this API may assume the following:

  1. That each instance of the FilterInvoker is associated with a single channel
  2. That shouldHandleRequest(ApiKeys, short) and onRequest(ApiKeys, short, RequestHeaderData, ApiMessage, FilterContext) (or on*Request as appropriate) will always be invoked on the same thread.
  3. That filters are applied in the order they were configured.

From 1. and 2. it follows that you can use member variables in your filter to store channel-local state.

Implementors should not assume:

  1. That invokers in the same chain execute on the same thread. Thus inter-filter communication/state transfer needs to be thread-safe
  • Method Summary

    Modifier and Type
    Method
    Description
    default CompletionStage<io.kroxylicious.proxy.filter.RequestFilterResult>
    onRequest(org.apache.kafka.common.protocol.ApiKeys apiKey, short apiVersion, org.apache.kafka.common.message.RequestHeaderData header, org.apache.kafka.common.protocol.ApiMessage body, io.kroxylicious.proxy.filter.FilterContext filterContext)
    Handle deserialized request data.
    default CompletionStage<io.kroxylicious.proxy.filter.ResponseFilterResult>
    onResponse(org.apache.kafka.common.protocol.ApiKeys apiKey, short apiVersion, org.apache.kafka.common.message.ResponseHeaderData header, org.apache.kafka.common.protocol.ApiMessage body, io.kroxylicious.proxy.filter.FilterContext filterContext)
    Handle deserialized response data.
    default boolean
    shouldHandleRequest(org.apache.kafka.common.protocol.ApiKeys apiKey, short apiVersion)
    Should this Invoker handle a request with a given api key and api version.
    default boolean
    shouldHandleResponse(org.apache.kafka.common.protocol.ApiKeys apiKey, short apiVersion)
    Should this Invoker handle a response with a given api key and api version.
  • Method Details

    • shouldHandleRequest

      default boolean shouldHandleRequest(org.apache.kafka.common.protocol.ApiKeys apiKey, short apiVersion)
      Should this Invoker handle a request with a given api key and api version. Returning true implies that this request will be deserialized and onRequest(ApiKeys, short, RequestHeaderData, ApiMessage, FilterContext) is eligible to be called with the deserialized data (if the message flows to that filter).
      Parameters:
      apiKey - the key of the message
      apiVersion - the version of the message
      Returns:
      true if the message should be deserialized
    • shouldHandleResponse

      default boolean shouldHandleResponse(org.apache.kafka.common.protocol.ApiKeys apiKey, short apiVersion)
      Should this Invoker handle a response with a given api key and api version. Returning true implies that this response will be deserialized and onResponse(ApiKeys, short, ResponseHeaderData, ApiMessage, FilterContext) is eligible to be called with the deserialized data (if the message flows to that filter).
      Parameters:
      apiKey - the key of the message
      apiVersion - the version of the message
      Returns:
      true if the message should be deserialized
    • onRequest

      default CompletionStage<io.kroxylicious.proxy.filter.RequestFilterResult> onRequest(org.apache.kafka.common.protocol.ApiKeys apiKey, short apiVersion, org.apache.kafka.common.message.RequestHeaderData header, org.apache.kafka.common.protocol.ApiMessage body, io.kroxylicious.proxy.filter.FilterContext filterContext)

      Handle deserialized request data. Implementations must tolerate being called with requests that this FilterInvoker is NOT interested in handling. If any FilterInvoker in the chain should handle a response, then all invokers in the chain are eligible to have their onRequest called. Implementations should forward requests they do not wish to operate on.

      Filters must return a CompletionStage<io.kroxylicious.proxy.filter.RequestFilterResult> object. This object encapsulates the request to be forwarded and, optionally, orders for actions such as closing the connection or dropping the request.

      The FilterContext is the factory for FilterResult objects. See FilterContext.forwardRequest(RequestHeaderData, ApiMessage) and FilterContext.requestFilterResultBuilder() for more details.

      Parameters:
      apiKey - the key of the message
      apiVersion - the apiVersion of the message
      header - the header of the message
      body - the body of the message
      filterContext - contains methods to continue the filter chain and other contextual data
      Returns:
      a CompletionStage<io.kroxylicious.proxy.filter.RequestFilterResult>, that when complete, will yield the request to forward.
    • onResponse

      default CompletionStage<io.kroxylicious.proxy.filter.ResponseFilterResult> onResponse(org.apache.kafka.common.protocol.ApiKeys apiKey, short apiVersion, org.apache.kafka.common.message.ResponseHeaderData header, org.apache.kafka.common.protocol.ApiMessage body, io.kroxylicious.proxy.filter.FilterContext filterContext)

      Handle deserialized response data. Implementations must tolerate being called with responses that this FilterInvoker is NOT interested in handling. If any FilterInvoker in the chain should handle a response, then all invokers in the chain are eligible to have their onResponse called. Implementations should forward responses they do not wish to operate on.

      Filters must return a CompletionStage<io.kroxylicious.proxy.filter.ResponseFilterResult> object. This object encapsulates the response to be forwarded and, optionally, orders for actions such as closing the connection or dropping the response.

      The FilterContext is the factory for FilterResult objects. See FilterContext.forwardResponse(ResponseHeaderData, ApiMessage) and FilterContext.responseFilterResultBuilder() for more details.

      Parameters:
      apiKey - the key of the message
      apiVersion - the apiVersion of the message
      header - the header of the message
      body - the body of the message
      filterContext - contains methods to continue the filter chain and other contextual data
      Returns:
      a CompletionStage<io.kroxylicious.proxy.filter.ResponseFilterResult>, that when complete, will yield the response to forward.