Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
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
33 changes: 20 additions & 13 deletions stubs/kafka-python/@tests/stubtest_allowlist.txt
Original file line number Diff line number Diff line change
Expand Up @@ -7,18 +7,25 @@ kafka.producer.__main__
# Benchmark modules are not included in type stubs.
kafka.benchmarks.*

# Concrete subclasses define these abstract properties as class attributes.
kafka.protocol.api.Request.API_KEY
kafka.protocol.api.Request.API_VERSION
kafka.protocol.api.Request.RESPONSE_TYPE
kafka.protocol.api.Request.SCHEMA
kafka.protocol.api.Response.API_KEY
kafka.protocol.api.Response.API_VERSION
kafka.protocol.api.Response.SCHEMA

# Vendored compatibility modules are implementation details.
kafka.vendor
kafka.vendor.enum34
kafka.vendor.selectors34
kafka.vendor.six
kafka.vendor.socketpair

# Undocumented and clearly not meant to be exposed.
kafka.coordinator.base.heartbeat_log
kafka.net.selector.log_trace
kafka.*\.log
kafka.*\.logger
kafka.*\.__license__

# The max_version param expects an int, even though it has as default value of float('inf').
kafka.protocol.broker_version_data.BrokerVersionData.api_version

# TODO: This has more attributes.
kafka.protocol.api_header.RequestHeader.client_id
kafka.protocol.api_header.RequestHeader.correlation_id
kafka.protocol.api_header.RequestHeader.request_api_key
kafka.protocol.api_header.RequestHeader.request_api_version
kafka.protocol.api_header.ResponseHeader.correlation_id
kafka.coordinator.assignors.sticky.user_data.StickyAssignorUserData.TopicPartition
kafka.coordinator.assignors.sticky.user_data.StickyAssignorUserData.generation
kafka.coordinator.assignors.sticky.user_data.StickyAssignorUserData.previous_assignment
2 changes: 1 addition & 1 deletion stubs/kafka-python/METADATA.toml
Original file line number Diff line number Diff line change
@@ -1,2 +1,2 @@
version = "2.3.2"
version = "3.0.9"
upstream-repository = "https://github.com/dpkp/kafka-python"
43 changes: 34 additions & 9 deletions stubs/kafka-python/kafka/__init__.pyi
Original file line number Diff line number Diff line change
@@ -1,13 +1,38 @@
import logging

from kafka.admin import KafkaAdminClient as KafkaAdminClient
from kafka.client_async import KafkaClient as KafkaClient
from kafka.conn import BrokerConnection as BrokerConnection
from kafka.consumer import KafkaConsumer as KafkaConsumer
from kafka.consumer.subscription_state import ConsumerRebalanceListener as ConsumerRebalanceListener
from kafka.consumer.subscription_state import (
AsyncConsumerRebalanceListener as AsyncConsumerRebalanceListener,
ConsumerRebalanceListener as ConsumerRebalanceListener,
)
from kafka.producer import KafkaProducer as KafkaProducer
from kafka.protocol.consumer import IsolationLevel as IsolationLevel, OffsetSpec as OffsetSpec
from kafka.serializer import (
DefaultSerializer as DefaultSerializer,
Deserializer as Deserializer,
JsonSerializer as JsonSerializer,
Serializer as Serializer,
)
from kafka.structs import (
ConsumerGroupMetadata as ConsumerGroupMetadata,
OffsetAndMetadata as OffsetAndMetadata,
TopicPartition as TopicPartition,
TopicPartitionReplica as TopicPartitionReplica,
)

__all__ = ["BrokerConnection", "ConsumerRebalanceListener", "KafkaAdminClient", "KafkaClient", "KafkaConsumer", "KafkaProducer"]

class NullHandler(logging.Handler):
def emit(self, record) -> None: ...
__all__ = [
"KafkaAdminClient",
"KafkaConsumer",
"KafkaProducer",
"AsyncConsumerRebalanceListener",
"ConsumerRebalanceListener",
"DefaultSerializer",
"JsonSerializer",
"Serializer",
"Deserializer",
"ConsumerGroupMetadata",
"OffsetAndMetadata",
"TopicPartition",
"TopicPartitionReplica",
"IsolationLevel",
"OffsetSpec",
]
60 changes: 49 additions & 11 deletions stubs/kafka-python/kafka/admin/__init__.pyi
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
from kafka.admin.acl_resource import (
from kafka.admin._acls import (
ACL as ACL,
ACLFilter as ACLFilter,
ACLOperation as ACLOperation,
Expand All @@ -8,23 +8,61 @@ from kafka.admin.acl_resource import (
ResourcePatternFilter as ResourcePatternFilter,
ResourceType as ResourceType,
)
from kafka.admin._cluster import UpdateFeatureType as UpdateFeatureType
from kafka.admin._configs import (
AlterConfigOp as AlterConfigOp,
ConfigFilterType as ConfigFilterType,
ConfigResource as ConfigResource,
ConfigResourceType as ConfigResourceType,
ConfigSourceType as ConfigSourceType,
ConfigType as ConfigType,
)
from kafka.admin._groups import GroupState as GroupState, GroupType as GroupType, MemberToRemove as MemberToRemove
from kafka.admin._partitions import NewPartitions as NewPartitions, OffsetSpec as OffsetSpec, OffsetTimestamp as OffsetTimestamp
from kafka.admin._topics import NewTopic as NewTopic
from kafka.admin._transactions import (
AbortTransactionSpec as AbortTransactionSpec,
PartitionProducerState as PartitionProducerState,
ProducerState as ProducerState,
TransactionDescription as TransactionDescription,
TransactionListing as TransactionListing,
TransactionState as TransactionState,
)
from kafka.admin._users import (
ScramMechanism as ScramMechanism,
UserScramCredentialDeletion as UserScramCredentialDeletion,
UserScramCredentialUpsertion as UserScramCredentialUpsertion,
)
from kafka.admin.client import KafkaAdminClient as KafkaAdminClient
from kafka.admin.config_resource import ConfigResource as ConfigResource, ConfigResourceType as ConfigResourceType
from kafka.admin.new_partitions import NewPartitions as NewPartitions
from kafka.admin.new_topic import NewTopic as NewTopic

__all__ = [
"ConfigResource",
"ConfigResourceType",
"KafkaAdminClient",
"NewTopic",
"NewPartitions",
"ACL",
"ACLFilter",
"ResourcePattern",
"ResourcePatternFilter",
"ACLOperation",
"ResourceType",
"ACLPermissionType",
"ACLResourcePatternType",
"ResourceType",
"ResourcePattern",
"ResourcePatternFilter",
"AlterConfigOp",
"ConfigResource",
"ConfigResourceType",
"ConfigType",
"ConfigSourceType",
"UpdateFeatureType",
"GroupState",
"GroupType",
"MemberToRemove",
"OffsetSpec",
"OffsetTimestamp",
"AbortTransactionSpec",
"PartitionProducerState",
"ProducerState",
"TransactionDescription",
"TransactionListing",
"TransactionState",
"ScramMechanism",
"UserScramCredentialDeletion",
"UserScramCredentialUpsertion",
]
Original file line number Diff line number Diff line change
@@ -1,13 +1,30 @@
from _typeshed import Incomplete
from collections.abc import Sequence
from enum import IntEnum
from typing import TypedDict, type_check_only

from kafka.errors import KafkaError

@type_check_only
class _CreateAclsResult(TypedDict):
succeeded: list[ACL]
failed: list[type[KafkaError]]

class ACLAdminMixin:
config: dict[Incomplete, Incomplete]
def describe_acls(self, acl_filter: ACLFilter) -> tuple[list[ACL], type[KafkaError]]: ...
def create_acls(self, acls: Sequence[ACL]) -> _CreateAclsResult: ...
def delete_acls(self, acl_filters: Sequence[ACLFilter]) -> list[tuple[ACLFilter, list[ACL], type[KafkaError]]]: ...

class ResourceType(IntEnum):
UNKNOWN = 0
ANY = 1
CLUSTER = 4
DELEGATION_TOKEN = 6
GROUP = 3
TOPIC = 2
GROUP = 3
CLUSTER = 4
TRANSACTIONAL_ID = 5
DELEGATION_TOKEN = 6
USER = 7

class ACLOperation(IntEnum):
UNKNOWN = 0
Expand All @@ -24,7 +41,7 @@ class ACLOperation(IntEnum):
ALTER_CONFIGS = 11
IDEMPOTENT_WRITE = 12
CREATE_TOKENS = 13
DESCRIBE_TOKENS = 13
DESCRIBE_TOKENS = 14

class ACLPermissionType(IntEnum):
UNKNOWN = 0
Expand All @@ -39,6 +56,25 @@ class ACLResourcePatternType(IntEnum):
LITERAL = 3
PREFIXED = 4

class ResourcePatternFilter:
resource_type: ResourceType
resource_name: str | None
pattern_type: ACLResourcePatternType
def __init__(self, resource_type: ResourceType, resource_name: str | None, pattern_type: ACLResourcePatternType) -> None: ...
def validate(self) -> None: ...
def __eq__(self, other: ResourcePatternFilter) -> bool: ... # type: ignore[override]
def __hash__(self) -> int: ...

class ResourcePattern(ResourcePatternFilter):
resource_name: str
def __init__(
self,
resource_type: ResourceType,
resource_name: str,
pattern_type: ACLResourcePatternType = ACLResourcePatternType.LITERAL,
) -> None: ...
def validate(self) -> None: ...

class ACLFilter:
principal: str | None
host: str | None
Expand All @@ -54,8 +90,8 @@ class ACLFilter:
resource_pattern: ResourcePatternFilter,
) -> None: ...
def validate(self) -> None: ...
def __eq__(self, other): ...
def __hash__(self): ...
def __eq__(self, other: ACLFilter) -> bool: ... # type: ignore[override]
def __hash__(self) -> int: ...

class ACL(ACLFilter):
resource_pattern: ResourcePattern
Expand All @@ -69,23 +105,4 @@ class ACL(ACLFilter):
) -> None: ...
def validate(self) -> None: ...

class ResourcePatternFilter:
resource_type: ResourceType
resource_name: str | None
pattern_type: ACLResourcePatternType
def __init__(self, resource_type: ResourceType, resource_name: str | None, pattern_type: ACLResourcePatternType) -> None: ...
def validate(self) -> None: ...
def __eq__(self, other): ...
def __hash__(self): ...

class ResourcePattern(ResourcePatternFilter):
resource_name: str
def __init__(
self,
resource_type: ResourceType,
resource_name: str,
pattern_type: ACLResourcePatternType = ACLResourcePatternType.LITERAL,
) -> None: ...
def validate(self) -> None: ...

def valid_acl_operations(int_vals) -> set[ACLOperation]: ...
27 changes: 27 additions & 0 deletions stubs/kafka-python/kafka/admin/_cluster.pyi
Original file line number Diff line number Diff line change
@@ -0,0 +1,27 @@
from _typeshed import Incomplete
from enum import IntEnum

from kafka.protocol.api_key import ApiKey
from kafka.util import EnumHelper

class ClusterAdminMixin:
def describe_cluster(self) -> dict[str, Incomplete]: ...
def describe_log_dirs(
self,
topic_partitions: dict[Incomplete, Incomplete] | list[Incomplete] | None = None,
brokers: list[Incomplete] | None = None,
) -> list[dict[str, Incomplete]]: ...
def alter_replica_log_dirs(self, replica_assignments): ...
def describe_metadata_quorum(self): ...
def get_broker_version_data(self, broker_id): ...
def api_versions(self) -> dict[ApiKey, tuple[int, int] | Incomplete]: ...
def describe_features(self, send_request_to_controller: bool = False) -> dict[str, Incomplete]: ...
def update_features(
self, feature_updates: dict[Incomplete, Incomplete], validate_only: bool = False, timeout_ms: int = 60000
) -> dict[Incomplete, Incomplete]: ...

class UpdateFeatureType(EnumHelper, IntEnum):
UNKNOWN = 0
UPGRADE = 1
SAFE_DOWNGRADE = 2
UNSAFE_DOWNGRADE = 3
83 changes: 83 additions & 0 deletions stubs/kafka-python/kafka/admin/_configs.pyi
Original file line number Diff line number Diff line change
@@ -0,0 +1,83 @@
from _typeshed import Incomplete
from collections.abc import Mapping, Sequence
from enum import IntEnum

from kafka.util import EnumHelper

class ConfigAdminMixin:
config: dict[Incomplete, Incomplete]
def describe_configs(
self,
config_resources: Sequence[ConfigResource],
include_synonyms: bool = False,
config_filter: ConfigFilterType | str = "modified",
): ...
def list_config_resources(self, resource_types: Sequence[ConfigResourceType | str] | None = None): ...
def alter_configs(
self,
config_resources: Sequence[ConfigResource],
validate_only: bool = False,
raise_on_unknown: bool = True,
incremental: bool | None = None,
): ...
def reset_configs(
self,
config_resources: Sequence[ConfigResource],
validate_only: bool = False,
raise_on_unknown: bool = True,
incremental: bool | None = None,
): ...

class AlterConfigOp(EnumHelper, IntEnum):
SET = 0
DELETE = 1
APPEND = 2
SUBTRACT = 3

class ConfigFilterType(EnumHelper, IntEnum):
ALL = 0
DYNAMIC = 1
MODIFIED = 2
DEFAULT = 3
STATIC = 4
def should_skip(self, config_source: ConfigSourceType) -> bool: ...

class ConfigResourceType(EnumHelper, IntEnum):
UNKNOWN = 0
TOPIC = 2
BROKER = 4
BROKER_LOGGER = 8
CLIENT_METRICS = 16
GROUP = 32

class ConfigResource:
resource_type: ConfigResourceType
name: str
configs: Mapping[str, str] | None
def __init__(self, resource_type: ConfigResourceType | str, name: str, configs: Mapping[str, str] | None = None) -> None: ...

class ConfigType(EnumHelper, IntEnum):
UNKNOWN = 0
BOOLEAN = 1
STRING = 2
INT = 3
SHORT = 4
LONG = 5
DOUBLE = 6
LIST = 7
CLASS = 8
PASSWORD = 9

class ConfigSourceType(EnumHelper, IntEnum):
UNKNOWN = 0
DYNAMIC_TOPIC_CONFIG = 1
DYNAMIC_BROKER_CONFIG = 2
DYNAMIC_DEFAULT_BROKER_CONFIG = 3
STATIC_BROKER_CONFIG = 4
DEFAULT_CONFIG = 5
DYNAMIC_BROKER_LOGGER_CONFIG = 6
DYNAMIC_CLIENT_METRICS_CONFIG = 7
DYNAMIC_GROUP_CONFIG = 8
def is_modified(self) -> bool: ...
@classmethod
def dynamic_for_resource_type(cls, resource_type: ConfigResourceType) -> ConfigSourceType: ...
Loading