-
Notifications
You must be signed in to change notification settings - Fork 38
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
- Loading branch information
Showing
11 changed files
with
177 additions
and
61 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
42 changes: 42 additions & 0 deletions
42
nflow-rest-api-spring-web/src/main/java/io/nflow/rest/config/springweb/SchedulerService.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,42 @@ | ||
package io.nflow.rest.config.springweb; | ||
|
||
import static java.lang.Math.max; | ||
import static org.slf4j.LoggerFactory.getLogger; | ||
import static reactor.core.publisher.Mono.fromCallable; | ||
import static reactor.core.scheduler.Schedulers.fromExecutor; | ||
|
||
import java.util.concurrent.Callable; | ||
import java.util.concurrent.Executors; | ||
|
||
import javax.inject.Inject; | ||
|
||
import org.slf4j.Logger; | ||
import org.springframework.core.env.Environment; | ||
import org.springframework.stereotype.Service; | ||
|
||
import io.nflow.engine.internal.executor.WorkflowInstanceExecutor; | ||
import reactor.core.publisher.Mono; | ||
import reactor.core.scheduler.Scheduler; | ||
|
||
/** | ||
* Service to hold a Webflux Scheduler in order to make blocking calls. | ||
*/ | ||
@Service | ||
public class SchedulerService { | ||
|
||
private static final Logger logger = getLogger(SchedulerService.class); | ||
private final Scheduler scheduler; | ||
|
||
@Inject | ||
public SchedulerService(WorkflowInstanceExecutor workflowInstanceExecutor, Environment env) { | ||
int dbPoolSize = env.getProperty("nflow.db.max_pool_size", Integer.class); | ||
int dispatcherCount = workflowInstanceExecutor.getThreadCount(); | ||
int threadPoolSize = max(dbPoolSize - dispatcherCount, 2); | ||
logger.info("Initializing REST API thread pool size to {}", threadPoolSize); | ||
this.scheduler = fromExecutor(Executors.newFixedThreadPool(threadPoolSize)); | ||
} | ||
|
||
public <T> Mono<T> callAsync(Callable<T> callable) { | ||
return fromCallable(callable).subscribeOn(this.scheduler); | ||
} | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
20 changes: 17 additions & 3 deletions
20
nflow-rest-api-spring-web/src/main/java/io/nflow/rest/v1/springweb/SpringWebResource.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -1,21 +1,35 @@ | ||
package io.nflow.rest.v1.springweb; | ||
|
||
import static org.springframework.http.ResponseEntity.status; | ||
import static reactor.core.publisher.Mono.just; | ||
|
||
import java.util.concurrent.Callable; | ||
import java.util.function.Supplier; | ||
|
||
import org.springframework.http.ResponseEntity; | ||
|
||
import io.nflow.rest.config.springweb.SchedulerService; | ||
import io.nflow.rest.v1.ResourceBase; | ||
import io.nflow.rest.v1.msg.ErrorResponse; | ||
import reactor.core.publisher.Mono; | ||
|
||
public abstract class SpringWebResource extends ResourceBase { | ||
|
||
protected ResponseEntity<?> handleExceptions(Supplier<ResponseEntity<?>> response) { | ||
private final SchedulerService scheduler; | ||
|
||
protected SpringWebResource(SchedulerService scheduler) { | ||
this.scheduler = scheduler; | ||
} | ||
|
||
protected Mono<ResponseEntity<?>> wrapBlocking(Callable<ResponseEntity<?>> callable) { | ||
return scheduler.callAsync(callable); | ||
} | ||
|
||
protected Mono<ResponseEntity<?>> handleExceptions(Supplier<Mono<ResponseEntity<?>>> response) { | ||
return handleExceptions(response::get, this::toErrorResponse); | ||
} | ||
|
||
private ResponseEntity<?> toErrorResponse(int statusCode, ErrorResponse body) { | ||
return status(statusCode).body(body); | ||
private Mono<ResponseEntity<?>> toErrorResponse(int statusCode, ErrorResponse body) { | ||
return just(status(statusCode).body(body)); | ||
} | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Oops, something went wrong.