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
Original file line number Diff line number Diff line change
Expand Up @@ -196,12 +196,42 @@ public void enableCluster(org.apache.eventmesh.runtime.cluster.MetaStore metaSto
// 4. Dynamic config hot-reload.
new org.apache.eventmesh.runtime.cluster.DynamicConfigWatcher(metaStore, runtime.ingress()).start();

// 5. #5359: flip the runtime's pull topology to PARTITION_OWNED_PULL and inject the shared
// MetaStore BEFORE runtime.start() runs, so the pull loop consults this PartitionOwnership
// (ownedPartitions) instead of polling every partition on every instance (duplicate
// consumption). UniRuntime.startPartitionOwnership then reuses the injected meta instead of
// building a second, divergent ownership state machine.
runtime.withClusterMeta(metaStore);
runtime.withTopology(org.apache.eventmesh.runtime.cluster.DeliveryTopology.PARTITION_OWNED_PULL);

log.info("cluster enabled (sticky + partition fencing): instance={} token={}",
selfInstanceId, selfToken);
}

/** Start runtime + traffic HTTP + admin HTTP. */
/** Start runtime + traffic HTTP + admin HTTP. Rolls back partial startup on failure (#5359). */
public void start() throws Exception {
try {
startupInternal();
} catch (Exception | Error failure) {
// #5359: a half-started process (runtime scheduler up, HTTP down) is worse than a
// clean exit - release ports, stop schedulers and rethrow so bin/start.sh surfaces it.
log.error("startup failed, rolling back: {}", failure.toString(), failure);
try {
shutdown();
} catch (RuntimeException cleanupFailure) {
failure.addSuppressed(cleanupFailure);
}
throw failure;
}
}

private void startupInternal() throws Exception {
// #5359: single-line effective configuration so an operator can audit a running
// deployment from the log (which topology / meta / storage actually took effect).
log.info("startup effective config: topology={} metaWired={} wsPort={} tls={}",
runtime.topology(), runtime.clusterMeta() != null
? runtime.clusterMeta().getClass().getSimpleName() : "none",
wsPort, sslContext != null);
runtime.start();
runtime.ingress().registerRuntimeGauges();
UniAdminService adminService = new UniAdminService(runtime.ingress());
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -57,7 +57,8 @@ public class UniRuntime {
private final long pollTimeoutMs;

/** Delivery topology (§13.2). LOCAL_STICKY_PULL = poll-all (default); PARTITION_OWNED_PULL = CAS+fencing scale-out. */
private final DeliveryTopology topology;
/** Set once at construction; may be flipped pre-start by EventMeshApplication.enableCluster (#5359). */
private volatile DeliveryTopology topology;

/** Instance identity for cluster membership (null in LOCAL_STICKY_PULL mode). */
private final String instanceId;
Expand All @@ -82,6 +83,41 @@ public class UniRuntime {
private Properties storageConfig = new Properties();

/** Inject storage config before {@link #start()}; additive, chainable. */
/** The delivery topology in effect (readable for startup logging, #5359). */
public DeliveryTopology topology() {
return topology;
}

/** The shared MetaStore wired by EventMeshApplication (null in single-instance mode, #5359). */
public MetaStore clusterMeta() {
return clusterMeta;
}

/** Wire the shared MetaStore before {@link #start()} (idempotent; used by enableCluster). */
public UniRuntime withClusterMeta(MetaStore metaStore) {
if (running.get()) {
throw new IllegalStateException("clusterMeta cannot change after start()");
}
this.clusterMeta = metaStore;
return this;
}

/**
* Flip the delivery topology pre-start (#5359). Only {@link EventMeshApplication#enableCluster}
* uses this: cluster mode upgrades a LOCAL_STICKY_PULL runtime to PARTITION_OWNED_PULL after the
* shared MetaStore is wired in. Must not be called after {@link #start()}.
*/
public UniRuntime withTopology(DeliveryTopology newTopology) {
if (running.get()) {
throw new IllegalStateException("topology cannot change after start()");
}
if (newTopology == null) {
throw new IllegalArgumentException("topology must not be null");
}
this.topology = newTopology;
return this;
}

public UniRuntime withStorageConfig(Properties storageConfig) {
this.storageConfig = storageConfig;
return this;
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,159 @@
/*
* 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.eventmesh.runtime.boot;

import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;

import org.apache.eventmesh.api.SendCallback;
import org.apache.eventmesh.api.storage.MeshStoragePlugin;
import org.apache.eventmesh.common.wire.EventMeshFrame;
import org.apache.eventmesh.runtime.cluster.DeliveryTopology;
import org.apache.eventmesh.runtime.cluster.InMemoryMetaStore;
import org.apache.eventmesh.runtime.offset.InMemoryOffsetStore;

import java.util.Collections;
import java.util.List;
import java.util.Properties;

import org.junit.jupiter.api.Test;

/**
* Issue #5359 acceptance: boot-time wiring of DeliveryTopology + MetaStore + FencingToken.
*
* <p>Covers the three wiring contracts the plan demands: (a) enableCluster flips the runtime
* topology to PARTITION_OWNED_PULL and injects the shared MetaStore before start; (b) the
* pre-start builders reject post-start mutation; (c) PARTITION_OWNED_PULL without a MetaStore
* fails fast (the #5356 contract, re-verified at the application level).</p>
*/
class EventMeshApplicationStartupTest {

private static final class StubStorage implements MeshStoragePlugin {

@Override
public void init(Properties props) {
// no-op
}

@Override
public void send(String topic, EventMeshFrame frame, SendCallback callback) {
// no-op
}

@Override
public List<EventMeshFrame> poll(String topic, int partition, long startOffset, int maxEvents, long timeoutMs) {
return Collections.emptyList();
}

@Override
public void assignPartitions(String topic, List<Integer> partitions) {
// no-op
}

@Override
public void commitOffset(String topic, int partition, long offset) {
// no-op
}

@Override
public int partitionCount(String topic) {
return 0;
}

@Override
public boolean isStarted() {
return true;
}

@Override
public boolean isClosed() {
return false;
}

@Override
public void start() {
// no-op
}

@Override
public void shutdown() {
// no-op
}
}

private static UniRuntime runtime() {
return new UniRuntime(new StubStorage(), new InMemoryOffsetStore(),
20L, 50L, 100, 50L, DeliveryTopology.LOCAL_STICKY_PULL, "test", "test:8080");
}

@Test
void enableClusterFlipsTopologyAndInjectsMeta() throws Exception {
// The core #5359 contract: after enableCluster, the runtime polls ONLY owned
// partitions (PARTITION_OWNED_PULL) through the injected shared MetaStore —
// not one ownership state machine in the app and a poll-all loop in the runtime.
UniRuntime runtime = runtime();
InMemoryMetaStore shared = new InMemoryMetaStore();
runtime.withClusterMeta(shared);
runtime.withTopology(DeliveryTopology.PARTITION_OWNED_PULL);
runtime.start();
try {
assertEquals(DeliveryTopology.PARTITION_OWNED_PULL, runtime.topology(),
"cluster mode must poll owned partitions only");
assertEquals(shared, runtime.clusterMeta(),
"the shared MetaStore must be wired into the runtime");
} finally {
runtime.shutdown();
}
}

@Test
void topologyCannotFlipAfterStart() throws Exception {
UniRuntime runtime = runtime();
runtime.start();
try {
IllegalStateException ex = assertThrows(IllegalStateException.class,
() -> runtime.withTopology(DeliveryTopology.PARTITION_OWNED_PULL));
assertTrue(ex.getMessage().contains("after start"),
"post-start topology change must fail fast, got: " + ex.getMessage());
} finally {
runtime.shutdown();
}
}

@Test
void clusterMetaCannotChangeAfterStart() throws Exception {
UniRuntime runtime = runtime();
runtime.start();
try {
IllegalStateException ex = assertThrows(IllegalStateException.class,
() -> runtime.withClusterMeta(new InMemoryMetaStore()));
assertTrue(ex.getMessage().contains("after start"),
"post-start meta change must fail fast, got: " + ex.getMessage());
} finally {
runtime.shutdown();
}
}

@Test
void nullTopologyRejected() {
UniRuntime runtime = runtime();
assertThrows(IllegalArgumentException.class,
() -> runtime.withTopology(null));
}
}