KnowHow

技術的なメモを中心にまとめます。
検索にて調べることができます。

Backend_Threadプログラム

登録日 :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/ (ログ保存場所)
  • start_clusters.sh
#!/bin/bash

PROGRAM_DIR=/home/app/thread_backend
$PROGRAM_DIR/main.py 2 2>&1 | tee -a $PROGRAM_DIR/logs/execute_clusters.log
  • start_hosts.sh
#!/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/*

  • src/thread_workers.py
#!/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
  • src/reporter.py
#!/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))
  • src/dataset.py
#!/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/*

  • target_hosts.ini
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
  • target_cluster.ini
[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/*

  • command.sh
#!/bin/bash

RES=$(systemctl status ypbind)
RPM=$(rpm -qa | grep postgres)
echo "$RES"
echo "$RPM"
exit 0

config/settings.py

  • 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