源码精读:PolicyEngineService — 三模型权限的统一评估
PolicyEngineServiceImpl 是 Control Layer 中的权限策略引擎,基于 Spring Boot 3.x + gRPC 实现,统一了 RBAC、ABAC、ReBAC 三种权限模型的评估。它通过 PolicyEvaluationService 抽象评估链,支持 8 种条件算子(EQ/NE/GT/GE/LT/LE/IN/CONTAINS),提供策略 CRUD、单条/批量评估、查询重写(行级安全 + 列级掩码)、策略测试(dry-run)和访问模拟五大功能。本文将逐行分析其分层架构、Cedar 策略引擎集成、评估链的执行流程、查询重写的字段级过滤实现,以及批量评估的部分失败语义。
源码精读:PolicyEngineService — 三模型权限的统一评估
“系列:S9 源码精读 · 第 4 篇 | 难度:高级 | 阅读时间:25 分钟
#TL;DR
PolicyEngineServiceImpl 是 Control Layer 中的权限策略引擎,基于 Spring Boot 3.x + gRPC 实现,统一了 RBAC、ABAC、ReBAC 三种权限模型的评估。它通过 PolicyEvaluationService 抽象评估链,支持 8 种条件算子(EQ/NE/GT/GE/LT/LE/IN/CONTAINS),提供策略 CRUD、单条/批量评估、查询重写(行级安全 + 列级掩码)、策略测试(dry-run)和访问模拟五大功能。本文将逐行分析其分层架构、Cedar 策略引擎集成、评估链的执行流程、查询重写的字段级过滤实现,以及批量评估的部分失败语义。
#目录
- 整体架构与三方协作者
- gRPC 服务层:PolicyEngineServiceImpl
- 策略 CRUD:Create / Get / List / Update / Delete
- 策略评估链:EvaluatePolicy
- Cedar 策略引擎:CedarPolicyEngine
- 开发模式:SimplePolicyEngine 的全放行策略
- 批量评估:BatchEvaluate 的部分失败语义
- 查询重写:RewriteQuery — 行级安全与列级掩码
- 策略测试:TestPolicy 的 dry-run 实现
- 访问模拟:SimulateAccess 的逐步追踪
- Key Takeaways
#1. 整体架构与三方协作者
PolicyEngineServiceImpl 的代码位于 Control Layer(Control Layer),遵循 Spring Boot gRPC 服务模式:
control-Layer/src/main/java/com/onto/control/policy/
├── api/
│ ├── PolicyEngineServiceImpl.java # gRPC 服务实现
│ └── PolicyProtoMapper.java # Proto ↔ Domain 转换
├── domain/
│ ├── Policy.java # 策略领域模型
│ ├── PolicyCondition.java # 条件定义
│ ├── PolicyRule.java # 规则定义
│ ├── PolicyDecision.java # 评估决策结果
│ ├── PolicyStatus.java # 策略状态枚举
│ └── EvaluationContext.java # 评估上下文
├── engine/
│ ├── CedarPolicyEngine.java # Cedar SDK 集成(生产)
│ ├── SimplePolicyEngine.java # 简化引擎(开发模式)
│ ├── CedarMapper.java # Cedar 类型映射
│ └── CedarPolicyConverter.java # 策略格式转换
├── service/
│ ├── PolicyStoreService.java # 策略存储接口
│ └── PolicyEvaluationService.java # 评估服务接口
└── proto/
├── PolicyEngineServiceGrpc.java # gRPC 生成代码
└── PolicyEngineProto.java # Proto 消息定义
三方协作者:
| 协作者 | 职责 | 注入方式 |
|---|---|---|
PolicyStoreService | 策略的 CRUD 持久化 | 构造器注入 |
PolicyEvaluationService | 策略评估链的执行 | 构造器注入 |
PolicyProtoMapper | Proto 与 Domain 的双向转换 | 构造器注入 |
@GrpcService
@Slf4j
@RequiredArgsConstructor
public class PolicyEngineServiceImpl
extends PolicyEngineServiceGrpc.PolicyEngineServiceImplBase {
private final PolicyStoreService policyStoreService;
private final PolicyEvaluationService policyEvaluationService;
private final PolicyProtoMapper mapper;
}
设计决策:使用 Lombok @RequiredArgsConstructor 自动生成构造器,避免手写 @Inject。@Slf4j 替代手动 Logger 声明,减少样板代码。
#2. gRPC 服务层:PolicyEngineServiceImpl
PolicyEngineServiceImpl 是一个典型的 gRPC 适配器层,暴露 8 个 RPC 方法:
| 方法 | 用途 | 请求类型 | 响应类型 |
|---|---|---|---|
createPolicy | 创建策略 | CreatePolicyRequest | PolicyResponse |
getPolicy | 获取策略 | GetPolicyRequest | PolicyResponse |
listPolicies | 列表查询 | ListPoliciesRequest | ListPoliciesResponse |
updatePolicy | 更新策略 | UpdatePolicyRequest | PolicyResponse |
deletePolicy | 删除策略 | DeletePolicyRequest | DeleteResponse |
evaluatePolicy | 单条评估 | EvaluateRequest | EvaluateResponse |
batchEvaluate | 批量评估 | BatchEvaluateRequest | BatchEvaluateResponse |
rewriteQuery | 查询重写 | RewriteQueryRequest | RewriteQueryResponse |
注解区别:与 Data Layer 的 @GrpcService(Quarkus)不同,Control Layer 使用 net.devh.boot.grpc.server.service.GrpcService——这是 grpc-spring-boot-starter 提供的 Spring Boot 适配注解。
#3. 策略 CRUD:Create / Get / List / Update / Delete
#3.1 创建策略
@Override
public void createPolicy(CreatePolicyRequest request,
StreamObserver<PolicyResponse> responseObserver) {
try {
// 1. 输入验证
if (!request.hasPolicy()) {
responseObserver.onError(Status.INVALID_ARGUMENT
.withDescription("Policy is required")
.asRuntimeException());
return;
}
// 2. Proto → Domain 转换
Policy policy = mapper.toDomain(request.getPolicy());
// 3. 从 WorldContext 提取创建者
String userId = WorldContextHolder.getUserIdOrDefault("system");
policy.setCreatedBy(userId);
// 4. 幂等性检查
if (policy.getPolicyId() != null
&& policyStoreService.existsById(policy.getPolicyId())) {
responseObserver.onError(Status.ALREADY_EXISTS
.withDescription("Policy already exists: " + policy.getPolicyId())
.asRuntimeException());
return;
}
// 5. 持久化
Policy saved = policyStoreService.savePolicy(policy);
// 6. 响应
responseObserver.onNext(PolicyResponse.newBuilder()
.setPolicy(mapper.toProto(saved))
.build());
responseObserver.onCompleted();
} catch (IllegalArgumentException e) {
responseObserver.onError(Status.INVALID_ARGUMENT
.withDescription(e.getMessage()).asRuntimeException());
} catch (Exception e) {
responseObserver.onError(Status.INTERNAL
.withDescription("Internal error: " + e.getMessage())
.asRuntimeException());
}
}
关键设计:
- 幂等性检查:在保存前检查
existsById,防止重复创建。这比依赖数据库唯一约束抛异常更优雅——客户端可以明确区分"已存在"和"内部错误" - WorldContextHolder:通过 ThreadLocal 从 gRPC Interceptor 中获取当前用户 ID,默认值
"system"用于无认证的内部调用
#3.2 更新策略的不可变字段保护
// Preserve immutable fields
updated.setPolicyId(existing.getPolicyId());
updated.setCreatedAt(existing.getCreatedAt());
updated.setCreatedBy(existing.getCreatedBy());
updated.setVersion(existing.getVersion());
// Set updater
String userId = WorldContextHolder.getUserIdOrDefault("system");
updated.setUpdatedBy(userId);
设计亮点:policyId、createdAt、createdBy、version 四个字段是不可变的——即使客户端在更新请求中发送了这些字段的新值,服务端也会用原始值覆盖。这是防御性编程的典范。
#3.3 列表查询的分页实现
// Apply pagination
int page = request.hasPage() ? request.getPage().getPage() : 0;
int size = request.hasPage() && request.getPage().getSize() > 0
? request.getPage().getSize() : 20;
int start = page * size;
int end = Math.min(start + size, policies.size());
List<Policy> pagedPolicies = start < policies.size()
? policies.subList(start, end)
: List.of();
分页采用偏移量模式(page * size),默认每页 20 条。注意 subList 的边界保护:当 start >= policies.size() 时返回空列表,避免 IndexOutOfBoundsException。
#3.4 资源类型映射
private String mapResourceTypeToString(ResourceType type) {
return switch (type) {
case RESOURCE_TYPE_WORLD -> "World";
case RESOURCE_TYPE_ONTOLOGY_TYPE -> "OntologyType";
case RESOURCE_TYPE_ONTOLOGY_INSTANCE -> "OntologyInstance";
case RESOURCE_TYPE_FIELD -> "Field";
case RESOURCE_TYPE_ACTION -> "Action";
case RESOURCE_TYPE_PIPELINE -> "Pipeline";
default -> null;
};
}
coomia-dip 支持 6 种资源类型的权限控制,覆盖了从数据世界(World)到具体字段(Field)的完整粒度。
#4. 策略评估链:EvaluatePolicy
@Override
public void evaluatePolicy(EvaluateRequest request,
StreamObserver<EvaluateResponse> responseObserver) {
try {
// 1. 验证三要素:Subject + Resource + Action
if (!request.hasSubject() || !request.hasResource()
|| request.getAction().isBlank()) {
responseObserver.onError(Status.INVALID_ARGUMENT
.withDescription("Subject, resource, and action are required")
.asRuntimeException());
return;
}
// 2. 构建评估上下文
EvaluationContext context = mapper.toEvaluationContext(request);
// 3. 执行评估链
PolicyDecision decision = policyEvaluationService.evaluate(context);
// 4. 构建响应
EvaluateResponse response = mapper.toEvaluateResponse(decision,
decision.policyId() != null ? List.of(decision.policyId()) : List.of());
responseObserver.onNext(response);
responseObserver.onCompleted();
} catch (Exception e) {
responseObserver.onError(Status.INTERNAL
.withDescription("Internal error: " + e.getMessage())
.asRuntimeException());
}
}
三要素模型:每次评估都必须提供 Subject(谁)、Resource(什么)、Action(做什么)——这是 ABAC/PBAC 的标准模型,也与 Cedar 策略语言的 principal/resource/action 三元组一致。
EvaluationContext 是三种权限模型的统一抽象:
EvaluationContext context = EvaluationContext.builder()
.subject(new Subject(id, type, attributes)) // RBAC: role; ABAC: attributes
.resource(new Resource(id, type, attributes)) // Resource attributes
.action(new Action(id, category)) // Action being performed
.build();
#5. Cedar 策略引擎:CedarPolicyEngine
Cedar 是 AWS 开源的策略语言和评估引擎,coomia-dip 在生产环境中使用它作为核心策略引擎:
@Service
@ConditionalOnProperty(name = "policy-engine.use-cedar",
havingValue = "true", matchIfMissing = false)
public class CedarPolicyEngine {
private final CedarMapper mapper;
private final CedarPolicyConverter converter;
private final AuthorizationEngine authorizationEngine;
public CedarPolicyEngine(CedarMapper mapper, CedarPolicyConverter converter) {
this.mapper = mapper;
this.converter = converter;
this.authorizationEngine = new BasicAuthorizationEngine();
}
}
条件激活:@ConditionalOnProperty 确保 Cedar 引擎仅在 policy-engine.use-cedar=true 时加载。这避免了在不支持 Cedar 本地库的开发环境中出现加载错误。
#5.1 评估流程五步法
public PolicyDecision evaluate(EvaluationContext context, List<Policy> policies) {
// Step 1: 过滤活跃策略
List<Policy> activePolicies = policies.stream()
.filter(Policy::isActive)
.collect(Collectors.toList());
// Step 2: 转换为 Cedar AuthorizationRequest
AuthorizationRequest request = mapper.toAuthorizationRequest(context);
// Step 3: 转换策略为 Cedar PolicySet
PolicySet policySet = converter.toPolicySet(activePolicies);
// Step 4: 构建 Cedar 实体集合
Set<Entity> entities = mapper.toEntities(context);
// Step 5: 调用 Cedar 评估引擎
AuthorizationResponse response =
authorizationEngine.isAuthorized(request, policySet, entities);
return mapper.toDecision(response, policyIds);
}
五步法的设计考量:
- 活跃策略过滤:避免将 DRAFT/DEPRECATED 状态的策略送入评估引擎
- 类型转换隔离:
CedarMapper和CedarPolicyConverter将 coomia-dip 领域模型与 Cedar SDK 类型完全解耦 - 错误兜底:
AuthException映射为DENY,确保评估失败时的安全默认行为
#5.2 默认拒绝策略
if (policies == null || policies.isEmpty()) {
return PolicyDecision.defaultDeny();
}
安全默认值:当没有匹配的策略时,返回 defaultDeny()——这是零信任安全模型的核心原则。
#6. 开发模式:SimplePolicyEngine 的全放行策略
@Service("cedarPolicyEngine")
@Profile({"dev", "default"})
public class SimplePolicyEngine {
public PolicyDecision evaluate(EvaluationContext context, List<Policy> policies) {
return PolicyDecision.permit(
"dev-mock-policy",
"Permitted by development mock policy engine",
List.of()
);
}
}
双 Profile 激活:@Profile({"dev", "default"}) 意味着在开发模式和未指定 Profile 的默认模式下,所有请求都被放行。这大幅降低了开发环境的配置复杂度。
Bean 名称冲突处理:@Service("cedarPolicyEngine") 使用与生产引擎相同的 Bean 名称,通过 Profile 切换实现无缝替换。
#7. 批量评估:BatchEvaluate 的部分失败语义
@Override
public void batchEvaluate(BatchEvaluateRequest request,
StreamObserver<BatchEvaluateResponse> responseObserver) {
try {
List<EvaluateResponse> responses = new java.util.ArrayList<>();
for (EvaluateRequest evalRequest : request.getRequestsList()) {
try {
// 逐个评估
EvaluationContext context = mapper.toEvaluationContext(evalRequest);
PolicyDecision decision = policyEvaluationService.evaluate(context);
responses.add(mapper.toEvaluateResponse(decision, ...));
} catch (Exception e) {
// 单个失败不影响整体
responses.add(EvaluateResponse.newBuilder()
.setAllowed(false)
.setDecision(PolicyEffect.POLICY_EFFECT_DENY)
.setReason("Evaluation error: " + e.getMessage())
.build());
}
}
responseObserver.onNext(BatchEvaluateResponse.newBuilder()
.addAllResponses(responses).build());
responseObserver.onCompleted();
} catch (Exception e) {
responseObserver.onError(Status.INTERNAL...);
}
}
部分失败语义:与 OntologyRuntimeService 的批量操作类似,单个评估失败时返回 DENY + 错误原因,而不是中止整个批量请求。这在权限预检(preflight check)场景中尤为重要——前端需要一次性知道用户对多个资源的访问权限。
#8. 查询重写:RewriteQuery — 行级安全与列级掩码
@Override
public void rewriteQuery(RewriteQueryRequest request,
StreamObserver<RewriteQueryResponse> responseObserver) {
try {
String originalQuery = request.getOriginalQuery();
List<String> requestedFields = request.getRequestedFieldsList();
// 获取目标 Schema 的活跃策略
List<Policy> policies =
policyStoreService.getActivePolicies(request.getTargetSchema());
List<AppliedFilter> appliedFilters = new ArrayList<>();
List<String> filteredFields = new ArrayList<>();
for (Policy policy : policies) {
// 匹配资源类型
if (policy.getMetadata().containsKey("resourceType")) {
String resourceType = policy.getMetadata()
.get("resourceType").toString();
if (resourceType.equals(request.getTargetSchema())) {
// 遍历规则条件,提取字段级过滤
for (PolicyRule rule : policy.getRules()) {
for (PolicyCondition condition : rule.getConditions()) {
if (condition.getAttribute().startsWith("field.")) {
String fieldName =
condition.getAttribute().substring(6);
filteredFields.add(fieldName);
appliedFilters.add(AppliedFilter.newBuilder()
.setField(fieldName)
.setFilterType("COLUMN_MASK")
.setFilterExpression(
condition.getValue().toString())
.setPolicyId(policy.getPolicyId())
.build());
}
}
}
}
}
}
RewriteQueryResponse response = RewriteQueryResponse.newBuilder()
.setRewrittenQuery(rewrittenQuery)
.addAllFilteredFields(filteredFields)
.addAllMaskedFields(maskedFields)
.addAllFilters(appliedFilters)
.build();
responseObserver.onNext(response);
responseObserver.onCompleted();
} catch (Exception e) { ... }
}
字段级过滤约定:策略条件的 attribute 以 field. 前缀开头时,表示这是一个字段级权限规则。例如 field.salary 表示对 salary 字段的掩码策略。
应用过滤器结构:
message AppliedFilter {
string field = 1; // 被过滤的字段名
string filter_type = 2; // COLUMN_MASK / ROW_FILTER
string filter_expression = 3; // 掩码表达式
string policy_id = 4; // 来源策略 ID
}
这个 AppliedFilter 结构允许 Data Layer 在查询执行时精确应用列级掩码——例如将 salary 字段替换为 **** 或 NULL。
#9. 策略测试:TestPolicy 的 dry-run 实现
@Override
public void testPolicy(TestPolicyRequest request,
StreamObserver<TestPolicyResponse> responseObserver) {
try {
Policy policy = mapper.toDomain(request.getPolicy());
List<TestResult> results = new ArrayList<>();
boolean allPassed = true;
for (TestCase testCase : request.getTestCasesList()) {
// 构建评估上下文
EvaluationContext context = EvaluationContext.builder()
.subject(new Subject(
testCase.getSubject().getId(),
testCase.getSubject().getType().name(),
new HashMap<>(testCase.getSubject().getAttributesMap())))
.resource(new Resource(...))
.action(new Action(testCase.getAction(), ...))
.build();
PolicyDecision decision = policyEvaluationService.evaluate(context);
boolean passed = decision.isPermitted() == testCase.getExpectedAllowed();
if (!passed) allPassed = false;
results.add(TestResult.newBuilder()
.setTestCaseName(testCase.getName())
.setPassed(passed)
.setActualAllowed(decision.isPermitted())
.setExpectedAllowed(testCase.getExpectedAllowed())
.setReason(decision.reason())
.build());
}
responseObserver.onNext(TestPolicyResponse.newBuilder()
.setAllPassed(allPassed)
.addAllResults(results)
.build());
responseObserver.onCompleted();
} catch (Exception e) { ... }
}
测试驱动策略开发:testPolicy 允许策略编写者在不部署到生产环境的情况下验证策略逻辑。每个 TestCase 包含预期结果(expectedAllowed),引擎执行后比较实际结果,输出 pass/fail 报告。
#10. 访问模拟:SimulateAccess 的逐步追踪
@Override
public void simulateAccess(SimulateAccessRequest request,
StreamObserver<SimulateAccessResponse> responseObserver) {
try {
List<Policy> policies = policyStoreService.getActivePolicies(resourceType);
List<PolicyEvaluationStep> steps = new ArrayList<>();
for (Policy policy : policies) {
PolicyDecision decision = policyEvaluationService.evaluate(context);
steps.add(PolicyEvaluationStep.newBuilder()
.setPolicyId(policy.getPolicyId())
.setPolicyName(policy.getName())
.setMatched(decision.policyId() != null
&& decision.policyId().equals(policy.getPolicyId()))
.setEffect(mapper.toProtoEffect(decision.effect()))
.setReason(decision.reason())
.build());
// verbose=false 时遇到首个匹配即停止
if (!request.getVerbose() && decision.policyId() != null) {
break;
}
}
boolean allowed = steps.stream()
.anyMatch(step -> step.getEffect() == PolicyEffect.POLICY_EFFECT_ALLOW);
responseObserver.onNext(SimulateAccessResponse.newBuilder()
.setAllowed(allowed)
.addAllSteps(steps)
.build());
responseObserver.onCompleted();
} catch (Exception e) { ... }
}
逐步追踪:SimulateAccess 的核心价值在于 steps 列表——它记录了每条策略的评估过程(匹配/不匹配、效果、原因),是权限调试的利器。
Verbose 模式:当 verbose=false 时采用短路求值(遇到首个匹配即停止),提高性能;verbose=true 时遍历所有策略,提供完整的评估链路。
#11. Key Takeaways
- 三模型统一:通过
EvaluationContext的 Subject/Resource/Action 三元组抽象,RBAC(角色)、ABAC(属性)、ReBAC(关系)三种模型共用一个评估入口 - Cedar 集成:生产环境使用 AWS Cedar 策略引擎,开发环境使用 SimplePolicyEngine 全放行——通过
@Profile+@ConditionalOnProperty无缝切换 - 默认拒绝:无匹配策略时返回
defaultDeny(),遵循零信任原则 - 查询重写:
field.前缀约定实现字段级权限控制,输出AppliedFilter供 Data Layer 执行列级掩码 - 批量部分失败:单个评估异常不阻断整个批量请求,失败项返回 DENY + 错误原因
- 策略测试:
testPolicy提供 dry-run 能力,支持测试驱动的策略开发 - 访问模拟:
simulateAccess的steps列表是权限调试的核心工具,支持 verbose/non-verbose 两种模式
#下一篇
S9-05:QueryFederationService — 多引擎查询路由。我们将深入 Data Layer 的查询联邦服务,看它如何根据查询计划将 OQL 查询路由到 Doris 或 DuckDB,以及代价估算和结果合并的实现。
Tags: #coomia-dip #source-code-reading #control-Layer #policy-engine #rbac #abac #rebac #cedar #query-rewriting #access-control