常時稼働エージェントの設計パターン
バックグラウンドプロセス、デーモンエージェント、永続的な接続管理を学びます。
「常時稼働エージェントの設計パターン」はCoddyKit上の無料AI Agentsレッスンです。 これはレッスン1/4です。 下記で完全なレッスンを無料で読むことができます。その後、ブラウザ内の組み込みコードエディタと24時間対応のAIチューターでハンズオン演習できます。 これはAI Agents学習パスの一部であり、ウェブとCoddyKitアプリ全体で進捗が同期されます。 AI Agentsコースには全4レッスンが含まれています。
常時稼働エージェントとは
常時稼働エージェントはバックグラウンドサービスとして継続的に動作し、イベントを待ち受けながら自律的にアクションを実行します。リクエストとレスポンスのやり取りごとに終了するエージェントとは異なり、インタラクションの間も動作し続け、時間の経過に伴って状態を維持します。
デーモンプロセスパターン
デーモンプロセスは、ターミナルセッションから独立してバックグラウンドで動作します。ターミナルを閉じた後もエージェントを稼働させるには、Pythonのデーモンスレッドまたはプロセススーパーバイザーを使用します。
import threading
import time
import signal
import sys
shutdown_flag = threading.Event()
def agent_main_loop():
print('Agent daemon started')
while not shutdown_flag.is_set():
try:
# Agent work: check for events, process tasks
perform_agent_cycle()
shutdown_flag.wait(timeout=60) # Sleep 60s, wakes on shutdown
except Exception as e:
print(f'Agent loop error: {e}')
shutdown_flag.wait(timeout=5) # Brief pause on error
print('Agent daemon stopped')
def perform_agent_cycle():
print(f'Agent cycle at {time.strftime("%H:%M:%S")}')
# Check emails, process queue, run scheduled tasks
def handle_signal(signum, frame):
print(f'Signal {signum} received, shutting down...')
shutdown_flag.set()
# Register signal handlers for graceful shutdown
signal.signal(signal.SIGTERM, handle_signal)
signal.signal(signal.SIGINT, handle_signal)
# Start as daemon thread
thread = threading.Thread(target=agent_main_loop, daemon=True)
thread.start()
print('Agent running in background')ウォッチドッグによるクラッシュ時の自動再起動
ウォッチドッグはエージェントプロセスを監視し、クラッシュした場合は再起動します。これは本番環境に不可欠です。エージェントは予期しないエラーに遭遇することが避けられないため、自動的に復旧できなければなりません。
import subprocess
import time
import logging
logger = logging.getLogger('watchdog')
class AgentWatchdog:
def __init__(self, agent_script: str, max_restarts: int = 10, restart_delay: float = 5.0):
self.agent_script = agent_script
self.max_restarts = max_restarts
self.restart_delay = restart_delay
self.restart_count = 0
self.process = None
def start(self):
self.process = subprocess.Popen(
['python', self.agent_script],
stdout=subprocess.PIPE,
stderr=subprocess.STDOUT
)
logger.info(f'Agent started (PID: {self.process.pid})')
def run_forever(self):
self.start()
while True:
return_code = self.process.wait()
logger.warning(f'Agent exited with code {return_code}')
if self.restart_count >= self.max_restarts:
logger.error(f'Max restarts ({self.max_restarts}) reached. Stopping watchdog.')
break
self.restart_count += 1
logger.info(f'Restarting agent (attempt {self.restart_count})...')
time.sleep(self.restart_delay * self.restart_count) # Backoff
self.start()
print('AgentWatchdog defined')永続 WebSocket 接続
リアルタイムでイベントを配信するには、永続的な WebSocket 接続を維持します。接続が切断された場合は自動的に再接続します。これは常時稼働エージェントにおける重要な課題です。
import asyncio
import websockets
import json
async def persistent_websocket_connection(uri: str, on_message):
backoff = 1
max_backoff = 60
while True: # Reconnect forever
try:
print(f'Connecting to {uri}')
async with websockets.connect(uri, ping_interval=30, ping_timeout=10) as ws:
print('WebSocket connected')
backoff = 1 # Reset backoff on successful connection
async for raw_message in ws:
try:
message = json.loads(raw_message)
await on_message(message)
except json.JSONDecodeError:
print(f'Invalid JSON message: {raw_message[:100]}')
except websockets.exceptions.ConnectionClosed as e:
print(f'WebSocket closed: {e}. Reconnecting in {backoff}s')
except Exception as e:
print(f'WebSocket error: {e}. Reconnecting in {backoff}s')
await asyncio.sleep(backoff)
backoff = min(backoff * 2, max_backoff) # Exponential backoff
async def handle_ws_message(message: dict):
print(f'Received: {message}')
print('Persistent WebSocket connection function defined')ハートビートチェック
ハートビートによって、エージェントが稼働し、処理を続けていることを確認できます。N秒ごとにハートビート信号を送信し、ハートビートが止まった場合は、ウォッチドッグがエージェントの停止またはハングアップを検知します。
import threading
import time
from datetime import datetime
class HeartbeatMonitor:
def __init__(self, max_silence_seconds: int = 300):
self.last_heartbeat = datetime.utcnow()
self.max_silence = max_silence_seconds
self.lock = threading.Lock()
def beat(self):
with self.lock:
self.last_heartbeat = datetime.utcnow()
def is_alive(self) -> bool:
with self.lock:
silence = (datetime.utcnow() - self.last_heartbeat).total_seconds()
return silence < self.max_silence
def silence_seconds(self) -> float:
with self.lock:
return (datetime.utcnow() - self.last_heartbeat).total_seconds()
monitor = HeartbeatMonitor(max_silence_seconds=60)
def agent_with_heartbeat():
while not shutdown_flag.is_set():
# Send heartbeat at start of each cycle
monitor.beat()
# Do agent work
perform_agent_cycle()
time.sleep(30)
# External watchdog checks the monitor
def watchdog_check():
while True:
if not monitor.is_alive():
print(f'ALERT: Agent silent for {monitor.silence_seconds():.0f}s')
# Restart agent here
time.sleep(30)
print('Heartbeat monitor defined')グレースフルシャットダウン
グレースフルシャットダウンでは、停止する前に進行中の処理を完了します。停止信号を受け取ると、新しい処理の開始を防ぎ、現在のタスクを完了し、状態を保存してから正常に終了します。
import signal
import threading
from contextlib import contextmanager
class GracefulShutdown:
def __init__(self, timeout: float = 30.0):
self.should_stop = threading.Event()
self.active_tasks = 0
self.lock = threading.Lock()
self.timeout = timeout
signal.signal(signal.SIGTERM, self._handle_signal)
signal.signal(signal.SIGINT, self._handle_signal)
def _handle_signal(self, signum, frame):
print(f'Shutdown signal received. Waiting for {self.active_tasks} active tasks...')
self.should_stop.set()
@contextmanager
def task(self):
if self.should_stop.is_set():
raise RuntimeError('Shutdown in progress, not accepting new tasks')
with self.lock:
self.active_tasks += 1
try:
yield
finally:
with self.lock:
self.active_tasks -= 1
def wait_for_all_tasks(self):
self.should_stop.wait()
deadline = time.time() + self.timeout
while self.active_tasks > 0 and time.time() < deadline:
time.sleep(0.1)
if self.active_tasks > 0:
print(f'WARNING: Forced shutdown with {self.active_tasks} tasks still active')
shutdown = GracefulShutdown(timeout=30)
print('Graceful shutdown manager created')切断と再接続への対応
常時稼働エージェントには、サービスの切断に対応する戦略が必要です。停止中はイベントをバッファーに蓄積し、再接続後に取りこぼしたイベントを再生し、切断中に到着したイベントが失われないようにします。
import asyncio
import redis
from datetime import datetime
r = redis.Redis(host='localhost', port=6379, decode_responses=True)
class DisconnectHandler:
def __init__(self, buffer_key: str = 'agent:offline_buffer'):
self.buffer_key = buffer_key
self.connected = True
def on_disconnect(self):
self.connected = False
print(f'Disconnected at {datetime.utcnow()}')
def on_reconnect(self):
self.connected = True
print(f'Reconnected at {datetime.utcnow()}')
self.replay_buffered_events()
def handle_event(self, event: dict):
if not self.connected:
# Buffer events for later replay
import json
r.lpush(self.buffer_key, json.dumps(event))
print(f'Event buffered (offline): {event["type"]}')
return
self.process_event(event)
def replay_buffered_events(self):
import json
replayed = 0
while True:
raw = r.rpop(self.buffer_key)
if not raw:
break
event = json.loads(raw)
self.process_event(event)
replayed += 1
if replayed:
print(f'Replayed {replayed} buffered events')
def process_event(self, event: dict):
print(f'Processing event: {event["type"]}')
handler = DisconnectHandler()
print('Disconnect handler created')systemdによるプロセス監視
Linuxの本番環境へのデプロイでは、systemdを使用してエージェントプロセスを管理します。systemdは自動再起動、journaldへのログ記録、起動時のエージェント開始を処理します。
# /etc/systemd/system/my-agent.service
SERVICE_FILE = '''
[Unit]
Description=My AI Agent Service
After=network.target
[Service]
Type=simple
User=ubuntu
WorkingDirectory=/home/ubuntu/agent
ExecStart=/home/ubuntu/venv/bin/python agent.py
Restart=always
RestartSec=10
StandardOutput=journal
StandardError=journal
# Environment variables
EnvironmentFile=/home/ubuntu/agent/.env
# Resource limits
MemoryLimit=1G
CPUQuota=50%
[Install]
WantedBy=multi-user.target
'''
# Deploy commands:
# sudo cp my-agent.service /etc/systemd/system/
# sudo systemctl daemon-reload
# sudo systemctl enable my-agent
# sudo systemctl start my-agent
# sudo systemctl status my-agent
# sudo journalctl -u my-agent -f # Follow logs
print('Systemd service configuration defined')
print('Enables: auto-start on boot, auto-restart on crash, centralized logging')再起動後も維持する状態の永続化
常時稼働エージェントは、再起動後に中断した場所から処理を再開できるよう、状態を保存する必要があります。チェックポイントデータを一定間隔で、また重要な状態変更の後にディスクまたはRedisへ保存します。
import json
import os
from datetime import datetime
CHECKPOINT_FILE = '/tmp/agent_checkpoint.json'
def save_checkpoint(state: dict):
state['last_saved'] = datetime.utcnow().isoformat()
with open(CHECKPOINT_FILE, 'w') as f:
json.dump(state, f, indent=2)
print(f'Checkpoint saved at {state["last_saved"]}')
def load_checkpoint() -> dict:
if not os.path.exists(CHECKPOINT_FILE):
print('No checkpoint found, starting fresh')
return {}
with open(CHECKPOINT_FILE) as f:
state = json.load(f)
print(f'Checkpoint loaded from {state.get("last_saved", "unknown")}')
return state
# Agent startup
agent_state = load_checkpoint()
last_processed_id = agent_state.get('last_processed_email_id', 0)
print(f'Resuming from email ID: {last_processed_id}')
# After processing each email
agent_state['last_processed_email_id'] = last_processed_id + 1
if agent_state['last_processed_email_id'] % 10 == 0: # Checkpoint every 10 items
save_checkpoint(agent_state)常時稼働エージェントの監視
常時稼働エージェントについて、稼働時間、1時間あたりの処理イベント数、エラー率、メモリ使用量、最終アクティブ時刻などの主要なメトリクスを追跡します。これらはヘルスエンドポイントで公開するか、監視サービスへ送信します。
from fastapi import FastAPI
from datetime import datetime
import psutil
import os
app = FastAPI()
start_time = datetime.utcnow()
events_processed = 0
last_event_time = None
@app.get('/health')
def health_check():
process = psutil.Process(os.getpid())
uptime_seconds = (datetime.utcnow() - start_time).total_seconds()
last_active = None
if last_event_time:
last_active = (datetime.utcnow() - last_event_time).total_seconds()
return {
'status': 'ok',
'uptime_seconds': round(uptime_seconds),
'events_processed': events_processed,
'memory_mb': round(process.memory_info().rss / 1024 / 1024, 1),
'cpu_percent': process.cpu_percent(interval=1),
'last_event_seconds_ago': round(last_active) if last_active else None,
'timestamp': datetime.utcnow().isoformat()
}外部依存関係向けサーキットブレーカー
常時稼働エージェントは、障害が発生する可能性のある外部サービスとやり取りします。サーキットブレーカーは、障害が発生しているサービスへの呼び出しを一定期間停止し、連鎖障害を防ぎます。
import time
from enum import Enum
class CircuitState(Enum):
CLOSED = 'closed' # Normal operation
OPEN = 'open' # Service down, not calling
HALF_OPEN = 'half_open' # Testing if service recovered
class CircuitBreaker:
def __init__(self, failure_threshold=5, recovery_timeout=60):
self.state = CircuitState.CLOSED
self.failure_count = 0
self.threshold = failure_threshold
self.recovery_timeout = recovery_timeout
self.last_failure_time = None
def call(self, fn, *args):
if self.state == CircuitState.OPEN:
if time.time() - self.last_failure_time > self.recovery_timeout:
self.state = CircuitState.HALF_OPEN
else:
raise RuntimeError('Circuit open: service unavailable')
try:
result = fn(*args)
if self.state == CircuitState.HALF_OPEN:
self.state = CircuitState.CLOSED
self.failure_count = 0
print('Circuit closed: service recovered')
return result
except Exception as e:
self.failure_count += 1
self.last_failure_time = time.time()
if self.failure_count >= self.threshold:
self.state = CircuitState.OPEN
print(f'Circuit opened after {self.failure_count} failures')
raise
circuit = CircuitBreaker(failure_threshold=3, recovery_timeout=30)
print('Circuit breaker created')理解度チェック:常時稼働エージェント
常時稼働エージェントの設計パターンについて理解度を確認します。
常時稼働エージェントの設計パターンのまとめ
信頼性の高い常時稼働エージェントは、正常なシャットダウンのためのシグナル処理を備えたデーモンプロセス、クラッシュ時の自動再起動を行うウォッチドッグ、指数バックオフによる再接続機能を備えた永続 WebSocket 接続、停止したプロセスを検知するハートビートチェック、処理中の作業を完了するグレースフルシャットダウン、再起動後に処理を再開するための状態チェックポイント保存を組み合わせます。本番環境でのプロセス管理には、systemdまたはsupervisordを使用します。
よくある質問
「常時稼働エージェントの設計パターン」レッスンは無料ですか?
はい。「常時稼働エージェントの設計パターン」の完全なテキストはこのウェブで無料で読めます。インタラクティブに演習し(組み込みコードエディタと24時間対応のAIチューター)、AI Agentsコースの残りをアンロックするには、CoddyKit PROにアップグレードしてください。 AI Agentsコースには全4レッスンが含まれています。
「常時稼働エージェントの設計パターン」で何を学びますか?
バックグラウンドプロセス、デーモンエージェント、永続的な接続管理を学びます。 ブラウザで直接実行するハンズオンコードでAI Agentsを演習し、24時間対応のAIチューターがレッスンを進める中での質問に答えます。
AI Agentsを始めるのに経験は必要ですか?
事前経験は必要ありません。CoddyKitのAI Agentsは初級者から上級者向けに構成されているため、ここから始めるか最初から始めて、自分のペースで進むことができます。 これはレッスン1/4です。
「常時稼働エージェントの設計パターン」レッスンにはどのくらい時間がかかりますか?
ほとんどのCoddyKitレッスンは約5~10分かかります。各レッスンはコンパクトでインタラクティブなので、着実に進歩し、ウェブとアプリ全体で正確に前回の場所から再開できます。
このAI Agentsレッスンでコードを書いて実行できますか?
はい。すべてのAI Agentsレッスンに組み込みコードエディタが含まれているため、ブラウザでリアルコードを書いて実行し、即座のAIフィードバックを取得できます。ローカル設定は不要です。
このコースのすべてのレッスン
- 常時稼働エージェントの設計パターン
- プロアクティブな通知・アラートシステム
- セッションをまたぐコンテキストの永続化
- 毎日のブリーフィングエージェントを構築する