Linuxプロセス間通信におけるパイプの仕組みと実装

プロセス間通信の基本概念

Linuxでは各プロセスが独立した仮想アドレス空間を持つため、別のプロセスとデータをやり取りするには何らかの共有リソースが必要です。すべてのIPC(Inter-Process Communication)手法は、まず「複数のプロセスが同一のリソースにアクセスできる環境を構築する」ことから始まります。その上で、一方がデータを書き込み、他方が読み出すことで通信が成立します。

無名パイプ

無名パイプは、カーネル内部のバッファを介してデータを流す単方向の通信路です。親プロセスがpipe()システムコールで読み取り用・書き込み用の2つのファイルディスクリプタを取得し、その後fork()で子プロセスを生成すると、子プロセスは親のファイルディスクリプタテーブルを引き継ぎます。親子それぞれで不要な端を閉じることで、単方向の通信チャネルが完成します。

pipe()システムコール

int pipe(int pipefd[2])は、配列の要素0に読み取り用、要素1に書き込み用のファイルディスクリプタを格納する出力型パラメータを持ちます。

無名パイプの特徴

  • 半二重通信:データは一方向にしか流れません。双方向通信が必要な場合は2本のパイプを作成します。
  • ファイルと同じライフサイクル:パイプの実体はカーネル内のファイルバッファであり、プロセス終了時に自動的に解放されます。
  • 親族関係が必要:親子・祖孫・兄弟など、血縁関係にあるプロセス間で利用されます。これはfork()によるファイルディスクリプタの継承が前提だからです。
  • バイトストリーム:書き込み回数と読み取り回数は1対1で対応しません。複数回の書き込みが1回の読み取りでまとめて取得されることもあります。
  • 同期機構の内蔵:パイプバッファが満杯のとき書き込み側はブロックし、空のとき読み取り側はブロックします。これにより読み書きの同期が自動的に取られます。

パイプ通信における4つの状況

  1. 読み取り側がバッファ内の全データを消費した後、書き込み側が新しいデータを送らない場合、読み取り側は待機(ブロック)します。
  2. 書き込み側がパイプバッファを満杯にした場合、それ以上書き込めずブロックします。
  3. 書き込み側のfdがすべて閉じられた後、読み取り側がバッファを空にして再度read()を呼ぶと、戻り値0が返り、EOFを検知したことを示します。
  4. 読み取り側のfdがすべて閉じられた状態で書き込みを続けると、OSは無駄を排除するためSIGPIPEシグナルを送信し、書き込みプロセスを強制終了します。

基本的なパイプ通信の実装例

#include <iostream>
#include <cstring>
#include <cerrno>
#include <unistd.h>
#include <sys/wait.h>
#include <sys/types.h>

int main()
{
    int fds[2] = {0};

    // パイプの作成
    if (pipe(fds) < 0) {
        std::cerr << "パイプ作成失敗: " << errno << " - " << strerror(errno) << std::endl;
        return 1;
    }

    std::cout << "読み取りfd: " << fds[0] << std::endl;
    std::cout << "書き込みfd: " << fds[1] << std::endl;

    pid_t child_pid = fork();
    if (child_pid < 0) {
        std::cerr << "fork失敗" << std::endl;
        return 1;
    }

    if (child_pid == 0) {
        // 子プロセス:書き込み側として動作
        close(fds[0]);

        int count = 0;
        while (count < 5) {
            char ch = 'A';
            ssize_t written = write(fds[1], &ch, 1);
            if (written > 0) {
                std::cout << "子: " << ++count << "文字目書き込み完了" << std::endl;
            }
            sleep(1);
        }

        close(fds[1]);
        exit(0);
    }

    // 親プロセス:読み取り側として動作
    close(fds[1]);

    char recv_buf[1024];
    while (true) {
        ssize_t bytes = read(fds[0], recv_buf, sizeof(recv_buf) - 1);
        if (bytes > 0) {
            recv_buf[bytes] = '\0';
            std::cout << "親: 受信データ[" << recv_buf << "]" << std::endl;
        } else if (bytes == 0) {
            std::cout << "親: 子プロセスが書き込み端を閉じました" << std::endl;
            break;
        } else {
            std::cerr << "親: 読み取りエラー " << strerror(errno) << std::endl;
            break;
        }
    }

    close(fds[0]);
    int status = 0;
    waitpid(child_pid, &status, 0);
    std::cout << "子プロセス終了コード: " << (status & 0x7F) << std::endl;

    return 0;
}

プロセスプールによるタスク分散の実装

複数の子プロセスを生成し、親プロセスからタスク指令をパイプ経由で各ワーカーに配送するパターンの実装例です。

// TaskManager.hpp
#pragma once
#include <iostream>
#include <vector>
#include <unistd.h>

typedef void (*TaskFunc)();

void LogTask() {
    std::cout << "[PID:" << getpid() << "] ログ出力タスク実行中..." << std::endl;
}

void DbTask() {
    std::cout << "[PID:" << getpid() << "] データベース操作タスク実行中..." << std::endl;
}

void HttpTask() {
    std::cout << "[PID:" << getpid() << "] HTTPリクエスト処理タスク実行中..." << std::endl;
}

class TaskManager {
private:
    std::vector<TaskFunc> tasks;
public:
    TaskManager() {
        tasks.push_back(LogTask);
        tasks.push_back(DbTask);
        tasks.push_back(HttpTask);
    }

    void run(int cmd) {
        if (cmd >= 0 && cmd < (int)tasks.size()) {
            tasks[cmd]();
        }
    }

    int size() const { return tasks.size(); }
};
// ProcessPool.cpp
#include <iostream>
#include <vector>
#include <cassert>
#include <unistd.h>
#include <sys/wait.h>
#include <sys/types.h>
#include "TaskManager.hpp"

const int WORKER_COUNT = 3;
TaskManager task_mgr;

class WorkerEndpoint {
private:
    static int seq;
public:
    pid_t pid;
    int write_fd;
    std::string alias;

    WorkerEndpoint(pid_t p, int fd) : pid(p), write_fd(fd) {
        char buf[64];
        snprintf(buf, sizeof(buf), "worker-%d[pid=%d,fd=%d]", seq++, pid, write_fd);
        alias = buf;
    }
};

int WorkerEndpoint::seq = 0;

void workerLoop() {
    while (true) {
        int cmd = 0;
        ssize_t n = read(STDIN_FILENO, &cmd, sizeof(int));
        if (n == sizeof(int)) {
            task_mgr.run(cmd);
        } else if (n == 0) {
            std::cout << "ワーカー " << getpid() << " は終了指令を受信" << std::endl;
            break;
        } else {
            break;
        }
    }
}

void spawnWorkers(std::vector<WorkerEndpoint>& workers) {
    std::vector<int> inherited_fds;

    for (int i = 0; i < WORKER_COUNT; i++) {
        int pipe_pair[2] = {0};
        if (pipe(pipe_pair) < 0) {
            perror("pipe");
            continue;
        }

        pid_t pid = fork();
        if (pid < 0) {
            perror("fork");
            continue;
        }

        if (pid == 0) {
            // 子プロセス:継承した書き込みfdをすべて閉じる
            for (int fd : inherited_fds) close(fd);
            close(pipe_pair[1]);
            dup2(pipe_pair[0], STDIN_FILENO);
            workerLoop();
            close(pipe_pair[0]);
            exit(0);
        }

        // 親プロセス
        close(pipe_pair[0]);
        workers.push_back(WorkerEndpoint(pid, pipe_pair[1]));
        inherited_fds.push_back(pipe_pair[1]);
    }
}

int showMenu() {
    std::cout << "=========================" << std::endl;
    std::cout << " 0: ログ出力" << std::endl;
    std::cout << " 1: DB操作" << std::endl;
    std::cout << " 2: HTTPリクエスト" << std::endl;
    std::cout << " 3: 終了" << std::endl;
    std::cout << "=========================" << std::endl;
    std::cout << "選択> ";
    int cmd;
    std::cin >> cmd;
    return cmd;
}

void dispatchTasks(const std::vector<WorkerEndpoint>& workers) {
    int round = 0;
    while (true) {
        int cmd = showMenu();
        if (cmd == 3) break;
        if (cmd < 0 || cmd >= task_mgr.size()) continue;

        int idx = round % workers.size();
        round++;

        std::cout << "ディスパッチ: " << workers[idx].alias
                  << " -> タスク番号 " << cmd << std::endl;

        write(workers[idx].write_fd, &cmd, sizeof(cmd));
        sleep(1);
    }
}

void cleanupWorkers(std::vector<WorkerEndpoint>& workers) {
    for (auto& w : workers) {
        std::cout << "終了通知送信 -> " << w.alias << std::endl;
        close(w.write_fd);
        waitpid(w.pid, nullptr, 0);
        std::cout << "回収完了: PID=" << w.pid << std::endl;
    }
}

int main() {
    std::vector<WorkerEndpoint> workers;

    spawnWorkers(workers);
    dispatchTasks(workers);
    cleanupWorkers(workers);

    return 0;
}

この実装では、子プロセス生成時に親が過去に開いたパイプの書き込みfdを子プロセス側で閉じることで、各ワーカーが自身専用のパイプからのみ読み取るようにしています。またdup2()でパイプの読み取りfdを標準入力に割り当てることで、ワーカー側はread(0, ...)で統一的に指令を受け取れます。

名前付きパイプ(FIFO)

無名パイプは血縁関係のあるプロセス間でしか使えませんが、名前付きパイプ(FIFO)はファイルシステム上に名前を持つ特殊ファイルとして存在するため、互いに無関係なプロセス間でも通信が可能です。

mkfifoコマンドまたはmkfifo()システムコールでFIFOファイルを作成します。このファイルはディスク上にinodeとパスを持ちますが、データブロックは持ちません。書き込まれたデータはカーネルのメモリバッファに留まり、ディスクへ書き込まれることはないため、ファイルサイズは常に0です。

mkfifo()システムコール

int mkfifo(const char *pathname, mode_t mode)は、指定したパスにFIFOファイルを作成します。pathnameはファイルのパス、modeはパーミッションを指定します。

FIFOの重要な特性として、open()が相手側のオープンを待つ点があります。読み取り専用で開こうとすると、他のプロセスが書き込み専用で開くまでブロックします(逆も同様)。これにより通信のタイミングを自動的に同期できます。

名前付きパイプによるサーバ・クライアント通信の実装

// ipc_common.hpp
#pragma once
#include <iostream>
#include <string>

const std::string FIFO_PATH = "/tmp/my_ipc_fifo";
const int BUF_SIZE = 1024;
const mode_t FIFO_MODE = 0666;
// fifo_server.cpp
#include <iostream>
#include <cerrno>
#include <cstring>
#include <sys/types.h>
#include <sys/stat.h>
#include <fcntl.h>
#include <unistd.h>
#include "ipc_common.hpp"

int main()
{
    umask(0);

    if (mkfifo(FIFO_PATH.c_str(), FIFO_MODE) != 0) {
        std::cerr << "FIFO作成失敗: " << errno << " - " << strerror(errno) << std::endl;
        return 1;
    }
    std::cout << "FIFOファイル作成完了: " << FIFO_PATH << std::endl;

    // 読み取り専用でオープン(クライアントが書き込みオープンするまでブロック)
    int read_fd = open(FIFO_PATH.c_str(), O_RDONLY);
    if (read_fd < 0) {
        std::cerr << "FIFOオープン失敗: " << strerror(errno) << std::endl;
        return 2;
    }
    std::cout << "クライアント接続確認、通信開始" << std::endl;

    char msg[BUF_SIZE];
    while (true) {
        ssize_t n = read(read_fd, msg, BUF_SIZE - 1);
        if (n > 0) {
            msg[n] = '\0';
            std::cout << msg << std::flush;
        } else if (n == 0) {
            std::cout << "\nクライアントが切断しました" << std::endl;
            break;
        } else {
            std::cerr << "読み取りエラー: " << strerror(errno) << std::endl;
            break;
        }
    }

    close(read_fd);
    unlink(FIFO_PATH.c_str());
    return 0;
}
// fifo_client.cpp
#include <iostream>
#include <cerrno>
#include <cstring>
#include <strings.h>
#include <sys/types.h>
#include <sys/stat.h>
#include <fcntl.h>
#include <unistd.h>
#include "ipc_common.hpp"

int main()
{
    int write_fd = open(FIFO_PATH.c_str(), O_WRONLY);
    if (write_fd < 0) {
        std::cerr << "FIFOオープン失敗: " << errno << " - " << strerror(errno) << std::endl;
        return 1;
    }

    std::cout << "メッセージを入力してください (quit で終了):" << std::endl;

    char input[BUF_SIZE];
    while (true) {
        std::cout << "> ";
        if (fgets(input, sizeof(input), stdin) == nullptr) break;

        if (strcasecmp(input, "quit\n") == 0) break;

        ssize_t n = write(write_fd, input, strlen(input));
        if (n < 0) {
            std::cerr << "書き込みエラー: " << strerror(errno) << std::endl;
            break;
        }
    }

    close(write_fd);
    return 0;
}

unlink()はファイルのリンクを削除するシステムコールです。参照カウントが0になると実際のファイルが削除されます。サーバ側で通信終了後にunlink()を呼ぶことで、FIFOファイルをクリーンアップします。

FIFOと無名パイプの根本的な違いは「関係のないプロセス間で通信できるか」にあります。無名パイプはfork()によるfd継承が前提ですが、FIFOはファイルパスという名前を使って独立したプロセスが同じリソースにアクセスできます。どちらもデータをカーネルバッファ上で受け渡すメモリレベルの通信であり、ディスクI/Oは発生しません。

タグ: linux IPC pipe fork mkfifo

9月13日 07:03 投稿