| 登録日 |
:2026/09/27 19:36 |
| カテゴリ |
:Python基礎 |
フォルダ構成
main.py
start_clusters.sh
start_hosts.sh
- config/
|- settings.py
- scripts/
| - command.sh
- src/
| - thread_workers.py
| - reporter.py
| - dataset.py
- targets/
| - target_cluster.ini
| - target_hosts.ini
- logs/ (ログ保存場所)
#!/bin/bash
PROGRAM_DIR=/home/app/thread_backend
$PROGRAM_DIR/main.py 2 2>&1 | tee -a $PROGRAM_DIR/logs/execute_clusters.log
#!/bin/bash
PROGRAM_DIR=/home/app/thread_backend
$PROGRAM_DIR/main.py 1 2>&1 | tee -a $PROGRAM_DIR/logs/execute_hosts.log
main.py
#!/usr/bin/python3
import os
import sys
import json
import queue
import datetime
import json
import logging
import time
import gc
sys_path = os.path.dirname(os.path.abspath(__file__))
sys.path.append(sys_path)
from config.settings import Config, setup_logger
from src.thread_workers import CommandThreadWorkers, ServerInfo, set_queue
from src.dataset import HostsIniDataset, ClusterIniDataset
from src.reporter import ReporterSample, RepoterJson
def import_hostlist() -> queue.Queue:
_datasets = HostsIniDataset()
#print(_datasets)
_q = set_queue(_datasets.targets_list)
return _q
def import_clusterlist() -> queue.Queue:
_datasets = ClusterIniDataset()
_q = set_queue(_datasets.targets_list)
return _q
def main_thread(_q: queue.Queue,
workers=1,
timeout=Config.TIMEOUT,
level=Config.LEVEL):
worker = CommandThreadWorkers(
_queue=_q,
workers=workers,
timeout=timeout,
level=level)
worker.run()
return worker.results
# -----------------------------
# CLI / main
# -----------------------------
def parse_args(argv):
today = datetime.datetime.now().strftime("%Y-%m-%d")
log_level = Config.LEVEL
debug = False
info = False
print(f"[INFO] {today}")
if len(argv) == 0:
print("[ERROR]Please input option (number)")
print("[INFO]Example: python main.py 1 -> imoprt hosts.ini and run")
print("[INFO]Example: python main.py 2 -> import cluster.ini and run")
print("[INFO]Example: python main.py <num> --info -> INFO mode")
print("[INFO]Example: python main.py <num> --debug -> DEBUG mode")
exit(1)
target = argv[0]
i = 1
while i < len(argv):
a = argv[i]
if a == "--debug":
log_level = logging.DEBUG
debug = True
i += 1
continue
if a == "--info":
log_level = logging.INFO
info = True
i += 1
continue
i += 1
return int(target), log_level, debug, info
def main(argv) -> None:
target, log_level, debug, info = parse_args(argv)
if target == 1:
print("[INFO] imoprt hosts.ini and run...")
q = import_hostlist()
if target == 2:
print("[INFO] import cluster.ini and run...")
q = import_clusterlist()
_results = main_thread(q,
workers=Config.MAX_WORKERS,
timeout=Config.TIMEOUT,
level=log_level)
return _results
if __name__ == '__main__':
dt_now = datetime.datetime.now().strftime('%Y/%m/%d %H:%M:%S')
start = time.time()
results = main(sys.argv[1:])
#RepoterJson.print_results(results)
ReporterSample.print_results(results)
end = time.time()
print(f"[INFO] Start: {dt_now}, Elapsed Time: {end - start:.2f} seconds")
gc.collect()
src/*
#!/usr/bin/python3
import os
import sys
from abc import ABC, abstractmethod
import subprocess
from subprocess import PIPE
import queue
import threading
import json
from typing import List, Type
import datetime
import gc
import socket
sys_path = os.path.dirname(os.path.dirname(os.path.abspath(__file__)))
sys.path.append(sys_path)
from config.settings import Config, setup_logger
#-----------------
#debug print
#print(sys_path)
#print(Config.TARGET_SCRIPT)
#print(Config.TARGET_HOSTS)
#print(Config.TARGET_CLUSTER)
#print(Config.TARGET_OUTPUT_DIR)
#----------------
"""
Date: 2026.09.26
Version: 1.0
Created by N.T
"""
#-----------------------------
# Thread Worker Interface
#-----------------------------
class IThreadWorkerInterface(ABC):
def __init__(self,
_queue: queue.Queue,
workers=1,
timeout=Config.TIMEOUT,
level=Config.LEVEL):
self.queue = _queue
self.workers = workers
self.timeout = timeout
self.logger = setup_logger(self.name, level)
self.result_lock = threading.Lock()
self.results = []
@property
def name(self) -> str:
return "ThreadWorkers"
def run(self):
ts = []
for _ in range(self.workers):
t = threading.Thread(target=self.worker)
t.start()
ts.append(t)
[self.queue.put(None) for _ in range(len(ts))]
[t.join() for t in ts]
def worker(self):
self.logger.debug('workers start')
while True:
item = self.queue.get()
if item is None:
break
#self.logger.debug({'thread': item})
res = self.executor(item)
with self.result_lock:
self.results.append({item.hostname: res})
self.logger.debug(type(res))
if type(res) == list:
self.logger.debug(json.dumps(res, indent=2, ensure_ascii=False))
else:
self.logger.debug(res)
self.queue.task_done()
self.logger.debug('workers end')
@abstractmethod
def executor(self, _item):
pass
#-----------------------------
# Thread Worker Concrete
#-----------------------------
class CommandThreadWorkers(IThreadWorkerInterface):
@property
def name(self) -> str:
return "CommandThreadWorkers"
def executor(self, _item):
self.logger.debug(f'{_item.hostname}:{_item.command}')
_res = None
try:
if _item.hostname and _item.command:
_command = "ssh "+ _item.hostname + " " + _item.command
self.logger.debug({'command': _command})
_res_temp = subprocess.run(
_command,
shell=True,
stdin=subprocess.DEVNULL,
stdout=PIPE,
stderr=PIPE,
timeout=self.timeout,
universal_newlines=True)
print(_res_temp.returncode)
if _res_temp.returncode != 0:
_res = _res_temp.stderr.split('\n')
else:
_res = _res_temp.stdout.split('\n')
messages = {_item.hostname: _res}
self.logger.info({
'status': 'success',
'messages': messages})
else:
self.logger.error({
'status': 'failed',
'message': 'hostname and command are required.',
'hostname': _item.hostname,
'_command': _item.command,
})
return _res
except Exception as e:
self.logger.error({
'status': 'failed',
'action': self.name,
'message': str(e)})
return _res
#-----------------------------
# Dataset class
#-----------------------------
class ServerInfo:
def __init__(
self,
ipaddr: str = None,
hostname: str = None,
username: str = None,
password: str = None,
command: str = Config.TARGET_SCRIPT,
port: int = Config.PORT,
timeout: int = Config.TIMEOUT,
):
self.ipaddr = ipaddr
self.hostname = hostname
self.username = username
self.password = password
self.command = command
self.port = port
self.timeout = timeout
#-----------------------------
# Set Queue
#-----------------------------
def set_queue(_targets: List[dict]) -> queue.Queue:
q = queue.Queue()
for t in _targets:
command=t.get("command", None)
if command == "" or command is None:
command = Config.TARGET_SCRIPT
server_info = ServerInfo(
ipaddr=t.get("ipaddr", None),
hostname=t.get("hostname", None),
username=t.get("username", None),
password=t.get("password", None),
command=command
)
q.put(server_info)
return q
#!/usr/bin/python3
# reporter.py
from abc import ABC, abstractmethod
from typing import List
import json
"""
Date: 2026.08.16
Version: 1.0
Created by N.T
"""
# -----------------------------
# Reporter
# -----------------------------
class IReporterInterface(ABC):
@staticmethod
@abstractmethod
def print_results(results: List[dict]) -> None:
pass
class ReporterSample(IReporterInterface):
@staticmethod
def print_results(results: List[dict]) -> None:
print('*' * 80)
print("[INFO]: Results Summary")
width = 12
for res in results:
for hostname, lines in res.items():
print("-" * 80)
hostname_str = str(hostname) if hostname is not None else "Unknown Host"
if lines == [] or lines is None:
print(f"{hostname_str:<{width}} | No data or error occurred.")
continue
for line in lines:
msg = line.replace("\\r", "").replace("\\n", "\n")
print(f"{hostname_str:<{width}} | {msg}")
class RepoterJson(IReporterInterface):
@staticmethod
def print_results(results: List[dict]) -> None:
print(json.dumps(results, indent=2, ensure_ascii=False))
#!/usr/bin/python3
# dataset.py
import os
import sys
import json
import configparser
from abc import ABC, abstractmethod
from typing import List, Dict
sys_path = os.path.dirname(os.path.dirname(os.path.abspath(__file__)))
sys.path.append(sys_path)
from config.settings import Config, setup_logger
"""
Date: 2026.08.16
Version: 1.0
Created by N.T
"""
class IDatasetInterface(ABC):
"""
ターゲットリスト(IPアドレス・ホスト名等)を読み込むデータセットの共通インターフェース。
"""
def __init__(self):
self.targets_list: List[Dict] = []
self.load()
@abstractmethod
def load(self) -> None:
"""self.targets_list を構築する"""
pass
def __str__(self):
return json.dumps(self.targets_list, indent=2, ensure_ascii=False)
class HostsIniDataset(IDatasetInterface):
"""
CSV形式(1行目ヘッダー:hosts.ini)からターゲットリストを読み込む。
fetch ip address and hostname list from config life.
- hosts.iniの内容を読み込んで、IPアドレスとホスト名のリストを作成する例。
- hosts.iniは以下のような形式を想定(1行目がヘッダー、2行目以降がデータ):
ipaddr,hostname,username,password,command
* username, password, commandは任意
"""
def __init__(self):
#self.targets_file = os.path.join(os.getcwd(), settings_dir, config_file)
self.targets_file = Config.TARGET_HOSTS
super().__init__()
def load(self) -> None:
with open(self.targets_file, 'r', encoding="utf-8") as f:
headers = []
for cnt, line in enumerate(f.readlines()):
line = line.rstrip("\n")
if not line:
continue
items = line.split(",")
if cnt == 0:
headers = items
continue
device = {headers[idx]: (item if item != "" else None)
for idx, item in enumerate(items)}
if device:
self.targets_list.append(device)
class ClusterIniDataset(IDatasetInterface):
"""
configparser形式(cluster.ini)からクラスタのヘッドノード情報を読み込む。
[cluster-a]
ipaddr = 192.168.10.1
hostname = cluster-a-headnode
command = df
not defined in cluster.ini, use default username and password.
"""
def __init__(self):
self.targets_file = Config.TARGET_CLUSTER
super().__init__()
def load(self) -> None:
parser = configparser.ConfigParser()
parser.read(self.targets_file, encoding="utf-8")
for section in parser.sections():
ipaddr = parser.get(section, "ipaddr", fallback=None)
hostname = parser.get(section, "hostname", fallback=section)
command = parser.get(section, "command", fallback=section)
if not ipaddr and not hostname:
continue
self.targets_list.append({
"ipaddr": ipaddr,
"hostname": hostname,
"command": command,
})
targets/*
hostname
ipa
client
ipa
client
hostname,ipaddr,username,password,command
ipa,192.168.56.104,,,
client,192.168.56.103,,,
client,192.168.56.103,,,/usr/bin/ls
[Xeon001]
ipaddr = 192.168.56.104
hostname = ipa
command = /home/app/thread_backend/scripts/command.sh
[Xeon002]
ipaddr = 192.168.56.103
hostname = client
command = /home/app/thread_backend/scripts/command.sh
[Xeon003]
ipaddr = 192.168.56.104
hostname = ipa
command = /home/app/thread_backend/scripts/command.sh
[MAC]
command = ls
; ヘッドノードが未確定/廃止予定のクラスタはコメントアウトで除外
;[cluster-l]
;ipaddr =
;hostname =
;command =
scripts/*
#!/bin/bash
RES=$(systemctl status ypbind)
RPM=$(rpm -qa | grep postgres)
echo "$RES"
echo "$RPM"
exit 0
config/settings.py
#!/usr/bin/python3
import logging
import os
class Config(object):
# SSH Connection Setting
USERNAME = "root"
PASSWORD = "rootroot"
PORT = 22
TIMEOUT = 10
CLUSTER_COMMAND_TIMEOUT = 3600
# Thead
MAX_WORKERS = 3
# File Settings
HOSTS_INI_FILE = "target_hosts.ini"
CLUSTER_INI_FILE = "target_cluster.ini"
SCRIPT = "command.sh"
# Target Settings
TARGET_PATH = "/home/app/thread_backend/"
TARGET_SCRIPT = os.path.join(TARGET_PATH + "scripts", SCRIPT)
TARGET_HOSTS = os.path.join(TARGET_PATH + "targets", HOSTS_INI_FILE)
TARGET_CLUSTER = os.path.join(TARGET_PATH + "targets", CLUSTER_INI_FILE)
TARGET_OUTPUT_DIR = os.path.join(TARGET_PATH + "out")
# Logger level
LEVEL = logging.WARN
#LEVEL = logging.INFO
#LEVEL = logging.DEBUG
# -----------------------------
# Logger
# -----------------------------
def setup_logger(name, level=logging.INFO):
logger = logging.getLogger(name)
logger.setLevel(level)
if not logger.handlers:
h = logging.StreamHandler()
fmt = logging.Formatter(
"%(asctime)s %(levelname)s %(name)s: %(threadName)s: %(message)s")
h.setFormatter(fmt)
logger.addHandler(h)
return logger