From ee9f7c15cd11d6be937d98991d7172cc2a72c6c2 Mon Sep 17 00:00:00 2001 From: Raymond Meester <23716892+skin27@users.noreply.github.com> Date: Fri, 7 Aug 2026 14:33:43 +0200 Subject: [PATCH] Add metadata to groprepository and add custom tenantrepository --- .../TenantVariableRepository.java | 202 ++++++++++++++++++ .../org/assimbly/dil/loader/FlowLoader.java | 26 +++ 2 files changed, 228 insertions(+) create mode 100644 dil/src/main/java/org/assimbly/dil/blocks/repositories/TenantVariableRepository.java diff --git a/dil/src/main/java/org/assimbly/dil/blocks/repositories/TenantVariableRepository.java b/dil/src/main/java/org/assimbly/dil/blocks/repositories/TenantVariableRepository.java new file mode 100644 index 00000000..922f7851 --- /dev/null +++ b/dil/src/main/java/org/assimbly/dil/blocks/repositories/TenantVariableRepository.java @@ -0,0 +1,202 @@ +package org.assimbly.dil.blocks.repositories; + +import java.util.ArrayList; +import java.util.List; +import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; +import java.util.stream.Stream; + +import org.apache.camel.CamelContext; +import org.apache.camel.CamelContextAware; +import org.apache.camel.StreamCache; +import org.apache.camel.StreamCacheException; +import org.apache.camel.spi.BrowsableVariableRepository; +import org.apache.camel.spi.StreamCachingStrategy; +import org.apache.camel.spi.VariableRepository; +import org.apache.camel.support.service.ServiceSupport; +import org.apache.camel.util.StringHelper; + +/** + * Tenant-scoped {@link VariableRepository}, stored in-memory per tenant id. + *

+ * Sits between {@code group} and {@code global} in scope: a tenant can span + * multiple named groups, but variables set here are not visible outside the + * tenant, and are wider than any single group or route. + *

+ * Variable names use the syntax {@code tenant::}, e.g. + *

+ * exchange.setVariable("tenant:acme:myKey", someValue);
+ * Object val = exchange.getVariable("tenant:acme:myKey");
+ * 
+ * Modeled directly on Camel's own {@code RouteVariableRepository} / + * {@code GroupVariableRepository} so it behaves consistently with the + * built-in repositories (including participating in stream-caching). + */ +public final class TenantVariableRepository extends ServiceSupport implements BrowsableVariableRepository, CamelContextAware { + + /** + * Registry bean id under which this repository must be bound so Camel's + * {@code VariableRepositoryFactory} can discover it (lookup is by id == + * repository id, i.e. "tenant"). + */ + public static final String TENANT_VARIABLE_REPOSITORY_ID = "tenant"; + + private final Map> tenants = new ConcurrentHashMap<>(); + private CamelContext camelContext; + private StreamCachingStrategy strategy; + + @Override + public CamelContext getCamelContext() { + return camelContext; + } + + @Override + public void setCamelContext(CamelContext camelContext) { + this.camelContext = camelContext; + } + + @Override + public String getId() { + return TENANT_VARIABLE_REPOSITORY_ID; + } + + @Override + public Object getVariable(String name) { + String tenantId = StringHelper.before(name, ":"); + String key = StringHelper.after(name, ":"); + if (tenantId == null || key == null) { + throw new IllegalArgumentException("Name must be tenantId:name syntax"); + } + Object answer = null; + Map variables = tenants.get(tenantId); + if (variables != null) { + answer = variables.get(key); + } + if (answer instanceof StreamCache sc) { + // reset so the cache is ready to be used as a variable + sc.reset(); + } + return answer; + } + + @Override + public void setVariable(String name, Object value) { + String tenantId = StringHelper.before(name, ":"); + String key = StringHelper.after(name, ":"); + if (tenantId == null || key == null) { + throw new IllegalArgumentException("Name must be tenantId:name syntax"); + } + + if (value != null && strategy != null) { + StreamCache sc = convertToStreamCache(value); + if (sc != null) { + value = sc; + } + } + if (value != null) { + Map variables = tenants.computeIfAbsent(tenantId, s -> new ConcurrentHashMap<>(8)); + variables.put(key, value); + } else { + // if the value is null, we just remove the key from the map + Map variables = tenants.get(tenantId); + if (variables != null) { + variables.remove(key); + } + } + } + + @Override + public Object removeVariable(String name) { + String tenantId = StringHelper.before(name, ":"); + String key = StringHelper.after(name, ":"); + if (tenantId == null || key == null) { + throw new IllegalArgumentException("Name must be tenantId:name syntax"); + } + + Map variables = tenants.get(tenantId); + if (variables != null) { + if ("*".equals(key)) { + variables.clear(); + return null; + } else { + return variables.remove(key); + } + } + return null; + } + + /** + * Removes all variables for a given tenant, e.g. when a tenant is + * offboarded. Not part of {@link VariableRepository}; call directly. + */ + public void removeTenant(String tenantId) { + tenants.remove(tenantId); + } + + public boolean hasVariables() { + for (var vars : tenants.values()) { + if (!vars.isEmpty()) { + return true; + } + } + return false; + } + + public int size() { + int size = 0; + for (var vars : tenants.values()) { + size += vars.size(); + } + return size; + } + + public Stream names() { + List answer = new ArrayList<>(); + for (var tenantEntry : tenants.entrySet()) { + for (var e : tenantEntry.getValue().entrySet()) { + answer.add(tenantEntry.getKey() + ":" + e.getKey()); + } + } + return answer.stream(); + } + + public Map getVariables() { + Map answer = new ConcurrentHashMap<>(); + for (var tenantEntry : tenants.entrySet()) { + for (var e : tenantEntry.getValue().entrySet()) { + answer.put(tenantEntry.getKey() + ":" + e.getKey(), e.getValue()); + } + } + return answer; + } + + public void clear() { + tenants.clear(); + } + + @Override + protected void doInit() throws Exception { + super.doInit(); + if (camelContext != null && camelContext.isStreamCaching()) { + strategy = camelContext.getStreamCachingStrategy(); + } + } + + private StreamCache convertToStreamCache(Object body) { + if (body == null) { + return null; + } else if (body instanceof StreamCache sc) { + sc.reset(); + return sc; + } + return tryStreamCache(body); + } + + private StreamCache tryStreamCache(Object body) { + try { + return strategy.cache(body); + } catch (Exception e) { + throw new StreamCacheException(body, e); + } + } +} \ No newline at end of file diff --git a/dil/src/main/java/org/assimbly/dil/loader/FlowLoader.java b/dil/src/main/java/org/assimbly/dil/loader/FlowLoader.java index 6e0cdb03..e4d02427 100644 --- a/dil/src/main/java/org/assimbly/dil/loader/FlowLoader.java +++ b/dil/src/main/java/org/assimbly/dil/loader/FlowLoader.java @@ -32,6 +32,14 @@ public class FlowLoader extends RouteBuilder { private final FlowLoaderReport flowLoaderReport; private final EncryptionUtil encryptionUtil; + // Define the fixed metadata key constants + public static final String METADATA_TENANT_NAME = "MetaData.TenantName"; + public static final String METADATA_ENVIRONMENT_NAME = "MetaData.EnvironmentName"; + public static final String METADATA_FLOW_NAME = "MetaData.FlowName"; + public static final String METADATA_FLOW_ID = "MetaData.FlowID"; + public static final String METADATA_FLOW_VERSION = "MetaData.FlowVersion"; + + public FlowLoader(final TreeMap props, FlowLoaderReport flowLoaderReport, EncryptionUtil encryptionUtil){ super(); this.props = props; @@ -51,6 +59,8 @@ public void configure() throws Exception { setResources(); + setMetadata(); + setErrorHandlers(); setRouteConfigurations(); @@ -302,4 +312,20 @@ private String decryptStepIfNeeded(String step) { return result.toString(); } + public void setMetadata() { + + setVariable(METADATA_FLOW_ID, "id"); + setVariable(METADATA_FLOW_NAME, "flow.name"); + setVariable(METADATA_FLOW_VERSION, "flow.version"); + setVariable(METADATA_ENVIRONMENT_NAME, "flow.environment"); + setVariable(METADATA_TENANT_NAME, "flow.tenant"); + + } + + private void setVariable(String type, String property){ + if(props.containsKey(property)) { + context.setVariable("group:" + flowId + ":" + type, props.get(property)); + } + } + } \ No newline at end of file