Examples in Java
FutoIn AsyncSteps Java API
The universal API is provided by org.futoin.api module, which is a dependency
of this Reference Implementation.
// Universal API
import org.futoin.api.AsyncSteps;
import org.futoin.api.AsyncTool;
// Reference Implementation
import org.futoin.ri.asyncsteps.AsyncStepsRI;
import org.futoin.ri.asyncsteps.AsyncToolRI;
class Examples {
boolean WITH_EXCEPTIONS = true;
void do_some_non_blocking_stuff() {}
void this_code_is_not_executed() {
assert (false);
}
void this_code_IS_executed() {}
Object schedule_external_callback(java.util.function.Consumer<Boolean> cb) {
cb.accept(false);
return null;
}
void external_cancel(Object handle) {}
final MutexRI mutex = new MutexRI();
AsyncSteps.ISync get_some_synchronization_object() {
return mutex;
}
void example_business_logic(AsyncSteps asi) {
// ---------------------------------------------------------------------
// 1. regular step example
//
// Prototypes:
// - asi.add(func(asi));
// - asi.add(func(asi), on_error(asi, code));
// - asi.add(func(asi, [A [,B [,C [,D]]]]));
// - asi.add(func(asi, [A [,B [,C [,D]]]]), on_error(asi, code));
//
asi.add((asi2) -> do_some_non_blocking_stuff());
// ---------------------------------------------------------------------
// 2. try {} catch {} block example
asi.add(
(asi2) -> {
// regular step example
do_some_non_blocking_stuff();
if (WITH_EXCEPTIONS) {
asi2.error("MyError"); // throw error
this_code_is_not_executed();
} else {
asi2.errorNoThrow("MyError"); // not exception
return; // ensure manual return without exceptions
}
},
(asi2, error_code) -> {
if (error_code.equals("MyError")) {
// Override error unwind with success
asi2.success();
} else {
asi2.error("OverrideErrorCode");
}
});
// ---------------------------------------------------------------------
// 3. Inner steps
asi.add(
(asi2) -> {
// regular step example
do_some_non_blocking_stuff();
asi2.add(
(asi3) -> {
asi3.error("MyError"); // throw error
});
// NOTE: inner steps are executed AFTER outer step body
this_code_IS_executed();
},
(asi2, error_code) -> {
if (error_code.equals("MyError")) {
asi2.success();
}
});
// ---------------------------------------------------------------------
// 4. Passing arbitrary parameters
asi.add(
(AsyncSteps asi2) -> {
asi2.success(123, true, "SomeString", List.of(1, 2, 3));
});
asi.<Integer, Boolean, String, List<Integer>>add(
(asi2, a, b, c, d) -> {
// Generic type inference example.
// NOTE: Maximum of 4 arguments is supported based on best practices.
assert (a == 123);
assert (b);
assert (c.equals("SomeString"));
assert (d.get(0) == 1);
asi2.success(a, b);
});
asi.add(
(AsyncSteps asi2, Integer a, Boolean b) -> {
// Explicit parameter types example.
assert (a == 123);
assert (b);
});
// ---------------------------------------------------------------------
// 5. state() - Thread Local Storage emulation
//
// Some predefined properties are set directly on State object for
// performance reasons.
//
// Business logic can use custom dynamic items as associative key-value map.
// Key is a String, value is of Object type with handy type casing getters.
{
// get reference to state object
var state = asi.state();
// Get-or-set-default variable
var some_var = state.set_default("SomeVar", 123);
// Same, but a shortcut
var some_var2 = asi.state("SomeVar", 123);
// Set
state.set("SomeVector", List.of(1, 2, 3));
}
asi.add(
(AsyncSteps asi2) -> {
// Get
var v = asi2.state().<List<Integer>>get("SomeVector");
// Get, but a shortcut
var i = asi2.<Integer>state("SomeVar");
try {
asi2.state().<Integer>get("SomeVector");
} catch (ClassCastException ex) {
// ...
}
});
// ---------------------------------------------------------------------
// 6. Advanced state handling
var sample_asi = asi.newInstance();
sample_asi
.state()
.set_catch_trace(
(AsyncSteps asi2, Throwable t) -> {
// Mostly a helper for debugging purposes
});
sample_asi
.state()
.set_unhandled_error(
(AsyncSteps asi2, String err) -> {
// Handle unexpected exception that cancels whole execution
});
sample_asi
.state()
.set_cancel_handler(
(AsyncSteps asi2) -> {
// Handle situations of AsyncSteps execution cancel
});
asi.add(
(asi2) -> {
// NOTE: error codes are associative, but not somes integers
// to be more network-friendly.
asi2.error("MyError", "Some arbitrary description of the error");
},
(asi2, error_code) -> {
// ErrorCode is wrapper around const char*
assert (error_code.equals("MyError"));
// Error info is stored in state
assert (asi2.state()
.error_info()
.equals("Some arbitrary description of the error"));
// Last exception thrown is also available in state
Throwable e = asi2.state().last_exception();
asi2.success(e);
});
// ---------------------------------------------------------------------
/// 7. Synchronization
//
// Unlike hardware race conditions, AsyncSteps synchronization serves
// logical purposes to limit concurrency of execution or rate of calls or
// both.
//
// FTN12 concept defines Mutex, Throttle and Limiter primitives which
// implement a single ISync interface.
var syncObj = get_some_synchronization_object();
// The same interface as asi.add(), but with extra synchronization object
// parameter.
asi.sync(
syncObj,
(asi2) -> {
// a critical section
},
(asi2, error_code) -> {
// an optional error handler for the critical section
});
// ---------------------------------------------------------------------
// 8. Loops
asi.loop(
(asi2) -> {
// infinite loop
if (WITH_EXCEPTIONS) {
asi2.breakLoop();
} else {
asi2.breakLoopNoThrow();
}
});
asi.repeat(
10,
(asi2, i) -> {
// range loop from i=0 till i=9 (inclusive)
});
asi.forEach(
List.of(1, 2, 3),
(asi2, index, value) -> {
// Iteration of arrays and sequences
});
asi.forEach(
new HashMap<String, String>(),
(asi2, key, value) -> {
// Iteration of map-like objects
});
// ---------------------------------------------------------------------
// 9. Timeout support
asi.add(
(asi2) -> {
// Raises Timeout error after specified period
asi2.setTimeout(Duration.ofMillis(10));
asi2.loop(
(asi3) -> {
// infinite loop
asi3.relinquish();
});
},
(asi2, error_code) -> {
if (error_code.equals(Error.Timeout)) {
asi2.success();
}
});
// ---------------------------------------------------------------------
// 10. External event integration
asi.add(
(asi2) -> {
var handle =
schedule_external_callback(
(isError) -> {
if (isError) {
asi2.errorNoThrow("ExternalError");
} else {
asi2.success();
}
});
// Handle an edge case when callbacks fires immediately
// and invalidates the current step.
if (asi2.state() != null) {
asi2.setCancel((asi3) -> external_cancel(handle));
}
});
asi.add(
(asi2) -> {
var handle =
schedule_external_callback(
(isError) -> {
if (asi2.state() == null) {
// AsyncSteps object is invalidated
// due to external cancel.
} else if (isError) {
asi2.errorNoThrow("ExternalError");
} else {
asi2.success();
}
});
// Handle an edge case when callbacks fires immediately
// and invalidates the current step.
if (asi2.state() != null) {
// alias for setCancel() with noop handler
asi2.waitExternal();
}
});
// ---------------------------------------------------------------------
// 11. Standard Promise/Await integration
asi.add(
(asi2) -> {
// The proper way to create new AsyncSteps instances
// without hard dependency on implementation.
var new_steps = asi2.newInstance();
new_steps.add(
(new_asi) -> {
new_asi.success(123);
});
// Proper way to wait for standard java.util.concurrent.Future
asi2.await(new_steps.<Integer>promise());
asi2.<Integer>add(
(asi3, i) -> {
assert (i == 123);
});
});
// ---------------------------------------------------------------------
// 12. Parallel execution
//
// It's designed for concurrent execution of sub-flows
// with shared state() on the same platform thread.
//
// Unhandled error in sub-flows lead to abort of all non-executed
// parallel steps.
asi.state().set("order", new ArrayList<Integer>());
var p =
asi.parallel(
(asi2, error_code) -> {
// Overall error handler
asi2.success();
});
p.add(
(asi2) -> {
// regular flow
asi2.state().<ArrayList<Integer>>get("order").add(1);
// Break execution burst - a CPU cache optimization, for demo
asi2.relinquish();
asi2.add(
(asi3) -> {
asi3.state().<ArrayList<Integer>>get("order").add(4);
});
});
p.add(
(asi2) -> {
asi2.state().<ArrayList<Integer>>get("order").add(2);
asi2.relinquish(); // Break execution burst
asi2.add(
(asi3) -> {
asi3.state().<ArrayList<Integer>>get("order").add(5);
asi3.error("SomeError");
});
});
p.add(
(asi2) -> {
asi2.state().<ArrayList<Integer>>get("order").add(3);
asi2.relinquish(); // Break execution burst
asi2.add(
(asi3) -> {
asi3.state().<ArrayList<Integer>>get("order").add(6);
});
});
asi.add(
(asi2) -> {
asi2.state().<ArrayList<Integer>>get("order"); // 1, 2, 3, 4, 5
});
// ---------------------------------------------------------------------
// 13. Control of AsyncSteps flow
{
// A new instance inherits the special handlers and use the same
// event loop.
var new_steps = asi.newInstance();
assert (new_steps.tool() == asi.tool());
// Add steps
new_steps.add((new_asi) -> {});
new_steps.loop((new_asi) -> new_asi.relinquish());
// Schedule execution of AsyncSteps flow
new_steps.execute();
// Cancel execution of AsyncSteps flow
new_steps.cancel();
}
}
void testAPIExample() throws Throwable {
var $as = new AsyncStepsRI();
$as.add(this::example_business_logic);
$as.promise().get();
}
}AsyncStepsRI - reference implementation
void example_AsyncStepsRI() throws Throwable {
// Use the default AsyncToolRI.shared() singleton
{
var $as = new AsyncStepsRI();
$as.add(
(asi) -> {
/* ... */
});
$as.execute();
}
// A dedicated event loop for some specific job
try (var async_tool = new AsyncToolRI()) {
for (int i = 0; i < 100; ++i) {
var $as = new AsyncStepsRI(async_tool);
$as.add(this::example_business_logic);
$as.execute();
}
}
}AsyncToolRI - reference implementation
void example_AsyncToolRI() throws Throwable {
// Create underlying event loop with own thread
try (var async_tool_own_loop = new AsyncToolRI()) {
async_tool_own_loop.immediate(
() -> {
// Executes first
async_tool_own_loop.immediate(
() -> {
// Likes executes second due to the race below.
});
});
async_tool_own_loop.immediate(
() -> {
// Likely executes third due to external scheduling race.
});
async_tool_own_loop.deferred(
Duration.ofMillis(100),
() -> {
// Executes not earlier than the time delay
});
(new AsyncStepsRI(async_tool_own_loop)).execute();
// Clean shutdown requires all scheduled jobs to complete.
}
// Create underlying event loop for integration into foreign event loop
try (var async_tool_foreign_loop =
new AsyncToolRI(
() -> {
// A callback to wake up foreign event loop to indicate new work
// available due to out-of-band API usage.
})) {
async_tool_foreign_loop.immediate(
() -> {
// Executes first
async_tool_foreign_loop.immediate(
() -> {
// Executes third
});
});
async_tool_foreign_loop.immediate(
() -> {
// Executes second
});
(new AsyncStepsRI(async_tool_foreign_loop))
.add(
(asi) -> {
// Executes fourth
})
.execute();
// Idiomatic foreing event loop logic
for (; ; ) {
var cycleResult = async_tool_foreign_loop.iterate();
if (cycleResult.haveWork()) {
// Delay may be zero, if new immediate() calls are scheduled
// during the current cycle.
Thread.sleep(Duration.ofNanos(cycleResult.delayNs()).toMillis());
} else {
// Wait for the wakeup callback, supplied to c-tor above.
break;
}
}
}
}Synchronization primitives
import org.futoin.api.Limiter;
import org.futoin.api.Mutex;
import org.futoin.api.Throttle;
import org.futoin.ri.asyncsteps.LimiterRI;
import org.futoin.ri.asyncsteps.MutexRI;
import org.futoin.ri.asyncsteps.ThrottleRI;
class Example {
// 1 concurrent, infinite queue
final Mutex mutex = new MutexRI();
// 10 concurrent, 1000 queue items
final Mutex mutex2 = new MutexRI(10, 1000);
// 100 entries per second. infinite queue
final Throttle throttle = new ThrottleRI(100);
// 100 entries per 10 seconds with maximum queue of 300
final Throttle throttle2 =
new ThrottleRI(AsyncToolRI.shared(), 100, Duration.ofSeconds(10), 300);
// 10 concurrent entries with queue of 20 with 100 entries over
// a period of 10 seconds with burst queue of 200
final Limiter limiter =
new LimiterRI(
(new Limiter.Options())
.withConcurrent(10)
.withMaxQueue(20)
.withRate(100)
.withPeriod(Duration.ofSeconds(10))
.withBurst(200));
}