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
2 changes: 1 addition & 1 deletion CONTRIBUTING.rst
Original file line number Diff line number Diff line change
Expand Up @@ -61,7 +61,7 @@ We follow pep8 rather stricly:
3. Line length is 80
4. Code must pass the default flake8 (version 2.5) tests (pep8 + pyflakes)

Code must be compatible with Python 3.8 to the newest released
Code must be compatible with Python 3.9 to the newest released
version of Python.

If external libraries must be used, they should be wrapped in a mechanism that
Expand Down
43 changes: 20 additions & 23 deletions py4j-python/src/py4j/clientserver.py
Original file line number Diff line number Diff line change
Expand Up @@ -39,10 +39,10 @@ class FinalizerWorker(Thread):

def __init__(self, deque):
self.deque = deque
super(FinalizerWorker, self).__init__()
super().__init__()

def run(self):
while(True):
while True:
try:
task = self.deque.pop()
if task == SHUTDOWN_FINALIZER_WORKER:
Expand Down Expand Up @@ -119,7 +119,7 @@ def __init__(
:param auth_token: if provided, an authentication that token clients
must provide to the server when connecting.
"""
super(JavaParameters, self).__init__(
super().__init__(
address, port, auto_field, auto_close, auto_convert, eager_load,
ssl_context, enable_memory_management, read_timeout, auth_token)
self.auto_gc = auto_gc
Expand Down Expand Up @@ -186,7 +186,7 @@ def __init__(
:param auth_token: if provided, an authentication token that clients
must provide to the server when connecting.
"""
super(PythonParameters, self).__init__(
super().__init__(
address, port, daemonize, daemonize_connections, eager_load,
ssl_context, accept_timeout, read_timeout,
propagate_java_exceptions, auth_token)
Expand Down Expand Up @@ -216,7 +216,7 @@ def __init__(
:param finalizer_deque: deque used to manage garbage collection
requests.
"""
super(JavaClient, self).__init__(
super().__init__(
java_parameters,
gateway_property=gateway_property)
self.java_parameters = java_parameters
Expand All @@ -233,7 +233,7 @@ def garbage_collect_object(self, target_id, enqueue=True):
if enqueue:
self.finalizer_deque.appendleft((self, target_id))
else:
super(JavaClient, self).garbage_collect_object(target_id)
super().garbage_collect_object(target_id)

def set_thread_connection(self, connection):
"""Associates a ClientServerConnection with the current thread.
Expand All @@ -248,7 +248,7 @@ def set_thread_connection(self, connection):

def shutdown_gateway(self):
try:
super(JavaClient, self).shutdown_gateway()
super().shutdown_gateway()
finally:
self.finalizer_deque.appendleft(SHUTDOWN_FINALIZER_WORKER)

Expand Down Expand Up @@ -291,7 +291,7 @@ def _create_new_connection(self):

def _should_retry(self, retry, connection, pne=None):
# Only retry if Python was driving the communication.
parent_retry = super(JavaClient, self)._should_retry(
parent_retry = super()._should_retry(
retry, connection, pne)
return parent_retry and retry and connection and\
connection.initiated_from_client
Expand Down Expand Up @@ -360,7 +360,7 @@ def __init__(

:param gateway_property: used to keep gateway preferences.
"""
super(PythonServer, self).__init__(
super().__init__(
pool=gateway_property.pool,
gateway_client=java_client,
callback_server_parameters=python_parameters)
Expand Down Expand Up @@ -450,10 +450,7 @@ def connect_to_java_server(self):

def _authenticate_connection(self):
if self.java_parameters.auth_token:
cmd = "{0}\n{1}\n".format(
proto.AUTH_COMMAND_NAME,
self.java_parameters.auth_token
)
cmd = f"{proto.AUTH_COMMAND_NAME}\n{self.java_parameters.auth_token}\n"
answer = self.send_command(cmd)
error, _ = proto.is_error(answer)
if error:
Expand Down Expand Up @@ -498,10 +495,10 @@ def shutdown_socket(self, remote_port, local_port):
logger.info(
"Send shutdown request for the Java socket {0}, remote port {1}, local port {2}".
format(address, remote_port, local_port))
self.socket.sendall("z\n".encode("utf-8"))
self.socket.sendall(("%s\n" % address).encode("utf-8"))
self.socket.sendall(("%s\n" % remote_port).encode("utf-8"))
self.socket.sendall(("%s\n" % local_port).encode("utf-8"))
self.socket.sendall(b"z\n")
self.socket.sendall(f"{address}\n".encode("utf-8"))
self.socket.sendall(f"{remote_port}\n".encode("utf-8"))
self.socket.sendall(f"{local_port}\n".encode("utf-8"))
logger.info("Close connection")
self.close()
self.is_connected = False
Expand All @@ -520,18 +517,18 @@ def run(self):

def send_command(self, command):
# TODO At some point extract common code from wait_for_commands
logger.debug("Command to send: {0}".format(command))
logger.debug(f"Command to send: {command}")
try:
self.socket.sendall(command.encode("utf-8"))
except Exception as e:
logger.info("Error while sending or receiving.", exc_info=True)
raise Py4JNetworkError(
"Error while sending", e, proto.ERROR_ON_SEND)
"Error while sending", e, proto.ERROR_ON_SEND) from e

try:
while True:
answer = self.stream.readline()[:-1].decode("utf-8")
logger.debug("Answer received: {0}".format(answer))
logger.debug(f"Answer received: {answer}")
# Happens when a the other end is dead. There might be an empty
# answer before the socket raises an error.
if answer.strip() == "":
Expand All @@ -552,7 +549,7 @@ def send_command(self, command):
self.socket.sendall(
proto.SUCCESS_RETURN_MESSAGE.encode("utf-8"))
else:
logger.error("Unknown command {0}".format(command))
logger.error(f"Unknown command {command}")
# We're sending something to prevent blocking,
# but at this point, the protocol is broken.
self.socket.sendall(
Expand Down Expand Up @@ -611,7 +608,7 @@ def wait_for_commands(self):
self.socket.sendall(
proto.SUCCESS_RETURN_MESSAGE.encode("utf-8"))
else:
logger.error("Unknown command {0}".format(command))
logger.error(f"Unknown command {command}")
# We're sending something to prevent blocking, but at this
# point, the protocol is broken.
self.socket.sendall(
Expand Down Expand Up @@ -701,7 +698,7 @@ def __init__(
python_parameters = PythonParameters()
self.java_parameters = java_parameters
self.python_parameters = python_parameters
super(ClientServer, self).__init__(
super().__init__(
gateway_parameters=java_parameters,
callback_server_parameters=python_parameters,
python_server_entry_point=python_server_entry_point
Expand Down
9 changes: 6 additions & 3 deletions py4j-python/src/py4j/finalizer.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,10 @@

:author: Barthelemy Dagenais
"""
from __future__ import annotations

from threading import RLock
from typing import Any


class ThreadSafeFinalizer(object):
Expand All @@ -28,7 +31,7 @@ class ThreadSafeFinalizer(object):
lock = RLock()

@classmethod
def add_finalizer(cls, id, weak_ref):
def add_finalizer(cls, id: Any, weak_ref: Any) -> None:
"""Registers a finalizer with an id.

:param id: The id of the object referenced by the weak reference.
Expand All @@ -38,7 +41,7 @@ def add_finalizer(cls, id, weak_ref):
cls.finalizers[id] = weak_ref

@classmethod
def remove_finalizer(cls, id):
def remove_finalizer(cls, id: Any) -> None:
"""Removes a finalizer associated with this id.

:param id: The id of the object for which the finalizer will be
Expand All @@ -48,7 +51,7 @@ def remove_finalizer(cls, id):
cls.finalizers.pop(id, None)

@classmethod
def clear_finalizers(cls, clear_all=False):
def clear_finalizers(cls, clear_all: bool = False) -> None:
"""Removes all registered finalizers.

:param clear_all: If `True`, all finalizers are deleted. Otherwise,
Expand Down
27 changes: 11 additions & 16 deletions py4j-python/src/py4j/java_collections.py
Original file line number Diff line number Diff line change
Expand Up @@ -88,7 +88,7 @@ def __str__(self):

def __repr__(self):
items = (
"{0}: {1}".format(repr(k), repr(v))
f"{repr(k)}: {repr(v)}"
for k, v in self.items())
return "{{{0}}}".format(", ".join(items))

Expand Down Expand Up @@ -190,8 +190,7 @@ def __getitem__(self, key):
elif isinstance(key, int):
return self.__compute_item(key)
else:
raise TypeError("array indices must be integers, not {0}".format(
key.__class__.__name__))
raise TypeError(f"array indices must be integers, not {key.__class__.__name__}")

def __repl_item_from_slice(self, range, iterable):
value_iter = iter(iterable)
Expand Down Expand Up @@ -220,15 +219,14 @@ def __setitem__(self, key, value):
if lenr != lenv:
raise ValueError(
"attempt to assign sequence of size "
"{0} to extended slice of size {1}".format(lenv, lenr))
f"{lenv} to extended slice of size {lenr}")
else:
return self.__repl_item_from_slice(self_range, value)

elif isinstance(key, int):
return self.__set_item(key, value)
else:
raise TypeError("list indices must be integers, not {0}".format(
key.__class__.__name__))
raise TypeError(f"list indices must be integers, not {key.__class__.__name__}")

def __len__(self):
command = proto.ARRAY_COMMAND_NAME +\
Expand Down Expand Up @@ -334,15 +332,14 @@ def __setitem__(self, key, value):
if lenr != lenv:
raise ValueError(
"attempt to assign sequence of size "
"{0} to extended slice of size {1}".format(lenv, lenr))
f"{lenv} to extended slice of size {lenr}")
else:
return self.__repl_item_from_slice(self_range, value)

elif isinstance(key, int):
return self.__set_item(key, value)
else:
raise TypeError("list indices must be integers, not {0}".format(
key.__class__.__name__))
raise TypeError(f"list indices must be integers, not {key.__class__.__name__}")

def __get_slice(self, indices):
command = proto.LIST_COMMAND_NAME +\
Expand All @@ -361,8 +358,7 @@ def __getitem__(self, key):
elif isinstance(key, int):
return self.__compute_item(key)
else:
raise TypeError("list indices must be integers, not {0}".format(
key.__class__.__name__))
raise TypeError(f"list indices must be integers, not {key.__class__.__name__}")

def __delitem__(self, key):
if isinstance(key, slice):
Expand All @@ -374,8 +370,7 @@ def __delitem__(self, key):
elif isinstance(key, int):
return self.__del_item(key)
else:
raise TypeError("list indices must be integers, not {0}".format(
key.__class__.__name__))
raise TypeError(f"list indices must be integers, not {key.__class__.__name__}")

def __contains__(self, item):
return self.contains(item)
Expand Down Expand Up @@ -421,8 +416,7 @@ def insert(self, key, value):
new_key = self.__compute_index(key, True)
return self.add(new_key, value)
else:
raise TypeError("list indices must be integers, not {0}".format(
key.__class__.__name__))
raise TypeError(f"list indices must be integers, not {key.__class__.__name__}")

def extend(self, other_list):
self.addAll(other_list)
Expand Down Expand Up @@ -472,7 +466,8 @@ def __str__(self):

def __repr__(self):
items = (repr(x) for x in self)
return "[{0}]".format(", ".join(items))
inside = ", ".join(items)
return f"[{inside}]"


class SetConverter(object):
Expand Down
Loading
Loading