oscerd commented on code in PR #26974: URL: https://github.com/apache/camel/pull/26974#discussion_r4142798325
########## components/camel-openfga/src/main/java/org/apache/camel/component/openfga/OpenFgaProducer.java: ########## @@ -0,0 +1,489 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.camel.component.openfga; + +import java.util.ArrayList; +import java.util.Collection; +import java.util.HashSet; +import java.util.List; +import java.util.Map; +import java.util.Set; +import java.util.function.Consumer; + +import dev.openfga.sdk.api.client.model.ClientBatchCheckClientResponse; +import dev.openfga.sdk.api.client.model.ClientCheckRequest; +import dev.openfga.sdk.api.client.model.ClientListObjectsRequest; +import dev.openfga.sdk.api.client.model.ClientListRelationsRequest; +import dev.openfga.sdk.api.client.model.ClientListUsersRequest; +import dev.openfga.sdk.api.client.model.ClientTupleKey; +import dev.openfga.sdk.api.client.model.ClientTupleKeyWithoutCondition; +import dev.openfga.sdk.api.configuration.ClientBatchCheckClientOptions; +import dev.openfga.sdk.api.configuration.ClientListObjectsOptions; +import dev.openfga.sdk.api.configuration.ClientListRelationsOptions; +import dev.openfga.sdk.api.configuration.ClientListUsersOptions; +import dev.openfga.sdk.api.model.ConsistencyPreference; +import dev.openfga.sdk.api.model.FgaObject; +import dev.openfga.sdk.api.model.User; +import dev.openfga.sdk.api.model.UserTypeFilter; +import org.apache.camel.Exchange; +import org.apache.camel.InvalidPayloadException; +import org.apache.camel.health.HealthCheckHelper; +import org.apache.camel.health.WritableHealthCheckRepository; +import org.apache.camel.support.DefaultProducer; +import org.apache.camel.util.ObjectHelper; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +public class OpenFgaProducer extends DefaultProducer { + + private static final Logger LOG = LoggerFactory.getLogger(OpenFgaProducer.class); + + private static final String DEFAULT_USER_FILTER = "user"; + + private OpenFgaProducerHealthCheck producerHealthCheck; + private WritableHealthCheckRepository healthCheckRepository; + + public OpenFgaProducer(final OpenFgaEndpoint endpoint) { + super(endpoint); + } + + @Override + public OpenFgaEndpoint getEndpoint() { + return (OpenFgaEndpoint) super.getEndpoint(); + } + + @Override + protected void doStart() throws Exception { + super.doStart(); + + OpenFgaConfiguration configuration = getEndpoint().getConfiguration(); + // an injected client can point anywhere and the endpoint has no way to ask it where, so there is no server + // this check would know the address of + if (configuration.getOpenFgaClient() != null || ObjectHelper.isEmpty(configuration.getApiUrl())) { + return; + } + + // health-check is optional so discover and resolve + healthCheckRepository = HealthCheckHelper.getHealthCheckRepository( + getEndpoint().getCamelContext(), + "producers", + WritableHealthCheckRepository.class); + + if (healthCheckRepository != null) { + producerHealthCheck = new OpenFgaProducerHealthCheck( + configuration.getApiUrl(), configuration.getApiToken(), configuration.getStoreId(), + // the endpoint URI is unique within the context, so two endpoints sharing a store but pointing at + // different servers get distinct health-check ids instead of colliding + getEndpoint().getEndpointUri(), getEndpoint().getSslContext()); + producerHealthCheck.setEnabled(getEndpoint().getComponent().isHealthCheckProducerEnabled()); + healthCheckRepository.addHealthCheck(producerHealthCheck); + } + } + + @Override + protected void doStop() throws Exception { + if (healthCheckRepository != null && producerHealthCheck != null) { + healthCheckRepository.removeHealthCheck(producerHealthCheck); + producerHealthCheck = null; + } + super.doStop(); + } + + @Override + public void process(Exchange exchange) throws Exception { + OpenFgaAuthorizer authorizer = getEndpoint().getAuthorizer(); + switch (getEndpoint().getOperation()) { + case check -> authorizer.check(exchange); + case batchCheck -> batchCheck(exchange, authorizer); + case listObjects -> listObjects(exchange, authorizer); + case listRelations -> listRelations(exchange, authorizer); + case listUsers -> listUsers(exchange, authorizer); + case writeTuples -> writeTuples(exchange, authorizer); + case deleteTuples -> deleteTuples(exchange, authorizer); + // unreachable today; here so that adding an operation to the enum without wiring it up fails loudly + // instead of silently letting the exchange through an endpoint that was asked to authorize it + default -> throw new IllegalArgumentException("Unsupported operation " + getEndpoint().getOperation()); + } + } + + /** + * Checks one subject against many objects, taken from the body, and replaces the body with the ones the check + * allowed - so a route can filter a collection down to what the caller may see in one step. + */ + private void batchCheck(Exchange exchange, OpenFgaAuthorizer authorizer) throws Exception { + OpenFgaAuthorizer.clearDecisionHeaders(exchange); + List<String> objects = bodyAsList(exchange); + String user = authorizer.resolveUser(exchange); + String relation = authorizer.resolveRelation(exchange); + if (user == null || relation == null) { + // the subject or the permission could not be established, so none of the objects is allowed. The deny + // reason is already on the exchange + exchange.getMessage().setBody(List.of()); + return; + } + + List<ClientCheckRequest> requests = new ArrayList<>(objects.size()); + for (String object : objects) { + if (OpenFgaIdentifiers.validate(object) != null) { + // an entry that cannot be an object of a check is not one the caller may have; leaving it out of the + // batch leaves it out of the result, which is the fail-closed answer + LOG.debug("Skipping '{}' in a batchCheck: not usable as an object identifier", object); + continue; + } + requests.add(new ClientCheckRequest().user(user).relation(relation)._object(object)); + } + if (requests.isEmpty()) { + // nothing in the body could be an object of a check, so nothing is allowed. Say so on the header as + // every other path does, rather than leaving the route to infer a verdict from an empty body + exchange.getMessage().setHeader(OpenFgaConstants.ALLOWED, false); + exchange.getMessage().setBody(List.of()); + return; + } + + ClientBatchCheckClientOptions options = new ClientBatchCheckClientOptions() + .maxParallelRequests(getEndpoint().getConfiguration().getMaxParallelRequests()); + if (ObjectHelper.isNotEmpty(getEndpoint().getConfiguration().getAuthorizationModelId())) { + options.authorizationModelId(getEndpoint().getConfiguration().getAuthorizationModelId()); + } + if (authorizer.getConsistency() != null) { + options.consistency(authorizer.getConsistency()); + } + + List<ClientBatchCheckClientResponse> responses = authorizer.await( + exchange, "batch check " + relation + " for " + user, + () -> authorizer.getClient().clientBatchCheck(requests, options)); + + Set<String> permitted = new HashSet<>(); + for (ClientBatchCheckClientResponse response : responses) { + if (response.getThrowable() != null) { + // one check in the batch never got an answer. Returning the rest would quietly present a partial + // filter as a complete one, so the whole operation fails - and unlike a single check, there is no + // failOpen here, because "proceed" for a filter would mean handing back everything unfiltered + throw new OpenFgaEvaluationException( + "OpenFGA failed to answer part of a batch check for " + user, exchange, + response.getThrowable()); + } + if (Boolean.TRUE.equals(response.getAllowed())) { + permitted.add(response.getRequest().getObject()); + } + } + // filter the body rather than collect the answers: the SDK issues the checks in parallel and hands them back + // in whatever order they completed, so building the result from the responses would reorder the caller's list + // between runs. A route that filters a list to render it wants the order it asked in + List<String> allowed = new ArrayList<>(); + for (String object : objects) { + if (permitted.contains(object)) { + allowed.add(object); + } + } + exchange.getMessage().setHeader(OpenFgaConstants.ALLOWED, !allowed.isEmpty()); + exchange.getMessage().setBody(allowed); + } + + /** + * Lists the objects of the configured type the subject can reach through the relation. + */ + private void listObjects(Exchange exchange, OpenFgaAuthorizer authorizer) throws Exception { + OpenFgaAuthorizer.clearDecisionHeaders(exchange); + String type = getEndpoint().getConfiguration().getType(); + String user = authorizer.resolveUser(exchange); + String relation = authorizer.resolveRelation(exchange); + if (user == null || relation == null) { + exchange.getMessage().setBody(List.of()); + return; + } + + ClientListObjectsRequest request = new ClientListObjectsRequest() + .user(user) + .relation(relation) + .type(type); + ClientListObjectsOptions options = new ClientListObjectsOptions(); + applyModelAndConsistency(authorizer, options::authorizationModelId, options::consistency); + + List<String> objects = authorizer.await(exchange, "list " + type + " objects " + user + " can " + relation, + () -> authorizer.getClient().listObjects(request, options)).getObjects(); + exchange.getMessage().setBody(objects != null ? objects : List.of()); + } + + /** + * Lists which of the configured relations the subject holds on the object. + */ + private void listRelations(Exchange exchange, OpenFgaAuthorizer authorizer) throws Exception { + OpenFgaAuthorizer.clearDecisionHeaders(exchange); + String user = authorizer.resolveUser(exchange); + String object = authorizer.resolveObject(exchange); + if (user == null || object == null) { + exchange.getMessage().setBody(List.of()); + return; + } + + ClientListRelationsRequest request = new ClientListRelationsRequest() + .user(user) + ._object(object) + .relations(splitToList(getEndpoint().getConfiguration().getRelations())); + ClientListRelationsOptions options = new ClientListRelationsOptions(); + applyModelAndConsistency(authorizer, options::authorizationModelId, options::consistency); + + List<String> relations = authorizer.await(exchange, "list the relations " + user + " has on " + object, + () -> authorizer.getClient().listRelations(request, options)).getRelations(); + exchange.getMessage().setBody(relations != null ? relations : List.of()); + } + + /** + * Lists the subjects that hold the relation on the object. + */ + private void listUsers(Exchange exchange, OpenFgaAuthorizer authorizer) throws Exception { + OpenFgaAuthorizer.clearDecisionHeaders(exchange); + String object = authorizer.resolveObject(exchange); + String relation = authorizer.resolveRelation(exchange); + if (object == null || relation == null) { + exchange.getMessage().setBody(List.of()); + return; + } + + ClientListUsersRequest request = new ClientListUsersRequest() + ._object(toFgaObject(object)) + .relation(relation) + .userFilters(userFilters()); + ClientListUsersOptions options = new ClientListUsersOptions(); + applyModelAndConsistency(authorizer, options::authorizationModelId, options::consistency); + + List<User> users = authorizer.await(exchange, "list the users with " + relation + " on " + object, + () -> authorizer.getClient().listUsers(request, options)).getUsers(); + List<String> identifiers = new ArrayList<>(); + if (users != null) { + for (User user : users) { + String identifier = toIdentifier(user); + if (identifier != null) { + identifiers.add(identifier); + } + } + } + exchange.getMessage().setBody(identifiers); + } + + /** + * Writes relationship tuples, granting access. + */ + private void writeTuples(Exchange exchange, OpenFgaAuthorizer authorizer) throws Exception { + List<ClientTupleKey> tuples = new ArrayList<>(); + for (Tuple tuple : resolveTuples(exchange, authorizer)) { + tuples.add(new ClientTupleKey().user(tuple.user).relation(tuple.relation)._object(tuple.object)); + } + // the write options carry no authorizationModelId - unlike every query's options - so the model this writes + // against is the one pinned on the client, which OpenFgaClientFactory already set from the configuration + authorizer.await(exchange, "write " + tuples.size() + " tuple(s)", + () -> authorizer.getClient().writeTuples(tuples)); + exchange.getMessage().setHeader(OpenFgaConstants.WRITTEN_TUPLES, tuples.size()); + } + + /** + * Deletes relationship tuples, revoking access. + */ + private void deleteTuples(Exchange exchange, OpenFgaAuthorizer authorizer) throws Exception { + List<ClientTupleKeyWithoutCondition> tuples = new ArrayList<>(); + for (Tuple tuple : resolveTuples(exchange, authorizer)) { + tuples.add(new ClientTupleKeyWithoutCondition() + .user(tuple.user).relation(tuple.relation)._object(tuple.object)); + } + authorizer.await(exchange, "delete " + tuples.size() + " tuple(s)", + () -> authorizer.getClient().deleteTuples(tuples)); + exchange.getMessage().setHeader(OpenFgaConstants.DELETED_TUPLES, tuples.size()); + } + + /** + * Collects the tuples to write or delete, from the body when it carries any and from the endpoint's + * {@code user}/{@code relation}/{@code object} otherwise - so a route that just created a resource can grant access + * to it without assembling a payload. + * <p/> + * A typed wildcard is accepted here, unlike on the check path: {@code user:*} is exactly how a resource is shared + * with everyone, and writing that tuple is a deliberate act by the route rather than something a caller chose. + */ + private List<Tuple> resolveTuples(Exchange exchange, OpenFgaAuthorizer authorizer) throws InvalidPayloadException { + // The endpoint wins whenever it names a tuple at all. This is an authorization component, so the same rule + // that governs the check governs the write: the configuration decides and the message does not. Reading the + // body in preference would mean a route that unmarshals an untrusted payload hands the caller the choice of + // which relationship to grant - and "user:attacker owner document:secret" is a legitimate-looking tuple. + Tuple configured = new Tuple( + authorizer.rawUser(exchange), authorizer.rawRelation(exchange), authorizer.rawObject(exchange)); Review Comment: Addressed in `fc044b07` — `resolveTuples` now asks `authorizer.hasConfiguredTuple()` (which is `user != null || relation != null || object != null` on the *compiled* expressions, OpenFgaAuthorizer:392), not the evaluated values. So `user=${header.u}&relation=...&object=...` with the headers missing still counts as configured, and `validated(..., true)` fails a part that resolved to null (naming which part) rather than quietly reading the tuple from the body. The body is only used when nothing is configured. _Claude Code on behalf of oscerd_ -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
