Not a member of Pastebin yet?
Sign Up,
it unlocks many cool features!
- #!/usr/bin/env python3
- import asyncio
- import os
- import sys
- import socket
- import subprocess
- from pathlib import Path
- from textwrap import dedent
- # =====================================================================
- # SYSTEMD TEMPLATE GENERATION FUNCTIONS (Modern Python Strings)
- # =====================================================================
- SERVICE_UNIT = """
- [Unit]
- Description=Resilient Async Python Service (FD-Store)
- Requires={socket_name}
- After=network.target
- [Service]
- Type=simple
- ExecStart={exec_path}
- Environment=PYTHONUNBUFFERED=1
- # FD-Store and Recovery configurations
- FileDescriptorStoreMax=100
- Restart=on-failure
- RestartSec=0.1s
- [Install]
- WantedBy=default.target
- """
- SOCKET_UNIT = """[Unit]
- Description=Resilient Async Python Service Socket
- [Socket]
- # Listens on localhost port {port} by default (ideal for user units)
- ListenStream={port}
- [Install]
- WantedBy=sockets.target
- """
- def create_service_unit(socket_name: str, exec_path: str) -> str:
- """Generates the systemd service configuration using modern multi-line formatting."""
- return SERVICE_UNIT.format(socket_name=socket_name, exec_path=exec_path)
- def create_socket_unit(port: int = 8080) -> str:
- """Generates the systemd socket configuration."""
- return SOCKET_UNIT.format(port=port)
- # =====================================================================
- # SYSTEMD COMMAND HELPER
- # =====================================================================
- def run_systemctl_user(args: list) -> bool:
- """Helper function to safely run systemctl commands in user space."""
- try:
- cmd = ["systemctl", "--user"] + args
- print(f"Running: {' '.join(cmd)}")
- subprocess.run(cmd, check=True, stdout=subprocess.PIPE, stderr=subprocess.PIPE)
- return True
- except subprocess.CalledProcessError as e:
- print(
- f"[ERROR] systemctl command failed: {e.stderr.decode().strip()}",
- file=sys.stderr,
- )
- return False
- def resilient_accept(conn):
- send_fd_to_systemd(conn.fileno())
- # =====================================================================
- # INSTALLATION & DEINSTALLATION LOGIC
- # =====================================================================
- def install_systemd_user_units():
- """Automatically installs, enables, and starts the systemd user units."""
- # human note:
- # Automatic installs are not good, if everything emerges from a single file.
- # The user doesn't know from where the file came and where it was installed.
- #
- print("Initializing systemd user unit installation...")
- # 1. Resolve structural paths
- script_path = Path(__file__).resolve()
- home_dir = Path.home()
- # Standard systemd user configuration directory
- systemd_user_dir = home_dir / ".config" / "systemd" / "user"
- service_file = systemd_user_dir / "async-resilient.service"
- socket_file = systemd_user_dir / "async-resilient.socket"
- # 2. Ensure the configuration directories exist
- systemd_user_dir.mkdir(parents=True, exist_ok=True)
- # 3. Generate unit contents via template functions
- service_content = create_service_unit(
- socket_name="async-resilient.socket", exec_path=str(script_path)
- )
- socket_content = create_socket_unit(port=8080)
- # 4. Write unit configurations to disk
- try:
- service_file.write_text(service_content, encoding="utf-8")
- socket_file.write_text(socket_content, encoding="utf-8")
- # Ensure the script itself remains executable
- # holy shit, the ai know octal notation!!!
- script_path.chmod(script_path.stat().st_mode | 0o755)
- print(f"[SUCCESS] Configurations written successfully.")
- # 5. Execute systemd lifecycle commands automatically
- print("\nExecuting systemd lifecycle commands...")
- if run_systemctl_user(["daemon-reload"]):
- if run_systemctl_user(["enable", "async-resilient.socket"]):
- if run_systemctl_user(["start", "async-resilient.socket"]):
- print("\n[SUCCESS] Server socket is up and running via systemd!")
- print("You can now connect to port 8080.")
- except Exception as e:
- print(f"[ERROR] Failed during systemd initialization: {e}", file=sys.stderr)
- sys.exit(1)
- def stop_and_disable_units():
- """Stops and disables the running systemd socket and service units."""
- print("Stopping and disabling systemd user units...")
- # Stop the units first so connections drop cleanly
- run_systemctl_user(["stop", "async-resilient.service"])
- run_systemctl_user(["stop", "async-resilient.socket"])
- # Disable them so they don't boot next time systemd starts
- run_systemctl_user(["disable", "async-resilient.socket"])
- run_systemctl_user(["disable", "async-resilient.service"])
- # Reload daemon to clear out the state
- run_systemctl_user(["daemon-reload"])
- print("[SUCCESS] Units have been successfully stopped and disabled.")
- # =====================================================================
- # CORE NETWORK DAEMON LOGIC
- # =====================================================================
- def send_fd_to_systemd(fd: int):
- """Sends an active client File Descriptor to the systemd FD-Store."""
- notify_socket_path = os.environ.get("NOTIFY_SOCKET")
- if not notify_socket_path:
- return
- if notify_socket_path.startswith("@"):
- notify_socket_path = "\x00" + notify_socket_path[1:]
- try:
- with socket.socket(socket.AF_UNIX, socket.SOCK_DGRAM) as sock:
- sock.connect(notify_socket_path)
- sock.sendmsg(
- [b"FDSTORE=1\n"],
- [
- (
- socket.SOL_SOCKET,
- socket.SCM_RIGHTS,
- int.to_bytes(fd, 4, sys.byteorder),
- )
- ],
- )
- except Exception as e:
- print(f"FD-Store Sync Error: {e}", file=sys.stderr)
- async def handle_client(
- reader: asyncio.StreamReader, writer: asyncio.StreamWriter, is_recovered=False
- ):
- """Asynchronously processes ongoing network clients."""
- try:
- addr = writer.get_extra_info("peername")
- print(f"Processing client from {addr} (Recovered: {is_recovered})")
- if is_recovered:
- writer.write(
- b"Your server crashed, but asyncio recovered your connection flawlessly!\n"
- )
- await writer.drain()
- else:
- resilient_accept(writer.get_extra_info("socket"))
- writer.write(b"Processing request... please wait 10 seconds.\n")
- await writer.drain()
- await asyncio.sleep(10)
- writer.write(b"Processing finished successfully.\n")
- await writer.drain()
- except Exception as e:
- print(f"Error handling client: {e}", file=sys.stderr)
- finally:
- writer.close()
- await writer.wait_closed()
- async def run_server():
- """Main daemon loop hooked into the systemd socket lifecycle."""
- listen_fds = os.environ.get("LISTEN_FDS")
- if not listen_fds:
- print(
- "Error: No systemd sockets passed down. Run 'init' or launch via systemd.socket.",
- file=sys.stderr,
- )
- sys.exit(1)
- num_fds = int(listen_fds)
- print(f"Received {num_fds} total File Descriptors from systemd.")
- # Main Listening Socket is always FD 3
- SYSTEMD_FIRST_FD = 3
- server_sock = socket.fromfd(SYSTEMD_FIRST_FD, socket.AF_INET, socket.SOCK_STREAM)
- server = await asyncio.start_server(
- lambda r, w: asyncio.create_task(handle_client(r, w, is_recovered=False)),
- sock=server_sock,
- )
- print("Async Server successfully bound and accepting connections...")
- # Rehydrate any existing client connections (FD 4+) rescued from a previous crash
- if num_fds > 1:
- print(f"Resurrecting {num_fds - 1} surviving client connections...")
- for fd_num in range(4, 3 + num_fds):
- try:
- recovered_sock = socket.fromfd(
- fd_num, socket.AF_INET, socket.SOCK_STREAM
- )
- reader, writer = await asyncio.open_connection(sock=recovered_sock)
- asyncio.create_task(handle_client(reader, writer, is_recovered=True))
- except Exception as e:
- print(
- f"Failed to recover existing connection on FD {fd_num}: {e}",
- file=sys.stderr,
- )
- async with server:
- await server.serve_forever()
- # =====================================================================
- # APPLICATION ENTRYPOINT
- # =====================================================================
- if __name__ == "__main__":
- if len(sys.argv) > 1:
- action = sys.argv[1].lower()
- if action == "init":
- install_systemd_user_units()
- elif action == "stop":
- stop_and_disable_units()
- else:
- print(
- f"Unknown argument: '{sys.argv[1]}'. Use 'init' or 'stop'.",
- file=sys.stderr,
- )
- else:
- try:
- asyncio.run(run_server())
- except KeyboardInterrupt:
- print("\nServer shutting down gracefully.")
Advertisement