Reformat code with black and isort.

This commit is contained in:
Felix Fontein
2025-10-06 18:34:59 +02:00
parent f45232635c
commit d65d37e9e9
132 changed files with 17581 additions and 14729 deletions
+179 -149
View File
@@ -6,6 +6,7 @@
from __future__ import annotations
DOCUMENTATION = r"""
module: docker_swarm
short_description: Manage Swarm cluster
@@ -292,8 +293,9 @@ actions:
import json
import traceback
try:
from docker.errors import DockerException, APIError
from docker.errors import APIError, DockerException
except ImportError:
# missing Docker SDK for Python handled in ansible.module_utils.docker.common
pass
@@ -302,13 +304,14 @@ from ansible_collections.community.docker.plugins.module_utils.common import (
DockerBaseClass,
RequestException,
)
from ansible_collections.community.docker.plugins.module_utils.swarm import (
AnsibleDockerSwarmClient,
)
from ansible_collections.community.docker.plugins.module_utils.util import (
DifferenceTracker,
sanitize_labels,
)
from ansible_collections.community.docker.plugins.module_utils.swarm import AnsibleDockerSwarmClient
class TaskParameters(DockerBaseClass):
def __init__(self):
@@ -353,68 +356,70 @@ class TaskParameters(DockerBaseClass):
return result
def update_from_swarm_info(self, swarm_info):
spec = swarm_info['Spec']
spec = swarm_info["Spec"]
ca_config = spec.get('CAConfig') or dict()
ca_config = spec.get("CAConfig") or dict()
if self.node_cert_expiry is None:
self.node_cert_expiry = ca_config.get('NodeCertExpiry')
self.node_cert_expiry = ca_config.get("NodeCertExpiry")
if self.ca_force_rotate is None:
self.ca_force_rotate = ca_config.get('ForceRotate')
self.ca_force_rotate = ca_config.get("ForceRotate")
dispatcher = spec.get('Dispatcher') or dict()
dispatcher = spec.get("Dispatcher") or dict()
if self.dispatcher_heartbeat_period is None:
self.dispatcher_heartbeat_period = dispatcher.get('HeartbeatPeriod')
self.dispatcher_heartbeat_period = dispatcher.get("HeartbeatPeriod")
raft = spec.get('Raft') or dict()
raft = spec.get("Raft") or dict()
if self.snapshot_interval is None:
self.snapshot_interval = raft.get('SnapshotInterval')
self.snapshot_interval = raft.get("SnapshotInterval")
if self.keep_old_snapshots is None:
self.keep_old_snapshots = raft.get('KeepOldSnapshots')
self.keep_old_snapshots = raft.get("KeepOldSnapshots")
if self.heartbeat_tick is None:
self.heartbeat_tick = raft.get('HeartbeatTick')
self.heartbeat_tick = raft.get("HeartbeatTick")
if self.log_entries_for_slow_followers is None:
self.log_entries_for_slow_followers = raft.get('LogEntriesForSlowFollowers')
self.log_entries_for_slow_followers = raft.get("LogEntriesForSlowFollowers")
if self.election_tick is None:
self.election_tick = raft.get('ElectionTick')
self.election_tick = raft.get("ElectionTick")
orchestration = spec.get('Orchestration') or dict()
orchestration = spec.get("Orchestration") or dict()
if self.task_history_retention_limit is None:
self.task_history_retention_limit = orchestration.get('TaskHistoryRetentionLimit')
self.task_history_retention_limit = orchestration.get(
"TaskHistoryRetentionLimit"
)
encryption_config = spec.get('EncryptionConfig') or dict()
encryption_config = spec.get("EncryptionConfig") or dict()
if self.autolock_managers is None:
self.autolock_managers = encryption_config.get('AutoLockManagers')
self.autolock_managers = encryption_config.get("AutoLockManagers")
if self.name is None:
self.name = spec['Name']
self.name = spec["Name"]
if self.labels is None:
self.labels = spec.get('Labels') or {}
self.labels = spec.get("Labels") or {}
if 'LogDriver' in spec['TaskDefaults']:
self.log_driver = spec['TaskDefaults']['LogDriver']
if "LogDriver" in spec["TaskDefaults"]:
self.log_driver = spec["TaskDefaults"]["LogDriver"]
def update_parameters(self, client):
assign = dict(
snapshot_interval='snapshot_interval',
task_history_retention_limit='task_history_retention_limit',
keep_old_snapshots='keep_old_snapshots',
log_entries_for_slow_followers='log_entries_for_slow_followers',
heartbeat_tick='heartbeat_tick',
election_tick='election_tick',
dispatcher_heartbeat_period='dispatcher_heartbeat_period',
node_cert_expiry='node_cert_expiry',
name='name',
labels='labels',
signing_ca_cert='signing_ca_cert',
signing_ca_key='signing_ca_key',
ca_force_rotate='ca_force_rotate',
autolock_managers='autolock_managers',
log_driver='log_driver',
snapshot_interval="snapshot_interval",
task_history_retention_limit="task_history_retention_limit",
keep_old_snapshots="keep_old_snapshots",
log_entries_for_slow_followers="log_entries_for_slow_followers",
heartbeat_tick="heartbeat_tick",
election_tick="election_tick",
dispatcher_heartbeat_period="dispatcher_heartbeat_period",
node_cert_expiry="node_cert_expiry",
name="name",
labels="labels",
signing_ca_cert="signing_ca_cert",
signing_ca_key="signing_ca_key",
ca_force_rotate="ca_force_rotate",
autolock_managers="autolock_managers",
log_driver="log_driver",
)
params = dict()
for dest, source in assign.items():
if not client.option_minimal_versions[source]['supported']:
if not client.option_minimal_versions[source]["supported"]:
continue
value = getattr(self, source)
if value is not None:
@@ -423,12 +428,21 @@ class TaskParameters(DockerBaseClass):
def compare_to_active(self, other, client, differences):
for k in self.__dict__:
if k in ('advertise_addr', 'listen_addr', 'remote_addrs', 'join_token',
'rotate_worker_token', 'rotate_manager_token', 'spec',
'default_addr_pool', 'subnet_size', 'data_path_addr',
'data_path_port'):
if k in (
"advertise_addr",
"listen_addr",
"remote_addrs",
"join_token",
"rotate_worker_token",
"rotate_manager_token",
"spec",
"default_addr_pool",
"subnet_size",
"data_path_addr",
"data_path_port",
):
continue
if not client.option_minimal_versions[k]['supported']:
if not client.option_minimal_versions[k]["supported"]:
continue
value = getattr(self, k)
if value is None:
@@ -437,9 +451,9 @@ class TaskParameters(DockerBaseClass):
if value != other_value:
differences.add(k, parameter=value, active=other_value)
if self.rotate_worker_token:
differences.add('rotate_worker_token', parameter=True, active=False)
differences.add("rotate_worker_token", parameter=True, active=False)
if self.rotate_manager_token:
differences.add('rotate_manager_token', parameter=True, active=False)
differences.add("rotate_manager_token", parameter=True, active=False)
return differences
@@ -454,9 +468,9 @@ class SwarmManager(DockerBaseClass):
self.check_mode = self.client.check_mode
self.swarm_info = {}
self.state = client.module.params['state']
self.force = client.module.params['force']
self.node_id = client.module.params['node_id']
self.state = client.module.params["state"]
self.force = client.module.params["force"]
self.node_id = client.module.params["node_id"]
self.differences = DifferenceTracker()
self.parameters = TaskParameters.from_ansible_params(client)
@@ -475,8 +489,8 @@ class SwarmManager(DockerBaseClass):
if self.client.module._diff or self.parameters.debug:
diff = dict()
diff['before'], diff['after'] = self.differences.get_before_after()
self.results['diff'] = diff
diff["before"], diff["after"] = self.differences.get_before_after()
self.results["diff"] = diff
def inspect_swarm(self):
try:
@@ -484,8 +498,8 @@ class SwarmManager(DockerBaseClass):
json_str = json.dumps(data, ensure_ascii=False)
self.swarm_info = json.loads(json_str)
self.results['changed'] = False
self.results['swarm_facts'] = self.swarm_info
self.results["changed"] = False
self.results["swarm_facts"] = self.swarm_info
unlock_key = self.get_unlock_key()
self.swarm_info.update(unlock_key)
@@ -493,7 +507,7 @@ class SwarmManager(DockerBaseClass):
return
def get_unlock_key(self):
default = {'UnlockKey': None}
default = {"UnlockKey": None}
if not self.has_swarm_lock_changed():
return default
try:
@@ -503,7 +517,7 @@ class SwarmManager(DockerBaseClass):
def has_swarm_lock_changed(self):
return self.parameters.autolock_managers and (
self.created or self.differences.has_difference_for('autolock_managers')
self.created or self.differences.has_difference_for("autolock_managers")
)
def init_swarm(self):
@@ -513,19 +527,19 @@ class SwarmManager(DockerBaseClass):
if not self.check_mode:
init_arguments = {
'advertise_addr': self.parameters.advertise_addr,
'listen_addr': self.parameters.listen_addr,
'force_new_cluster': self.force,
'swarm_spec': self.parameters.spec,
"advertise_addr": self.parameters.advertise_addr,
"listen_addr": self.parameters.listen_addr,
"force_new_cluster": self.force,
"swarm_spec": self.parameters.spec,
}
if self.parameters.default_addr_pool is not None:
init_arguments['default_addr_pool'] = self.parameters.default_addr_pool
init_arguments["default_addr_pool"] = self.parameters.default_addr_pool
if self.parameters.subnet_size is not None:
init_arguments['subnet_size'] = self.parameters.subnet_size
init_arguments["subnet_size"] = self.parameters.subnet_size
if self.parameters.data_path_addr is not None:
init_arguments['data_path_addr'] = self.parameters.data_path_addr
init_arguments["data_path_addr"] = self.parameters.data_path_addr
if self.parameters.data_path_port is not None:
init_arguments['data_path_port'] = self.parameters.data_path_port
init_arguments["data_path_port"] = self.parameters.data_path_port
try:
self.client.init_swarm(**init_arguments)
except APIError as exc:
@@ -537,180 +551,196 @@ class SwarmManager(DockerBaseClass):
self.created = True
self.inspect_swarm()
self.results['actions'].append(f"New Swarm cluster created: {self.swarm_info.get('ID')}")
self.differences.add('state', parameter='present', active='absent')
self.results['changed'] = True
self.results['swarm_facts'] = {
'JoinTokens': self.swarm_info.get('JoinTokens'),
'UnlockKey': self.swarm_info.get('UnlockKey')
self.results["actions"].append(
f"New Swarm cluster created: {self.swarm_info.get('ID')}"
)
self.differences.add("state", parameter="present", active="absent")
self.results["changed"] = True
self.results["swarm_facts"] = {
"JoinTokens": self.swarm_info.get("JoinTokens"),
"UnlockKey": self.swarm_info.get("UnlockKey"),
}
def __update_swarm(self):
try:
self.inspect_swarm()
version = self.swarm_info['Version']['Index']
version = self.swarm_info["Version"]["Index"]
self.parameters.update_from_swarm_info(self.swarm_info)
old_parameters = TaskParameters()
old_parameters.update_from_swarm_info(self.swarm_info)
self.parameters.compare_to_active(old_parameters, self.client, self.differences)
self.parameters.compare_to_active(
old_parameters, self.client, self.differences
)
if self.differences.empty:
self.results['actions'].append("No modification")
self.results['changed'] = False
self.results["actions"].append("No modification")
self.results["changed"] = False
return
update_parameters = TaskParameters.from_ansible_params(self.client)
update_parameters.update_parameters(self.client)
if not self.check_mode:
self.client.update_swarm(
version=version, swarm_spec=update_parameters.spec,
version=version,
swarm_spec=update_parameters.spec,
rotate_worker_token=self.parameters.rotate_worker_token,
rotate_manager_token=self.parameters.rotate_manager_token)
rotate_manager_token=self.parameters.rotate_manager_token,
)
except APIError as exc:
self.client.fail(f"Can not update a Swarm Cluster: {exc}")
return
self.inspect_swarm()
self.results['actions'].append("Swarm cluster updated")
self.results['changed'] = True
self.results["actions"].append("Swarm cluster updated")
self.results["changed"] = True
def join(self):
if self.client.check_if_swarm_node():
self.results['actions'].append("This node is already part of a swarm.")
self.results["actions"].append("This node is already part of a swarm.")
return
if not self.check_mode:
join_arguments = {
'remote_addrs': self.parameters.remote_addrs,
'join_token': self.parameters.join_token,
'listen_addr': self.parameters.listen_addr,
'advertise_addr': self.parameters.advertise_addr,
"remote_addrs": self.parameters.remote_addrs,
"join_token": self.parameters.join_token,
"listen_addr": self.parameters.listen_addr,
"advertise_addr": self.parameters.advertise_addr,
}
if self.parameters.data_path_addr is not None:
join_arguments['data_path_addr'] = self.parameters.data_path_addr
join_arguments["data_path_addr"] = self.parameters.data_path_addr
try:
self.client.join_swarm(**join_arguments)
except APIError as exc:
self.client.fail(f"Can not join the Swarm Cluster: {exc}")
self.results['actions'].append("New node is added to swarm cluster")
self.differences.add('joined', parameter=True, active=False)
self.results['changed'] = True
self.results["actions"].append("New node is added to swarm cluster")
self.differences.add("joined", parameter=True, active=False)
self.results["changed"] = True
def leave(self):
if not self.client.check_if_swarm_node():
self.results['actions'].append("This node is not part of a swarm.")
self.results["actions"].append("This node is not part of a swarm.")
return
if not self.check_mode:
try:
self.client.leave_swarm(force=self.force)
except APIError as exc:
self.client.fail(f"This node can not leave the Swarm Cluster: {exc}")
self.results['actions'].append("Node has left the swarm cluster")
self.differences.add('joined', parameter='absent', active='present')
self.results['changed'] = True
self.results["actions"].append("Node has left the swarm cluster")
self.differences.add("joined", parameter="absent", active="present")
self.results["changed"] = True
def remove(self):
if not self.client.check_if_swarm_manager():
self.client.fail("This node is not a manager.")
try:
status_down = self.client.check_if_swarm_node_is_down(node_id=self.node_id, repeat_check=5)
status_down = self.client.check_if_swarm_node_is_down(
node_id=self.node_id, repeat_check=5
)
except APIError:
return
if not status_down:
self.client.fail("Can not remove the node. The status node is ready and not down.")
self.client.fail(
"Can not remove the node. The status node is ready and not down."
)
if not self.check_mode:
try:
self.client.remove_node(node_id=self.node_id, force=self.force)
except APIError as exc:
self.client.fail(f"Can not remove the node from the Swarm Cluster: {exc}")
self.results['actions'].append("Node is removed from swarm cluster.")
self.differences.add('joined', parameter=False, active=True)
self.results['changed'] = True
self.client.fail(
f"Can not remove the node from the Swarm Cluster: {exc}"
)
self.results["actions"].append("Node is removed from swarm cluster.")
self.differences.add("joined", parameter=False, active=True)
self.results["changed"] = True
def _detect_remove_operation(client):
return client.module.params['state'] == 'remove'
return client.module.params["state"] == "remove"
def main():
argument_spec = dict(
advertise_addr=dict(type='str'),
data_path_addr=dict(type='str'),
data_path_port=dict(type='int'),
state=dict(type='str', default='present', choices=['present', 'join', 'absent', 'remove']),
force=dict(type='bool', default=False),
listen_addr=dict(type='str', default='0.0.0.0:2377'),
remote_addrs=dict(type='list', elements='str'),
join_token=dict(type='str', no_log=True),
snapshot_interval=dict(type='int'),
task_history_retention_limit=dict(type='int'),
keep_old_snapshots=dict(type='int'),
log_entries_for_slow_followers=dict(type='int'),
heartbeat_tick=dict(type='int'),
election_tick=dict(type='int'),
dispatcher_heartbeat_period=dict(type='int'),
node_cert_expiry=dict(type='int'),
name=dict(type='str'),
labels=dict(type='dict'),
signing_ca_cert=dict(type='str'),
signing_ca_key=dict(type='str', no_log=True),
ca_force_rotate=dict(type='int'),
autolock_managers=dict(type='bool'),
node_id=dict(type='str'),
rotate_worker_token=dict(type='bool', default=False),
rotate_manager_token=dict(type='bool', default=False),
default_addr_pool=dict(type='list', elements='str'),
subnet_size=dict(type='int'),
advertise_addr=dict(type="str"),
data_path_addr=dict(type="str"),
data_path_port=dict(type="int"),
state=dict(
type="str",
default="present",
choices=["present", "join", "absent", "remove"],
),
force=dict(type="bool", default=False),
listen_addr=dict(type="str", default="0.0.0.0:2377"),
remote_addrs=dict(type="list", elements="str"),
join_token=dict(type="str", no_log=True),
snapshot_interval=dict(type="int"),
task_history_retention_limit=dict(type="int"),
keep_old_snapshots=dict(type="int"),
log_entries_for_slow_followers=dict(type="int"),
heartbeat_tick=dict(type="int"),
election_tick=dict(type="int"),
dispatcher_heartbeat_period=dict(type="int"),
node_cert_expiry=dict(type="int"),
name=dict(type="str"),
labels=dict(type="dict"),
signing_ca_cert=dict(type="str"),
signing_ca_key=dict(type="str", no_log=True),
ca_force_rotate=dict(type="int"),
autolock_managers=dict(type="bool"),
node_id=dict(type="str"),
rotate_worker_token=dict(type="bool", default=False),
rotate_manager_token=dict(type="bool", default=False),
default_addr_pool=dict(type="list", elements="str"),
subnet_size=dict(type="int"),
)
required_if = [
('state', 'join', ['remote_addrs', 'join_token']),
('state', 'remove', ['node_id'])
("state", "join", ["remote_addrs", "join_token"]),
("state", "remove", ["node_id"]),
]
option_minimal_versions = dict(
labels=dict(docker_py_version='2.6.0', docker_api_version='1.32'),
signing_ca_cert=dict(docker_py_version='2.6.0', docker_api_version='1.30'),
signing_ca_key=dict(docker_py_version='2.6.0', docker_api_version='1.30'),
ca_force_rotate=dict(docker_py_version='2.6.0', docker_api_version='1.30'),
autolock_managers=dict(docker_py_version='2.6.0'),
log_driver=dict(docker_py_version='2.6.0'),
labels=dict(docker_py_version="2.6.0", docker_api_version="1.32"),
signing_ca_cert=dict(docker_py_version="2.6.0", docker_api_version="1.30"),
signing_ca_key=dict(docker_py_version="2.6.0", docker_api_version="1.30"),
ca_force_rotate=dict(docker_py_version="2.6.0", docker_api_version="1.30"),
autolock_managers=dict(docker_py_version="2.6.0"),
log_driver=dict(docker_py_version="2.6.0"),
remove_operation=dict(
docker_py_version='2.4.0',
docker_py_version="2.4.0",
detect_usage=_detect_remove_operation,
usage_msg='remove swarm nodes'
usage_msg="remove swarm nodes",
),
default_addr_pool=dict(docker_py_version='4.0.0', docker_api_version='1.39'),
subnet_size=dict(docker_py_version='4.0.0', docker_api_version='1.39'),
data_path_addr=dict(docker_py_version='4.0.0', docker_api_version='1.30'),
data_path_port=dict(docker_py_version='6.0.0', docker_api_version='1.40'),
default_addr_pool=dict(docker_py_version="4.0.0", docker_api_version="1.39"),
subnet_size=dict(docker_py_version="4.0.0", docker_api_version="1.39"),
data_path_addr=dict(docker_py_version="4.0.0", docker_api_version="1.30"),
data_path_port=dict(docker_py_version="6.0.0", docker_api_version="1.40"),
)
client = AnsibleDockerSwarmClient(
argument_spec=argument_spec,
supports_check_mode=True,
required_if=required_if,
min_docker_version='1.10.0',
min_docker_version="1.10.0",
option_minimal_versions=option_minimal_versions,
)
sanitize_labels(client.module.params['labels'], 'labels', client)
sanitize_labels(client.module.params["labels"], "labels", client)
try:
results = dict(
changed=False,
result='',
actions=[]
)
results = dict(changed=False, result="", actions=[])
SwarmManager(client, results)()
client.module.exit_json(**results)
except DockerException as e:
client.fail(f'An unexpected Docker error occurred: {e}', exception=traceback.format_exc())
client.fail(
f"An unexpected Docker error occurred: {e}",
exception=traceback.format_exc(),
)
except RequestException as e:
client.fail(
f'An unexpected requests error occurred when Docker SDK for Python tried to talk to the docker daemon: {e}',
exception=traceback.format_exc())
f"An unexpected requests error occurred when Docker SDK for Python tried to talk to the docker daemon: {e}",
exception=traceback.format_exc(),
)
if __name__ == '__main__':
if __name__ == "__main__":
main()