非同期
`AsyncOpenEmail` は `OpenEmail` のすべてのメソッドを備え、asyncio または trio の上で await して使う。
非同期クライアント
AsyncOpenEmail は OpenEmail が持つすべてのメソッドを同じ引数と同じ戻り値の型で備えており、それぞれが await するコルーチンである。同じキーワード引数で作成し、省略したものは同じ環境変数から読み、同じエラーを送出する。
import asyncio from openemail import AsyncOpenEmail async def main() -> None: async with AsyncOpenEmail() as client: sent = await client.emails.send({ 'from': 'Acme Billing <[email protected]>', 'to': '[email protected]', 'subject': 'Your September invoice', 'text': 'Your invoice is attached.', }) email = await client.emails.get(sent['id']) print(email['status'], email['sentAt']) asyncio.run(main())async with は、ブロックが正常に終わっても例外で終わっても、終了時にコネクションプールを閉じる。プロセスと同じだけ生きるクライアントは、起動時に一度作成し、終了時に await client.aclose() で閉じる。プログラム全体でクライアントは 1 つあれば足りる。同じイベントループ上のコルーチンであれば、いくつでも同時にそれを使える。
既製の openemail クライアントと init() は同期型であり、それらの非同期版は存在しない。AsyncOpenEmail はプログラムの開始時に作成して必要なコードに渡すか、フレームワークのアプリケーション状態に保持すること。
パッケージのパリティチェックは 2 つのクライアントをメソッドごとに比較し、OpenEmail と AsyncOpenEmail で引数が異なるメソッドがあれば失敗する。そのため、各メソッドのページの説明は両方に当てはまる。
async for によるページング
list と list_all は、他のすべてのメソッドと同じく await で呼び出す。iterate はそうではない。すぐに非同期イテレーターを返し、async for はループがそのページに到達したときに各ページを取得するので、ループを抜ければリクエストも止まる。
import asyncio from openemail import AsyncOpenEmail async def main() -> None: async with AsyncOpenEmail() as client: page = await client.emails.list(status='failed', limit=50) print(len(page['items']), page['nextCursor']) complaints = await client.suppressions.list_all(reason='complaint') print(len(complaints)) async for thread in client.threads.iterate(folder='inbox'): print(thread['id']) asyncio.run(main())多数の呼び出しを同時に行う
1 つのクライアントで、いくつでもリクエストを同時に扱える。asyncio.gather はそれらをまとめて開始し、セマフォは実行中のリクエスト数を自分で選んだ数に抑える。
import asyncio from openemail import AsyncOpenEmailfrom openemail.types import SentEmailResource async def main() -> None: recipients = ['[email protected]', '[email protected]', '[email protected]'] gate = asyncio.Semaphore(8) async with AsyncOpenEmail() as client: async def welcome(address: str) -> SentEmailResource: async with gate: return await client.emails.send({ 'from': 'Acme <[email protected]>', 'to': address, 'subject': 'Welcome to Acme', 'text': 'Your workspace is ready.', }) results = await asyncio.gather( *(welcome(address) for address in recipients), return_exceptions=True, ) for address, result in zip(recipients, results): if isinstance(result, BaseException): print(address, 'failed:', result) else: print(address, result['status']) asyncio.run(main())現在の API は、通常の読み取りと書き込みの頻度に全般的な上限を設けていないので、集中したリクエストを代わりに減速させてくれるものは何もない。こうした上限は今後追加される可能性がある。API が実際に数えているものは 429 を返す。ワークスペースの月間送信枠、1 日あたりの AI アクション、1 時間あたり 500 件のファイルアップロードなどである。いずれも Retry-After を付けないので、クライアントはリトライせず、is_rate_limited が true の OpenEmailApiError を即座に送出する。return_exceptions=True を指定すると、gather は最初の拒否で例外を送出せず、すべての結果を返す(拒否された呼び出しの結果はその例外になる)。ループはその一つひとつを読む。
タスクをキャンセルするとそのリクエストもキャンセルされるが、送信中にキャンセルされた送信はすでに API に届いている可能性がある。キャンセルしてから繰り返すかもしれない送信には、独自の idempotency_key= を付けること。そうすれば、繰り返しは 2 通目のメッセージを送るのではなく、1 回目の送信を再生する。
asyncio と trio
クライアントは待機とタイムアウトを anyio で、送信を httpx で行い、どちらもいずれのイベントループでも動くので、同じ main() が asyncio.run(main()) でも trio.run(main) でも動作する。trio はパッケージの依存関係ではないので、使う場合は自分でインストールすること。
上の例のような asyncio 自身の gather と Semaphore は asyncio でしか動かない。両方で動かす必要があるコードには、パッケージがすでに依存している anyio の create_task_group と Semaphore を使うこと。
非同期のアクセストークン
本人が OAuth で接続したアプリは、API キーではなくアクセストークンを持ち、それを access_token= として渡す。トークンそのものか、トークンを返す関数のどちらかである。関数は毎回のリクエストの前に実行されるので、期限が近づいたらトークンを更新でき、クライアントを作り直す必要はない。AsyncOpenEmail では async 関数にでき、クライアントはその戻り値を await する。
import asyncioimport timefrom dataclasses import dataclass from openemail import AsyncOpenEmail from acme.auth import refresh_access_token @dataclassclass CachedToken: value: str = '' expires_at: float = 0.0 cached = CachedToken() async def access_token() -> str: if cached.expires_at - time.time() < 60: cached.value, lifetime = await refresh_access_token() cached.expires_at = time.time() + lifetime return cached.value async def main() -> None: async with AsyncOpenEmail(access_token=access_token) as client: me = await client.me.get() print(me['object']) asyncio.run(main())すべてのリクエストがこの関数を待つので、関数は軽く保つこと。上の例のように、キャッシュしたトークンを返し、期限が近いときだけ更新すればよい。OpenEmail はコルーチンを待てないので、async 関数を渡すと最初のリクエストで ValueError を送出する。
使い捨て受信トレイ
create_async_temp_mail() は create_temp_mail() の非同期版である。API キーは持たない。create と list_domains は認証情報をまったく送らず、それ以外のすべてのメソッドは create が返したトークンを inbox_token= として受け取る。1 つの受信箱に結び付いたクライアントにするには、create_async_temp_mail 自体に inbox_token= を渡すこと。
import asyncio import httpxfrom openemail import create_async_temp_mail async def main() -> None: async with httpx.AsyncClient(follow_redirects=True) as http: temp = create_async_temp_mail(http_client=http) inbox = await temp.create({'ttlMinutes': 60}) print(inbox['address'], inbox['expiresAt']) async for message in temp.iterate_messages(inbox['id'], inbox_token=inbox['token']): print(message['from']['email'], message['subject']) asyncio.run(main())これが返すクライアントには独自の aclose() がない。作業が終わったときに接続を閉じるには、上の例のように async with で httpx.AsyncClient を開き、それを http_client= として渡すこと。
独自の httpx クライアント
http_client= は、プロキシ、接続数の上限、独自の証明書、テストでのモックトランスポートのために自分で作った httpx.AsyncClient を受け取る。httpx.Client は OpenEmail 用なので、渡すと TypeError が送出される。
import asyncio import httpxfrom openemail import AsyncOpenEmail async def main() -> None: async with httpx.AsyncClient( proxy='http://proxy.internal:3128', limits=httpx.Limits(max_connections=20), follow_redirects=True, ) as http: client = AsyncOpenEmail(http_client=http, timeout=20) page = await client.threads.list(folder='inbox', limit=10) print(len(page['items'])) asyncio.run(main())渡したクライアントは引き続き自分の管理下にある。aclose() と async with の終了が閉じるのは SDK が開いたプールだけなので、httpx.AsyncClient は自分で閉じること。ここではそれ自身の async with で閉じている。SDK が開くプールはリダイレクトに従うので、それに合わせて自分のクライアントにも follow_redirects=True を設定すること。クライアントや個々の呼び出しに指定した timeout= は、httpx.AsyncClient がどんなタイムアウトを持っていても、引き続きすべての試行の上限となる。
テストでは、httpx.MockTransport が自分で用意した関数からすべてのリクエストに応答するので、ネットワークには何も届かない。
import asyncio import httpxfrom openemail import AsyncOpenEmail def answer(request: httpx.Request) -> httpx.Response: return httpx.Response(200, json={'object': 'list', 'data': [], 'hasMore': False, 'nextCursor': None}) async def main() -> None: async with httpx.AsyncClient(transport=httpx.MockTransport(answer)) as http: client = AsyncOpenEmail('oe_test_fixture', http_client=http) page = await client.suppressions.list() assert page['items'] == [] asyncio.run(main())