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
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:
- That each instance of the FilterInvoker is associated with a single channel
- That
shouldHandleRequest(ApiKeys, short)andonRequest(ApiKeys, short, RequestHeaderData, ApiMessage, FilterContext)(oron*Requestas appropriate) will always be invoked on the same thread. - 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:
- 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 TypeMethodDescriptiondefault 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 booleanshouldHandleRequest(org.apache.kafka.common.protocol.ApiKeys apiKey, short apiVersion) Should this Invoker handle a request with a given api key and api version.default booleanshouldHandleResponse(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 andonRequest(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 messageapiVersion- 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 andonResponse(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 messageapiVersion- 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
FilterContextis the factory forFilterResultobjects. SeeFilterContext.forwardRequest(RequestHeaderData, ApiMessage)andFilterContext.requestFilterResultBuilder()for more details.- Parameters:
apiKey- the key of the messageapiVersion- the apiVersion of the messageheader- the header of the messagebody- the body of the messagefilterContext- 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
FilterContextis the factory forFilterResultobjects. SeeFilterContext.forwardResponse(ResponseHeaderData, ApiMessage)andFilterContext.responseFilterResultBuilder()for more details.- Parameters:
apiKey- the key of the messageapiVersion- the apiVersion of the messageheader- the header of the messagebody- the body of the messagefilterContext- 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.
-