Akar Masalah: Streaming Subprocess pada Asynchronous Queue

Menjalankan task komputasi berat, manipulasi biner (seperti FFmpeg, ImageMagick), atau eksekusi isolated script di worker backend sering kali memerlukan subprocess OS. Pada arsitektur asinkron berbasis event loop, interaksi antarmuka pipa I/O (stdout/stderr) antara parent process dan child process sering memicu tiga kegagalan utama:

  • Buffer Desynchronization: Streaming output biner atau framed-protocol melalui asynchronous stream reader memecah payload menjadi chunk arbitrer. Menulis chunk ini secara langsung ke shared cache (misalnya Redis hash atau memcached) tanpa batasan sekuensial atau pemisah frame presisi menimbulkan race condition dan korupsi data stream.
  • Event Loop Starvation: Output streaming bervolume tinggi membanjiri buffer pipa OS. Jika handler loop membaca stream tanpa yield point kooperatif atau mencoba memproses deserialisasi secara intensif di loop thread utama, event loop kehabisan siklus CPU untuk melayani task lain.
  • Zombie dan Orphan Processes: Task timeout atau worker thread crash yang gagal mengirim SIGTERM/SIGKILL dan tidak memanggil sistem waitpid() meninggalkan zombie process pada tabel proses OS, menguras PID pool dan file descriptor.

Pola Desain: Process Sentinel dan Stream Filter

GNU Emacs menangani isolasi proses eksternal melalui pemisahan tegas antara penanganan aliran data dan siklus hidup proses: Filter Functions dan Process Sentinels. Arsitektur backend modern dapat mengadopsi prinsip ini untuk task queue worker.

1. Process Sentinel (Rekonsiliasi Status)

Sentinel adalah fungsi deterministik yang hanya dieksekusi saat child process mengalami transisi status (exit, crash, terbunuh sinyal). Sentinel bertanggung jawab mutlak atas:

  • Membaca exit code subprocess dan status pemutusan pipa.
  • Menjalankan rekonsiliasi state ke task broker (acknowledging message, requeue, atau logging failed state).
  • Memastikan resource cleanup lokal: penutupan pipe file descriptor dan flushes ring buffer.

2. Process Filter & Bounded Ring Buffer

Stream filter bertindak sebagai middleware antara stdout subprocess dan storage tujuan. Alih-alih melakukan broadcast langsung, chunk biner diarahkan ke ring buffer berukuran tetap (bounded buffer). Ketika laju stream child process melampaui kemampuan konsumsi downstream (misal upload S3 atau write Redis), backpressure aktif: worker berhenti memanggil read() dari pipe, memaksa buffer pipa level kernel (pipe buffer) penuh dan menahan eksekusi child process (sistem blokir write() di level kernel).

Implementasi Subprocess Worker Minimalis

Kode Python berikut mengimplementasikan subprocess worker dengan backpressure, queue bounded buffer, dan pola process sentinel menggunakan pustaka standar asyncio.

import asyncio
import os
import signal
import sys
from typing import Callable, Coroutine, Optional

# ponytail: fixed queue capacity to 32 chunks. Upgrade to dynamic byte-size watermarking if chunk sizes vary widely.
MAX_BUFFER_CHUNKS = 32

class SubprocessTaskWorker:
    def __init__(self, cmd: list[str], timeout: float = 30.0):
        self.cmd = cmd
        self.timeout = timeout
        self.buffer: asyncio.Queue[Optional[bytes]] = asyncio.Queue(maxsize=MAX_BUFFER_CHUNKS)
        self.process: Optional[asyncio.subprocess.Process] = None

    async def _stream_filter(self, stream: asyncio.StreamReader) -> None:
        """Membaca stdout subprocess dan mengalirkan ke bounded queue dengan backpressure."""
        try:
            while not stream.at_eof():
                chunk = await stream.read(65536)
                if not chunk:
                    break
                # Backpressure: put() akan pause jika buffer penuh
                await self.buffer.put(chunk)
        finally:
            await self.buffer.put(None)  # Sinyal sentinel stream EOF

    async def _process_sentinel(self, return_code: int) -> None:
        """Rekonsiliasi status akhir task dan pembersihan resource."""
        if return_code == 0:
            # Sukses: finalisasi state task
            return
        # Tangani error state jika process exit dengan code selain 0
        raise RuntimeError(f"Subprocess gagal dieksekusi dengan return code {return_code}")

    async def execute(self, sink_consumer: Callable[[bytes], Coroutine]) -> None:
        """Eksekusi subprocess, hubungkan consumer, dan pastikan cleanup."""
        self.process = await asyncio.create_subprocess_exec(
            *self.cmd,
            stdout=asyncio.subprocess.PIPE,
            stderr=asyncio.subprocess.PIPE,
            preexec_fn=os.setsid  # Grup proses baru untuk isolasi sinyal
        )

        stream_task = asyncio.create_task(self._stream_filter(self.process.stdout))
        
        async def drain_buffer():
            while True:
                chunk = await self.buffer.get()
                if chunk is None:
                    self.buffer.task_done()
                    break
                await sink_consumer(chunk)
                self.buffer.task_done()

        drain_task = asyncio.create_task(drain_buffer())

        try:
            # Jalankan dengan timeout batas atas
            await asyncio.wait_for(self.process.wait(), timeout=self.timeout)
            await asyncio.gather(stream_task, drain_task)
            await self._process_sentinel(self.process.returncode)
        except (asyncio.TimeoutError, asyncio.CancelledError, Exception):
            # Mitigasi Zombie Process: Bunuh seluruh process group jika crash/timeout
            if self.process and self.process.returncode is None:
                try:
                    os.killpg(os.getpgid(self.process.pid), signal.SIGKILL)
                except ProcessLookupError:
                    pass
                await self.process.wait()  # Hindari status defunct/zombie
            stream_task.cancel()
            drain_task.cancel()
            raise

# --- Self-check run ---
async def _test():
    collector = []
    async def sink(chunk: bytes):
        collector.append(chunk)

    worker = SubprocessTaskWorker([sys.executable, "-c", "import sys; [sys.stdout.write('data') for _ in range(3)]"])
    await worker.execute(sink)
    assert b"".join(collector) == b"datadatadata", "Buffer stream mismatch"

if __name__ == "__main__":
    asyncio.run(_test())

Implementasi dilewati: penulisan chunk biner ke storage persistent langsung; tambahkan adapter stream saat integrasi S3 multi-part upload atau framing Redis Streams.

Manajemen Siklus Hidup dan Pembasmian Zombie

Kegagalan paling umum dalam task worker asinkron adalah pengabaian child process saat timeout tercapai. Pemanggilan process.kill() bawaan modul async hanya menargetkan parent PID dari command yang dijalankan, meninggalkan child process turunan (misal command berupa shell script yang memanggil binary lain) berjalan tanpa pengawasan (orphan process).

  • Process Grouping: Gunakan os.setsid pada spawn hook POSIX. Hal ini memposisikan subprocess dan seluruh turunannya ke dalam process group unik.
  • Signal Broadcasting: Ketika task timeout terpicu, kirim sinyal langsung ke seluruh group menggunakan os.killpg(os.getpgid(pid), signal.SIGKILL).
  • Deterministic Reaping: Wajib jalankan await process.wait() di blok finally atau except. Tanpa eksekusi wait, kernel tidak akan menghapus entri dari process table, menyebabkan zombie accumulation.

Evaluasi Arsitektur

Gunakan pendekatan subprocess worker terisolasi ini apabila payload task tidak kompatibel dengan multi-threading (terkendala thread safety pustaka native biner atau GIL). Jika overhead alokasi proses (forking latency) terlalu membebani throughput broker, gantikan proses spawning ad-hoc dengan pre-warmed persistent worker pools menggunakan protokol IPC berbasis Unix Domain Socket dengan framing ukuran pasti.