package com.xfestudio.xfeservermanager.core.spi;
import com.xfestudio.xfeservermanager.api.event.ManagementEvent;
import com.xfestudio.xfeservermanager.api.spi.ManagementEventSink;
import java.util.Collection;
import java.util.List;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.CompletionStage;
/** Fans one event out to all configured sinks without blocking its caller. */
public final class CompositeManagementEventSink implements ManagementEventSink {
private final List<ManagementEventSink> sinks;
public CompositeManagementEventSink(Collection<? extends ManagementEventSink> sinks) {
this.sinks = List.copyOf(sinks);
}
@Override
public CompletionStage<Void> publish(ManagementEvent event) {
CompletableFuture<?>[] pending = sinks.stream()
.map(sink -> publishSafely(sink, event))
.toArray(CompletableFuture[]::new);
return CompletableFuture.allOf(pending);
}
private static CompletableFuture<Void> publishSafely(ManagementEventSink sink, ManagementEvent event) {
try {
return sink.publish(event).toCompletableFuture();
} catch (RuntimeException exception) {
return CompletableFuture.failedFuture(exception);
}
}
}
package com.xfestudio.xfeservermanager.core.spi;
import com.xfestudio.xfeservermanager.api.event.ManagementEvent;
import com.xfestudio.xfeservermanager.api.spi.ManagementEventSink;
import java.util.Collection;
import java.util.List;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.CompletionStage;
/** Fans one event out to all configured sinks without blocking its caller. */
public final class CompositeManagementEventSink implements ManagementEventSink {
private final List<ManagementEventSink> sinks;
public CompositeManagementEventSink(Collection<? extends ManagementEventSink> sinks) {
this.sinks = List.copyOf(sinks);
}
@Override
public CompletionStage<Void> publish(ManagementEvent event) {
CompletableFuture<?>[] pending = sinks.stream()
.map(sink -> publishSafely(sink, event))
.toArray(CompletableFuture[]::new);
return CompletableFuture.allOf(pending);
}
private static CompletableFuture<Void> publishSafely(ManagementEventSink sink, ManagementEvent event) {
try {
return sink.publish(event).toCompletableFuture();
} catch (RuntimeException exception) {
return CompletableFuture.failedFuture(exception);
}
}
}