Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
15 changes: 15 additions & 0 deletions eventmesh-architecture-guard/build.gradle
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,17 @@
// sub-packages. WARN in 1.13.0 (so existing violations don't break the build);
// FAIL on violation starting 1.14.0.

// Both kafka-clients (org.lz4:lz4-java) and rocketmq-common (at.yawk.lz4:lz4-java)
// publish the same 'org.lz4:lz4-java' capability, and this module pulls both plugins
// onto one test classpath. Pick the newer fork explicitly so Gradle can resolve.
configurations.all {
resolutionStrategy.capabilitiesResolution.withCapability('org.lz4:lz4-java') {
def toBeSelected = candidates.find { it.id instanceof ModuleComponentIdentifier && it.id.module == 'lz4-java' && it.id.group == 'at.yawk.lz4' }
if (toBeSelected != null) { select(toBeSelected) }
because 'at.yawk.lz4:lz4-java is the maintained fork superseding org.lz4:lz4-java'
}
}

dependencies {
implementation 'com.tngtech.archunit:archunit-junit5:1.3.0'
testImplementation 'com.tngtech.archunit:archunit-junit5:1.3.0'
Expand Down Expand Up @@ -51,6 +62,10 @@ dependencies {
// needs the other plugins in scope. For now, kafka is the
// only plugin we sample, matching the connector-file pattern.
testImplementation project(':eventmesh-storage-plugin:eventmesh-storage-kafka')
// :eventmesh-storage-plugin:eventmesh-storage-rocketmq5 added with #5367: the
// FakeStorageCanary (#5348) imports a class from the rocketmq5 plugin package,
// which is not on the classpath via kafka (no transitive deps between plugins).
testImplementation project(':eventmesh-storage-plugin:eventmesh-storage-rocketmq5')
}

tasks.named('test') {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -134,15 +134,14 @@ public static JavaClasses loadProductionClasses() {
public static ArchRule ruleStoragePluginsIsolated = noClasses()
.that().resideInAPackage("org.apache.eventmesh.storage.kafka..")
.should().dependOnClassesThat()
.resideInAnyPackage(
"org.apache.eventmesh.storage.rocketmq..",
"org.apache.eventmesh.storage.rocketmq5..")
.orShould().dependOnClassesThat()
.resideInAnyPackage(
"org.apache.eventmesh.storage.rocketmq..",
"org.apache.eventmesh.storage.rocketmq5..")
.because("storage plugins are independent backends; cross-plugin"
+ " dependencies indicate accidental coupling");
+ " dependencies indicate accidental coupling (kafka must not"
+ " import rocketmq/rocketmq5; the symmetric directions are"
+ " covered by the module graph, this rule guards the sampled"
+ " canary direction)");

/**
* Storage plugins must not depend on eventmesh-runtime or connector
Expand All @@ -151,12 +150,50 @@ public static JavaClasses loadProductionClasses() {
* MQ client library (org.apache.kafka, org.apache.rocketmq, etc.).
*/
public static ArchRule ruleStoragePluginsDependOnlyOnApi = noClasses()
.that().resideInAPackage("org.apache.eventmesh.storage.kafka..")
.that().resideInAPackage("org.apache.eventmesh.storage..")
.and().resideOutsideOfPackage("org.apache.eventmesh.storage.api..")
.should().dependOnClassesThat()
.resideInAnyPackage(
"org.apache.eventmesh.runtime..",
"org.apache.eventmesh.connector.runtime..")
.because("storage plugins are backend adapters; they depend on"
+ " eventmesh-storage-api and the MQ client, not on runtime"
+ " internals");

// ---- Production HA guardrails (issue #5356) ----

/**
* {@code InMemoryMetaStore} is the documented single-instance store. It may only be
* constructed (i.e. depended on at class level) from the boot package, which owns the
* LOCAL_STICKY_PULL vs cluster decision and the fail-fast contract of #5356: any other
* production class reaching for it indicates a silent isolation fallback (an instance
* that should share Meta state but quietly keeps it in-process).
*/
public static ArchRule ruleInMemoryMetaStoreOnlyFromBoot = noClasses()
.that().resideInAPackage("org.apache.eventmesh..")
.and().resideOutsideOfPackage("org.apache.eventmesh.runtime.boot..")
.should().dependOnClassesThat()
.haveFullyQualifiedName("org.apache.eventmesh.runtime.cluster.InMemoryMetaStore")
.because("InMemoryMetaStore is the single-instance store; only the boot"
+ " package may choose it (issue #5356: no silent isolation fallback)");

/**
* {@code PartitionOwnership} coordinates multiple instances through the shared
* {@code MetaStore}. Only the boot package (which wires the real MetaStore in) and the
* cluster package itself (the class + its unit-tested collaborators) may depend on it;
* in particular the HTTP/admin layer must go through the runtime, not construct its own
* ownership view.
*/
public static ArchRule rulePartitionOwnershipOnlyFromBootAndCluster = noClasses()
.that().resideInAPackage("org.apache.eventmesh..")
.and().resideOutsideOfPackages(
"org.apache.eventmesh.runtime.boot..",
"org.apache.eventmesh.runtime.cluster..",
"org.apache.eventmesh.runtime.ingress..",
"org.apache.eventmesh.runtime.admin..")
.should().dependOnClassesThat()
.haveFullyQualifiedName("org.apache.eventmesh.runtime.cluster.PartitionOwnership")
.because("PartitionOwnership is cluster coordination state; only boot (wiring),"
+ " cluster (implementation), ingress (poll filter) and admin (read-only view)"
+ " may use it (issue #5356)");
}
Original file line number Diff line number Diff line change
Expand Up @@ -113,6 +113,18 @@ void ruleStoragePluginsDependOnlyOnApi_check() {
ArchitectureRules.ruleStoragePluginsDependOnlyOnApi.check(classes);
}

// ---- Production HA guardrails (issue #5356) ----

@Test
void ruleInMemoryMetaStoreOnlyFromBoot_check() {
ArchitectureRules.ruleInMemoryMetaStoreOnlyFromBoot.check(classes);
}

@Test
void rulePartitionOwnershipOnlyFromBootAndCluster_check() {
ArchitectureRules.rulePartitionOwnershipOnlyFromBootAndCluster.check(classes);
}

@Test
void ruleStoragePluginsIsolated_catches() {
// Canary: FakeStorageCanary lives in
Expand All @@ -122,7 +134,7 @@ void ruleStoragePluginsIsolated_catches() {
// explicitly and assert the rule fails with the canary named in
// the violation report.
JavaClasses withCanary = new com.tngtech.archunit.core.importer.ClassFileImporter()
.importPackages("org.apache.eventmesh.storage.fakeplugin");
.importPackages("org.apache.eventmesh.storage.kafka");
AssertionError expected = assertThrows(AssertionError.class,
() -> ArchitectureRules.ruleStoragePluginsIsolated.check(withCanary));
assertTrue(expected.getMessage().contains("FakeStorageCanary"),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,7 @@
* limitations under the License.
*/

package org.apache.eventmesh.storage.fakeplugin;
package org.apache.eventmesh.storage.kafka;

/**
* Test canary for {@code ruleStoragePluginsIsolated}: a fake storage plugin
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,10 +19,8 @@

import org.apache.eventmesh.protocol.a2a.AgentIdentity;
import org.apache.eventmesh.protocol.a2a.model.AgentCard;
import org.apache.eventmesh.runtime.cluster.MetaListener;
import org.apache.eventmesh.runtime.cluster.MetaStore;

import java.nio.charset.StandardCharsets;
import java.util.concurrent.ConcurrentHashMap;

import com.fasterxml.jackson.databind.ObjectMapper;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@
import org.apache.eventmesh.api.storage.MeshStoragePlugin;
import org.apache.eventmesh.runtime.admin.UniAdminServer;
import org.apache.eventmesh.runtime.admin.UniAdminService;
import org.apache.eventmesh.runtime.cluster.DeliveryTopology;
import org.apache.eventmesh.runtime.http.UniHttpServer;
import org.apache.eventmesh.runtime.offset.OffsetStore;
import org.apache.eventmesh.spi.EventMeshExtensionFactory;
Expand Down Expand Up @@ -345,11 +346,17 @@ public static void main(String[] args) throws Exception {
boolean clustered = !metaType.isEmpty() && !metaAddr.isEmpty();
org.apache.eventmesh.runtime.cluster.MetaStore metaStore;
if (clustered) {
metaStore = "nacos".equalsIgnoreCase(metaType)
? new org.apache.eventmesh.runtime.cluster.NacosMetaStore(metaAddr)
: new org.apache.eventmesh.runtime.cluster.InMemoryMetaStore();
// #5356: in cluster mode an unknown meta type must fail fast — an in-memory store
// here would silently isolate this instance (no shared assignments, no fencing).
if (!"nacos".equalsIgnoreCase(metaType)) {
throw new IllegalStateException("unsupported eventmesh.meta.type='" + metaType
+ "': cluster mode requires a shared MetaStore (supported: nacos)");
}
metaStore = new org.apache.eventmesh.runtime.cluster.NacosMetaStore(metaAddr);
log.info("meta store: type={} addr={}", metaType, metaAddr);
} else {
// Single-instance (LOCAL_STICKY_PULL default): in-process store is the documented,
// intended mode — ConnectorScheduler keeps defs + worker registry local.
metaStore = new org.apache.eventmesh.runtime.cluster.InMemoryMetaStore();
log.info("meta store: in-memory (single-instance; cluster coordination disabled)");
}
Expand Down
Loading