返回博客

源码精读: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 策略引擎集成、评估链的执行流程、查询重写的字段级过滤实现,以及批量评估的部分失败语义。

Coomia发布于 2025年12月3日13 分钟阅读
分享本文Twitter / X

源码精读: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 策略引擎集成、评估链的执行流程、查询重写的字段级过滤实现,以及批量评估的部分失败语义。

#目录

  1. 整体架构与三方协作者
  2. gRPC 服务层:PolicyEngineServiceImpl
  3. 策略 CRUD:Create / Get / List / Update / Delete
  4. 策略评估链:EvaluatePolicy
  5. Cedar 策略引擎:CedarPolicyEngine
  6. 开发模式:SimplePolicyEngine 的全放行策略
  7. 批量评估:BatchEvaluate 的部分失败语义
  8. 查询重写:RewriteQuery — 行级安全与列级掩码
  9. 策略测试:TestPolicy 的 dry-run 实现
  10. 访问模拟:SimulateAccess 的逐步追踪
  11. Key Takeaways

#1. 整体架构与三方协作者

PolicyEngineServiceImpl 的代码位于 Control Layer(Control Layer),遵循 Spring Boot gRPC 服务模式:

Code
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策略评估链的执行构造器注入
PolicyProtoMapperProto 与 Domain 的双向转换构造器注入
Java
@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创建策略CreatePolicyRequestPolicyResponse
getPolicy获取策略GetPolicyRequestPolicyResponse
listPolicies列表查询ListPoliciesRequestListPoliciesResponse
updatePolicy更新策略UpdatePolicyRequestPolicyResponse
deletePolicy删除策略DeletePolicyRequestDeleteResponse
evaluatePolicy单条评估EvaluateRequestEvaluateResponse
batchEvaluate批量评估BatchEvaluateRequestBatchEvaluateResponse
rewriteQuery查询重写RewriteQueryRequestRewriteQueryResponse

注解区别:与 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 创建策略

Java
@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());
    }
}

关键设计

  1. 幂等性检查:在保存前检查 existsById,防止重复创建。这比依赖数据库唯一约束抛异常更优雅——客户端可以明确区分"已存在"和"内部错误"
  2. WorldContextHolder:通过 ThreadLocal 从 gRPC Interceptor 中获取当前用户 ID,默认值 "system" 用于无认证的内部调用

#3.2 更新策略的不可变字段保护

Java
// 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);

设计亮点policyIdcreatedAtcreatedByversion 四个字段是不可变的——即使客户端在更新请求中发送了这些字段的新值,服务端也会用原始值覆盖。这是防御性编程的典范。

#3.3 列表查询的分页实现

Java
// 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 资源类型映射

Java
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

Java
@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 是三种权限模型的统一抽象:

Java
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 在生产环境中使用它作为核心策略引擎:

Java
@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 评估流程五步法

Java
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);
}

五步法的设计考量

  1. 活跃策略过滤:避免将 DRAFT/DEPRECATED 状态的策略送入评估引擎
  2. 类型转换隔离CedarMapperCedarPolicyConverter 将 coomia-dip 领域模型与 Cedar SDK 类型完全解耦
  3. 错误兜底AuthException 映射为 DENY,确保评估失败时的安全默认行为

#5.2 默认拒绝策略

Java
if (policies == null || policies.isEmpty()) {
    return PolicyDecision.defaultDeny();
}

安全默认值:当没有匹配的策略时,返回 defaultDeny()——这是零信任安全模型的核心原则。

#6. 开发模式:SimplePolicyEngine 的全放行策略

Java
@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 的部分失败语义

Java
@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 — 行级安全与列级掩码

Java
@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) { ... }
}

字段级过滤约定:策略条件的 attributefield. 前缀开头时,表示这是一个字段级权限规则。例如 field.salary 表示对 salary 字段的掩码策略。

应用过滤器结构

PROTOBUF
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 实现

Java
@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 的逐步追踪

Java
@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

  1. 三模型统一:通过 EvaluationContext 的 Subject/Resource/Action 三元组抽象,RBAC(角色)、ABAC(属性)、ReBAC(关系)三种模型共用一个评估入口
  2. Cedar 集成:生产环境使用 AWS Cedar 策略引擎,开发环境使用 SimplePolicyEngine 全放行——通过 @Profile + @ConditionalOnProperty 无缝切换
  3. 默认拒绝:无匹配策略时返回 defaultDeny(),遵循零信任原则
  4. 查询重写field. 前缀约定实现字段级权限控制,输出 AppliedFilter 供 Data Layer 执行列级掩码
  5. 批量部分失败:单个评估异常不阻断整个批量请求,失败项返回 DENY + 错误原因
  6. 策略测试testPolicy 提供 dry-run 能力,支持测试驱动的策略开发
  7. 访问模拟simulateAccesssteps 列表是权限调试的核心工具,支持 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