フレームワーク
Django、Flask、FastAPI、Webhook エンドポイント、そして二重に送信しないバックグラウンドジョブ。
Django
キーは環境変数から読み込んで settings に置き、ビューが import する専用のモジュールで OpenEmail を 1 つ作成する。クライアントはスレッド間で安全に共有できるので、1 つのインスタンスがすべてのリクエストを処理し、プロセスにつき 1 つのコネクションプールを保つ。モジュール名は openemail.py 以外にすること。その名前だとパッケージが隠れてしまうことがある。
import os OPENEMAIL_API_KEY = os.environ['OPENEMAIL_API_KEY']OPENEMAIL_TIMEOUT = 10.0OPENEMAIL_SENDER = 'Acme <[email protected]>'from django.conf import settingsfrom openemail import OpenEmail mailer = OpenEmail(settings.OPENEMAIL_API_KEY, timeout=settings.OPENEMAIL_TIMEOUT)from django.conf import settingsfrom django.http import HttpRequest, JsonResponsefrom django.views.decorators.http import require_POSTfrom openemail import OpenEmailApiError from acme.mailer import mailer @require_POSTdef invite(request: HttpRequest) -> JsonResponse: try: sent = mailer.emails.send({ 'from': settings.OPENEMAIL_SENDER, 'to': request.POST['email'], 'subject': 'You are invited to Acme', 'text': 'Accept the invitation to join the workspace.', }) except OpenEmailApiError as error: if error.is_validation: return JsonResponse({'error': error.message, 'field': error.param}, status=422) raise return JsonResponse({'id': sent['id']}, status=202)同梱の openemail クライアントもここで使える。アプリ設定の ready() メソッドから init(settings.OPENEMAIL_API_KEY) を一度だけ呼び、送信する場所ではどこでも openemail を import すればよい。
Flask
アプリケーションファクトリーがアプリと一緒にクライアントを作成して app.extensions に保持し、ビューは current_app を通じてそれにアクセスする。クライアントも受け取れるファクトリーにしておけば、テストから httpx.MockTransport で作ったクライアントを渡せる。
import os from flask import Flask, current_app, requestfrom openemail import OpenEmail def create_app(mailer: OpenEmail | None = None) -> Flask: app = Flask(__name__) app.extensions['openemail'] = mailer or OpenEmail(os.environ['OPENEMAIL_API_KEY'], timeout=10) @app.post('/invites') def invite() -> tuple[dict[str, str], int]: client: OpenEmail = current_app.extensions['openemail'] sent = client.emails.send({ 'from': 'Acme <[email protected]>', 'to': request.form['email'], 'subject': 'You are invited to Acme', 'text': 'Accept the invitation to join the workspace.', }) return {'id': sent['id']}, 202 return appFastAPI
lifespan の中で AsyncOpenEmail を 1 つ作成する。こうすると、クライアントはリクエストを処理するイベントループに属し、サーバーが停止すると閉じられる。ルートには依存関係を使って渡す。
from collections.abc import AsyncIteratorfrom contextlib import asynccontextmanagerfrom typing import Annotated from fastapi import Body, Depends, FastAPI, Requestfrom openemail import AsyncOpenEmail @asynccontextmanagerasync def lifespan(app: FastAPI) -> AsyncIterator[None]: async with AsyncOpenEmail(timeout=10) as mailer: app.state.openemail = mailer yield app = FastAPI(lifespan=lifespan) def get_openemail(request: Request) -> AsyncOpenEmail: mailer: AsyncOpenEmail = request.app.state.openemail return mailer Mailer = Annotated[AsyncOpenEmail, Depends(get_openemail)] @app.post('/invites', status_code=202)async def invite(email: Annotated[str, Body(embed=True)], mailer: Mailer) -> dict[str, str]: sent = await mailer.emails.send({ 'from': 'Acme <[email protected]>', 'to': email, 'subject': 'You are invited to Acme', 'text': 'Accept the invitation to join the workspace.', }) return {'id': sent['id'], 'status': sent['status']}テストでは app.dependency_overrides[get_openemail] を通じてクライアントを差し替えるので、どのルートも API に届かない。
Webhook エンドポイント
すべての配信は、処理する前に生のボディとリクエストヘッダーで検証すること。FastAPI では await request.body() と request.headers、Django では request.body と request.headers、Flask では request.get_data() と request.headers を使う。ヘッダーの検索は大文字と小文字を区別しないので、各フレームワーク自身のヘッダーオブジェクトをそのまま使える。
import os from fastapi import FastAPI, HTTPException, Request, Responsefrom openemail import WEBHOOK_EVENTS, WebhookVerificationError, verify_webhook_signature from acme.jobs import handle_bounce app = FastAPI() @app.post('/webhooks/openemail', status_code=204)async def openemail_webhook(request: Request) -> Response: try: event = verify_webhook_signature( payload=await request.body(), headers=request.headers, secret=os.environ['OPENEMAIL_WEBHOOK_SECRET'], ) except WebhookVerificationError as error: raise HTTPException(status_code=400, detail='bad signature') from error if event['type'] == WEBHOOK_EVENTS.EMAIL_BOUNCED: handle_bounce.delay(event['id'], event['data']) return Response(status_code=204)import os from django.http import HttpRequest, HttpResponsefrom django.views.decorators.csrf import csrf_exemptfrom django.views.decorators.http import require_POSTfrom openemail import WebhookVerificationError, verify_webhook_signature from acme.models import ReceivedEvent @csrf_exempt@require_POSTdef openemail_webhook(request: HttpRequest) -> HttpResponse: try: event = verify_webhook_signature( payload=request.body, headers=request.headers, secret=os.environ['OPENEMAIL_WEBHOOK_SECRET'], ) except WebhookVerificationError: return HttpResponse('bad signature', status=400) ReceivedEvent.objects.get_or_create( id=event['id'], defaults={'type': event['type'], 'data': event['data']}, ) return HttpResponse(status=204)Django は CSRF トークンのない POST を拒否するが、配信には CSRF トークンが付いていないので、ビューは csrf_exempt にする。リクエストが OpenEmail から来たことを証明するのは署名である。request.META はキーの名前が HTTP_X_OPENEMAIL_SIGNATURE の形式に変えられているので、代わりに request.headers を読むこと。
すぐに 2xx で応答し、処理はその後で行うこと。応答のない配信や、408、425、429、5xx が返った配信は、約 27 時間半のうちに最大 8 回まで再試行される。また再送では同じ id のイベントがもう一度送られるので、処理済みの id を保存しておき、重複はスキップすること。
バックグラウンドジョブ
ジョブキューは失敗したタスクを再実行するが、タスクはメールが送り出された後に失敗することもある。応答が失われた場合や、ワーカーが完了前に停止した場合である。送信が必要になった原因から導出した idempotency_key= を渡すこと。そうすればタスクのどの実行も同じキーを持つので、繰り返しは 2 通目を送るのではなく元のメッセージを再生する。
from celery import Task, shared_taskfrom openemail import OpenEmail, OpenEmailApiError, OpenEmailNetworkError mailer = OpenEmail() @shared_task(bind=True, acks_late=True, max_retries=5)def send_receipt(self: Task, order_id: str, email: str) -> str: try: sent = mailer.emails.send( { 'from': 'Acme Billing <[email protected]>', 'to': email, 'subject': f'Receipt for order {order_id}', 'text': f'Thank you for order {order_id}.', }, idempotency_key=f'receipt:{order_id}', ) except OpenEmailNetworkError as error: raise self.retry(exc=error, countdown=30) except OpenEmailApiError as error: if error.is_server_error: raise self.retry(exc=error, countdown=30) raise return sent['id']acks_late=True を指定すると、Celery はタスクの実行が終わってから初めてそれを確認応答する。そのため、ワーカーの停止で中断されたタスクは再度配送されることがあるが、キーがその 2 回目の実行を再生に変えるので、ここでは安全である。RQ の Retry は失敗したジョブを再実行し、同じ効果がある。
from redis import Redisfrom rq import Queue, Retry from openemail import OpenEmail mailer = OpenEmail() def send_receipt(order_id: str, email: str) -> str: sent = mailer.emails.send( { 'from': 'Acme Billing <[email protected]>', 'to': email, 'subject': f'Receipt for order {order_id}', 'text': f'Thank you for order {order_id}.', }, idempotency_key=f'receipt:{order_id}', ) return sent['id'] queue = Queue(connection=Redis())queue.enqueue(send_receipt, 'AC-4192', '[email protected]', retry=Retry(max=5, interval=[10, 60, 300]))キーは、その送信が必要になった原因から導出すること。決して時計から作ってはならない。キーは英字、数字、アンダースコア、ドット、コロン、ハイフンからなる 1〜255 文字なので、メールアドレスではなく id から組み立てること。ボディも同様にタスクの引数だけから組み立てること。同じキーで異なるボディを送った繰り返しは、再生されるのではなく 422 idempotency_key_reuse で拒否される。
クライアントはモジュールレベルで作成すること。最初のリクエストまで接続を開かないので、親からフォークされた各ワーカープロセスはそれぞれ自分の接続を開く。