Interface StreamBridge

All Superinterfaces:
AutoCloseable
All Known Implementing Classes:
DisruptorStreamBridge, KafkaStreamEngine

public interface StreamBridge extends AutoCloseable
High-performance streaming SPI bridge for in-process event ingestion, rule evaluation, and result dispatching.
  • Method Details

    • subscribe

      void subscribe(String topic, RuleStreamHandler handler)
      Subscribes a functional rule stream handler to a given topic.
      Parameters:
      topic - stream topic
      handler - handler invoked for incoming events on the topic
    • subscribeListener

      void subscribeListener(String topic, StreamListener listener)
      Subscribes an asynchronous result listener to a given topic.
      Parameters:
      topic - stream topic
      listener - listener invoked when events on the topic complete evaluation
    • addListener

      default void addListener(String topic, StreamListener listener)
      Parameters:
      topic - stream topic
      listener - listener invoked when events on the topic complete evaluation
    • publish

      Publishes a rule evaluation event into the stream and returns a CompletableFuture for the result.
      Parameters:
      event - the rule event to publish
      Returns:
      CompletableFuture completing with the evaluation result
    • publish

      <T> CompletableFuture<StreamResult> publish(StreamRecord<T> record)
      Publishes a generic stream record into the stream and returns a CompletableFuture for the result.
      Type Parameters:
      T - payload type
      Parameters:
      record - the stream record to publish
      Returns:
      CompletableFuture completing with the evaluation result
    • tryPublish

      boolean tryPublish(RuleEvent event, StreamListener listener)
      Fast-path non-blocking publish of a rule event with a completion listener, avoiding Future allocation.
      Parameters:
      event - rule event
      listener - optional listener for evaluation result
      Returns:
      true if published onto the ring buffer, false if saturated
    • tryPublish

      default boolean tryPublish(RuleEvent event)
      Fast-path fire-and-forget publish of an event.
      Parameters:
      event - rule event
      Returns:
      true if accepted, false if saturated
    • stats

      StreamStats stats()
      Retrieves current streaming engine metrics and capacity stats.
      Returns:
      runtime stats snapshot
    • isRunning

      boolean isRunning()
      Checks if the stream bridge is active and accepting events.
      Returns:
      true if running, false if stopped or shut down
    • shutdown

      void shutdown(Duration timeout)
      Gracefully shuts down the stream bridge within the specified timeout duration.
      Parameters:
      timeout - maximum time to await in-flight events draining
    • close

      default void close()
      Specified by:
      close in interface AutoCloseable