From 68c8d7fd0ed23b2a4705f58ceb0e6a9af7b075b9 Mon Sep 17 00:00:00 2001 From: "igor.petrenko" Date: Tue, 21 Jul 2026 09:13:16 +0300 Subject: [PATCH 01/23] CE-178 oap-notification ( oap-notification-mqtt ) --- .../oap-notification-client/pom.xml | 27 ++++ .../java/oap/notification/Notification.java | 23 ++++ .../oap/notification/NotificationService.java | 22 ++++ .../notification/NotificationTransport.java | 9 ++ .../src/main/java/oap/notification/Qos.java | 7 ++ .../oap-notification-mqtt/pom.xml | 38 ++++++ .../mqtt/HivemqNotificationTransport.java | 118 ++++++++++++++++++ .../oap-notification-test/pom.xml | 32 +++++ .../notification/mqtt/MosquittoFixture.java | 47 +++++++ .../notification/TestNotificationMessage.java | 18 +++ .../MosquittoNotificationServiceTest.java | 45 +++++++ .../test/resources/META-INF/oap-module.oap | 14 +++ oap-notification/pom.xml | 22 ++++ oap-stdlib/src/main/java/oap/json/Binder.java | 15 +++ pom.xml | 5 +- 15 files changed, 441 insertions(+), 1 deletion(-) create mode 100644 oap-notification/oap-notification-client/pom.xml create mode 100644 oap-notification/oap-notification-client/src/main/java/oap/notification/Notification.java create mode 100644 oap-notification/oap-notification-client/src/main/java/oap/notification/NotificationService.java create mode 100644 oap-notification/oap-notification-client/src/main/java/oap/notification/NotificationTransport.java create mode 100644 oap-notification/oap-notification-client/src/main/java/oap/notification/Qos.java create mode 100644 oap-notification/oap-notification-mqtt/pom.xml create mode 100644 oap-notification/oap-notification-mqtt/src/main/java/oap/notification/mqtt/HivemqNotificationTransport.java create mode 100644 oap-notification/oap-notification-test/pom.xml create mode 100644 oap-notification/oap-notification-test/src/main/java/oap/notification/mqtt/MosquittoFixture.java create mode 100644 oap-notification/oap-notification-test/src/test/java/oap/notification/TestNotificationMessage.java create mode 100644 oap-notification/oap-notification-test/src/test/java/oap/notification/mqtt/MosquittoNotificationServiceTest.java create mode 100644 oap-notification/oap-notification-test/src/test/resources/META-INF/oap-module.oap create mode 100644 oap-notification/pom.xml diff --git a/oap-notification/oap-notification-client/pom.xml b/oap-notification/oap-notification-client/pom.xml new file mode 100644 index 0000000000..2bc20ca37f --- /dev/null +++ b/oap-notification/oap-notification-client/pom.xml @@ -0,0 +1,27 @@ + + + 4.0.0 + + + oap + oap-notification + ${oap.project.version} + + + oap-notification-client + + + + oap + oap-stdlib + ${project.version} + + + + org.projectlombok + lombok + + + diff --git a/oap-notification/oap-notification-client/src/main/java/oap/notification/Notification.java b/oap-notification/oap-notification-client/src/main/java/oap/notification/Notification.java new file mode 100644 index 0000000000..7803ec865e --- /dev/null +++ b/oap-notification/oap-notification-client/src/main/java/oap/notification/Notification.java @@ -0,0 +1,23 @@ +package oap.notification; + +import com.fasterxml.jackson.annotation.JsonTypeInfo; +import com.fasterxml.jackson.databind.annotation.JsonTypeIdResolver; +import oap.json.TypeIdFactory; + +import java.io.Serial; +import java.io.Serializable; + +public class Notification implements Serializable { + @Serial + private static final long serialVersionUID = -1730908173571715179L; + + public final String sender; + @JsonTypeIdResolver( TypeIdFactory.class ) + @JsonTypeInfo( use = JsonTypeInfo.Id.CUSTOM, include = JsonTypeInfo.As.EXTERNAL_PROPERTY, property = "object:type" ) + public final Serializable message; + + public Notification( String sender, Serializable message ) { + this.sender = sender; + this.message = message; + } +} diff --git a/oap-notification/oap-notification-client/src/main/java/oap/notification/NotificationService.java b/oap-notification/oap-notification-client/src/main/java/oap/notification/NotificationService.java new file mode 100644 index 0000000000..5b1dab17b2 --- /dev/null +++ b/oap-notification/oap-notification-client/src/main/java/oap/notification/NotificationService.java @@ -0,0 +1,22 @@ +package oap.notification; + +import java.io.Serializable; +import java.util.function.Consumer; + +public class NotificationService { + public final String id; + public final NotificationTransport notificationTransport; + + public NotificationService( String id, NotificationTransport notificationTransport ) { + this.id = id; + this.notificationTransport = notificationTransport; + } + + public void sendNotification( String topic, Qos qos, TMessage message ) { + notificationTransport.publish( topic, qos, new Notification( id, message ) ); + } + + public void subscribeToTopic( String topic, Consumer notificationSupplier ) { + notificationTransport.subscribe( topic, notificationSupplier ); + } +} diff --git a/oap-notification/oap-notification-client/src/main/java/oap/notification/NotificationTransport.java b/oap-notification/oap-notification-client/src/main/java/oap/notification/NotificationTransport.java new file mode 100644 index 0000000000..d040e4dc31 --- /dev/null +++ b/oap-notification/oap-notification-client/src/main/java/oap/notification/NotificationTransport.java @@ -0,0 +1,9 @@ +package oap.notification; + +import java.util.function.Consumer; + +public interface NotificationTransport { + void publish( String topic, Qos qos, Notification notification ); + + void subscribe( String topic, Consumer notificationConsumer ); +} diff --git a/oap-notification/oap-notification-client/src/main/java/oap/notification/Qos.java b/oap-notification/oap-notification-client/src/main/java/oap/notification/Qos.java new file mode 100644 index 0000000000..e79d5fa1bc --- /dev/null +++ b/oap-notification/oap-notification-client/src/main/java/oap/notification/Qos.java @@ -0,0 +1,7 @@ +package oap.notification; + +public enum Qos { + AT_MOST_ONCE, + AT_LEAST_ONCE, + EXACTLY_ONCE +} diff --git a/oap-notification/oap-notification-mqtt/pom.xml b/oap-notification/oap-notification-mqtt/pom.xml new file mode 100644 index 0000000000..31417ba9aa --- /dev/null +++ b/oap-notification/oap-notification-mqtt/pom.xml @@ -0,0 +1,38 @@ + + + 4.0.0 + + + oap + oap-notification + ${oap.project.version} + + + oap-notification-mqtt + + + + com.hivemq + hivemq-mqtt-client + ${oap.deps.hivemq-mqtt-client.version} + + + + oap + oap-notification-client + ${parent.version} + + + oap + oap-application-annotation + ${parent.version} + + + + org.projectlombok + lombok + + + diff --git a/oap-notification/oap-notification-mqtt/src/main/java/oap/notification/mqtt/HivemqNotificationTransport.java b/oap-notification/oap-notification-mqtt/src/main/java/oap/notification/mqtt/HivemqNotificationTransport.java new file mode 100644 index 0000000000..a16b02230d --- /dev/null +++ b/oap-notification/oap-notification-mqtt/src/main/java/oap/notification/mqtt/HivemqNotificationTransport.java @@ -0,0 +1,118 @@ +package oap.notification.mqtt; + +import com.hivemq.client.mqtt.MqttClient; +import com.hivemq.client.mqtt.datatypes.MqttQos; +import com.hivemq.client.mqtt.mqtt5.Mqtt5AsyncClient; +import com.hivemq.client.mqtt.mqtt5.message.connect.connack.Mqtt5ConnAck; +import com.hivemq.client.mqtt.mqtt5.message.publish.Mqtt5Publish; +import com.hivemq.client.mqtt.mqtt5.message.publish.Mqtt5PublishResult; +import com.hivemq.client.mqtt.mqtt5.message.subscribe.suback.Mqtt5SubAck; +import lombok.extern.slf4j.Slf4j; +import oap.application.annotation.Start; +import oap.application.annotation.Stop; +import oap.json.Binder; +import oap.notification.Notification; +import oap.notification.NotificationTransport; +import oap.notification.Qos; +import oap.util.Dates; + +import java.util.concurrent.CompletionException; +import java.util.concurrent.TimeUnit; +import java.util.function.Consumer; + +@Slf4j +public class HivemqNotificationTransport implements NotificationTransport, AutoCloseable { + private final String identifier; + private final String host; + private final int port; + public long connectTimeout = Dates.s( 10 ); + public long publishTimeout = Dates.s( 1 ); + private Mqtt5AsyncClient client; + + public HivemqNotificationTransport( String identifier, String host, int port ) { + this.identifier = identifier; + this.host = host; + this.port = port; + } + + @Start + public void start() { + client = MqttClient + .builder() + .useMqttVersion5() + .identifier( identifier ) + .serverHost( host ) + .serverPort( port ) + + .automaticReconnect() + .applyAutomaticReconnect() + + .buildAsync(); + + Mqtt5ConnAck ack = client + .connectWith() + .send() + .orTimeout( connectTimeout, TimeUnit.MILLISECONDS ) + .join(); + + log.debug( "Connected to MQTT server at {}:{} response {}", host, port, ack ); + } + + @Stop + public void close() { + if( client != null && client.getState().isConnectedOrReconnect() ) { + try { + client + .disconnectWith() + .send() + .orTimeout( connectTimeout, TimeUnit.MILLISECONDS ) + .join(); + } catch( CompletionException e ) { + log.error( e.getCause().getMessage() ); + } + } + } + + @Override + public void publish( String topic, Qos qos, Notification notification ) { + log.trace( "publish topic {} qos {} notification {}", topic, qos, Binder.json.marshal( notification ) ); + + Mqtt5PublishResult result = client + .publishWith() + .topic( topic ) + .qos( convertQos( qos ) ) + .payload( Binder.json.marshal( notification ).getBytes() ) + .send() + .orTimeout( publishTimeout, TimeUnit.MILLISECONDS ) + .join(); + + log.trace( "publish topic {} qos {} result {}", topic, qos, result ); + } + + @Override + public void subscribe( String topic, Consumer notificationConsumer ) { + Mqtt5SubAck ack = client + .subscribeWith() + .topicFilter( topic ) + .callback( new Consumer() { + @Override + public void accept( Mqtt5Publish mqtt5Publish ) { + byte[] payloadAsBytes = mqtt5Publish.getPayloadAsBytes(); + notificationConsumer.accept( Binder.json.unmarshal( Notification.class, payloadAsBytes ) ); + } + } ) + .send() + .orTimeout( publishTimeout, TimeUnit.MILLISECONDS ) + .join(); + + log.trace( "publish topic {} result {}", topic, ack ); + } + + private MqttQos convertQos( Qos qos ) { + return switch( qos ) { + case AT_MOST_ONCE -> MqttQos.AT_MOST_ONCE; + case EXACTLY_ONCE -> MqttQos.EXACTLY_ONCE; + case AT_LEAST_ONCE -> MqttQos.AT_LEAST_ONCE; + }; + } +} diff --git a/oap-notification/oap-notification-test/pom.xml b/oap-notification/oap-notification-test/pom.xml new file mode 100644 index 0000000000..aa05b9d742 --- /dev/null +++ b/oap-notification/oap-notification-test/pom.xml @@ -0,0 +1,32 @@ + + + 4.0.0 + + + oap + oap-notification + ${oap.project.version} + + + oap-notification-test + + + + oap + oap-notification-mqtt + ${project.version} + + + oap + oap-stdlib-test + ${project.version} + + + + org.projectlombok + lombok + + + diff --git a/oap-notification/oap-notification-test/src/main/java/oap/notification/mqtt/MosquittoFixture.java b/oap-notification/oap-notification-test/src/main/java/oap/notification/mqtt/MosquittoFixture.java new file mode 100644 index 0000000000..5d6d867d06 --- /dev/null +++ b/oap-notification/oap-notification-test/src/main/java/oap/notification/mqtt/MosquittoFixture.java @@ -0,0 +1,47 @@ +package oap.notification.mqtt; + +import com.github.dockerjava.api.model.ExposedPort; +import com.github.dockerjava.api.model.PortBinding; +import com.github.dockerjava.api.model.Ports; +import lombok.Getter; +import lombok.extern.slf4j.Slf4j; +import oap.testng.AbstractFixture; +import org.testcontainers.containers.GenericContainer; +import org.testcontainers.containers.output.Slf4jLogConsumer; +import org.testcontainers.utility.DockerImageName; + +@Slf4j +public class MosquittoFixture extends AbstractFixture { + private static final String VERSION = "2.1.2-alpine"; + @Getter + private final int port; + private GenericContainer container; + + public MosquittoFixture() { + port = definePort( "MQTT_PORT" ); + } + + @Override + protected void before() { + super.before(); + + PortBinding portBinding = new PortBinding( + Ports.Binding.bindPort( port ), + new ExposedPort( 1883 ) ); + + container = new GenericContainer<>( DockerImageName.parse( "eclipse-mosquitto:" + VERSION ) ) + .withExposedPorts( 1883 ) + .withCreateContainerCmdModifier( cmd -> cmd.getHostConfig().withPortBindings( portBinding ) ) + .withLogConsumer( new Slf4jLogConsumer( log ) ); + container.start(); + } + + @Override + protected void after() { + if( container != null ) { + container.stop(); + } + + super.after(); + } +} diff --git a/oap-notification/oap-notification-test/src/test/java/oap/notification/TestNotificationMessage.java b/oap-notification/oap-notification-test/src/test/java/oap/notification/TestNotificationMessage.java new file mode 100644 index 0000000000..d23f4e5444 --- /dev/null +++ b/oap-notification/oap-notification-test/src/test/java/oap/notification/TestNotificationMessage.java @@ -0,0 +1,18 @@ +package oap.notification; + +import lombok.AllArgsConstructor; +import lombok.EqualsAndHashCode; +import lombok.ToString; + +import java.io.Serial; +import java.io.Serializable; + +@EqualsAndHashCode +@ToString +@AllArgsConstructor +public class TestNotificationMessage implements Serializable { + @Serial + private static final long serialVersionUID = -4166315788080834194L; + + public final String value; +} diff --git a/oap-notification/oap-notification-test/src/test/java/oap/notification/mqtt/MosquittoNotificationServiceTest.java b/oap-notification/oap-notification-test/src/test/java/oap/notification/mqtt/MosquittoNotificationServiceTest.java new file mode 100644 index 0000000000..f86905271c --- /dev/null +++ b/oap-notification/oap-notification-test/src/test/java/oap/notification/mqtt/MosquittoNotificationServiceTest.java @@ -0,0 +1,45 @@ +package oap.notification.mqtt; + +import oap.notification.NotificationService; +import oap.notification.Qos; +import oap.notification.TestNotificationMessage; +import oap.testng.Fixtures; +import org.testng.annotations.Test; + +import java.util.StringJoiner; + +import static org.assertj.core.api.Assertions.assertThat; + +public class MosquittoNotificationServiceTest extends Fixtures { + private final MosquittoFixture mosquittoFixture; + + public MosquittoNotificationServiceTest() { + mosquittoFixture = fixture( new MosquittoFixture() ); + } + + @Test + public void testMessages() { + StringJoiner msg = new StringJoiner( " / " ); + + try( HivemqNotificationTransport notificationTransportClient1 = new HivemqNotificationTransport( "client1", "127.0.0.1", mosquittoFixture.getPort() ); + HivemqNotificationTransport notificationTransportClient2 = new HivemqNotificationTransport( "client2", "127.0.0.1", mosquittoFixture.getPort() ) ) { + + notificationTransportClient1.start(); + notificationTransportClient2.start(); + + NotificationService notificationService1 = new NotificationService( "c1", notificationTransportClient1 ); + NotificationService notificationService2 = new NotificationService( "c2", notificationTransportClient2 ); + + notificationService1.sendNotification( "/test", Qos.AT_LEAST_ONCE, new TestNotificationMessage( "val1" ) ); + + notificationService2.subscribeToTopic( "/test", notification -> { + TestNotificationMessage notificationMessage = ( TestNotificationMessage ) notification.message; + msg.add( notificationMessage.value ); + } ); + + notificationService1.sendNotification( "/test", Qos.AT_LEAST_ONCE, new TestNotificationMessage( "val2" ) ); + + assertThat( msg ).hasToString( "val1 / val2" ); + } + } +} diff --git a/oap-notification/oap-notification-test/src/test/resources/META-INF/oap-module.oap b/oap-notification/oap-notification-test/src/test/resources/META-INF/oap-module.oap new file mode 100644 index 0000000000..5a5e1f0526 --- /dev/null +++ b/oap-notification/oap-notification-test/src/test/resources/META-INF/oap-module.oap @@ -0,0 +1,14 @@ +name = oap-notification-test + +services { + +} + +configurations = [ + { + loader = oap.json.TypeIdFactory + config { + oap-test-notification = oap.notification.TestNotificationMessage + } + } +] \ No newline at end of file diff --git a/oap-notification/pom.xml b/oap-notification/pom.xml new file mode 100644 index 0000000000..a7f46d032d --- /dev/null +++ b/oap-notification/pom.xml @@ -0,0 +1,22 @@ + + + 4.0.0 + + pom + + + oap + oap + ${oap.project.version} + + + oap-notification + + + oap-notification-client + oap-notification-mqtt + oap-notification-test + + diff --git a/oap-stdlib/src/main/java/oap/json/Binder.java b/oap-stdlib/src/main/java/oap/json/Binder.java index 4274099177..9b6c79704c 100644 --- a/oap-stdlib/src/main/java/oap/json/Binder.java +++ b/oap-stdlib/src/main/java/oap/json/Binder.java @@ -89,6 +89,7 @@ import java.util.Optional; import java.util.Set; +import static java.nio.charset.StandardCharsets.UTF_8; import static oap.io.IoStreams.DEFAULT_BUFFER; import static oap.io.IoStreams.Encoding.from; @@ -277,6 +278,11 @@ private static String getLimitation( String json ) { return json; } + private static String getLimitation( byte[] json ) { + if( json != null && json.length > 20 ) return new String( json, UTF_8 ).substring( 0, 20 ) + "..."; + return json != null ? new String( json, UTF_8 ) : null; + } + public ObjectMapper getMapper() { return mapper; } @@ -560,6 +566,15 @@ public T unmarshal( Class clazz, String json ) throws JsonException { } } + public T unmarshal( Class clazz, byte[] json ) throws JsonException { + try { + return mapper.readValue( json, clazz ); + } catch( Exception e ) { + log.trace( "Cannot deserialize [{}] into {}", json, clazz.getCanonicalName() ); + throw new JsonException( "Cannot deserialize [" + getLimitation( json ) + "] to class: " + clazz.getCanonicalName(), e ); + } + } + public T unmarshal( Class clazz, Map map ) throws JsonException { try { return mapper.convertValue( map, clazz ); diff --git a/pom.xml b/pom.xml index 8a4c07f637..49fbf74d6f 100644 --- a/pom.xml +++ b/pom.xml @@ -33,6 +33,7 @@ oap-mcp oap-jpath oap-message + oap-notification oap-application oap-formats oap-storage @@ -66,7 +67,7 @@ - 25.9.2 + 25.9.3 25.0.1 25.0.0 @@ -123,6 +124,8 @@ 3.6.4 12.1.6 + 1.3.14 + 4.9.8 From 65c7d5e6298e7dc8d14f97de6d7d90aedf6101aa Mon Sep 17 00:00:00 2001 From: "igor.petrenko" Date: Tue, 21 Jul 2026 10:07:03 +0300 Subject: [PATCH 02/23] CE-178 oap-notification ( oap-notification-mqtt ) --- .../oap/notification/mqtt/MosquittoNotificationServiceTest.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/oap-notification/oap-notification-test/src/test/java/oap/notification/mqtt/MosquittoNotificationServiceTest.java b/oap-notification/oap-notification-test/src/test/java/oap/notification/mqtt/MosquittoNotificationServiceTest.java index f86905271c..77c981440f 100644 --- a/oap-notification/oap-notification-test/src/test/java/oap/notification/mqtt/MosquittoNotificationServiceTest.java +++ b/oap-notification/oap-notification-test/src/test/java/oap/notification/mqtt/MosquittoNotificationServiceTest.java @@ -39,7 +39,7 @@ public void testMessages() { notificationService1.sendNotification( "/test", Qos.AT_LEAST_ONCE, new TestNotificationMessage( "val2" ) ); - assertThat( msg ).hasToString( "val1 / val2" ); + assertThat( msg ).hasToString( "val2" ); } } } From 82d805909616a5a1d4abe95a921c8811303688e5 Mon Sep 17 00:00:00 2001 From: "igor.petrenko" Date: Tue, 21 Jul 2026 10:24:23 +0300 Subject: [PATCH 03/23] CE-178 oap-notification ( oap-notification-mqtt ) --- .../java/oap/notification/Notification.java | 6 +++++ .../oap/notification/NotificationPublish.java | 26 +++++++++++++++++++ .../oap/notification/NotificationService.java | 9 +++++-- .../notification/NotificationTransport.java | 7 ++++- .../mqtt/HivemqNotificationTransport.java | 24 ++++++++++------- 5 files changed, 59 insertions(+), 13 deletions(-) create mode 100644 oap-notification/oap-notification-client/src/main/java/oap/notification/NotificationPublish.java diff --git a/oap-notification/oap-notification-client/src/main/java/oap/notification/Notification.java b/oap-notification/oap-notification-client/src/main/java/oap/notification/Notification.java index 7803ec865e..5fb0f1b70b 100644 --- a/oap-notification/oap-notification-client/src/main/java/oap/notification/Notification.java +++ b/oap-notification/oap-notification-client/src/main/java/oap/notification/Notification.java @@ -1,5 +1,6 @@ package oap.notification; +import com.fasterxml.jackson.annotation.JsonCreator; import com.fasterxml.jackson.annotation.JsonTypeInfo; import com.fasterxml.jackson.databind.annotation.JsonTypeIdResolver; import oap.json.TypeIdFactory; @@ -16,8 +17,13 @@ public class Notification implements Serializable { @JsonTypeInfo( use = JsonTypeInfo.Id.CUSTOM, include = JsonTypeInfo.As.EXTERNAL_PROPERTY, property = "object:type" ) public final Serializable message; + @JsonCreator public Notification( String sender, Serializable message ) { this.sender = sender; this.message = message; } + + public Notification( Notification notification ) { + this( notification.sender, notification.message ); + } } diff --git a/oap-notification/oap-notification-client/src/main/java/oap/notification/NotificationPublish.java b/oap-notification/oap-notification-client/src/main/java/oap/notification/NotificationPublish.java new file mode 100644 index 0000000000..f92abb6119 --- /dev/null +++ b/oap-notification/oap-notification-client/src/main/java/oap/notification/NotificationPublish.java @@ -0,0 +1,26 @@ +package oap.notification; + +import lombok.ToString; + +import java.io.Serial; +import java.io.Serializable; + +@ToString +public class NotificationPublish extends Notification { + @Serial + private static final long serialVersionUID = 8509736862218143643L; + + public final String topic; + + public NotificationPublish( String topic, Notification notification ) { + super( notification ); + + this.topic = topic; + } + + public NotificationPublish( String topic, String sender, Serializable message ) { + super( sender, message ); + + this.topic = topic; + } +} diff --git a/oap-notification/oap-notification-client/src/main/java/oap/notification/NotificationService.java b/oap-notification/oap-notification-client/src/main/java/oap/notification/NotificationService.java index 5b1dab17b2..d6042e1768 100644 --- a/oap-notification/oap-notification-client/src/main/java/oap/notification/NotificationService.java +++ b/oap-notification/oap-notification-client/src/main/java/oap/notification/NotificationService.java @@ -1,6 +1,7 @@ package oap.notification; import java.io.Serializable; +import java.util.List; import java.util.function.Consumer; public class NotificationService { @@ -16,7 +17,11 @@ public void sendNotification( String topic, Qos notificationTransport.publish( topic, qos, new Notification( id, message ) ); } - public void subscribeToTopic( String topic, Consumer notificationSupplier ) { - notificationTransport.subscribe( topic, notificationSupplier ); + public void subscribeToTopic( String topic, Consumer notificationConsumer ) { + notificationTransport.subscribe( topic, notificationConsumer ); + } + + public void subscribeToTopic( List topics, Consumer notificationConsumer ) { + notificationTransport.subscribe( topics, notificationConsumer ); } } diff --git a/oap-notification/oap-notification-client/src/main/java/oap/notification/NotificationTransport.java b/oap-notification/oap-notification-client/src/main/java/oap/notification/NotificationTransport.java index d040e4dc31..1617197cfd 100644 --- a/oap-notification/oap-notification-client/src/main/java/oap/notification/NotificationTransport.java +++ b/oap-notification/oap-notification-client/src/main/java/oap/notification/NotificationTransport.java @@ -1,9 +1,14 @@ package oap.notification; +import java.util.List; import java.util.function.Consumer; public interface NotificationTransport { void publish( String topic, Qos qos, Notification notification ); - void subscribe( String topic, Consumer notificationConsumer ); + default void subscribe( String topic, Consumer notificationConsumer ) { + subscribe( List.of( topic ), notificationConsumer ); + } + + void subscribe( List topics, Consumer notificationConsumer ); } diff --git a/oap-notification/oap-notification-mqtt/src/main/java/oap/notification/mqtt/HivemqNotificationTransport.java b/oap-notification/oap-notification-mqtt/src/main/java/oap/notification/mqtt/HivemqNotificationTransport.java index a16b02230d..c47510ab03 100644 --- a/oap-notification/oap-notification-mqtt/src/main/java/oap/notification/mqtt/HivemqNotificationTransport.java +++ b/oap-notification/oap-notification-mqtt/src/main/java/oap/notification/mqtt/HivemqNotificationTransport.java @@ -4,18 +4,21 @@ import com.hivemq.client.mqtt.datatypes.MqttQos; import com.hivemq.client.mqtt.mqtt5.Mqtt5AsyncClient; import com.hivemq.client.mqtt.mqtt5.message.connect.connack.Mqtt5ConnAck; -import com.hivemq.client.mqtt.mqtt5.message.publish.Mqtt5Publish; import com.hivemq.client.mqtt.mqtt5.message.publish.Mqtt5PublishResult; +import com.hivemq.client.mqtt.mqtt5.message.subscribe.Mqtt5Subscription; import com.hivemq.client.mqtt.mqtt5.message.subscribe.suback.Mqtt5SubAck; import lombok.extern.slf4j.Slf4j; import oap.application.annotation.Start; import oap.application.annotation.Stop; import oap.json.Binder; import oap.notification.Notification; +import oap.notification.NotificationPublish; import oap.notification.NotificationTransport; import oap.notification.Qos; import oap.util.Dates; +import oap.util.Lists; +import java.util.List; import java.util.concurrent.CompletionException; import java.util.concurrent.TimeUnit; import java.util.function.Consumer; @@ -90,22 +93,23 @@ public void publish( String topic, Qos qos, Notification notification ) { } @Override - public void subscribe( String topic, Consumer notificationConsumer ) { + public void subscribe( List topics, Consumer notificationConsumer ) { Mqtt5SubAck ack = client .subscribeWith() - .topicFilter( topic ) - .callback( new Consumer() { - @Override - public void accept( Mqtt5Publish mqtt5Publish ) { - byte[] payloadAsBytes = mqtt5Publish.getPayloadAsBytes(); - notificationConsumer.accept( Binder.json.unmarshal( Notification.class, payloadAsBytes ) ); - } + .addSubscriptions( Lists.map( topics, topic -> Mqtt5Subscription.builder().topicFilter( topic ).build() ) ) + .callback( mqtt5Publish -> { + byte[] payloadAsBytes = mqtt5Publish.getPayloadAsBytes(); + + log.trace( "topic {} payload {}", mqtt5Publish.getTopic(), + payloadAsBytes.length > 0 ? new String( payloadAsBytes ) : "" ); + + notificationConsumer.accept( new NotificationPublish( mqtt5Publish.getTopic().toString(), Binder.json.unmarshal( Notification.class, payloadAsBytes ) ) ); } ) .send() .orTimeout( publishTimeout, TimeUnit.MILLISECONDS ) .join(); - log.trace( "publish topic {} result {}", topic, ack ); + log.trace( "publish topics {} result {}", topics, ack ); } private MqttQos convertQos( Qos qos ) { From 1fb943b07c5c4155ed4fdd8827756f587998f38e Mon Sep 17 00:00:00 2001 From: "igor.petrenko" Date: Tue, 21 Jul 2026 10:40:46 +0300 Subject: [PATCH 04/23] CE-178 oap-notification ( oap-notification-mqtt ) --- .../MockNotificationTransport.java | 20 ++++++++++++++++++ .../src/resources/META-INF/oap-module.oap | 21 +++++++++++++++++++ .../main/resources/META-INF/oap-module.oap | 9 ++++++++ 3 files changed, 50 insertions(+) create mode 100644 oap-notification/oap-notification-client/src/main/java/oap/notification/MockNotificationTransport.java create mode 100644 oap-notification/oap-notification-client/src/resources/META-INF/oap-module.oap create mode 100644 oap-notification/oap-notification-mqtt/src/main/resources/META-INF/oap-module.oap diff --git a/oap-notification/oap-notification-client/src/main/java/oap/notification/MockNotificationTransport.java b/oap-notification/oap-notification-client/src/main/java/oap/notification/MockNotificationTransport.java new file mode 100644 index 0000000000..edf0c87eb1 --- /dev/null +++ b/oap-notification/oap-notification-client/src/main/java/oap/notification/MockNotificationTransport.java @@ -0,0 +1,20 @@ +package oap.notification; + +import lombok.extern.slf4j.Slf4j; +import oap.json.Binder; + +import java.util.List; +import java.util.function.Consumer; + +@Slf4j +public class MockNotificationTransport implements NotificationTransport { + @Override + public void publish( String topic, Qos qos, Notification notification ) { + log.trace( "publish topic {} qos {} notification {}", topic, qos, Binder.json.marshal( notification ) ); + } + + @Override + public void subscribe( List topics, Consumer notificationConsumer ) { + log.trace( "subscribe topics {}", topics ); + } +} diff --git a/oap-notification/oap-notification-client/src/resources/META-INF/oap-module.oap b/oap-notification/oap-notification-client/src/resources/META-INF/oap-module.oap new file mode 100644 index 0000000000..d742075a04 --- /dev/null +++ b/oap-notification/oap-notification-client/src/resources/META-INF/oap-module.oap @@ -0,0 +1,21 @@ +name = oap-notification + +services { + mock-notification-transport { + implementation = oap.notification.MockNotificationTransport + } + + notification-transport { + abstract = true + implementation = oap.notification.NotificationTransport + default = + } + + notification { + implementation = oap.notification.NotificationService + parameters { + id = "" + notificationTransport = + } + } +} diff --git a/oap-notification/oap-notification-mqtt/src/main/resources/META-INF/oap-module.oap b/oap-notification/oap-notification-mqtt/src/main/resources/META-INF/oap-module.oap new file mode 100644 index 0000000000..8cb6183126 --- /dev/null +++ b/oap-notification/oap-notification-mqtt/src/main/resources/META-INF/oap-module.oap @@ -0,0 +1,9 @@ +name = oap-notification-mqtt + +dependsOn = oap-notification + +services { + mqtt-notification-transport { + implementation = oap.notification.mqtt.HivemqNotificationTransport + } +} From 2f15d6764408e2e7e2eabfa697d0efde0443c52a Mon Sep 17 00:00:00 2001 From: "igor.petrenko" Date: Tue, 21 Jul 2026 10:55:32 +0300 Subject: [PATCH 05/23] CE-178 oap-notification ( oap-notification-mqtt ) --- .../src/main/resources/META-INF/oap-module.oap | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/oap-notification/oap-notification-mqtt/src/main/resources/META-INF/oap-module.oap b/oap-notification/oap-notification-mqtt/src/main/resources/META-INF/oap-module.oap index 8cb6183126..516fb8b100 100644 --- a/oap-notification/oap-notification-mqtt/src/main/resources/META-INF/oap-module.oap +++ b/oap-notification/oap-notification-mqtt/src/main/resources/META-INF/oap-module.oap @@ -5,5 +5,10 @@ dependsOn = oap-notification services { mqtt-notification-transport { implementation = oap.notification.mqtt.HivemqNotificationTransport + parameters { + identifier = OAP + host = "" + port = 1883 + } } } From de9487a9ef7bcaba6dd991df0bc052ee09fc653f Mon Sep 17 00:00:00 2001 From: "igor.petrenko" Date: Tue, 21 Jul 2026 11:25:17 +0300 Subject: [PATCH 06/23] CE-178 oap-notification ( oap-notification-mqtt ) --- .../src/{ => main}/resources/META-INF/oap-module.oap | 0 1 file changed, 0 insertions(+), 0 deletions(-) rename oap-notification/oap-notification-client/src/{ => main}/resources/META-INF/oap-module.oap (100%) diff --git a/oap-notification/oap-notification-client/src/resources/META-INF/oap-module.oap b/oap-notification/oap-notification-client/src/main/resources/META-INF/oap-module.oap similarity index 100% rename from oap-notification/oap-notification-client/src/resources/META-INF/oap-module.oap rename to oap-notification/oap-notification-client/src/main/resources/META-INF/oap-module.oap From 6684d3f9dc2ba2cfc4eea31a773171ca679effa9 Mon Sep 17 00:00:00 2001 From: "igor.petrenko" Date: Tue, 21 Jul 2026 11:36:40 +0300 Subject: [PATCH 07/23] CE-178 oap-notification ( oap-notification-mqtt ) --- .../src/main/resources/META-INF/oap-module.oap | 1 + 1 file changed, 1 insertion(+) diff --git a/oap-notification/oap-notification-mqtt/src/main/resources/META-INF/oap-module.oap b/oap-notification/oap-notification-mqtt/src/main/resources/META-INF/oap-module.oap index 516fb8b100..d0844afe57 100644 --- a/oap-notification/oap-notification-mqtt/src/main/resources/META-INF/oap-module.oap +++ b/oap-notification/oap-notification-mqtt/src/main/resources/META-INF/oap-module.oap @@ -10,5 +10,6 @@ services { host = "" port = 1883 } + supervision.supervise = true } } From e088e1cd728e2995a729f12d2b1daddc39e4f6b1 Mon Sep 17 00:00:00 2001 From: "igor.petrenko" Date: Tue, 21 Jul 2026 12:26:23 +0300 Subject: [PATCH 08/23] CE-178 oap-notification ( oap-notification-mqtt ) --- .../src/main/resources/META-INF/oap-module.oap | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/oap-notification/oap-notification-client/src/main/resources/META-INF/oap-module.oap b/oap-notification/oap-notification-client/src/main/resources/META-INF/oap-module.oap index d742075a04..ba3c3e599b 100644 --- a/oap-notification/oap-notification-client/src/main/resources/META-INF/oap-module.oap +++ b/oap-notification/oap-notification-client/src/main/resources/META-INF/oap-module.oap @@ -14,7 +14,7 @@ services { notification { implementation = oap.notification.NotificationService parameters { - id = "" +// id = "" notificationTransport = } } From 55c8ce0bc4e174aee5556a7653f8c022e44fb327 Mon Sep 17 00:00:00 2001 From: "igor.petrenko" Date: Tue, 21 Jul 2026 14:37:28 +0300 Subject: [PATCH 09/23] CE-178 oap-notification ( oap-notification-mqtt ) --- .../src/main/java/oap/application/Kernel.java | 6 +- .../main/java/oap/application/ModuleItem.java | 8 +-- .../application/supervision/Supervisor.java | 63 ++++++++++--------- 3 files changed, 39 insertions(+), 38 deletions(-) diff --git a/oap-application/oap-application/src/main/java/oap/application/Kernel.java b/oap-application/oap-application/src/main/java/oap/application/Kernel.java index 7252551488..ea788c5b2d 100644 --- a/oap-application/oap-application/src/main/java/oap/application/Kernel.java +++ b/oap-application/oap-application/src/main/java/oap/application/Kernel.java @@ -488,13 +488,13 @@ private void startService( Supervisor supervisor, ModuleItem.ServiceItem si ) { } if( service.supervision.thread ) { - supervisor.startThread( si.serviceName, instance, applicationConfiguration.shutdown ); + supervisor.startThread( si, instance, applicationConfiguration.shutdown ); } else { if( service.supervision.schedule && service.supervision.cron != null ) - supervisor.scheduleCron( si.serviceName, ( Runnable ) instance, + supervisor.scheduleCron( si, ( Runnable ) instance, service.supervision.cron ); else if( service.supervision.schedule && service.supervision.delay != 0 ) - supervisor.scheduleWithFixedDelay( si.serviceName, ( Runnable ) instance, + supervisor.scheduleWithFixedDelay( si, ( Runnable ) instance, service.supervision.delay, MILLISECONDS ); } } diff --git a/oap-application/oap-application/src/main/java/oap/application/ModuleItem.java b/oap-application/oap-application/src/main/java/oap/application/ModuleItem.java index 04e3c03f7f..f070d0e601 100644 --- a/oap-application/oap-application/src/main/java/oap/application/ModuleItem.java +++ b/oap-application/oap-application/src/main/java/oap/application/ModuleItem.java @@ -76,7 +76,7 @@ public boolean equals( Object o ) { if( this == o ) return true; if( o == null || getClass() != o.getClass() ) return false; - var that = ( ModuleItem ) o; + ModuleItem that = ( ModuleItem ) o; return module.name.equals( that.module.name ); } @@ -158,7 +158,7 @@ public boolean equals( Object o ) { if( this == o ) return true; if( o == null || getClass() != o.getClass() ) return false; - var that = ( ServiceItem ) o; + ServiceItem that = ( ServiceItem ) o; if( !moduleItem.module.name.equals( that.moduleItem.module.name ) ) return false; return serviceName.equals( that.serviceName ); @@ -172,7 +172,7 @@ public int hashCode() { } public void addDependsOn( ServiceReference serviceReference ) { - var found = Lists.find2( dependsOn, d -> d.equals( serviceReference ) ); + ServiceReference found = Lists.find2( dependsOn, d -> d.equals( serviceReference ) ); if( found == null || found.required ) { if( found != null ) dependsOn.remove( found ); dependsOn.add( serviceReference ); @@ -201,7 +201,7 @@ public boolean equals( Object o ) { if( this == o ) return true; if( o == null || getClass() != o.getClass() ) return false; - var that = ( ServiceReference ) o; + ServiceReference that = ( ServiceReference ) o; return serviceItem.serviceName.equals( that.serviceItem.serviceName ); } diff --git a/oap-application/oap-application/src/main/java/oap/application/supervision/Supervisor.java b/oap-application/oap-application/src/main/java/oap/application/supervision/Supervisor.java index 7ddff9e540..7ab3827d0e 100644 --- a/oap-application/oap-application/src/main/java/oap/application/supervision/Supervisor.java +++ b/oap-application/oap-application/src/main/java/oap/application/supervision/Supervisor.java @@ -26,6 +26,7 @@ import lombok.extern.slf4j.Slf4j; import oap.application.ApplicationConfiguration; import oap.application.KernelHelper; +import oap.application.ModuleItem; import oap.concurrent.Executors; import oap.util.BiStream; import oap.util.Dates; @@ -45,7 +46,7 @@ @Slf4j public class Supervisor { private final LinkedHashMap supervised = new LinkedHashMap<>(); - private final LinkedHashMap> wrappers = new LinkedHashMap<>(); + private final LinkedHashMap> wrappers = new LinkedHashMap<>(); private boolean stopped = false; @@ -92,16 +93,16 @@ public synchronized void startSupervised( String name, Object service, // this.wrappers.put( name, new ThreadService( name, ( Runnable ) instance, this ) ); // } - public synchronized void startThread( String name, Object instance, ApplicationConfiguration.ModuleShutdown shutdown ) { - this.wrappers.put( name, new ThreadService( name, ( Runnable ) instance, this, shutdown ) ); + public synchronized void startThread( ModuleItem.ServiceItem si, Object instance, ApplicationConfiguration.ModuleShutdown shutdown ) { + this.wrappers.put( si, new ThreadService( si.toString(), ( Runnable ) instance, this, shutdown ) ); } - public synchronized void scheduleWithFixedDelay( String name, Runnable service, long delay, TimeUnit unit ) { - this.wrappers.put( name, new DelayScheduledService( service, delay, unit ) ); + public synchronized void scheduleWithFixedDelay( ModuleItem.ServiceItem si, Runnable service, long delay, TimeUnit unit ) { + this.wrappers.put( si, new DelayScheduledService( service, delay, unit ) ); } - public synchronized void scheduleCron( String name, Runnable service, String cron ) { - this.wrappers.put( name, new CronScheduledService( service, cron ) ); + public synchronized void scheduleCron( ModuleItem.ServiceItem si, Runnable service, String cron ) { + this.wrappers.put( si, new CronScheduledService( service, cron ) ); } public synchronized void preStart() { @@ -119,15 +120,15 @@ public synchronized void preStart() { BiStream.of( this.wrappers ) .reversed() - .forEach( ( name, service ) -> { - log.debug( "[{}] pre starting {}...", service.type(), name ); - KernelHelper.setThreadNameSuffix( name ); + .forEach( ( si, service ) -> { + log.debug( "[{}] pre starting {}...", service.type(), si ); + KernelHelper.setThreadNameSuffix( si.toString() ); try { service.preStart(); } finally { KernelHelper.restoreThreadName(); } - log.debug( "[{}] pre starting {}... Done.", service.type(), name ); + log.debug( "[{}] pre starting {}... Done.", service.type(), si ); } ); } @@ -147,17 +148,17 @@ public synchronized void start() { log.debug( "starting {}... Done. ({}ms)", name, end - start ); } ); - this.wrappers.forEach( ( name, service ) -> { - log.debug( "[{}] starting {}...", service.type(), name ); + this.wrappers.forEach( ( si, service ) -> { + log.debug( "[{}] starting {}...", service.type(), si ); long start = System.currentTimeMillis(); - KernelHelper.setThreadNameSuffix( name ); + KernelHelper.setThreadNameSuffix( si.toString() ); try { service.start(); } finally { KernelHelper.restoreThreadName(); } long end = System.currentTimeMillis(); - log.debug( "[{}] starting {}... Done. ({}ms)", service.type(), name, end - start ); + log.debug( "[{}] starting {}... Done. ({}ms)", service.type(), si, end - start ); } ); } @@ -168,19 +169,19 @@ public synchronized void preStop( ApplicationConfiguration.ModuleShutdown shutdo BiStream.of( this.wrappers ) .reversed() - .forEach( ( name, service ) -> { + .forEach( ( si, service ) -> { Runnable func = () -> { - log.debug( "[{}] pre stopping {}...", service.type(), name ); - KernelHelper.setThreadNameSuffix( name ); + log.debug( "[{}] pre stopping {}...", service.type(), si ); + KernelHelper.setThreadNameSuffix( si.toString() ); try { service.preStop(); } finally { KernelHelper.restoreThreadName(); } - log.debug( "[{}] pre stopping {}... Done.", service.type(), name ); + log.debug( "[{}] pre stopping {}... Done.", service.type(), si ); }; - runAndDetectTimeout( name, shutdownConfiguration, func ); + runAndDetectTimeout( si.toString(), shutdownConfiguration, func ); } ); BiStream.of( this.supervised ) @@ -211,19 +212,19 @@ public synchronized void stop( ApplicationConfiguration.ModuleShutdown shutdown BiStream.of( this.wrappers ) .reversed() - .forEach( ( name, service ) -> { + .forEach( ( si, service ) -> { Runnable func = () -> { - log.debug( "[{}] stopping {}...", service.type(), name ); - KernelHelper.setThreadNameSuffix( name ); + log.debug( "[{}] stopping {}...", service.type(), si ); + KernelHelper.setThreadNameSuffix( si.toString() ); try { service.stop(); } finally { KernelHelper.restoreThreadName(); } - log.debug( "[{}] stopping {}... Done.", service.type(), name ); + log.debug( "[{}] stopping {}... Done.", service.type(), si ); }; - runAndDetectTimeout( name, shutdownConfiguration, func ); + runAndDetectTimeout( si.toString(), shutdownConfiguration, func ); } ); this.wrappers.clear(); @@ -255,21 +256,21 @@ public synchronized void stop( String serviceName, ApplicationConfiguration.Modu try( ShutdownConfiguration shutdownConfiguration = new ShutdownConfiguration( shutdown ) ) { BiStream.of( this.wrappers ) - .filter( ( name, _ ) -> name.equals( serviceName ) ) - .forEach( ( name, service ) -> { + .filter( ( si, _ ) -> si.serviceName.equals( serviceName ) ) + .forEach( ( si, service ) -> { Runnable func = () -> { - log.debug( "[{}] stopping {}...", service.type(), name ); - KernelHelper.setThreadNameSuffix( name ); + log.debug( "[{}] stopping {}...", service.type(), si ); + KernelHelper.setThreadNameSuffix( si.toString() ); try { service.preStop(); service.stop(); } finally { KernelHelper.restoreThreadName(); } - log.debug( "[{}] stopping {}... Done.", service.type(), name ); + log.debug( "[{}] stopping {}... Done.", service.type(), si ); }; - runAndDetectTimeout( name, shutdownConfiguration, func ); + runAndDetectTimeout( si.toString(), shutdownConfiguration, func ); } ); this.wrappers.clear(); From 0c8d244c71328b3df5d42d47536109ee90447df7 Mon Sep 17 00:00:00 2001 From: "igor.petrenko" Date: Tue, 21 Jul 2026 15:28:27 +0300 Subject: [PATCH 10/23] CE-178 oap-notification ( oap-notification-mqtt ) --- .../notification/NotificationException.java | 18 +++++++++++ .../oap/notification/NotificationService.java | 2 +- .../notification/NotificationTransport.java | 2 +- .../mqtt/HivemqNotificationTransport.java | 31 +++++++++++-------- 4 files changed, 38 insertions(+), 15 deletions(-) create mode 100644 oap-notification/oap-notification-client/src/main/java/oap/notification/NotificationException.java diff --git a/oap-notification/oap-notification-client/src/main/java/oap/notification/NotificationException.java b/oap-notification/oap-notification-client/src/main/java/oap/notification/NotificationException.java new file mode 100644 index 0000000000..84e3da9a95 --- /dev/null +++ b/oap-notification/oap-notification-client/src/main/java/oap/notification/NotificationException.java @@ -0,0 +1,18 @@ +package oap.notification; + +public class NotificationException extends RuntimeException { + public NotificationException() { + } + + public NotificationException( String message ) { + super( message ); + } + + public NotificationException( String message, Throwable cause ) { + super( message, cause ); + } + + public NotificationException( Throwable cause ) { + super( cause ); + } +} diff --git a/oap-notification/oap-notification-client/src/main/java/oap/notification/NotificationService.java b/oap-notification/oap-notification-client/src/main/java/oap/notification/NotificationService.java index d6042e1768..d971d42b2c 100644 --- a/oap-notification/oap-notification-client/src/main/java/oap/notification/NotificationService.java +++ b/oap-notification/oap-notification-client/src/main/java/oap/notification/NotificationService.java @@ -13,7 +13,7 @@ public NotificationService( String id, NotificationTransport notificationTranspo this.notificationTransport = notificationTransport; } - public void sendNotification( String topic, Qos qos, TMessage message ) { + public void sendNotification( String topic, Qos qos, TMessage message ) throws NotificationException { notificationTransport.publish( topic, qos, new Notification( id, message ) ); } diff --git a/oap-notification/oap-notification-client/src/main/java/oap/notification/NotificationTransport.java b/oap-notification/oap-notification-client/src/main/java/oap/notification/NotificationTransport.java index 1617197cfd..184714001b 100644 --- a/oap-notification/oap-notification-client/src/main/java/oap/notification/NotificationTransport.java +++ b/oap-notification/oap-notification-client/src/main/java/oap/notification/NotificationTransport.java @@ -4,7 +4,7 @@ import java.util.function.Consumer; public interface NotificationTransport { - void publish( String topic, Qos qos, Notification notification ); + void publish( String topic, Qos qos, Notification notification ) throws NotificationException; default void subscribe( String topic, Consumer notificationConsumer ) { subscribe( List.of( topic ), notificationConsumer ); diff --git a/oap-notification/oap-notification-mqtt/src/main/java/oap/notification/mqtt/HivemqNotificationTransport.java b/oap-notification/oap-notification-mqtt/src/main/java/oap/notification/mqtt/HivemqNotificationTransport.java index c47510ab03..d0d4038b7d 100644 --- a/oap-notification/oap-notification-mqtt/src/main/java/oap/notification/mqtt/HivemqNotificationTransport.java +++ b/oap-notification/oap-notification-mqtt/src/main/java/oap/notification/mqtt/HivemqNotificationTransport.java @@ -12,6 +12,7 @@ import oap.application.annotation.Stop; import oap.json.Binder; import oap.notification.Notification; +import oap.notification.NotificationException; import oap.notification.NotificationPublish; import oap.notification.NotificationTransport; import oap.notification.Qos; @@ -77,19 +78,23 @@ public void close() { } @Override - public void publish( String topic, Qos qos, Notification notification ) { - log.trace( "publish topic {} qos {} notification {}", topic, qos, Binder.json.marshal( notification ) ); - - Mqtt5PublishResult result = client - .publishWith() - .topic( topic ) - .qos( convertQos( qos ) ) - .payload( Binder.json.marshal( notification ).getBytes() ) - .send() - .orTimeout( publishTimeout, TimeUnit.MILLISECONDS ) - .join(); - - log.trace( "publish topic {} qos {} result {}", topic, qos, result ); + public void publish( String topic, Qos qos, Notification notification ) throws NotificationException { + try { + log.trace( "publish topic {} qos {} notification {}", topic, qos, Binder.json.marshal( notification ) ); + + Mqtt5PublishResult result = client + .publishWith() + .topic( topic ) + .qos( convertQos( qos ) ) + .payload( Binder.json.marshal( notification ).getBytes() ) + .send() + .orTimeout( publishTimeout, TimeUnit.MILLISECONDS ) + .join(); + + log.trace( "publish topic {} qos {} result {}", topic, qos, result ); + } catch( CompletionException e ) { + throw new NotificationException( e.getCause() ); + } } @Override From 1c2db8c1c67d884b6cdf65360b2c24ebfc56ec2f Mon Sep 17 00:00:00 2001 From: "igor.petrenko" Date: Tue, 21 Jul 2026 16:06:38 +0300 Subject: [PATCH 11/23] CE-178 oap-notification ( oap-notification-mqtt ) --- .../src/main/java/oap/application/Kernel.java | 2 +- .../application/supervision/Supervisor.java | 52 +++++++++---------- 2 files changed, 27 insertions(+), 27 deletions(-) diff --git a/oap-application/oap-application/src/main/java/oap/application/Kernel.java b/oap-application/oap-application/src/main/java/oap/application/Kernel.java index ea788c5b2d..46f41ca9e9 100644 --- a/oap-application/oap-application/src/main/java/oap/application/Kernel.java +++ b/oap-application/oap-application/src/main/java/oap/application/Kernel.java @@ -479,7 +479,7 @@ private void startService( Supervisor supervisor, ModuleItem.ServiceItem si ) { Service service = si.service; Object instance = si.instance; if( service.supervision.supervise ) { - supervisor.startSupervised( si.serviceName, instance, + supervisor.startSupervised( si, instance, service.supervision.preStartWith, service.supervision.startWith, service.supervision.preStopWith, diff --git a/oap-application/oap-application/src/main/java/oap/application/supervision/Supervisor.java b/oap-application/oap-application/src/main/java/oap/application/supervision/Supervisor.java index 7ab3827d0e..664ba1611a 100644 --- a/oap-application/oap-application/src/main/java/oap/application/supervision/Supervisor.java +++ b/oap-application/oap-application/src/main/java/oap/application/supervision/Supervisor.java @@ -45,7 +45,7 @@ @Slf4j public class Supervisor { - private final LinkedHashMap supervised = new LinkedHashMap<>(); + private final LinkedHashMap supervised = new LinkedHashMap<>(); private final LinkedHashMap> wrappers = new LinkedHashMap<>(); private boolean stopped = false; @@ -83,10 +83,10 @@ private static void runAndDetectTimeout( String name, ShutdownConfiguration shut } } - public synchronized void startSupervised( String name, Object service, + public synchronized void startSupervised( ModuleItem.ServiceItem si, Object service, List preStartWith, List startWith, List preStopWith, List stopWith ) { - this.supervised.put( name, new StartableService( service, preStartWith, startWith, preStopWith, stopWith ) ); + this.supervised.put( si, new StartableService( service, preStartWith, startWith, preStopWith, stopWith ) ); } // public synchronized void startScheduledThread( String name, Object instance, long delay, TimeUnit milliseconds ) { @@ -108,9 +108,9 @@ public synchronized void scheduleCron( ModuleItem.ServiceItem si, Runnable servi public synchronized void preStart() { log.debug( "pre starting..." ); - this.supervised.forEach( ( name, service ) -> { - log.debug( "pre starting {}...", name ); - KernelHelper.setThreadNameSuffix( name ); + this.supervised.forEach( ( si, service ) -> { + log.debug( "pre starting {}...", si ); + KernelHelper.setThreadNameSuffix( si.toString() ); try { service.preStart(); } finally { @@ -135,17 +135,17 @@ public synchronized void preStart() { public synchronized void start() { log.debug( "starting..." ); this.stopped = false; - this.supervised.forEach( ( name, service ) -> { - log.debug( "starting {}...", name ); + this.supervised.forEach( ( si, service ) -> { + log.debug( "starting {}...", si ); long start = System.currentTimeMillis(); - KernelHelper.setThreadNameSuffix( name ); + KernelHelper.setThreadNameSuffix( si.toString() ); try { service.start(); } finally { KernelHelper.restoreThreadName(); } long end = System.currentTimeMillis(); - log.debug( "starting {}... Done. ({}ms)", name, end - start ); + log.debug( "starting {}... Done. ({}ms)", si, end - start ); } ); this.wrappers.forEach( ( si, service ) -> { @@ -186,19 +186,19 @@ public synchronized void preStop( ApplicationConfiguration.ModuleShutdown shutdo BiStream.of( this.supervised ) .reversed() - .forEach( ( name, service ) -> { + .forEach( ( si, service ) -> { Runnable func = () -> { - log.debug( "pre stopping {}...", name ); - KernelHelper.setThreadNameSuffix( name ); + log.debug( "pre stopping {}...", si ); + KernelHelper.setThreadNameSuffix( si.toString() ); try { service.preStop(); } finally { KernelHelper.restoreThreadName(); } - log.debug( "pre stopping {}... Done.", name ); + log.debug( "pre stopping {}... Done.", si ); }; - runAndDetectTimeout( name, shutdownConfiguration, func ); + runAndDetectTimeout( si.toString(), shutdownConfiguration, func ); } ); } } @@ -230,19 +230,19 @@ public synchronized void stop( ApplicationConfiguration.ModuleShutdown shutdown BiStream.of( this.supervised ) .reversed() - .forEach( ( name, service ) -> { + .forEach( ( si, service ) -> { Runnable func = () -> { - log.debug( "stopping {}...", name ); - KernelHelper.setThreadNameSuffix( name ); + log.debug( "stopping {}...", si ); + KernelHelper.setThreadNameSuffix( si.toString() ); try { service.stop(); } finally { KernelHelper.restoreThreadName(); } - log.debug( "stopping {}... Done.", name ); + log.debug( "stopping {}... Done.", si ); }; - runAndDetectTimeout( name, shutdownConfiguration, func ); + runAndDetectTimeout( si.toString(), shutdownConfiguration, func ); } ); this.supervised.clear(); } @@ -275,21 +275,21 @@ public synchronized void stop( String serviceName, ApplicationConfiguration.Modu this.wrappers.clear(); BiStream.of( this.supervised ) - .filter( ( name, _ ) -> name.equals( serviceName ) ) - .forEach( ( name, service ) -> { + .filter( ( si, _ ) -> si.serviceName.equals( serviceName ) ) + .forEach( ( si, service ) -> { Runnable func = () -> { - log.debug( "stopping {}...", name ); - KernelHelper.setThreadNameSuffix( name ); + log.debug( "stopping {}...", si ); + KernelHelper.setThreadNameSuffix( si.toString() ); try { service.preStop(); service.stop(); } finally { KernelHelper.restoreThreadName(); } - log.debug( "stopping {}... Done.", name ); + log.debug( "stopping {}... Done.", si ); }; - runAndDetectTimeout( name, shutdownConfiguration, func ); + runAndDetectTimeout( si.toString(), shutdownConfiguration, func ); } ); } } From f35b8f2eaab6f64e74216c58efd8a309ba2c7c25 Mon Sep 17 00:00:00 2001 From: "igor.petrenko" Date: Tue, 21 Jul 2026 16:37:40 +0300 Subject: [PATCH 12/23] CE-178 oap-notification ( oap-notification-mqtt ) --- .../src/main/java/oap/notification/Notification.java | 6 ++---- .../src/main/java/oap/notification/NotificationPublish.java | 4 ++-- .../src/main/java/oap/notification/NotificationService.java | 6 ++---- .../src/main/resources/META-INF/oap-module.oap | 1 - .../notification/mqtt/MosquittoNotificationServiceTest.java | 4 ++-- 5 files changed, 8 insertions(+), 13 deletions(-) diff --git a/oap-notification/oap-notification-client/src/main/java/oap/notification/Notification.java b/oap-notification/oap-notification-client/src/main/java/oap/notification/Notification.java index 5fb0f1b70b..31aeb0046a 100644 --- a/oap-notification/oap-notification-client/src/main/java/oap/notification/Notification.java +++ b/oap-notification/oap-notification-client/src/main/java/oap/notification/Notification.java @@ -12,18 +12,16 @@ public class Notification implements Serializable { @Serial private static final long serialVersionUID = -1730908173571715179L; - public final String sender; @JsonTypeIdResolver( TypeIdFactory.class ) @JsonTypeInfo( use = JsonTypeInfo.Id.CUSTOM, include = JsonTypeInfo.As.EXTERNAL_PROPERTY, property = "object:type" ) public final Serializable message; @JsonCreator - public Notification( String sender, Serializable message ) { - this.sender = sender; + public Notification( Serializable message ) { this.message = message; } public Notification( Notification notification ) { - this( notification.sender, notification.message ); + this( notification.message ); } } diff --git a/oap-notification/oap-notification-client/src/main/java/oap/notification/NotificationPublish.java b/oap-notification/oap-notification-client/src/main/java/oap/notification/NotificationPublish.java index f92abb6119..47a14212cc 100644 --- a/oap-notification/oap-notification-client/src/main/java/oap/notification/NotificationPublish.java +++ b/oap-notification/oap-notification-client/src/main/java/oap/notification/NotificationPublish.java @@ -18,8 +18,8 @@ public NotificationPublish( String topic, Notification notification ) { this.topic = topic; } - public NotificationPublish( String topic, String sender, Serializable message ) { - super( sender, message ); + public NotificationPublish( String topic, Serializable message ) { + super( message ); this.topic = topic; } diff --git a/oap-notification/oap-notification-client/src/main/java/oap/notification/NotificationService.java b/oap-notification/oap-notification-client/src/main/java/oap/notification/NotificationService.java index d971d42b2c..16d1eff012 100644 --- a/oap-notification/oap-notification-client/src/main/java/oap/notification/NotificationService.java +++ b/oap-notification/oap-notification-client/src/main/java/oap/notification/NotificationService.java @@ -5,16 +5,14 @@ import java.util.function.Consumer; public class NotificationService { - public final String id; public final NotificationTransport notificationTransport; - public NotificationService( String id, NotificationTransport notificationTransport ) { - this.id = id; + public NotificationService( NotificationTransport notificationTransport ) { this.notificationTransport = notificationTransport; } public void sendNotification( String topic, Qos qos, TMessage message ) throws NotificationException { - notificationTransport.publish( topic, qos, new Notification( id, message ) ); + notificationTransport.publish( topic, qos, new Notification( message ) ); } public void subscribeToTopic( String topic, Consumer notificationConsumer ) { diff --git a/oap-notification/oap-notification-client/src/main/resources/META-INF/oap-module.oap b/oap-notification/oap-notification-client/src/main/resources/META-INF/oap-module.oap index ba3c3e599b..534190fdd9 100644 --- a/oap-notification/oap-notification-client/src/main/resources/META-INF/oap-module.oap +++ b/oap-notification/oap-notification-client/src/main/resources/META-INF/oap-module.oap @@ -14,7 +14,6 @@ services { notification { implementation = oap.notification.NotificationService parameters { -// id = "" notificationTransport = } } diff --git a/oap-notification/oap-notification-test/src/test/java/oap/notification/mqtt/MosquittoNotificationServiceTest.java b/oap-notification/oap-notification-test/src/test/java/oap/notification/mqtt/MosquittoNotificationServiceTest.java index 77c981440f..ff611ace6f 100644 --- a/oap-notification/oap-notification-test/src/test/java/oap/notification/mqtt/MosquittoNotificationServiceTest.java +++ b/oap-notification/oap-notification-test/src/test/java/oap/notification/mqtt/MosquittoNotificationServiceTest.java @@ -27,8 +27,8 @@ public void testMessages() { notificationTransportClient1.start(); notificationTransportClient2.start(); - NotificationService notificationService1 = new NotificationService( "c1", notificationTransportClient1 ); - NotificationService notificationService2 = new NotificationService( "c2", notificationTransportClient2 ); + NotificationService notificationService1 = new NotificationService( notificationTransportClient1 ); + NotificationService notificationService2 = new NotificationService( notificationTransportClient2 ); notificationService1.sendNotification( "/test", Qos.AT_LEAST_ONCE, new TestNotificationMessage( "val1" ) ); From c38b986b12fba5c4a9af2cdb9d8ee5371be5083c Mon Sep 17 00:00:00 2001 From: "igor.petrenko" Date: Tue, 21 Jul 2026 17:06:54 +0300 Subject: [PATCH 13/23] CE-178 oap-notification ( oap-notification-mqtt ) --- .../src/main/resources/META-INF/oap-module.oap | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/oap-notification/oap-notification-mqtt/src/main/resources/META-INF/oap-module.oap b/oap-notification/oap-notification-mqtt/src/main/resources/META-INF/oap-module.oap index d0844afe57..d86aba6dbc 100644 --- a/oap-notification/oap-notification-mqtt/src/main/resources/META-INF/oap-module.oap +++ b/oap-notification/oap-notification-mqtt/src/main/resources/META-INF/oap-module.oap @@ -6,7 +6,7 @@ services { mqtt-notification-transport { implementation = oap.notification.mqtt.HivemqNotificationTransport parameters { - identifier = OAP +// identifier = OAP host = "" port = 1883 } From e2391d5b74356b5850051ddd2d7b241bfdb702bf Mon Sep 17 00:00:00 2001 From: "igor.petrenko" Date: Tue, 21 Jul 2026 18:22:47 +0300 Subject: [PATCH 14/23] CE-178 oap-notification ( oap-notification-mqtt ) --- .../src/main/java/oap/testng/Asserts.java | 19 ++++++++++++++++--- 1 file changed, 16 insertions(+), 3 deletions(-) diff --git a/oap-stdlib-test/src/main/java/oap/testng/Asserts.java b/oap-stdlib-test/src/main/java/oap/testng/Asserts.java index ac05aa156e..81146d6044 100644 --- a/oap-stdlib-test/src/main/java/oap/testng/Asserts.java +++ b/oap-stdlib-test/src/main/java/oap/testng/Asserts.java @@ -42,6 +42,7 @@ import org.testng.Assert; import javax.annotation.Nonnull; +import javax.annotation.Nullable; import java.lang.reflect.InvocationTargetException; import java.lang.reflect.Method; import java.net.URL; @@ -58,7 +59,7 @@ public final class Asserts { @SneakyThrows - public static void eventually( long retryTimeout, int retries, Try.ThrowingRunnable asserts ) { + public static void eventually( long retryTimeout, int retries, Try.ThrowingRunnable asserts, @Nullable Runnable onFailure ) { boolean passed = false; Throwable exception = null; int count = retries; @@ -71,6 +72,10 @@ public static void eventually( long retryTimeout, int retries, Try.ThrowingRunna exception = e; Threads.sleepSafely( retryTimeout ); count--; + + if( onFailure != null ) { + onFailure.run(); + } } } if( !passed ) @@ -79,12 +84,20 @@ public static void eventually( long retryTimeout, int retries, Try.ThrowingRunna } public static void assertEventually( long retryTimeout, int retries, oap.util.function.Try.ThrowingRunnable asserts ) { - eventually( retryTimeout, retries, asserts ); + assertEventually( retryTimeout, retries, asserts, null ); + } + + public static void assertEventually( long retryTimeout, int retries, oap.util.function.Try.ThrowingRunnable asserts, @Nullable Runnable onFailure ) { + eventually( retryTimeout, retries, asserts, onFailure ); } public static void assertEventually( Duration duration, Duration retryInterval, oap.util.function.Try.ThrowingRunnable asserts ) { + assertEventually( duration, retryInterval, asserts, null ); + } + + public static void assertEventually( Duration duration, Duration retryInterval, oap.util.function.Try.ThrowingRunnable asserts, @Nullable Runnable onFailure ) { int retries = ( int ) ( duration.toMillis() / retryInterval.toMillis() ); - eventually( retryInterval.toMillis(), retries, asserts ); + eventually( retryInterval.toMillis(), retries, asserts, onFailure ); } @Deprecated From 0420787e2a421d272e1c06b5b88bd5500ae3b70d Mon Sep 17 00:00:00 2001 From: "igor.petrenko" Date: Tue, 21 Jul 2026 20:42:21 +0300 Subject: [PATCH 15/23] CE-178 oap-notification ( oap-notification-mqtt ) --- .../src/main/java/oap/logstream/disk/AbstractWriter.java | 5 ++++- .../main/java/oap/logstream/disk/DiskLoggerBackend.java | 8 ++++++-- .../java/oap/logstream/disk/FileWriterNotification.java | 7 +++++++ .../src/main/java/oap/logstream/disk/RowBinaryWriter.java | 5 +++-- 4 files changed, 20 insertions(+), 5 deletions(-) create mode 100644 oap-formats/oap-logstream/oap-logstream/src/main/java/oap/logstream/disk/FileWriterNotification.java diff --git a/oap-formats/oap-logstream/oap-logstream/src/main/java/oap/logstream/disk/AbstractWriter.java b/oap-formats/oap-logstream/oap-logstream/src/main/java/oap/logstream/disk/AbstractWriter.java index 65abfd06cd..f20cb14ba4 100644 --- a/oap-formats/oap-logstream/oap-logstream/src/main/java/oap/logstream/disk/AbstractWriter.java +++ b/oap-formats/oap-logstream/oap-logstream/src/main/java/oap/logstream/disk/AbstractWriter.java @@ -55,6 +55,7 @@ public abstract class AbstractWriter implements Closeable { protected final Stopwatch stopwatch = new Stopwatch(); protected final int maxVersions; protected final String hostname; + protected final FileWriterNotification notification; protected final ReentrantLock lock = new ReentrantLock(); protected LogFile logFile; protected String lastPattern; @@ -62,13 +63,14 @@ public abstract class AbstractWriter implements Closeable { protected boolean closed = false; protected AbstractWriter( TemplateEngine templateEngine, LogFormat logFormat, Path logDirectory, String filePattern, LogId logId, int bufferSize, Timestamp timestamp, - int maxVersions, String hostname ) { + int maxVersions, String hostname, FileWriterNotification notification ) { this.templateEngine = templateEngine; this.logFormat = logFormat; this.logDirectory = logDirectory; this.filePattern = filePattern; this.maxVersions = maxVersions; this.hostname = hostname; + this.notification = notification; log.trace( "filePattern {}", filePattern ); Preconditions.checkArgument( filePattern.matches( ".*[${]\\{\\s*LOG_VERSION\\s*}}?.*" ), "file pattern must contains LOG_VERSION variable" ); @@ -161,6 +163,7 @@ protected void closeOutput() throws LoggerException { Metrics.summary( "logstream_logging_server_bucket_time_seconds" ).record( Dates.nanosToSeconds( stopwatch.elapsed() ) ); logFile.readyForUpload(); + notification.fileClosed( logFile.outFilename ); } finally { logFile = null; } diff --git a/oap-formats/oap-logstream/oap-logstream/src/main/java/oap/logstream/disk/DiskLoggerBackend.java b/oap-formats/oap-logstream/oap-logstream/src/main/java/oap/logstream/disk/DiskLoggerBackend.java index ce50188250..68196ad98c 100644 --- a/oap-formats/oap-logstream/oap-logstream/src/main/java/oap/logstream/disk/DiskLoggerBackend.java +++ b/oap-formats/oap-logstream/oap-logstream/src/main/java/oap/logstream/disk/DiskLoggerBackend.java @@ -75,7 +75,7 @@ *
  • POD_NAME
  • */ @Slf4j -public class DiskLoggerBackend extends AbstractLoggerBackend implements Cloneable, AutoCloseable { +public class DiskLoggerBackend extends AbstractLoggerBackend implements FileWriterNotification, Cloneable, AutoCloseable { public static final int DEFAULT_BUFFER = 1024 * 100; public static final long DEFAULT_FREE_SPACE_REQUIRED = 2000000000L; public final LinkedHashMap filePatternByType = new LinkedHashMap<>(); @@ -126,7 +126,7 @@ public AbstractWriter load( LogId id ) { log.trace( "new writer id '{}' filePattern '{}'", id, fp ); - return new RowBinaryWriter( templateEngine, DiskLoggerBackend.this.logDirectory, fp.path, id, bufferSize, timestamp, maxVersions, hostname ); + return new RowBinaryWriter( templateEngine, DiskLoggerBackend.this.logDirectory, fp.path, id, bufferSize, timestamp, maxVersions, hostname, DiskLoggerBackend.this ); } } ); Metrics.gauge( "logstream_logging_disk_writers", List.of( Tag.of( "path", this.logDirectory.toString() ) ), @@ -243,6 +243,10 @@ public String toString() { .toString(); } + @Override + public void fileClosed( Path outFilename ) { + } + @ToString @EqualsAndHashCode public static class FilePatternConfiguration { diff --git a/oap-formats/oap-logstream/oap-logstream/src/main/java/oap/logstream/disk/FileWriterNotification.java b/oap-formats/oap-logstream/oap-logstream/src/main/java/oap/logstream/disk/FileWriterNotification.java new file mode 100644 index 0000000000..eb245f114c --- /dev/null +++ b/oap-formats/oap-logstream/oap-logstream/src/main/java/oap/logstream/disk/FileWriterNotification.java @@ -0,0 +1,7 @@ +package oap.logstream.disk; + +import java.nio.file.Path; + +public interface FileWriterNotification { + void fileClosed( Path outFilename ); +} diff --git a/oap-formats/oap-logstream/oap-logstream/src/main/java/oap/logstream/disk/RowBinaryWriter.java b/oap-formats/oap-logstream/oap-logstream/src/main/java/oap/logstream/disk/RowBinaryWriter.java index 9aeda25d3c..0bf8336964 100644 --- a/oap-formats/oap-logstream/oap-logstream/src/main/java/oap/logstream/disk/RowBinaryWriter.java +++ b/oap-formats/oap-logstream/oap-logstream/src/main/java/oap/logstream/disk/RowBinaryWriter.java @@ -18,8 +18,9 @@ @Slf4j public class RowBinaryWriter extends AbstractWriter { - public RowBinaryWriter( TemplateEngine templateEngine, Path logDirectory, String filePattern, LogId logId, int bufferSize, Timestamp timestamp, int maxVersions, String hostname ) { - super( templateEngine, LogFormat.ROW_BINARY_GZ, logDirectory, filePattern, logId, bufferSize, timestamp, maxVersions, hostname ); + public RowBinaryWriter( TemplateEngine templateEngine, Path logDirectory, String filePattern, LogId logId, + int bufferSize, Timestamp timestamp, int maxVersions, String hostname, FileWriterNotification notification ) { + super( templateEngine, LogFormat.ROW_BINARY_GZ, logDirectory, filePattern, logId, bufferSize, timestamp, maxVersions, hostname, notification ); } @Override From d36f296e30c650cbb55755cf0683188d2f2981d9 Mon Sep 17 00:00:00 2001 From: "igor.petrenko" Date: Tue, 21 Jul 2026 20:45:42 +0300 Subject: [PATCH 16/23] CE-178 oap-notification ( oap-notification-mqtt ) --- .../main/java/oap/logstream/disk/DiskLoggerBackend.java | 8 ++++++-- 1 file changed, 6 insertions(+), 2 deletions(-) diff --git a/oap-formats/oap-logstream/oap-logstream/src/main/java/oap/logstream/disk/DiskLoggerBackend.java b/oap-formats/oap-logstream/oap-logstream/src/main/java/oap/logstream/disk/DiskLoggerBackend.java index 68196ad98c..9c8585511a 100644 --- a/oap-formats/oap-logstream/oap-logstream/src/main/java/oap/logstream/disk/DiskLoggerBackend.java +++ b/oap-formats/oap-logstream/oap-logstream/src/main/java/oap/logstream/disk/DiskLoggerBackend.java @@ -93,6 +93,7 @@ public class DiskLoggerBackend extends AbstractLoggerBackend implements FileWrit public long refreshInitDelay = Dates.s( 10 ); public long refreshPeriod = Dates.s( 10 ); public volatile boolean closed; + protected ArrayList notifications = new ArrayList<>(); public DiskLoggerBackend( TemplateEngine templateEngine, Path logDirectory, Timestamp timestamp, int bufferSize, String hostname ) { this( templateEngine, logDirectory, new WriterConfiguration(), timestamp, bufferSize, hostname ); @@ -116,8 +117,8 @@ public DiskLoggerBackend( TemplateEngine templateEngine, Path logDirectory, Writ this.writers = CacheBuilder.newBuilder() .ticker( JodaTicker.JODA_TICKER ) .expireAfterAccess( 60 / timestamp.bucketsPerHour * 3, TimeUnit.MINUTES ) - .removalListener( notification -> { - Closeables.close( ( Closeable ) notification.getValue() ); + .removalListener( l -> { + Closeables.close( ( Closeable ) l.getValue() ); } ) .build( new CacheLoader<>() { @Override @@ -245,6 +246,9 @@ public String toString() { @Override public void fileClosed( Path outFilename ) { + for( FileWriterNotification notification : notifications ) { + notification.fileClosed( outFilename ); + } } @ToString From befb804f18ee37643c113cebdb49687625aa1b2d Mon Sep 17 00:00:00 2001 From: "igor.petrenko" Date: Tue, 21 Jul 2026 20:54:21 +0300 Subject: [PATCH 17/23] CE-178 oap-notification ( oap-notification-mqtt ) --- .../test/java/oap/logstream/disk/RowBinaryWriterTest.java | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/oap-formats/oap-logstream/oap-logstream-test/src/test/java/oap/logstream/disk/RowBinaryWriterTest.java b/oap-formats/oap-logstream/oap-logstream-test/src/test/java/oap/logstream/disk/RowBinaryWriterTest.java index 5221ce919d..4147a04fcc 100644 --- a/oap-formats/oap-logstream/oap-logstream-test/src/test/java/oap/logstream/disk/RowBinaryWriterTest.java +++ b/oap-formats/oap-logstream/oap-logstream-test/src/test/java/oap/logstream/disk/RowBinaryWriterTest.java @@ -88,7 +88,7 @@ public void testWrite() throws IOException { LogId logId = new LogId( "", "log", "log", Map.of( "p", "1" ), headers, types ); Path logs = testDirectoryFixture.testPath( "logs" ); - try( RowBinaryWriter writer = new RowBinaryWriter( templateEngineFixture.templateEngine, logs, FILE_PATTERN, logId, 1024, BPH_12, 20, "localhost" ) ) { + try( RowBinaryWriter writer = new RowBinaryWriter( templateEngineFixture.templateEngine, logs, FILE_PATTERN, logId, 1024, BPH_12, 20, "localhost", null ) ) { writer.write( CURRENT_PROTOCOL_VERSION, content1 ); writer.write( CURRENT_PROTOCOL_VERSION, content2 ); } @@ -124,7 +124,7 @@ public void testConcurrency() throws IOException { int count = 10; - try( RowBinaryWriter writer = new RowBinaryWriter( templateEngineFixture.templateEngine, logs, FILE_PATTERN, logId, 1024, BPH_12, 20, "localhost" ) ) { + try( RowBinaryWriter writer = new RowBinaryWriter( templateEngineFixture.templateEngine, logs, FILE_PATTERN, logId, 1024, BPH_12, 20, "localhost", null ) ) { try( ExecutorService executorService = Executors.newVirtualThreadPerTaskExecutor() ) { for( long i = 0; i < count; i++ ) { @@ -168,7 +168,7 @@ public void testWriteToNewVersionWhenCompleted() throws IOException { Path v1 = logs.resolve( "1-file-02-47b82ddc0-1.rb.gz.rb.gz" ); Path v2 = logs.resolve( "1-file-02-47b82ddc0-2.rb.gz.rb.gz" ); - try( RowBinaryWriter writer = new RowBinaryWriter( templateEngineFixture.templateEngine, logs, FILE_PATTERN, logId, 1024, BPH_12, 20, "localhost" ) ) { + try( RowBinaryWriter writer = new RowBinaryWriter( templateEngineFixture.templateEngine, logs, FILE_PATTERN, logId, 1024, BPH_12, 20, "localhost", null ) ) { writer.write( CURRENT_PROTOCOL_VERSION, content1 ); writer.refresh(); From e8e387da719ddc444526716191aa83e88c77f200 Mon Sep 17 00:00:00 2001 From: "igor.petrenko" Date: Tue, 21 Jul 2026 21:01:19 +0300 Subject: [PATCH 18/23] CE-178 oap-notification ( oap-notification-mqtt ) --- .../main/java/oap/logstream/disk/DiskLoggerBackend.java | 7 ++++++- 1 file changed, 6 insertions(+), 1 deletion(-) diff --git a/oap-formats/oap-logstream/oap-logstream/src/main/java/oap/logstream/disk/DiskLoggerBackend.java b/oap-formats/oap-logstream/oap-logstream/src/main/java/oap/logstream/disk/DiskLoggerBackend.java index 9c8585511a..94b6277571 100644 --- a/oap-formats/oap-logstream/oap-logstream/src/main/java/oap/logstream/disk/DiskLoggerBackend.java +++ b/oap-formats/oap-logstream/oap-logstream/src/main/java/oap/logstream/disk/DiskLoggerBackend.java @@ -63,6 +63,7 @@ import java.util.LinkedHashMap; import java.util.List; import java.util.Map; +import java.util.concurrent.CopyOnWriteArrayList; import java.util.concurrent.TimeUnit; import static java.util.concurrent.TimeUnit.MILLISECONDS; @@ -93,7 +94,7 @@ public class DiskLoggerBackend extends AbstractLoggerBackend implements FileWrit public long refreshInitDelay = Dates.s( 10 ); public long refreshPeriod = Dates.s( 10 ); public volatile boolean closed; - protected ArrayList notifications = new ArrayList<>(); + protected CopyOnWriteArrayList notifications = new CopyOnWriteArrayList<>(); public DiskLoggerBackend( TemplateEngine templateEngine, Path logDirectory, Timestamp timestamp, int bufferSize, String hostname ) { this( templateEngine, logDirectory, new WriterConfiguration(), timestamp, bufferSize, hostname ); @@ -251,6 +252,10 @@ public void fileClosed( Path outFilename ) { } } + public void addNotification( FileWriterNotification notification ) { + this.notifications.add( notification ); + } + @ToString @EqualsAndHashCode public static class FilePatternConfiguration { From 0298b23be9de8b0f6527dbe26dbd8742b32e5c38 Mon Sep 17 00:00:00 2001 From: "igor.petrenko" Date: Tue, 21 Jul 2026 21:18:33 +0300 Subject: [PATCH 19/23] CE-178 oap-notification ( oap-notification-mqtt ) --- .../src/main/java/oap/logstream/disk/DiskLoggerBackend.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/oap-formats/oap-logstream/oap-logstream/src/main/java/oap/logstream/disk/DiskLoggerBackend.java b/oap-formats/oap-logstream/oap-logstream/src/main/java/oap/logstream/disk/DiskLoggerBackend.java index 94b6277571..d46804b502 100644 --- a/oap-formats/oap-logstream/oap-logstream/src/main/java/oap/logstream/disk/DiskLoggerBackend.java +++ b/oap-formats/oap-logstream/oap-logstream/src/main/java/oap/logstream/disk/DiskLoggerBackend.java @@ -252,7 +252,7 @@ public void fileClosed( Path outFilename ) { } } - public void addNotification( FileWriterNotification notification ) { + public void addNotificationListener( FileWriterNotification notification ) { this.notifications.add( notification ); } From a23bcc9309bd1aa0ac17a83076bafcce5d5d2137 Mon Sep 17 00:00:00 2001 From: "igor.petrenko" Date: Tue, 21 Jul 2026 22:04:47 +0300 Subject: [PATCH 20/23] CE-178 oap-notification ( oap-notification-mqtt ) --- .../src/main/java/oap/logstream/disk/AbstractWriter.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/oap-formats/oap-logstream/oap-logstream/src/main/java/oap/logstream/disk/AbstractWriter.java b/oap-formats/oap-logstream/oap-logstream/src/main/java/oap/logstream/disk/AbstractWriter.java index f20cb14ba4..a4e77a6fab 100644 --- a/oap-formats/oap-logstream/oap-logstream/src/main/java/oap/logstream/disk/AbstractWriter.java +++ b/oap-formats/oap-logstream/oap-logstream/src/main/java/oap/logstream/disk/AbstractWriter.java @@ -140,7 +140,7 @@ public void refresh( boolean forceSync ) { lastPattern = currentPattern; } else { - log.debug( "refresh {}... SKIP", lastPattern ); + log.trace( "refresh {}... SKIP", lastPattern ); } } finally { lock.unlock(); From 018e14b6a50ac36468e3b5fa63fc6fec0795803b Mon Sep 17 00:00:00 2001 From: "igor.petrenko" Date: Tue, 21 Jul 2026 22:12:03 +0300 Subject: [PATCH 21/23] CE-178 oap-notification ( oap-notification-mqtt ) --- .../src/main/java/oap/logstream/disk/AbstractWriter.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/oap-formats/oap-logstream/oap-logstream/src/main/java/oap/logstream/disk/AbstractWriter.java b/oap-formats/oap-logstream/oap-logstream/src/main/java/oap/logstream/disk/AbstractWriter.java index a4e77a6fab..bfa5590dc1 100644 --- a/oap-formats/oap-logstream/oap-logstream/src/main/java/oap/logstream/disk/AbstractWriter.java +++ b/oap-formats/oap-logstream/oap-logstream/src/main/java/oap/logstream/disk/AbstractWriter.java @@ -127,7 +127,7 @@ public void refresh( boolean forceSync ) { String currentPattern = currentPattern(); if( forceSync || !Objects.equals( this.lastPattern, currentPattern ) ) { - log.debug( "lastPattern {} currentPattern {} version {}", lastPattern, currentPattern, fileVersion ); + log.trace( "lastPattern {} currentPattern {} version {}", lastPattern, currentPattern, fileVersion ); String patternWithPreviousVersion = currentPattern( fileVersion - 1 ); if( !Objects.equals( patternWithPreviousVersion, this.lastPattern ) ) { From 2179d60cc2f2c813da4306a2cf159593e18d58bd Mon Sep 17 00:00:00 2001 From: "igor.petrenko" Date: Tue, 21 Jul 2026 22:25:23 +0300 Subject: [PATCH 22/23] CE-178 oap-notification ( oap-notification-mqtt ) --- .../src/main/java/oap/logstream/disk/AbstractWriter.java | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/oap-formats/oap-logstream/oap-logstream/src/main/java/oap/logstream/disk/AbstractWriter.java b/oap-formats/oap-logstream/oap-logstream/src/main/java/oap/logstream/disk/AbstractWriter.java index bfa5590dc1..ab86055cec 100644 --- a/oap-formats/oap-logstream/oap-logstream/src/main/java/oap/logstream/disk/AbstractWriter.java +++ b/oap-formats/oap-logstream/oap-logstream/src/main/java/oap/logstream/disk/AbstractWriter.java @@ -163,7 +163,10 @@ protected void closeOutput() throws LoggerException { Metrics.summary( "logstream_logging_server_bucket_time_seconds" ).record( Dates.nanosToSeconds( stopwatch.elapsed() ) ); logFile.readyForUpload(); - notification.fileClosed( logFile.outFilename ); + + if( notification != null ) { + notification.fileClosed( logFile.outFilename ); + } } finally { logFile = null; } From af71cccfff133aea8d711a645bc88972c98a1a8f Mon Sep 17 00:00:00 2001 From: "igor.petrenko" Date: Wed, 22 Jul 2026 09:05:40 +0300 Subject: [PATCH 23/23] CE-178 oap-notification ( oap-notification-mqtt ) --- .../java/oap/logstream/disk/AbstractWriter.java | 13 +++++++------ 1 file changed, 7 insertions(+), 6 deletions(-) diff --git a/oap-formats/oap-logstream/oap-logstream/src/main/java/oap/logstream/disk/AbstractWriter.java b/oap-formats/oap-logstream/oap-logstream/src/main/java/oap/logstream/disk/AbstractWriter.java index ab86055cec..d685a121fb 100644 --- a/oap-formats/oap-logstream/oap-logstream/src/main/java/oap/logstream/disk/AbstractWriter.java +++ b/oap-formats/oap-logstream/oap-logstream/src/main/java/oap/logstream/disk/AbstractWriter.java @@ -28,7 +28,6 @@ import io.micrometer.core.instrument.Metrics; import lombok.SneakyThrows; import lombok.extern.slf4j.Slf4j; -import oap.concurrent.Stopwatch; import oap.logstream.LogId; import oap.logstream.LogIdTemplate; import oap.logstream.LogStreamProtocol.ProtocolVersion; @@ -52,7 +51,6 @@ public abstract class AbstractWriter implements Closeable { protected final LogId logId; protected final Timestamp timestamp; protected final int bufferSize; - protected final Stopwatch stopwatch = new Stopwatch(); protected final int maxVersions; protected final String hostname; protected final FileWriterNotification notification; @@ -122,7 +120,9 @@ public void refresh() { public void refresh( boolean forceSync ) { lock.lock(); try { - log.debug( "refresh {}...", lastPattern ); + if( logFile != null ) { + log.debug( "refresh {}...", lastPattern ); + } String currentPattern = currentPattern(); @@ -155,12 +155,11 @@ protected void closeOutput() throws LoggerException { lock.lock(); try { if( logFile != null ) try { - stopwatch.count( logFile::close ); + logFile.close(); long fileSize = logFile.getDataSize(); log.trace( "closing output {} ({} bytes)", this, fileSize ); Metrics.summary( "logstream_logging_server_bucket_size" ).record( fileSize ); - Metrics.summary( "logstream_logging_server_bucket_time_seconds" ).record( Dates.nanosToSeconds( stopwatch.elapsed() ) ); logFile.readyForUpload(); @@ -179,7 +178,9 @@ protected void closeOutput() throws LoggerException { public void close() { lock.lock(); try { - log.debug( "closing {}", this ); + if( logFile != null ) { + log.debug( "closing {}", this ); + } closed = true; closeOutput(); } finally {