概要
asynccontextmanager + asyncio.timeout で分散ロックを宣言的に
@asynccontextmanager は yield の前後に「取得・解放」を書くことで、async with distributed_lock(redis, product_id): という直感的な API を実現する。asyncio.timeout() は Python 3.11+ の宣言的タイムアウト API で、asyncio.wait_for を置き換え、TaskGroup 内でも自然に動作する。finally でロック解放を保証するため、例外発生時もデッドロックにならない。
model_validator(mode="after") で cross-field バリデーション
Pydantic v2 の field_validator は単一フィールドの検証だが、model_validator(mode="after") は全フィールドが設定された後に複数フィールドを参照したバリデーションを行う。「qty > 0 かつ product_id が PROD- 形式」のような依存関係をここで表現し、エラーメッセージに product_id を含めることで運用時のデバッグを容易にする。
Workload Identity Federation — SA キーレス設計
GKE Pod が SA キー JSON を使う設計は「キーローテーション負債・漏洩リスク・Secret 管理コスト」が発生する。Workload Identity では KSA の iam.gke.io/gcp-service-account Annotation + GSA の workloadIdentityUser バインディングの2つを設定するだけで、Pod が自動的に GSA として振る舞える。キーファイルは存在しないのでローテーション不要。
Cloud Armor evaluatePreconfiguredExpr で OWASP WAF を宣言的に
google_compute_security_policy の evaluatePreconfiguredExpr('sqli-v33-stable') は Google が管理する OWASP CRS SQLi ルールセットを適用する。-stable は本番向けに false positive が低い設定。Cloud Run v2 は直接 Cloud Armor を受け付けず、Cloud Run → NEG → Backend Service → GLB 経由でポリシーをアタッチする。
問題 A: コーディング — asynccontextmanager × asyncio.timeout × Pydantic v2 model_validator(分散ロック付き在庫予約 Bad→Good)
以下の「悪いコード」は、ECサイト MOps チームの在庫予約バッチです。問題点を全て洗い出し、asynccontextmanager・asyncio.timeout(Python 3.11+ 新 API)・Pydantic v2 model_validator(mode="after")・frozen dataclass・redis.asyncio を使って Bad→Good にリファクタリングしてください。
制約・前提条件
- Python 3.12+、asyncio、redis.asyncio、Pydantic v2 を使うこと
@asynccontextmanagerで分散ロックをコンテキストマネージャとして実装することasyncio.timeout(LOCK_ACQUIRE_TIMEOUT_SECONDS)でロック取得タイムアウトを制御すること- Pydantic v2 の
model_validator(mode="after")で qty > 0 かつ product_id 形式を cross-field バリデーションすること - 予約結果を
frozen dataclassのReservationResultで返すこと - Google スタイル docstring・インラインコメント・名前付き定数を含めること
悪いコード (Before) — カテゴリ A
import asyncio
import redis # 問題6: 同期 redis-py を asyncio 内で使用
async def reserve_inventory(items, redis_client):
results = []
for item in items:
# 問題1: TTL なし・ロック解放なし(デッドロックリスク)
lock_key = "lock:" + item["product_id"]
redis_client.set(lock_key, "1")
# 問題2: タイムアウトなし(Redis ダウン時に無限待機)
stock = int(redis_client.get("stock:" + item["product_id"]) or 0)
# 問題3: バリデーションなし(マイナス数量・型エラーが素通り)
if stock >= item["qty"]:
redis_client.decrby("stock:" + item["product_id"], item["qty"])
results.append({"product_id": item["product_id"], "reserved": True})
else:
results.append({"product_id": item["product_id"], "reserved": False})
# 問題4: ロック解放が finally ブロックにない(例外時にロック残留)
redis_client.delete(lock_key)
# 問題5: dict で返す(型安全性なし)
return results
# 問題7: マジックナンバー(TTL 値・タイムアウト値が即値)
redis.set(lock_key, "1") は TTL 無しで永続化。例外時に finally がないためロックが残留しデッドロックになる。ex=LOCK_TTL_SECONDS, nx=True + finally: await redis.delete() で修正asyncio.timeout(LOCK_ACQUIRE_TIMEOUT_SECONDS) で取得タイムアウトを設定item["qty"] が負数・ゼロ・文字列でも素通り。Pydantic v2 model_validator(mode="after") で cross-field バリデーション@asynccontextmanager + finally で解放を保証{"product_id": ..., "reserved": True} は型安全でない。frozen dataclass ReservationResult で型安全にredis.get() はイベントループをブロックする。redis.asyncio に切り替えるLOCK_TTL_SECONDS / LOCK_ACQUIRE_TIMEOUT_SECONDS に切り出すヒント A(段階的開示)
ヒント1 — 方向性
asynccontextmanager は yield の前後に「取得・解放」を書くことで、async with distributed_lock(): という直感的な API を実現する。asyncio.timeout() は Python 3.11+ の宣言的タイムアウト API で、asyncio.wait_for を置き換える。model_validator(mode="after") は Pydantic v2 で全フィールド設定後にフィールド間依存バリデーションを行う。frozen dataclass ReservationResult は予約結果の型安全な値オブジェクト。
ヒント2 — アプローチ
LOCK_TTL_SECONDS: Final[int] = 30/LOCK_ACQUIRE_TIMEOUT_SECONDS: Final[float] = 5.0で定数化@asynccontextmanagerの中でawait redis.set(lock_key, "1", ex=LOCK_TTL_SECONDS, nx=True)→yield→finally: await redis.delete(lock_key)async with asyncio.timeout(LOCK_ACQUIRE_TIMEOUT_SECONDS):でロック取得ループを囲むmodel_validator(mode="after")でqty > 0・product_id.startswith("PROD-")をチェック@dataclass(frozen=True) ReservationResultにremaining_stock: intとsummary: strプロパティを追加
ヒント3 — コードの骨格
LOCK_TTL_SECONDS: Final[int] = 30
LOCK_ACQUIRE_TIMEOUT_SECONDS: Final[float] = 5.0
LOCK_POLL_INTERVAL_SECONDS: Final[float] = 0.05
@asynccontextmanager
async def distributed_lock(redis, product_id: str) -> AsyncGenerator[None, None]:
lock_key = f"lock:{product_id}"
acquired = False
try:
async with asyncio.timeout(LOCK_ACQUIRE_TIMEOUT_SECONDS):
while not acquired:
acquired = bool(
await redis.set(lock_key, "1", ex=LOCK_TTL_SECONDS, nx=True)
)
if not acquired:
await asyncio.sleep(LOCK_POLL_INTERVAL_SECONDS)
yield
finally:
if acquired:
await redis.delete(lock_key)
class ReservationRequest(BaseModel):
product_id: str
qty: int
@model_validator(mode="after")
def validate_request(self) -> "ReservationRequest":
if self.qty <= 0:
raise ValueError(f"qty は 1 以上必要: {self.qty}")
if not self.product_id.startswith("PROD-"):
raise ValueError(f"product_id は 'PROD-' で始まる必要があります: {self.product_id}")
return self
問題点分析 — カテゴリ A
| # | 問題点 | 分類 | 改善方法 |
|---|---|---|---|
| 1 | ロック TTL なし・解放なし | エラー処理 Ch10 | ex=TTL, nx=True + finally: delete() |
| 2 | タイムアウトなし | エラー処理 Ch10 | asyncio.timeout() |
| 3 | バリデーションなし | 型の活用 Ch2 | Pydantic v2 model_validator(mode="after") |
| 4 | 例外時ロック残留 | エラー処理 Ch10 | @asynccontextmanager + finally |
| 5 | dict で返す | 型の活用 Ch4 | frozen dataclass ReservationResult |
| 6 | 同期 Redis をイベントループ内で使用 | パフォーマンス | redis.asyncio に切り替え |
| 7 | マジックナンバー | 可読性 Ch7 | LOCK_TTL_SECONDS / LOCK_ACQUIRE_TIMEOUT_SECONDS |
模範解答 A
import asyncio
import redis
async def reserve_inventory(items, redis_client):
results = []
for item in items:
lock_key = "lock:" + item["product_id"]
redis_client.set(lock_key, "1") # TTL なし
stock = int(redis_client.get("stock:" + item["product_id"]) or 0)
if stock >= item["qty"]: # qty バリデーションなし
redis_client.decrby("stock:" + item["product_id"], item["qty"])
results.append({"product_id": item["product_id"], "reserved": True})
else:
results.append({"product_id": item["product_id"], "reserved": False})
redis_client.delete(lock_key) # finally なし → 例外時ロック残留
return results # dict 返し → 型安全性なし
"""inventory_reserve.py — 分散ロック付き在庫予約バッチ。
Ch2: 型の活用(Pydantic v2 ReservationRequest + model_validator)
Ch4: コレクション(frozen dataclass ReservationResult)
Ch7: 名前付き定数(LOCK_TTL_SECONDS / LOCK_ACQUIRE_TIMEOUT_SECONDS)
Ch10: エラー処理(asyncio.timeout × finally ロック解放保証)
Ch11: テスト容易性(asynccontextmanager で分散ロックを分離)
"""
from __future__ import annotations
import asyncio
import logging
from collections.abc import AsyncGenerator
from contextlib import asynccontextmanager
from dataclasses import dataclass
from typing import Final
import redis.asyncio as aioredis
from pydantic import BaseModel, field_validator, model_validator
logger = logging.getLogger(__name__)
# 名前付き定数(マジックナンバー禁止)
LOCK_TTL_SECONDS: Final[int] = 30 # Redis ロックの TTL
LOCK_ACQUIRE_TIMEOUT_SECONDS: Final[float] = 5.0
LOCK_POLL_INTERVAL_SECONDS: Final[float] = 0.05
PRODUCT_ID_PREFIX: Final[str] = "PROD-"
class ReservationRequest(BaseModel):
"""在庫予約リクエストの値オブジェクト(Pydantic v2)。"""
product_id: str
qty: int
@field_validator("product_id")
@classmethod
def validate_product_id_format(cls, v: str) -> str:
if not v.startswith(PRODUCT_ID_PREFIX):
raise ValueError(f"product_id は '{PRODUCT_ID_PREFIX}' で始まる必要: {v}")
if not (8 <= len(v) <= 32):
raise ValueError(f"product_id の長さは 8〜32 文字: {len(v)}")
return v
@model_validator(mode="after")
def validate_reservation_logic(self) -> "ReservationRequest":
"""フィールド間依存バリデーション: qty は 1 以上必要。"""
if self.qty <= 0:
raise ValueError(
f"qty は 1 以上の整数が必要: {self.qty} "
f"(product_id={self.product_id})" # 関連フィールドも含めて報告
)
return self
@dataclass(frozen=True)
class ReservationResult:
"""在庫予約結果のイミュータブル値オブジェクト。"""
product_id: str
reserved: bool
qty: int
remaining_stock: int
@property
def summary(self) -> str:
status = "成功" if self.reserved else "在庫不足"
return (
f"[{status}] product_id={self.product_id} "
f"qty={self.qty} remaining={self.remaining_stock}"
)
@asynccontextmanager
async def distributed_lock(
redis: aioredis.Redis,
product_id: str,
) -> AsyncGenerator[None, None]:
"""Redis 分散ロックを Context Manager として提供する。
Raises:
asyncio.TimeoutError: タイムアウト以内にロック取得できない場合。
"""
lock_key = f"lock:{product_id}"
acquired = False
try:
# asyncio.timeout: Python 3.11+ の宣言的タイムアウト API
async with asyncio.timeout(LOCK_ACQUIRE_TIMEOUT_SECONDS):
while not acquired:
# SET key 1 EX {ttl} NX — アトミックなロック取得
acquired = bool(
await redis.set(lock_key, "1", ex=LOCK_TTL_SECONDS, nx=True)
)
if not acquired:
await asyncio.sleep(LOCK_POLL_INTERVAL_SECONDS)
logger.debug("分散ロック取得: key=%s", lock_key)
yield # ← コンテキストブロック内のコードが実行される
finally:
# 例外発生時・正常終了時の両方でロック解放(デッドロック防止)
if acquired:
await redis.delete(lock_key)
logger.debug("分散ロック解放: key=%s", lock_key)
async def reserve_inventory(
raw_requests: list[dict],
redis: aioredis.Redis,
) -> list[ReservationResult]:
"""在庫予約を実行し ReservationResult のリストで返す。"""
# Step 1: 入力バリデーション(不正データを早期に弾く)
requests: list[ReservationRequest] = []
for raw in raw_requests:
try:
requests.append(ReservationRequest.model_validate(raw))
except Exception as e:
logger.warning("バリデーションエラー: %s / 入力: %s", e, raw)
results: list[ReservationResult] = []
for req in requests:
try:
async with distributed_lock(redis, req.product_id): # 分散ロック取得
stock_key = f"stock:{req.product_id}"
raw_stock = await redis.get(stock_key) # 非同期 Redis(ループブロックなし)
current_stock = int(raw_stock or 0)
if current_stock >= req.qty:
remaining = await redis.decrby(stock_key, req.qty)
results.append(ReservationResult(
product_id=req.product_id, reserved=True,
qty=req.qty, remaining_stock=remaining,
))
else:
results.append(ReservationResult(
product_id=req.product_id, reserved=False,
qty=req.qty, remaining_stock=current_stock,
))
except asyncio.TimeoutError:
logger.error("ロック取得タイムアウト: product_id=%s", req.product_id)
return results
import asyncio
import redis.asyncio as aioredis
async def main():
r = aioredis.from_url("redis://localhost:6379")
await r.set("stock:PROD-001-AB", 10)
await r.set("stock:PROD-002-CD", 2)
requests = [
{"product_id": "PROD-001-AB", "qty": 3}, # 在庫 10 → 成功
{"product_id": "PROD-002-CD", "qty": 5}, # 在庫 2 → 在庫不足
{"product_id": "PROD-999-XX", "qty": -1}, # qty < 0 → バリデーションエラー(スキップ)
]
results = await reserve_inventory(requests, r)
for result in results:
print(result.summary)
# [成功] product_id=PROD-001-AB qty=3 remaining=7
# [在庫不足] product_id=PROD-002-CD qty=5 remaining=2
# frozen dataclass なので変更不可
# results[0].reserved = False → FrozenInstanceError
asyncio.run(main())
| ポイント | 適用した設計原則/パターン | 書籍対応章 |
|---|---|---|
@asynccontextmanager で分散ロックを分離 | テスト容易性・単一責任原則 | Ch11 |
asyncio.timeout() でタイムアウト制御 | エラー処理・Fail-Fast | Ch10 |
finally ブロックでロック解放保証 | エラー処理・リソース管理 | Ch10 |
model_validator(mode="after") cross-field バリデーション | 型の活用・入力検証 | Ch2 |
@dataclass(frozen=True) ReservationResult | 値オブジェクト・不変性 | Ch4 / Ch2 |
LOCK_TTL_SECONDS / LOCK_ACQUIRE_TIMEOUT_SECONDS | マジックナンバー排除 | Ch7 |
redis.asyncio で非同期 Redis 操作 | パフォーマンス・イベントループ設計 | Ch11 |
# tests/test_inventory_reserve.py
import asyncio
import pytest
from unittest.mock import AsyncMock, MagicMock
from inventory_reserve import (
reserve_inventory, distributed_lock,
ReservationResult, ReservationRequest,
LOCK_TTL_SECONDS,
)
class TestReservationRequest:
def test_valid_request(self):
req = ReservationRequest(product_id="PROD-001-AB", qty=5)
assert req.qty == 5
assert req.product_id == "PROD-001-AB"
def test_invalid_qty_zero(self):
with pytest.raises(Exception, match="qty は 1 以上"):
ReservationRequest(product_id="PROD-001-AB", qty=0)
def test_invalid_qty_negative(self):
with pytest.raises(Exception, match="qty は 1 以上"):
ReservationRequest(product_id="PROD-001-AB", qty=-3)
def test_invalid_product_id_prefix(self):
with pytest.raises(Exception, match="PROD-"):
ReservationRequest(product_id="ITEM-001-AB", qty=1)
class TestReservationResult:
def test_summary_reserved(self):
r = ReservationResult(product_id="PROD-001-AB", reserved=True, qty=3, remaining_stock=7)
assert "成功" in r.summary
assert "remaining=7" in r.summary
def test_summary_insufficient(self):
r = ReservationResult(product_id="PROD-002-CD", reserved=False, qty=5, remaining_stock=2)
assert "在庫不足" in r.summary
def test_frozen_immutability(self):
r = ReservationResult(product_id="PROD-001-AB", reserved=True, qty=3, remaining_stock=7)
with pytest.raises(Exception):
r.reserved = False # type: ignore[misc]
class TestReserveInventory:
@pytest.mark.asyncio
async def test_successful_reservation(self):
redis_mock = AsyncMock()
redis_mock.set = AsyncMock(return_value=True) # ロック取得成功
redis_mock.get = AsyncMock(return_value=b"10") # 在庫 10
redis_mock.decrby = AsyncMock(return_value=7) # 3 減算後 7
results = await reserve_inventory(
[{"product_id": "PROD-001-AB", "qty": 3}], redis_mock
)
assert len(results) == 1
assert results[0].reserved is True
assert results[0].remaining_stock == 7
@pytest.mark.asyncio
async def test_insufficient_stock(self):
redis_mock = AsyncMock()
redis_mock.set = AsyncMock(return_value=True)
redis_mock.get = AsyncMock(return_value=b"2") # 在庫 2 < qty 5
results = await reserve_inventory(
[{"product_id": "PROD-002-CD", "qty": 5}], redis_mock
)
assert results[0].reserved is False
assert results[0].remaining_stock == 2
問題 B: インフラ — GKE Workload Identity Federation × Cloud Run v2 最小権限 × Cloud Armor WAF(Terraform)
ECサイト MOps チームの GKE Autopilot クラスタと Cloud Run v2 に、以下の IAM・セキュリティ課題があります。
- Python バッチ Pod が SA キー JSON ファイル(
GOOGLE_APPLICATION_CREDENTIALS)を使っている(キー漏洩リスク・ローテーション負債) - Cloud Run v2 配信 API が
roles/bigquery.admin(過剰権限)で稼働している - GKE KSA と GSA の紐付けが Terraform 管理外で手動設定
- Cloud Run v2 が
--no-allow-unauthenticatedなのに Cloud Armor(WAF)未設定
要件
| # | 要件 |
|---|---|
| 1 | GKE Pod が SA キーなし(キーレス)で BigQuery・Pub/Sub・GCS にアクセスできるよう Workload Identity Federation を Terraform で設定すること |
| 2 | Cloud Run v2 の IAM を roles/bigquery.dataEditor + roles/bigquery.jobUser(最小権限)に変更し、roles/bigquery.admin を削除すること |
| 3 | Cloud Run v2 に Cloud Armor セキュリティポリシー(SQLインジェクション OWASP ルールセット)を Terraform でアタッチすること |
| 4 | GKE KSA と GSA の紐付けを Terraform の google_service_account_iam_member で管理すること |
ヒント B(段階的開示)
ヒント1 — 方向性
iam.gke.io/gcp-service-account Annotation を付け、GSA に workloadIdentityUser バインディングを追加」する2ステップ。Cloud Armor は google_compute_security_policy リソースで OWASP ルールセットを適用し、Cloud Run の NEG(Serverless Network Endpoint Group)経由でバックエンドサービスにアタッチする。roles/bigquery.admin を削除して dataEditor + jobUser に絞ると「テーブル読み書き + クエリ実行」のみに権限が限定される。
ヒント2 — Terraform リソース構成
google_service_account→ GKE バッチ用 GSA(mops-batch-sa)google_service_account_iam_member→roles/iam.workloadIdentityUserをserviceAccount:{project}.svc.id.goog[mops/mops-batch-ksa]にバインドkubernetes_service_account→iam.gke.io/gcp-service-accountAnnotation を付与google_project_iam_member→dataEditor + jobUser(admin は使わない)google_compute_security_policy→evaluatePreconfiguredExpr('sqli-v33-stable')+xss-v33-stable+ レートリミットgoogle_compute_region_network_endpoint_group→ Cloud Run v2 の SERVERLESS NEGgoogle_compute_backend_service→security_policyで Cloud Armor アタッチ
ヒント3 — Workload Identity + Cloud Armor の骨格
resource "google_service_account_iam_member" "workload_identity_binding" {
service_account_id = google_service_account.mops_batch_gsa.name
role = "roles/iam.workloadIdentityUser"
# 形式: serviceAccount:{project_id}.svc.id.goog[{namespace}/{ksa_name}]
member = "serviceAccount:${var.project_id}.svc.id.goog[mops/mops-batch-ksa]"
}
resource "kubernetes_service_account" "mops_batch_ksa" {
metadata {
name = "mops-batch-ksa"
namespace = "mops"
annotations = {
"iam.gke.io/gcp-service-account" = google_service_account.mops_batch_gsa.email
}
}
}
resource "google_compute_security_policy" "mops_waf_policy" {
name = "mops-waf-policy"
rule {
action = "deny(403)"
priority = 1000
match {
expr {
expression = "evaluatePreconfiguredExpr('sqli-v33-stable')"
}
}
}
rule {
action = "allow"
priority = 2147483647
match {
versioned_expr = "SRC_IPS_V1"
config { src_ip_ranges = ["*"] }
}
}
}
resource "google_compute_backend_service" "delivery_api_backend" {
security_policy = google_compute_security_policy.mops_waf_policy.id
backend {
group = google_compute_region_network_endpoint_group.delivery_api_neg.id
}
}
アーキテクチャ図 — GKE Workload Identity × Cloud Run v2 × Cloud Armor 設計
模範解答 B
# terraform/modules/mops-iam/main.tf
# GKE Workload Identity Federation × Cloud Run v2 最小権限 × Cloud Armor OWASP WAF
locals {
project_id = var.project_id
namespace = "mops"
ksa_name = "mops-batch-ksa"
}
# ── GKE バッチ用 GSA(Workload Identity で使用)────────────────────────────
resource "google_service_account" "mops_batch_gsa" {
account_id = "mops-batch-sa"
display_name = "MOps Batch GKE Workload Identity SA"
description = "GKE Pod が SA キーなし(キーレス)で GCP API にアクセスするための GSA"
project = local.project_id
}
# BigQuery・Pub/Sub・GCS の最小権限を付与(admin は付与しない)
resource "google_project_iam_member" "mops_batch_bq_data_editor" {
project = local.project_id
role = "roles/bigquery.dataEditor" # テーブルの読み書き
member = "serviceAccount:${google_service_account.mops_batch_gsa.email}"
}
resource "google_project_iam_member" "mops_batch_bq_job_user" {
project = local.project_id
role = "roles/bigquery.jobUser" # クエリジョブの実行(dataEditor とセット)
member = "serviceAccount:${google_service_account.mops_batch_gsa.email}"
}
resource "google_project_iam_member" "mops_batch_pubsub" {
project = local.project_id
role = "roles/pubsub.publisher" # メッセージ発行のみ
member = "serviceAccount:${google_service_account.mops_batch_gsa.email}"
}
resource "google_project_iam_member" "mops_batch_gcs" {
project = local.project_id
role = "roles/storage.objectCreator" # GCS 書き込みのみ(削除不可)
member = "serviceAccount:${google_service_account.mops_batch_gsa.email}"
}
# ── Workload Identity バインディング(KSA → GSA の紐付け)────────────────
# GKE Pod の KSA が GSA を impersonate できるようバインド
resource "google_service_account_iam_member" "workload_identity_binding" {
service_account_id = google_service_account.mops_batch_gsa.name
role = "roles/iam.workloadIdentityUser"
# 形式: serviceAccount:{project_id}.svc.id.goog[{namespace}/{ksa_name}]
member = "serviceAccount:${local.project_id}.svc.id.goog[${local.namespace}/${local.ksa_name}]"
}
# ── Kubernetes ServiceAccount(GKE 側)──────────────────────────────────────
resource "kubernetes_service_account" "mops_batch_ksa" {
metadata {
name = local.ksa_name
namespace = local.namespace
annotations = {
# この Annotation で KSA と GSA が紐付く(Workload Identity の核心)
"iam.gke.io/gcp-service-account" = google_service_account.mops_batch_gsa.email
}
}
depends_on = [google_service_account_iam_member.workload_identity_binding]
}
# ── Cloud Run v2 用 GSA(最小権限)───────────────────────────────────────────
resource "google_service_account" "delivery_api_gsa" {
account_id = "delivery-api-sa"
display_name = "MOps Delivery API Cloud Run SA"
description = "Cloud Run v2 専用 SA(roles/bigquery.admin は使わない)"
project = local.project_id
}
# roles/bigquery.admin を削除 → dataEditor + jobUser に絞る(最小権限原則)
resource "google_project_iam_member" "delivery_api_bq_data_editor" {
project = local.project_id
role = "roles/bigquery.dataEditor"
member = "serviceAccount:${google_service_account.delivery_api_gsa.email}"
}
resource "google_project_iam_member" "delivery_api_bq_job_user" {
project = local.project_id
role = "roles/bigquery.jobUser"
member = "serviceAccount:${google_service_account.delivery_api_gsa.email}"
}
# ── Cloud Armor セキュリティポリシー(OWASP SQLi / XSS + レートリミット)──
resource "google_compute_security_policy" "mops_waf_policy" {
name = "mops-waf-policy"
description = "MOps 配信 API 用 WAF(OWASP CRS SQLi / XSS + DDoS 緩和)"
project = local.project_id
# デフォルトルール: 全トラフィックを許可
rule {
action = "allow"
priority = 2147483647
match {
versioned_expr = "SRC_IPS_V1"
config { src_ip_ranges = ["*"] }
}
description = "デフォルト: 全許可"
}
# SQLインジェクション検知(OWASP CRS sqli-v33-stable)
rule {
action = "deny(403)"
priority = 1000
match {
expr {
# -stable は本番向けに false positive が低いルールセット
expression = "evaluatePreconfiguredExpr('sqli-v33-stable')"
}
}
description = "OWASP SQLi: SQLインジェクションを 403 で拒否"
}
# XSS 検知
rule {
action = "deny(403)"
priority = 1001
match {
expr {
expression = "evaluatePreconfiguredExpr('xss-v33-stable')"
}
}
description = "OWASP XSS: クロスサイトスクリプティングを 403 で拒否"
}
# レートリミット(DDoS 緩和): IP あたり 100 req/分 を超えたら 429
rule {
action = "throttle"
priority = 900
match {
versioned_expr = "SRC_IPS_V1"
config { src_ip_ranges = ["*"] }
}
rate_limit_options {
rate_limit_threshold {
count = 100
interval_sec = 60
}
conform_action = "allow"
exceed_action = "deny(429)"
enforce_on_key = "IP"
}
description = "DDoS 緩和: 100 req/min 超えで 429"
}
}
# ── Cloud Run v2 デプロイ(最小権限 SA)─────────────────────────────────────
resource "google_cloud_run_v2_service" "delivery_api" {
name = "delivery-api"
location = var.region
project = local.project_id
template {
service_account = google_service_account.delivery_api_gsa.email # 最小権限 SA
containers {
image = "gcr.io/${local.project_id}/delivery-api:latest"
resources {
limits = { cpu = "1", memory = "512Mi" }
cpu_idle = true # アイドル時 CPU 解放(コスト削減)
startup_cpu_boost = true # コールドスタート高速化
}
}
}
ingress = "INGRESS_TRAFFIC_INTERNAL_LOAD_BALANCER"
}
# ── Cloud Run v2 の Serverless NEG(Cloud Armor アタッチ用)────────────────
resource "google_compute_region_network_endpoint_group" "delivery_api_neg" {
name = "delivery-api-neg"
region = var.region
project = local.project_id
network_endpoint_type = "SERVERLESS"
cloud_run {
service = google_cloud_run_v2_service.delivery_api.name
}
}
# ── Backend Service(Cloud Armor をここでアタッチ)──────────────────────────
resource "google_compute_backend_service" "delivery_api_backend" {
name = "delivery-api-backend"
project = local.project_id
security_policy = google_compute_security_policy.mops_waf_policy.id # Cloud Armor アタッチ
backend {
group = google_compute_region_network_endpoint_group.delivery_api_neg.id
}
log_config {
enable = true
sample_rate = 1.0 # 全リクエストをロギング(WAF チューニング用)
}
}
Bad vs Good 設計比較
| 観点 | Bad(現状) | Good(改善後) |
|---|---|---|
| GKE Pod の認証方式 | SA キー JSON(GOOGLE_APPLICATION_CREDENTIALS)を Secret に格納 → ローテーション負債・漏洩リスク |
Workload Identity Federation — KSA Annotation + GSA バインディングのみ → キーファイルなし・ローテーション不要 |
| Cloud Run v2 IAM 権限 | roles/bigquery.admin(データセット削除・IAM 変更まで可能) |
roles/bigquery.dataEditor + jobUser(テーブル読み書き + クエリ実行のみ) |
| WAF / アプリ保護 | Cloud Armor 未設定(SQLi・XSS が直接 Cloud Run に届く) | Cloud Armor evaluatePreconfiguredExpr で OWASP CRS SQLi / XSS を 403 拒否 |
| IAM リソース管理 | 手動で gcloud コマンド実行 → Terraform state と乖離 | google_service_account_iam_member で Terraform 管理 → drift 検知・レビュー可能 |
| 侵害時の影響範囲 | admin 権限で全データセット削除・エクスポートが可能(影響: 壊滅的) | dataEditor のみ → テーブル書き込みまで(DDL・IAM 変更は不可) |
| DDoS 対策 | なし(トラフィックが制限なく Cloud Run に到達) | Cloud Armor レートリミット(100 req/min/IP → 429) |
Workload Identity 確認コマンド
# 1. GKE クラスタの Workload Identity Pool 確認
gcloud container clusters describe mops-cluster \
--region asia-northeast1 \
--format="value(workloadIdentityConfig.workloadPool)"
# 出力例: {project_id}.svc.id.goog
# 2. GSA の IAM バインディング確認
gcloud iam service-accounts get-iam-policy mops-batch-sa@{project_id}.iam.gserviceaccount.com
# roles/iam.workloadIdentityUser が
# serviceAccount:{project_id}.svc.id.goog[mops/mops-batch-ksa] に付与されていることを確認
# 3. KSA の Annotation 確認
kubectl get serviceaccount mops-batch-ksa -n mops -o yaml
# annotations:
# iam.gke.io/gcp-service-account: mops-batch-sa@{project_id}.iam.gserviceaccount.com
# 4. Pod 内から認証が通ることを確認
kubectl run -it test-pod --image=google/cloud-sdk:slim \
--serviceaccount=mops-batch-ksa -n mops --rm \
-- bash -c "gcloud auth list && bq ls"
# ACTIVE ACCOUNT が mops-batch-sa@... になっていれば Workload Identity 設定成功
# 5. Cloud Armor ポリシーの確認
gcloud compute security-policies describe mops-waf-policy --format=json | \
jq '.rules[] | {priority: .priority, action: .action}'
# priority 900 (throttle), 1000 (deny sqli), 1001 (deny xss), 2147483647 (allow) が確認できれば OK
# 6. SA キー削除(Workload Identity 移行完了後に実施)
# 既存の SA キーを全て削除してローテーション負債を解消
gcloud iam service-accounts keys list \
--iam-account=mops-batch-sa@{project_id}.iam.gserviceaccount.com
# 上記で USER_MANAGED のキーが 0 件になっていることを確認
# (SYSTEM_MANAGED は GCP 内部用なので問題なし)
ポイント解説
カテゴリ A
@asynccontextmanager と yield の組み合わせで「取得・使用・解放」を1つの Context Manager に封じ込める。async with distributed_lock(redis, product_id): という呼び出し側のコードから「ロックの実装詳細」が完全に隠蔽される。これにより reserve_inventory はビジネスロジックのみに集中でき、単体テストでも distributed_lock を Mock しやすくなる(Ch11: テスト容易性)。
Python 3.11+ の
asyncio.timeout(seconds) は async with 構文で使える宣言的タイムアウト API。asyncio.wait_for(coro, timeout=5.0) はコルーチンをラップするため TaskGroup 内では使いづらかった。timeout() はキャンセル伝播が CancelledError ベースで正確に動作し、TimeoutError(asyncio.TimeoutError)として呼び出し元に伝わる。
Pydantic v2 の
field_validator は単一フィールドの検証に使い、model_validator(mode="after") は全フィールドが設定された後に複数フィールドを参照したバリデーションを行う。「qty > 0 かつ product_id が PROD- 形式でなければならない」というビジネスルールをモデルに閉じ込めることで、呼び出し元がバリデーションロジックを書かなくて済む。
カテゴリ B
GKE の Workload Identity Pool は
{project_id}.svc.id.goog という形式。KSA の iam.gke.io/gcp-service-account Annotation と GSA の workloadIdentityUser バインディングが揃うと、Pod が OIDC トークンを自動取得して GSA に impersonate できる。SA キーファイルは一切存在しないため、Secret 管理・ローテーションのコストがゼロになる。
bigquery.admin はデータセット削除・IAM 設定変更・全テーブルへのアクセスを含む。配信 API に必要なのは「テーブル読み書き(dataEditor)」と「クエリ実行(jobUser)」のみ。侵害時に admin 権限が悪用されると全データセット削除が可能だが、dataEditor のみでは DDL(CREATE/DROP)・IAM 変更は不可。影響範囲を最小化する「最小権限原則(Principle of Least Privilege)」の実践。
Cloud Run v2 は SERVERLESS インフラのため Cloud Armor を直接アタッチできない。
Cloud Run → Serverless NEG → Backend Service → GLB → Cloud Armor というルートを構築する。google_compute_region_network_endpoint_group の network_endpoint_type = "SERVERLESS" と cloud_run.service でエンドポイントを定義し、google_compute_backend_service の security_policy でポリシーをアタッチする。
実務への応用
- asynccontextmanager + Redis ロックは MOps Flash Sale の必須パターン: ECサイトの Flash Sale(瞬間大量アクセス)で同一商品への同時予約リクエストが殺到する。Redis NX ロック + TTL で「在庫チェックと減算をアトミック」にする設計が MOps の Cart / Checkout で基本になる。
asyncio.timeoutでロック待機の上限を設けることで、Redis ダウン時の無限待機を防ぎ SLA を守れる - frozen dataclass → pytest が読みやすくなる:
ReservationResult(product_id="PROD-001", reserved=True, qty=3, remaining_stock=7)のようにテストデータを直接コンストラクタで作れる。assert result.success_rate == 0.9と書けるアサーションは仕様として読め、dict のresult["remaining_stock"]より型安全 - Workload Identity への移行は GKE プロジェクトの最優先 IAM 作業: SA キー JSON を ConfigMap や Secret に置いている既存システムは即移行すべき。移行コストは「GSA 作成 + バインディング + KSA Annotation 追加 + GOOGLE_APPLICATION_CREDENTIALS 環境変数削除」のみ。キー削除後はローテーション作業が消える
- Cloud Armor WAF は Cloud Run v2 の GLB 経路で必須: PCI DSS・SOC2 を意識する EC サイトでは WAF は必須要件。Cloud Armor の
evaluatePreconfiguredExprは Google が管理するため、OWASP ルールの更新を自分でメンテナンスしなくて良い。log_config.sample_rate = 1.0で全リクエストをロギングし、false positive チューニングに使う - 証券マン視点: IAM 過剰権限の財務リスク: IBM 2024 データ漏洩コスト調査では平均 4.88M USD(約 7.5億円)。
bigquery.admin漏洩なら全データセット削除・エクスポートが可能。最小権限への移行工数 0.5 人日(3万円)に対してリスク低減は数億円規模。ROI は数万倍以上でやらない理由がない
今日のまとめ
@asynccontextmanager + asyncio.timeout() + finally の組み合わせで「取得タイムアウト付き・例外時解放保証付き」分散ロックを宣言的に実装でき、Pydantic v2 model_validator(mode="after") でフィールド間依存バリデーションを型安全に表現し、frozen dataclass ReservationResult で予約結果を値オブジェクトとして返す設計が MOps 在庫バッチの基盤パターンになる(Ch2/Ch4/Ch7/Ch10/Ch11)。インフラ側では Workload Identity Federation で GKE Pod を SA キーレスにし(KSA Annotation + GSA
workloadIdentityUser バインディングの2ステップ)、Cloud Run v2 の IAM を roles/bigquery.admin から dataEditor + jobUser 最小権限に絞り、Cloud Armor の evaluatePreconfiguredExpr('sqli-v33-stable') で OWASP CRS WAF を宣言的に有効化する「キーレス × 最小権限 × WAF 防護」のセキュリティ設計が 2026 年の GCP デファクトスタンダードになっている。