Class OrchestratorEngine

java.lang.Object
com.iantapply.orchestra.engine.OrchestratorEngine
All Implemented Interfaces:
AutoCloseable

public final class OrchestratorEngine extends Object implements AutoCloseable
Durable, at-least-once orchestration engine. Action implementations should use ActionContext.idempotencyKey() when calling systems that support deduplication.
  • Constructor Summary

    Constructors
    Constructor
    Description
    OrchestratorEngine(com.iantapply.orchestra.port.DefinitionRepository definitions, com.iantapply.orchestra.port.ExecutionRepository executions, com.iantapply.orchestra.port.DistributedLock locks, com.iantapply.orchestra.port.TargetResolver targets, ActionRegistry registry, Clock clock, int workerCount, int queueCapacity)
    Creates an engine with bounded worker concurrency and queue capacity.
    OrchestratorEngine(com.iantapply.orchestra.port.DefinitionRepository definitions, com.iantapply.orchestra.port.ExecutionRepository executions, com.iantapply.orchestra.port.DistributedLock locks, com.iantapply.orchestra.port.TargetResolver targets, ActionRegistry registry, Clock clock, EngineOptions options)
    Creates an engine with explicit operational timing and capacity options.
    OrchestratorEngine(com.iantapply.orchestra.port.DefinitionRepository definitions, com.iantapply.orchestra.port.ExecutionRepository executions, com.iantapply.orchestra.port.DistributedLock locks, com.iantapply.orchestra.port.TargetResolver targets, ActionRegistry registry, Clock clock, EngineOptions options, Consumer<String> failureCounter)
    Creates an engine whose operational failure paths increment named counters.
  • Method Summary

    Modifier and Type
    Method
    Description
    int
    Returns the number of workers currently executing a task.
    void
    addListener(com.iantapply.orchestra.api.EventLifecycleListener listener)
    Registers a lifecycle listener.
    boolean
    Cancels an execution.
    void
    Stops polling and shuts down engine workers.
    boolean
    pause(UUID id)
    Pauses an execution.
    int
    Returns the number of execution tasks currently waiting for a worker.
    void
    Immediately scans persisted state; called on startup to recover interrupted work.
    boolean
    Resumes an execution.
    boolean
    retry(UUID id)
    Reschedules a failed execution from its first stage.
    schedule(String definitionId, Instant startAt, Map<String,Object> variables)
    Creates a durable execution for a known definition.
    boolean
    setVariable(UUID id, String key, Object value)
    Atomically sets or removes an execution variable.
    void
    Starts the due-execution polling loop once.
    startNow(String definitionId)
    Schedules an event immediately with no initial variables.

    Methods inherited from class Object

    clone, equals, finalize, getClass, hashCode, notify, notifyAll, toString, wait, wait, wait
  • Constructor Details

    • OrchestratorEngine

      public OrchestratorEngine(com.iantapply.orchestra.port.DefinitionRepository definitions, com.iantapply.orchestra.port.ExecutionRepository executions, com.iantapply.orchestra.port.DistributedLock locks, com.iantapply.orchestra.port.TargetResolver targets, ActionRegistry registry, Clock clock, int workerCount, int queueCapacity)
      Creates an engine with bounded worker concurrency and queue capacity.
      Parameters:
      definitions - event definition store
      executions - durable execution store
      locks - distributed execution lease provider
      targets - target resolver
      registry - action and condition registry
      clock - engine clock
      workerCount - number of concurrent execution workers
      queueCapacity - maximum queued execution tasks
    • OrchestratorEngine

      public OrchestratorEngine(com.iantapply.orchestra.port.DefinitionRepository definitions, com.iantapply.orchestra.port.ExecutionRepository executions, com.iantapply.orchestra.port.DistributedLock locks, com.iantapply.orchestra.port.TargetResolver targets, ActionRegistry registry, Clock clock, EngineOptions options)
      Creates an engine with explicit operational timing and capacity options.
      Parameters:
      definitions - event definition store
      executions - durable execution store
      locks - distributed execution lease provider
      targets - target resolver
      registry - action and condition registry
      clock - engine clock
      options - operational timing and capacity options
    • OrchestratorEngine

      public OrchestratorEngine(com.iantapply.orchestra.port.DefinitionRepository definitions, com.iantapply.orchestra.port.ExecutionRepository executions, com.iantapply.orchestra.port.DistributedLock locks, com.iantapply.orchestra.port.TargetResolver targets, ActionRegistry registry, Clock clock, EngineOptions options, Consumer<String> failureCounter)
      Creates an engine whose operational failure paths increment named counters.
      Parameters:
      definitions - event definition store
      executions - durable execution store
      locks - distributed execution lease provider
      targets - target resolver
      registry - action and condition registry
      clock - engine clock
      options - operational timing and capacity options
      failureCounter - consumer of metric names for failures and rejected work
  • Method Details

    • start

      public void start()
      Starts the due-execution polling loop once.
    • schedule

      public UUID schedule(String definitionId, Instant startAt, Map<String,Object> variables)
      Creates a durable execution for a known definition.
      Parameters:
      definitionId - event definition identifier
      startAt - requested start time
      variables - initial execution variables
      Returns:
      new execution identifier
    • startNow

      public UUID startNow(String definitionId)
      Schedules an event immediately with no initial variables.
      Parameters:
      definitionId - event definition to schedule immediately
      Returns:
      new execution identifier
    • addListener

      public void addListener(com.iantapply.orchestra.api.EventLifecycleListener listener)
      Registers a lifecycle listener.
      Parameters:
      listener - lifecycle listener retained until the engine closes
    • pause

      public boolean pause(UUID id)
      Pauses an execution.
      Parameters:
      id - execution identifier
      Returns:
      whether the execution was found and atomically paused
    • resume

      public boolean resume(UUID id)
      Resumes an execution.
      Parameters:
      id - execution identifier
      Returns:
      whether the execution was found and atomically resumed
    • cancel

      public boolean cancel(UUID id)
      Cancels an execution.
      Parameters:
      id - execution identifier
      Returns:
      whether the execution was found and atomically cancelled
    • retry

      public boolean retry(UUID id)
      Reschedules a failed execution from its first stage.
      Parameters:
      id - execution identifier
      Returns:
      whether the failed execution was atomically rescheduled
    • setVariable

      public boolean setVariable(UUID id, String key, Object value)
      Atomically sets or removes an execution variable.
      Parameters:
      id - execution identifier
      key - variable name
      value - new value, or null to remove it
      Returns:
      whether the update succeeded
    • recover

      public void recover()
      Immediately scans persisted state; called on startup to recover interrupted work.
    • queuedTaskCount

      public int queuedTaskCount()
      Returns the number of execution tasks currently waiting for a worker.
      Returns:
      queued task count
    • activeWorkerCount

      public int activeWorkerCount()
      Returns the number of workers currently executing a task.
      Returns:
      active worker count
    • close

      public void close()
      Stops polling and shuts down engine workers.
      Specified by:
      close in interface AutoCloseable