Compare commits

...
14 Commits
Author SHA1 Message Date
brian.long 0bad7d4f0e fix NPE in comparator 2026-08-17 13:25:19 -04:00
brian.long ae14f183c8 updated various artifact versions 2026-08-17 13:24:08 -04:00
brian.long 2ed236b3cd fix copy/paste compile issue 2026-08-17 12:48:17 -04:00
brian.long 0685bf9e76 add reconcile throttling waits 2026-03-24 22:30:58 -04:00
brian.long c0b02b9004 basic cleanup
(cherry picked from commit 727a566ad55100023d613b98bb69f7ee21891be3)
2026-03-24 22:22:20 -04:00
brian.long 06d16eb223 various minor improvements
(cherry picked from commit f2cf774bad5987c91b4e8acef94fb97fdb921063)
2026-03-24 22:21:10 -04:00
brian.long 83796054e1 remove cxf-jaxrs-platform-module dups from packaging 2026-03-24 21:33:16 -04:00
brian.long f7b546d640 readability cleanup 2026-03-24 21:28:34 -04:00
brian.long 00c2ddd125 fix date/time regex hashes 2026-03-24 21:28:08 -04:00
brian.long 259dd72e85 fix ASIE core query param: nodeId -> nodeid 2026-03-24 21:26:56 -04:00
brian.long 37351974d3 POM reorg 2026-03-24 21:26:34 -04:00
brian.long 66933e587d separate reconcile threads from index/reindex threads 2026-02-02 18:48:33 -05:00
brian.long 83efe3a679 move spring value defaults to module alfresco-global 2026-02-02 18:47:13 -05:00
brian.long fbf6c17206 moving shard determination to SolrShardHashService 2026-01-12 15:29:03 -05:00
26 changed files with 457 additions and 329 deletions
+23 -22
View File
@@ -10,17 +10,12 @@
<relativePath>../</relativePath>
</parent>
<groupId>com.inteligr8.alfresco</groupId>
<artifactId>asie-api</artifactId>
<version>1.1-SNAPSHOT-asie2</version>
<version>1.2-SNAPSHOT-asie2</version>
<packaging>jar</packaging>
<name>ASIE Jakarta RS API</name>
<description>Alfresco Search &amp; Insight Engine Jakarta RS API</description>
<properties>
<alfresco.platform.version>23.2.0</alfresco.platform.version>
</properties>
<dependencyManagement>
<dependencies>
@@ -38,22 +33,11 @@
<dependency>
<groupId>com.inteligr8</groupId>
<artifactId>solr-api</artifactId>
<version>1.1-SNAPSHOT-solr6</version>
<version>1.2-SNAPSHOT-solr6</version>
</dependency>
<dependency>
<groupId>org.alfresco</groupId>
<artifactId>alfresco-data-model</artifactId>
<exclusions>
<exclusion>
<groupId>*</groupId>
<artifactId>*</artifactId>
</exclusion>
</exclusions>
</dependency>
<dependency>
<groupId>org.apache.commons</groupId>
<artifactId>commons-lang3</artifactId>
<version>3.17.0</version>
</dependency>
<dependency>
<groupId>org.apache.logging.log4j</groupId>
@@ -63,27 +47,44 @@
<dependency>
<groupId>com.inteligr8</groupId>
<artifactId>common-rest-client</artifactId>
<version>3.0.2-jersey</version>
<version>${commom-rest-client.base.version}-jersey</version>
<scope>test</scope>
</dependency>
<dependency>
<groupId>org.glassfish.jersey.inject</groupId>
<artifactId>jersey-hk2</artifactId>
<version>3.1.10</version>
<version>4.0.2</version>
<scope>test</scope>
</dependency>
<dependency>
<groupId>org.junit.jupiter</groupId>
<artifactId>junit-jupiter-api</artifactId>
<version>5.11.2</version>
<scope>test</scope>
</dependency>
</dependencies>
<build>
<plugins>
<plugin>
<groupId>io.repaint.maven</groupId>
<artifactId>tiles-maven-plugin</artifactId>
<extensions>true</extensions>
<configuration>
<tiles>
<!-- Documentation: https://git.inteligr8.com/inteligr8/ootbee-beedk/src/stable/beedk-acs-platform-webapp-tile -->
<!-- TODO spin up just ASIE
<tile>com.inteligr8.ootbee:beedk-acs-platform-webapp-tile:[1.1.6,2.0.0)</tile>
-->
</tiles>
</configuration>
</plugin>
</plugins>
</build>
<repositories>
<repository>
<id>alfresco-public</id>
<url>https://artifacts.alfresco.com/nexus/repository/public/</url>
<url>https://artifacts.alfresco.com/nexus/repository/releases/</url>
</repository>
</repositories>
</project>
@@ -24,7 +24,7 @@ public class IndexRequest extends JsonFormattedResponseRequest<IndexRequest> {
@QueryParam("acltxid")
private Long aclTransactionId;
@QueryParam("nodeId")
@QueryParam("nodeid")
private Long nodeId;
@QueryParam("aclid")
@@ -3,7 +3,7 @@ package com.inteligr8.alfresco.asie.model.core;
import java.util.Collection;
import org.alfresco.service.cmr.repository.StoreRef;
import org.apache.commons.lang3.StringUtils;
import org.springframework.util.StringUtils;
import com.inteligr8.solr.model.JsonFormattedResponseRequest;
@@ -89,7 +89,7 @@ public class NewCoreRequest extends JsonFormattedResponseRequest<NewCoreRequest>
}
public NewCoreRequest withShardIds(Collection<String> shardIds) {
this.shardIds = StringUtils.join(shardIds, ",");
this.shardIds = StringUtils.collectionToDelimitedString(shardIds, ",");
return this;
}
@@ -24,7 +24,7 @@ public class PurgeRequest extends JsonFormattedResponseRequest<PurgeRequest> {
@QueryParam("acltxid")
private Long aclTransactionId;
@QueryParam("nodeId")
@QueryParam("nodeid")
private Long nodeId;
@QueryParam("aclid")
@@ -24,7 +24,7 @@ public class ReindexRequest extends JsonFormattedResponseRequest<ReindexRequest>
@QueryParam("acltxid")
private Long aclTransactionId;
@QueryParam("nodeId")
@QueryParam("nodeid")
private Long nodeId;
@QueryParam("aclid")
@@ -8,7 +8,7 @@ import org.slf4j.LoggerFactory;
import com.inteligr8.alfresco.asie.AsieClient;
public class AbstractApiUnitTest {
public class AbstractApiIT {
protected Logger logger = LoggerFactory.getLogger(this.getClass());
@@ -11,7 +11,7 @@ import com.inteligr8.solr.model.Action.Status;
import com.inteligr8.solr.model.Cores;
import com.inteligr8.solr.model.ResponseHeader;
public class CoreAdminReindexUnitTest extends AbstractApiUnitTest {
public class CoreAdminReindexIT extends AbstractApiIT {
@Test
public void reindex() {
@@ -19,7 +19,7 @@ import com.inteligr8.solr.model.core.StatusResponse;
import jakarta.ws.rs.ProcessingException;
public class CoreAdminStatusUnitTest extends AbstractApiUnitTest {
public class CoreAdminStatusIT extends AbstractApiIT {
@Test
public void noHost() {
+6 -7
View File
@@ -16,10 +16,10 @@
<name>ASIE Platform Module for ACS Community</name>
<properties>
<alfresco.sdk.version>4.9.0</alfresco.sdk.version>
<alfresco.platform.version>23.3.0</alfresco.platform.version>
<alfresco.platform.war.version>23.3.0.98</alfresco.platform.war.version>
<tomcat-rad.version>10-2.1</tomcat-rad.version>
<alfresco.sdk.version>4.16.0</alfresco.sdk.version>
<alfresco.platform.version>26.1.0</alfresco.platform.version>
<alfresco.platform.war.version>26.1.0.61</alfresco.platform.war.version>
<tomcat-rad.version>2.3-tomcat-11.0.22</tomcat-rad.version>
<beedk.rad.acs-search.enabled>true</beedk.rad.acs-search.enabled>
</properties>
@@ -59,7 +59,7 @@
<dependency>
<groupId>com.inteligr8.alfresco</groupId>
<artifactId>cxf-jaxrs-platform-module</artifactId>
<version>1.3.1-acs-v23.3</version>
<version>1.4.0-acs-v26.1</version>
<type>amp</type>
</dependency>
@@ -81,7 +81,6 @@
<plugin>
<groupId>io.repaint.maven</groupId>
<artifactId>tiles-maven-plugin</artifactId>
<version>2.40</version>
<extensions>true</extensions>
<configuration>
<tiles>
@@ -100,7 +99,7 @@
<repositories>
<repository>
<id>alfresco-public</id>
<url>https://artifacts.alfresco.com/nexus/content/groups/public</url>
<url>https://artifacts.alfresco.com/nexus/repository/releases/</url>
</repository>
</repositories>
</project>
+14 -20
View File
@@ -16,10 +16,10 @@
<name>ASIE Platform Module for ACS Enterprise</name>
<properties>
<alfresco.sdk.version>4.9.0</alfresco.sdk.version>
<alfresco.platform.version>23.3.0</alfresco.platform.version>
<alfresco.platform.war.version>23.3.0.98</alfresco.platform.war.version>
<tomcat-rad.version>10-2.1</tomcat-rad.version>
<alfresco.platform.version>26.1.0</alfresco.platform.version>
<alfresco.platform.war.version>26.1.0.61</alfresco.platform.war.version>
<tomcat-rad.version>2.3-tomcat-11.0.22</tomcat-rad.version>
<jackson.version>2.17.2</jackson.version>
<beedk.rad.acs-search.enabled>true</beedk.rad.acs-search.enabled>
</properties>
@@ -37,10 +37,13 @@
</dependencyManagement>
<dependencies>
<!-- Alfresco Modules required to use this module -->
<dependency>
<groupId>com.inteligr8.alfresco</groupId>
<artifactId>asie-shared</artifactId>
<version>${project.version}</version>
<artifactId>cxf-jaxrs-platform-module</artifactId>
<version>1.4.0-acs-v26.1</version>
<type>amp</type>
<scope>provided</scope>
</dependency>
<!-- Needed by this module, but provided by ACS -->
@@ -50,19 +53,11 @@
<scope>provided</scope>
</dependency>
<!-- Alfresco Modules required to use this module -->
<!-- core dependency, but dependencies of provided ones above should override -->
<dependency>
<groupId>com.inteligr8.alfresco</groupId>
<artifactId>cxf-jaxrs-platform-module</artifactId>
<version>1.3.2-acs-v23.3</version>
<type>amp</type>
</dependency>
<!-- Provided by cxf-jaxrs-platform-module, but packaged due to shared -->
<dependency>
<groupId>com.inteligr8</groupId>
<artifactId>common-rest-client</artifactId>
<scope>provided</scope>
<artifactId>asie-shared</artifactId>
<version>${project.version}</version>
</dependency>
<!-- Provided by cxf-jaxrs-platform-module, but packaged due to solr-api -->
@@ -80,13 +75,13 @@
<dependency>
<groupId>com.fasterxml.jackson.datatype</groupId>
<artifactId>jackson-datatype-jsr310</artifactId>
<version>2.17.3</version>
<version>${jackson.version}</version>
<scope>provided</scope>
</dependency>
<dependency>
<groupId>com.fasterxml.jackson.module</groupId>
<artifactId>jackson-module-jakarta-xmlbind-annotations</artifactId>
<version>2.17.2</version>
<version>${jackson.version}</version>
<scope>provided</scope>
</dependency>
@@ -108,7 +103,6 @@
<plugin>
<groupId>io.repaint.maven</groupId>
<artifactId>tiles-maven-plugin</artifactId>
<version>2.40</version>
<extensions>true</extensions>
<configuration>
<tiles>
@@ -264,7 +264,7 @@ public class ShardDiscoveryService implements com.inteligr8.alfresco.asie.spi.Sh
ShardInstanceState nodeShardState = ShardInstanceState.from(shardState);
Pair<SolrHost, ShardInstanceState> pair = new Pair<>(node, nodeShardState);
if (comparator.compare(pair, shardNodeStates.get(shardId)) < 0)
if (!shardNodeStates.containsKey(shardId) || comparator.compare(pair, shardNodeStates.get(shardId)) < 0)
shardNodeStates.put(shardId, pair);
}
}
+30 -19
View File
@@ -43,15 +43,21 @@
<maven.compiler.target>17</maven.compiler.target>
<maven.compiler.release>17</maven.compiler.release>
<maven.deploy.skip>true</maven.deploy.skip>
<!-- must be aligned with the target ACS platform junit -->
<junit.version>5.12.2</junit.version>
<commom-rest-client.base.version>3.0.4</commom-rest-client.base.version>
<alfresco.platform.version>26.1.0</alfresco.platform.version>
</properties>
<dependencyManagement>
<dependencies>
<!-- Provided by cxf-jaxrs-platform-module, but packaged due to shared -->
<dependency>
<groupId>com.inteligr8</groupId>
<artifactId>common-rest-client</artifactId>
<version>3.0.3-cxf</version>
<groupId>org.junit.jupiter</groupId>
<artifactId>junit-jupiter-api</artifactId>
<version>${junit.version}</version>
</dependency>
</dependencies>
</dependencyManagement>
@@ -59,43 +65,48 @@
<build>
<pluginManagement>
<plugins>
<!-- avoids log4j dependency -->
<plugin>
<artifactId>maven-compiler-plugin</artifactId>
<version>3.14.1</version>
</plugin>
<!-- avoids struts dependency -->
<!-- helps avoid vulnerable dependencies -->
<plugin>
<artifactId>maven-site-plugin</artifactId>
<version>3.21.0</version>
<version>3.22.0</version>
</plugin>
<plugin>
<artifactId>maven-compiler-plugin</artifactId>
<version>3.15.0</version>
</plugin>
<!-- Force use of a new maven-dependency-plugin that doesn't download struts dependency -->
<plugin>
<artifactId>maven-dependency-plugin</artifactId>
<version>3.9.0</version>
<version>3.11.0</version>
</plugin>
<plugin>
<artifactId>maven-surefire-plugin</artifactId>
<version>3.5.4</version>
<version>3.5.6</version>
<dependencies>
<dependency>
<groupId>org.junit.jupiter</groupId>
<artifactId>junit-jupiter-engine</artifactId>
<version>5.14.0</version>
<version>${junit.version}</version>
</dependency>
</dependencies>
</plugin>
<plugin>
<artifactId>maven-failsafe-plugin</artifactId>
<version>3.5.4</version>
<version>3.5.6</version>
<dependencies>
<dependency>
<groupId>org.junit.jupiter</groupId>
<artifactId>junit-jupiter-engine</artifactId>
<version>5.14.0</version>
<version>${junit.version}</version>
</dependency>
</dependencies>
</plugin>
<plugin>
<groupId>io.repaint.maven</groupId>
<artifactId>tiles-maven-plugin</artifactId>
<version>2.45</version>
</plugin>
</plugins>
</pluginManagement>
</build>
@@ -151,7 +162,7 @@
<plugin>
<groupId>org.sonatype.central</groupId>
<artifactId>central-publishing-maven-plugin</artifactId>
<version>0.8.0</version>
<version>0.10.0</version>
<extensions>true</extensions>
<configuration>
<publishingServerId>central</publishingServerId>
+6 -4
View File
@@ -16,8 +16,9 @@
<name>ASIE Shared Library for Platform Modules</name>
<properties>
<alfresco.sdk.version>4.9.0</alfresco.sdk.version>
<alfresco.platform.version>23.3.0</alfresco.platform.version>
<alfresco.sdk.version>4.16.0</alfresco.sdk.version>
<common-rest-client.version>${commom-rest-client.base.version}-cxf</common-rest-client.version>
</properties>
<dependencyManagement>
@@ -36,11 +37,12 @@
<dependency>
<groupId>com.inteligr8.alfresco</groupId>
<artifactId>asie-api</artifactId>
<version>1.1-SNAPSHOT-asie2</version>
<version>1.2-SNAPSHOT-asie2</version>
</dependency>
<dependency>
<groupId>com.inteligr8</groupId>
<artifactId>common-rest-client</artifactId>
<version>${common-rest-client.version}</version>
</dependency>
<!-- Needed by this module, but provided by ACS -->
@@ -71,7 +73,7 @@
<repositories>
<repository>
<id>alfresco-public</id>
<url>https://artifacts.alfresco.com/nexus/content/groups/public</url>
<url>https://artifacts.alfresco.com/nexus/repository/releases/</url>
</repository>
</repositories>
</project>
@@ -9,14 +9,8 @@ import java.util.Map;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
import org.alfresco.model.ContentModel;
import org.alfresco.service.cmr.repository.InvalidNodeRefException;
import org.alfresco.service.cmr.repository.NodeRef;
import org.alfresco.service.cmr.repository.NodeService;
import org.alfresco.service.cmr.repository.StoreRef;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.extensions.webscripts.WebScriptException;
import org.springframework.extensions.webscripts.WebScriptRequest;
import org.springframework.extensions.webscripts.WebScriptResponse;
@@ -44,10 +44,10 @@ public class ReconcileAcsNodesWebScript extends AbstractAsieWebScript {
public void reconciled(long nodeDbId) {
if (includeReconciled) {
@SuppressWarnings("unchecked")
List<Long> unreconciledNodeDbIds = (List<Long>) responseMap.get("reconciled");
if (unreconciledNodeDbIds == null)
responseMap.put("reconciled", unreconciledNodeDbIds = new LinkedList<>());
unreconciledNodeDbIds.add(nodeDbId);
List<Long> reconciledNodeDbIds = (List<Long>) responseMap.get("reconciled");
if (reconciledNodeDbIds == null)
responseMap.put("reconciled", reconciledNodeDbIds = new LinkedList<>());
reconciledNodeDbIds.add(nodeDbId);
}
}
@@ -12,17 +12,12 @@ import java.util.concurrent.Future;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
import org.alfresco.model.ContentModel;
import org.alfresco.repo.index.shard.Floc;
import org.alfresco.repo.index.shard.Shard;
import org.alfresco.repo.index.shard.ShardInstance;
import org.alfresco.repo.index.shard.ShardRegistry;
import org.alfresco.repo.index.shard.ShardState;
import org.alfresco.service.cmr.repository.StoreRef;
import org.alfresco.service.cmr.search.SearchParameters;
import org.alfresco.service.cmr.search.SearchService;
import org.alfresco.service.namespace.NamespaceService;
import org.alfresco.service.namespace.QName;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
@@ -45,9 +40,6 @@ public abstract class AbstractActionService {
private final Logger logger = LoggerFactory.getLogger(this.getClass());
@Autowired
private NamespaceService namespaceService;
@Autowired
private ApiService apiService;
@@ -58,10 +50,10 @@ public abstract class AbstractActionService {
@Qualifier(Constants.QUALIFIER_ASIE)
private ShardRegistry shardRegistry;
@Value("${inteligr8.asie.default.concurrentQueueSize:64}")
@Value("${inteligr8.asie.default.concurrentQueueSize}")
private int concurrentQueueSize;
@Value("${inteligr8.asie.default.concurrency:16}")
@Value("${inteligr8.asie.default.concurrency}")
private int concurrency;
protected int getConcurrency() {
@@ -8,26 +8,22 @@ import java.util.Map.Entry;
import java.util.Set;
import java.util.concurrent.Callable;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Future;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
import java.util.regex.Matcher;
import org.alfresco.model.ContentModel;
import org.alfresco.repo.index.shard.Floc;
import org.alfresco.repo.index.shard.Shard;
import org.alfresco.repo.index.shard.ShardInstance;
import org.alfresco.repo.index.shard.ShardRegistry;
import org.alfresco.repo.index.shard.ShardState;
import org.alfresco.service.cmr.repository.NodeRef;
import org.alfresco.service.cmr.repository.NodeService;
import org.alfresco.service.cmr.repository.StoreRef;
import org.alfresco.service.cmr.search.SearchParameters;
import org.alfresco.service.cmr.search.SearchService;
import org.alfresco.service.namespace.NamespaceService;
import org.alfresco.service.namespace.QName;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.DisposableBean;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.beans.factory.annotation.Value;
@@ -44,16 +40,13 @@ import com.inteligr8.solr.model.Action;
import com.inteligr8.solr.model.ActionResponse;
import com.inteligr8.solr.model.BaseResponse;
public abstract class AbstractNodeActionService {
public abstract class AbstractNodeActionService implements DisposableBean {
private final Logger logger = LoggerFactory.getLogger(this.getClass());
@Autowired
private NamespaceService namespaceService;
@Autowired
private NodeService nodeService;
@Autowired
private ApiService apiService;
@@ -67,10 +60,10 @@ public abstract class AbstractNodeActionService {
@Qualifier(Constants.QUALIFIER_ASIE)
private ShardRegistry shardRegistry;
@Value("${inteligr8.asie.default.concurrentQueueSize:64}")
@Value("${inteligr8.asie.default.concurrentQueueSize}")
private int concurrentQueueSize;
@Value("${inteligr8.asie.default.concurrency:16}")
@Value("${inteligr8.asie.default.concurrency}")
private int concurrency;
protected int getConcurrency() {
@@ -85,6 +78,22 @@ public abstract class AbstractNodeActionService {
protected abstract String getActionName();
@Override
public void destroy() {
ExecutorService executor = this.executorManager.get(this.getThreadNamePrefix());
if (executor != null) {
this.logger.info("Shutting down throttled thread pool executor: {}", this.getThreadNamePrefix());
executor.shutdown();
}
}
private ThrottledThreadPoolExecutor getExecutor() {
return this.executorManager.createThrottled(
this.getThreadNamePrefix(),
this.getConcurrency(), this.getConcurrency(), this.getConcurrentQueueSize(),
1L, TimeUnit.MINUTES);
}
/**
* This method executes an action on the specified node in Solr using its
* ACS unique database identifier. The callback handles all the return
@@ -139,16 +148,15 @@ public abstract class AbstractNodeActionService {
}
}
private Future<Void> _action(long nodeDbId, ActionCallback callback, Long fullQueueExpireTimeMillis) throws TimeoutException, InterruptedException {
private Future<Void> _action(
final long nodeDbId,
final ActionCallback callback,
Long fullQueueExpireTimeMillis) throws TimeoutException, InterruptedException {
List<com.inteligr8.alfresco.asie.model.ShardInstance> eligibleInstances = this.findPossibleShardInstances(nodeDbId);
this.logger.debug("Will attempt to {} ACS node against {} shard instances: {}", this.getActionName(), eligibleInstances.size(), nodeDbId);
CompositeFuture<Void> future = new CompositeFuture<>();
ThrottledThreadPoolExecutor executor = this.executorManager.createThrottled(
this.getThreadNamePrefix(),
this.getConcurrency(), this.getConcurrency(), this.getConcurrentQueueSize(),
1L, TimeUnit.MINUTES);
ThrottledThreadPoolExecutor executor = this.getExecutor();
for (final com.inteligr8.alfresco.asie.model.ShardInstance instance : eligibleInstances) {
this.logger.trace("Will attempt to {} ACS node against shard instance: {}: {}", this.getActionName(), nodeDbId, instance);
@@ -156,49 +164,12 @@ public abstract class AbstractNodeActionService {
Callable<Void> callable = new Callable<>() {
@Override
public Void call() {
String core = instance.extractShard().getCoreName();
SolrHost host = instance.extractNode();
URL url = host.toUrl(apiService.isSecure() ? "https" : "http");
CoreAdminApi api = apiService.createApi(url.toString(), CoreAdminApi.class);
try {
logger.debug("Performing {} of ACS node against shard instance: {}: {}", getActionName(), nodeDbId, instance);
BaseResponse apiResponse = execute(api, core, nodeDbId);
logger.trace("Performed {} of ACS node against shard instance: {}: {}", getActionName(), nodeDbId, instance);
Action action = null;
if (apiResponse instanceof ActionCoreResponse<?>) {
action = ((ActionCoreResponse<Action>) apiResponse).getCores().getByCore(core);
} else if (apiResponse instanceof ActionResponse<?>) {
action = ((ActionResponse<Action>) apiResponse).getAction();
}
if (action == null) {
callback.unknownResult(instance);
} else {
switch (action.getStatus()) {
case Scheduled:
callback.scheduled(instance);
break;
case Success:
callback.success(instance);
break;
default:
if (apiResponse instanceof com.inteligr8.alfresco.asie.model.BaseResponse) {
com.inteligr8.alfresco.asie.model.BaseResponse asieResponse = (com.inteligr8.alfresco.asie.model.BaseResponse) apiResponse;
logger.debug("Performance of {} of ACS node against shard instance failed: {}: {}: {}", getActionName(), nodeDbId, instance, asieResponse.getException());
callback.error(instance, asieResponse.getException());
} else {
logger.debug("Performance of {} of ACS node against shard instance failed: {}: {}: {}", getActionName(), nodeDbId, instance, apiResponse.getResponseHeader().getStatus());
callback.error(instance, String.valueOf(apiResponse.getResponseHeader().getStatus()));
}
}
}
actionToShard(nodeDbId, callback, instance);
} catch (Exception e) {
logger.error("An exception occurred", e);
logger.error("An unexpected exception occurred", e);
callback.error(instance, e.getMessage());
}
return null;
}
};
@@ -213,15 +184,54 @@ public abstract class AbstractNodeActionService {
return future;
}
@SuppressWarnings("unchecked")
protected void actionToShard(long nodeDbId, ActionCallback callback, com.inteligr8.alfresco.asie.model.ShardInstance instance) {
String core = instance.extractShard().getCoreName();
SolrHost host = instance.extractNode();
URL url = host.toUrl(this.apiService.isSecure() ? "https" : "http");
CoreAdminApi api = this.apiService.createApi(url.toString(), CoreAdminApi.class);
this.logger.debug("Performing {} of ACS node against shard instance: {}: {}", this.getActionName(), nodeDbId, instance);
BaseResponse apiResponse = execute(api, core, nodeDbId);
this.logger.trace("Performed {} of ACS node against shard instance: {}: {}", this.getActionName(), nodeDbId, instance);
Action action = null;
if (apiResponse instanceof ActionCoreResponse<?>) {
action = ((ActionCoreResponse<Action>) apiResponse).getCores().getByCore(core);
} else if (apiResponse instanceof ActionResponse<?>) {
action = ((ActionResponse<Action>) apiResponse).getAction();
}
if (action == null) {
callback.unknownResult(instance);
} else {
switch (action.getStatus()) {
case Scheduled:
callback.scheduled(instance);
break;
case Success:
callback.success(instance);
break;
default:
if (apiResponse instanceof com.inteligr8.alfresco.asie.model.BaseResponse) {
com.inteligr8.alfresco.asie.model.BaseResponse asieResponse = (com.inteligr8.alfresco.asie.model.BaseResponse) apiResponse;
this.logger.debug("Performance of {} of ACS node against shard instance failed: {}: {}: {}",
this.getActionName(), nodeDbId, instance, asieResponse.getException());
callback.error(instance, asieResponse.getException());
} else {
this.logger.debug("Performance of {} of ACS node against shard instance failed: {}: {}: {}",
this.getActionName(), nodeDbId, instance, apiResponse.getResponseHeader().getStatus());
callback.error(instance, String.valueOf(apiResponse.getResponseHeader().getStatus()));
}
}
}
}
protected abstract BaseResponse execute(CoreAdminApi api, String core, long nodeDbId);
private List<com.inteligr8.alfresco.asie.model.ShardInstance> findPossibleShardInstances(long nodeDbId) {
if (this.shardRegistry == null)
throw new UnsupportedOperationException("ACS instances without a sharding configuration are not yet implemented");
SearchParameters searchParams = new SearchParameters();
searchParams.setLanguage(SearchService.LANGUAGE_FTS_ALFRESCO);
searchParams.setQuery("@" + this.formatForFts(ContentModel.PROP_NODE_DBID) + ":" + nodeDbId);
List<com.inteligr8.alfresco.asie.model.ShardInstance> instances = new LinkedList<>();
@@ -233,7 +243,7 @@ public abstract class AbstractNodeActionService {
if (!floc.getKey().getStoreRefs().contains(StoreRef.STORE_REF_WORKSPACE_SPACESSTORE))
continue;
Integer shardHash = null;
int shardHash = -1;
ShardSet shardset = null;
for (Entry<Shard, Set<ShardState>> shard : floc.getValue().entrySet()) {
for (ShardState shardState : shard.getValue()) {
@@ -242,26 +252,10 @@ public abstract class AbstractNodeActionService {
// so we are computing the shardHash once and caching it
if (shardset == null) {
shardset = ShardSet.from(floc.getKey(), shardState);
switch (shardset.getMethod()) {
case PROPERTY:
this.logger.trace("Using property-based sharding method; discovering target shard ...");
NodeRef nodeRef = this.nodeService.getNodeRef(nodeDbId);
Object propValue = this.nodeService.getProperty(nodeRef, QName.createQName(shardset.getPrefixedProperty(), this.namespaceService));
if (propValue != null) {
this.logger.trace("Discovered node property for sharding: {} <=> {} => {}", nodeDbId, nodeRef, propValue);
Matcher matcher = shardset.getRegex().matcher(propValue.toString());
String hashable = matcher.group(1);
this.logger.trace("Extracted shardable value from node: {} <=> {} => {}", nodeDbId, propValue, hashable);
shardHash = this.shardHashService.hash(hashable, shardset.getShards().intValue());
this.logger.debug("Hash shardable value to shard instance ID: {} <=> {} => {}", nodeDbId, hashable, shardHash);
}
break;
default:
this.logger.trace("Despite sharding, considering all shards and nodes without optimization: {}: {}", nodeDbId, instances);
}
shardHash = this.shardHashService.computeShardInstanceId(shardset, nodeDbId);
}
if (shardHash != null) {
if (shardHash >= 0) {
if (shard.getKey().getInstance() == shardHash)
instances.add(this.toModel(shardState.getShardInstance(), shardState));
} else {
@@ -1,11 +1,10 @@
package com.inteligr8.alfresco.asie.service;
import java.util.HashMap;
import java.util.HashSet;
import java.util.Map;
import java.util.Collections;
import java.util.Set;
import java.util.concurrent.Callable;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Future;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
@@ -26,7 +25,6 @@ import org.apache.commons.collections4.SetUtils;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.DisposableBean;
import org.springframework.beans.factory.InitializingBean;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Component;
@@ -39,7 +37,7 @@ import com.inteligr8.alfresco.asie.util.CompositeFuture;
import com.inteligr8.alfresco.asie.util.ThrottledThreadPoolExecutor;
@Component
public class AcsReconcileService implements InitializingBean, DisposableBean {
public class AcsReconcileService implements DisposableBean {
private final Logger logger = LoggerFactory.getLogger(this.getClass());
private final Logger reconcileLogger = LoggerFactory.getLogger("inteligr8.asie.reconcile");
@@ -61,30 +59,42 @@ public class AcsReconcileService implements InitializingBean, DisposableBean {
@Autowired
private ReindexService reindexService;
@Autowired
private ExecutorManager executorManager;
@Value("${inteligr8.asie.reconciliation.nodesChunkSize:250}")
@Value("${inteligr8.asie.reconciliation.nodesChunkSize}")
private int nodesChunkSize;
@Value("${inteligr8.asie.reconciliation.nodeTimeoutSeconds:10}")
@Value("${inteligr8.asie.reconciliation.nodeTimeoutSeconds}")
private int nodeTimeoutSeconds;
@Value("${inteligr8.asie.reconciliation.concurrentQueueSize:64}")
@Value("${inteligr8.asie.reconciliation.concurrentQueueSize}")
private int concurrentQueueSize;
@Value("${inteligr8.asie.reconciliation.concurrency:2}")
@Value("${inteligr8.asie.reconciliation.concurrency}")
private int concurrency;
private ThrottledThreadPoolExecutor executor;
@Override
public void afterPropertiesSet() {
this.executor = new ThrottledThreadPoolExecutor(this.concurrency, this.concurrency, this.concurrentQueueSize, 1L, TimeUnit.MINUTES, "solr-reconcile");
this.executor.prestartAllCoreThreads();
}
@Value("${inteligr8.asie.reconciliation.waitAfterSolrNodeActionMillis}")
private long waitAfterSolrNodeActionMillis;
@Value("${inteligr8.asie.reconciliation.waitAfterSolrNodeReconcileMillis}")
private long waitAfterSolrNodeReconcileMillis;
@Override
public void destroy() {
this.executor.shutdown();
ExecutorService executor = this.executorManager.get("solr-reconcile");
if (executor != null) {
this.logger.info("Shutting down throttled thread pool executor: {}", "solr-reconcile");
executor.shutdown();
}
}
private ThrottledThreadPoolExecutor getExecutor() {
return this.executorManager.createThrottled(
"solr-reconcile",
this.concurrency, this.concurrency, this.concurrentQueueSize,
10L, TimeUnit.SECONDS);
}
/**
@@ -96,7 +106,7 @@ public class AcsReconcileService implements InitializingBean, DisposableBean {
*
* There are two sets of parameters regarding timeouts. The queue timeouts
* are for how long the requesting thread should wait for a full queue to
* open up space for new re-index executions. The execution timeouts are
* open up space for new reconcile executions. The execution timeouts are
* for how long the execution should be allowed to take once dequeued.
* There is no timeout for how long the execution is queued.
*
@@ -107,10 +117,10 @@ public class AcsReconcileService implements InitializingBean, DisposableBean {
* @param callback A callback to process multiple returned values from the re-index.
* @param queueTimeout A timeout for how long the calling thread should wait for space on the queue.
* @param queueUnit The time units for the `queueTimeout`.
* @param execTimeout A timeout for the elapsed time the reindex execution should take when dequeued.
* @param execTimeout A timeout for the elapsed time the reconcile execution should take when dequeued.
* @param execUnit The time units for the `execTimeout`.
* @throws TimeoutException Either the queue or execution timeout lapsed.
* @throws InterruptedException The re-index was interrupted (server shutdown).
* @throws InterruptedException The reconciliation was interrupted (server shutdown).
*/
public void reconcile(
long fromDbId, long toDbId, Integer nodesChunkSize,
@@ -119,18 +129,10 @@ public class AcsReconcileService implements InitializingBean, DisposableBean {
ReconcileCallback callback,
long queueTimeout, TimeUnit queueUnit,
long execTimeout, TimeUnit execUnit) throws InterruptedException, TimeoutException {
if (nodesChunkSize == null)
nodesChunkSize = this.nodesChunkSize;
if (this.logger.isTraceEnabled())
this.logger.trace("reconcile({}, {}, {}, {}, {}, {}, {})", fromDbId, toDbId, nodesChunkSize, indexUnreconciled, reindexReconciled, queueUnit.toMillis(queueTimeout), execUnit.toMillis(execTimeout));
CompositeFuture<Void> future = new CompositeFuture<>();
for (long startDbId = fromDbId; startDbId < toDbId; startDbId += nodesChunkSize) {
long endDbId = Math.min(toDbId, startDbId + nodesChunkSize);
future.combine(this.reconcileChunk(startDbId, endDbId, indexUnreconciled, reindexReconciled, callback, queueTimeout, queueUnit, execTimeout, execUnit));
future.purge(true);
}
Future<Void> future = this._reconcile(fromDbId, toDbId, nodesChunkSize, indexUnreconciled, reindexReconciled, callback, queueTimeout, queueUnit, execTimeout, execUnit);
try {
future.get(execTimeout, execUnit);
@@ -139,27 +141,52 @@ public class AcsReconcileService implements InitializingBean, DisposableBean {
}
}
/**
* This method reconciles the specified node range between ACS and Solr.
* The node range is specified using the ACS unique database identifiers.
* There is no other reasonably efficient attack vector. The callback
* handles all the return values. This is the synchronous alternative to
* the other `reconcile` method.
*
* @param fromDbId A node database ID, inclusive.
* @param toDbId A node database ID, exclusive.
* @param indexUnreconciled For nodes not found in Solr, attempt to index against all applicable Solr instances.
* @param reindexReconciled For nodes found in Solr, attempt to re-index against all applicable Solr instances.
* @param callback A callback to process multiple returned values from the re-index.
* @throws InterruptedException The reconciliation was interrupted (server shutdown).
*/
public Future<Void> reconcile(
long fromDbId, long toDbId, Integer nodesChunkSize,
boolean indexUnreconciled,
boolean reindexReconciled,
ReconcileCallback callback) throws InterruptedException {
if (nodesChunkSize == null)
nodesChunkSize = this.nodesChunkSize;
this.logger.trace("reconcile({}, {}, {}, {}, {})", fromDbId, toDbId, nodesChunkSize, indexUnreconciled, reindexReconciled);
CompositeFuture<Void> future = new CompositeFuture<>();
try {
for (long startDbId = fromDbId; startDbId < toDbId; startDbId += nodesChunkSize) {
long endDbId = Math.min(toDbId, startDbId + nodesChunkSize);
future.combine(this.reconcileChunk(startDbId, endDbId, indexUnreconciled, reindexReconciled, callback, -1L, null, -1L, null));
future.purge(true);
}
return this._reconcile(fromDbId, toDbId, nodesChunkSize, indexUnreconciled, reindexReconciled, callback, -1L, null, -1L, null);
} catch (TimeoutException te) {
throw new RuntimeException("This should never happen: " + te.getMessage(), te);
}
}
protected Future<Void> _reconcile(
long fromDbId, long toDbId, Integer nodesChunkSize,
boolean indexUnreconciled,
boolean reindexReconciled,
ReconcileCallback callback,
long queueTimeout, TimeUnit queueUnit,
long execTimeout, TimeUnit execUnit) throws InterruptedException, TimeoutException {
if (nodesChunkSize == null)
nodesChunkSize = this.nodesChunkSize;
CompositeFuture<Void> future = new CompositeFuture<>();
for (long startDbId = fromDbId; startDbId < toDbId; startDbId += nodesChunkSize) {
long endDbId = Math.min(toDbId, startDbId + nodesChunkSize);
future.combine(this.reconcileChunk(startDbId, endDbId, indexUnreconciled, reindexReconciled, callback, queueTimeout, queueUnit, execTimeout, execUnit));
future.purge(true);
}
return future;
}
@@ -205,12 +232,14 @@ public class AcsReconcileService implements InitializingBean, DisposableBean {
}
CompositeFuture<Void> future = new CompositeFuture<>();
ThrottledThreadPoolExecutor executor = this.getExecutor();
for (long _nodeDbId = fromDbId; _nodeDbId < toDbId; _nodeDbId++) {
final long nodeDbId = _nodeDbId;
this.logger.trace("Attempting to reconcile ACS node: {}", nodeDbId);
Callable<Void> callable;
boolean callingSolr = false;
final int dbIdIndex = (int) (nodeDbId - fromDbId);
if (nodeRefs[dbIdIndex] != null) {
@@ -222,144 +251,157 @@ public class AcsReconcileService implements InitializingBean, DisposableBean {
callback.reconciled(nodeDbId);
if (reindexReconciled)
reindex(nodeDbId, nodeRefs[dbIdIndex], callback, execTimeout, execUnit);
reindex(nodeDbId, nodeRefs[dbIdIndex], callback);
// purposefully forgetting about the returned future
// its results will be logged
// the reconcile thread will continue independently
// the callback will lag
return null;
}
};
if (reindexReconciled)
callingSolr = true;
} else {
callable = new Callable<Void>() {
@Override
public Void call() throws InterruptedException, TimeoutException {
reconcile(nodeDbId, indexUnreconciled, callback, execTimeout, execUnit);
reconcile(nodeDbId, indexUnreconciled, callback);
return null;
}
};
callingSolr = true;
}
if (queueTimeout < 0L) {
future.combine(this.executor.submit(callable, -1L, null));
future.combine(executor.submit(callable, -1L, null));
} else {
future.combine(this.executor.submit(callable, queueTimeout, queueUnit));
future.combine(executor.submit(callable, queueTimeout, queueUnit));
}
if (callingSolr && this.waitAfterSolrNodeReconcileMillis > 0L) {
this.logger.trace("Waiting between each node reconcile");
Thread.sleep(this.waitAfterSolrNodeReconcileMillis);
}
}
return future;
}
public void reconcile(long nodeDbId,
public boolean reconcile(long nodeDbId,
boolean index,
ReconcileCallback callback,
long execTimeout, TimeUnit execUnit) throws InterruptedException, TimeoutException {
ReconcileCallback callback) throws InterruptedException, TimeoutException {
NodeRef nodeRef = this.nodeService.getNodeRef(nodeDbId);
if (nodeRef == null) {
this.logger.trace("No such ACS node: {}; skipping ...", nodeDbId);
return;
return false;
}
if (!StoreRef.STORE_REF_WORKSPACE_SPACESSTORE.equals(nodeRef.getStoreRef())) {
this.logger.trace("A deliberately ignored store in the DB is not indexed in Solr: {}: {}", nodeDbId, nodeRef);
return;
return false;
}
Set<QName> aspects = this.nodeService.getAspects(nodeRef);
aspects.retainAll(this.ignoreNodesWithAspects);
if (!aspects.isEmpty()) {
this.logger.trace("A deliberately ignored node in the DB is not indexed in Solr: {}: {}: {}", nodeDbId, nodeRef, aspects);
return;
return false;
}
if (!index) {
this.logger.debug("A node in the DB is not indexed in Solr: {}: {}", nodeDbId, nodeRef);
this.reconcileLogger.info("UNRECONCILED: {} <=> {}", nodeDbId, nodeRef);
callback.unreconciled(nodeDbId);
return false;
} else {
logger.debug("A node in the DB is not indexed in Solr; attempt to index: {}: {}", nodeDbId, nodeRef);
this.index(nodeDbId, nodeRef, callback, execTimeout, execUnit);
this.logger.debug("A node in the DB is not indexed in Solr; attempt to index: {}: {}", nodeDbId, nodeRef);
this.index(nodeDbId, nodeRef, callback);
// purposefully forgetting about the returned future
// its results will be logged
// the reconcile thread will continue independently
// the callback will lag
return true;
}
}
public void index(long nodeDbId, NodeRef nodeRef,
ReconcileCallback callback,
long execTimeout, TimeUnit execUnit) throws InterruptedException, TimeoutException {
Set<ShardInstance> syncHosts = new HashSet<>();
Set<ShardInstance> asyncHosts = new HashSet<>();
Map<ShardInstance, String> errorHosts = new HashMap<>();
public Future<Void> index(long nodeDbId, NodeRef nodeRef, ReconcileCallback callback) throws InterruptedException, TimeoutException {
IndexCallback indexCallback = new IndexCallback() {
@Override
public void success(ShardInstance instance) {
reconcileLogger.info("INDEXED: {} <=> {} in {}", nodeDbId, nodeRef, instance);
syncHosts.add(instance);
if (callback != null)
callback.processed(nodeDbId,
Collections.singleton(instance), Collections.emptySet(),
Collections.emptyMap());
}
@Override
public void scheduled(ShardInstance instance) {
reconcileLogger.info("INDEXING: {} <=> {} in {}", nodeDbId, nodeRef, instance);
asyncHosts.add(instance);
if (callback != null)
callback.processed(nodeDbId,
Collections.emptySet(), Collections.singleton(instance),
Collections.emptyMap());
}
@Override
public void error(ShardInstance instance, String message) {
reconcileLogger.info("FAILED INDEX: {} <=> {} in {}", nodeDbId, nodeRef, instance);
errorHosts.put(instance, message);
if (callback != null)
callback.processed(nodeDbId,
Collections.emptySet(), Collections.emptySet(),
Collections.singletonMap(instance, message));
}
};
try {
if (execTimeout < 0L) {
this.indexService.index(nodeDbId, indexCallback).get();
} else {
this.indexService.index(nodeDbId, indexCallback).get(execTimeout, execUnit);
}
} catch (ExecutionException ee) {
throw new RuntimeException("An unexpected exception occurred: " + ee.getMessage(), ee);
}
Future<Void> future = this.indexService.index(nodeDbId, indexCallback);
if (callback != null)
callback.processed(nodeDbId, syncHosts, asyncHosts, errorHosts);
if (this.waitAfterSolrNodeActionMillis > 0L) {
Thread.sleep(this.waitAfterSolrNodeActionMillis);
}
return future;
}
public void reindex(long nodeDbId, NodeRef nodeRef,
ReconcileCallback callback,
long execTimeout, TimeUnit execUnit) throws InterruptedException, TimeoutException {
Set<ShardInstance> syncHosts = new HashSet<>();
Set<ShardInstance> asyncHosts = new HashSet<>();
Map<ShardInstance, String> errorHosts = new HashMap<>();
public Future<Void> reindex(long nodeDbId, NodeRef nodeRef, ReconcileCallback callback) throws InterruptedException {
ReindexCallback reindexCallback = new ReindexCallback() {
@Override
public void success(ShardInstance instance) {
reconcileLogger.info("REINDEXED: {} <=> {} in {}", nodeDbId, nodeRef, instance);
syncHosts.add(instance);
if (callback != null)
callback.processed(nodeDbId,
Collections.singleton(instance), Collections.emptySet(),
Collections.emptyMap());
}
@Override
public void scheduled(ShardInstance instance) {
reconcileLogger.info("REINDEXING: {} <=> {} in {}", nodeDbId, nodeRef, instance);
asyncHosts.add(instance);
if (callback != null)
callback.processed(nodeDbId,
Collections.emptySet(), Collections.singleton(instance),
Collections.emptyMap());
}
@Override
public void error(ShardInstance instance, String message) {
reconcileLogger.info("FAILED REINDEX: {} <=> {} in {}", nodeDbId, nodeRef, instance);
errorHosts.put(instance, message);
if (callback != null)
callback.processed(nodeDbId,
Collections.emptySet(), Collections.emptySet(),
Collections.singletonMap(instance, message));
}
};
Future<Void> future = this.reindexService.reindex(nodeDbId, reindexCallback);
try {
if (execTimeout < 0L) {
this.reindexService.reindex(nodeDbId, reindexCallback).get();
} else {
this.reindexService.reindex(nodeDbId, reindexCallback).get(execTimeout, execUnit);
}
} catch (ExecutionException ee) {
throw new RuntimeException("An unexpected exception occurred: " + ee.getMessage(), ee);
if (this.waitAfterSolrNodeActionMillis > 0L) {
Thread.sleep(this.waitAfterSolrNodeActionMillis);
}
if (callback != null)
callback.processed(nodeDbId, syncHosts, asyncHosts, errorHosts);
return future;
}
private String formatForFts(QName qname) {
@@ -43,7 +43,7 @@ public class ApiService implements InitializingBean {
@Value("${inteligr8.asie.basePath}")
private String solrBaseUrl;
@Value("${inteligr8.asie.reconciliation.nodesChunkSize:250}")
@Value("${inteligr8.asie.reconciliation.nodesChunkSize}")
private int nodesChunkSize;
@Override
@@ -6,6 +6,8 @@ import java.util.concurrent.ExecutorService;
import java.util.concurrent.RejectedExecutionHandler;
import java.util.concurrent.TimeUnit;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.DisposableBean;
import org.springframework.beans.factory.InitializingBean;
import org.springframework.beans.factory.annotation.Value;
@@ -30,42 +32,36 @@ import com.inteligr8.alfresco.asie.util.ThrottledThreadPoolExecutor;
*/
@Component
public class ExecutorManager implements InitializingBean, DisposableBean, RemovalListener<String, ExecutorService> {
@Value("${inteligr8.asie.executors.expireTimeInMinutes:30}")
private final Logger logger = LoggerFactory.getLogger(this.getClass());
@Value("${inteligr8.asie.executors.expireTimeInMinutes}")
private int expireTimeInMinutes;
private Cache<String, ExecutorService> refCache;
private Cache<String, ExecutorService> expiringCache;
private Cache<String, ThrottledThreadPoolExecutor> cache;
@Override
public void afterPropertiesSet() throws Exception {
// a weak value happens when the executor is no longer referenced
// the possible references are by the caller temporarily using and the `expiringCache` (below; so it expired)
// this keeps the pool from being shutdown after it expires if the caller is still referencing it
// ultimately, if it is cached, it will be in this cache and MAY be in the `expiringCache`.
this.refCache = CacheBuilder.newBuilder()
.initialCapacity(8)
.weakValues()
.removalListener(this)
.build();
this.expiringCache = CacheBuilder.newBuilder()
this.cache = CacheBuilder.newBuilder()
.initialCapacity(8)
.expireAfterAccess(this.expireTimeInMinutes, TimeUnit.MINUTES)
.removalListener(this)
.build();
}
@Override
public void destroy() throws Exception {
this.refCache.invalidateAll();
this.refCache.cleanUp();
this.expiringCache.invalidateAll();
this.expiringCache.cleanUp();
this.cache.invalidateAll();
this.cache.cleanUp();
}
@Override
public void onRemoval(RemovalNotification<String, ExecutorService> notification) {
notification.getValue().shutdown();
this.logger.debug("Throttled thread pool removed/expired from cache: {}", notification.getKey());
if (!notification.getValue().isShutdown()) {
notification.getValue().shutdown();
this.logger.info("Throttled thread pool shut down: {}", notification.getKey());
}
}
public ThrottledThreadPoolExecutor createThrottled(
@@ -88,9 +84,10 @@ public class ExecutorManager implements InitializingBean, DisposableBean, Remova
final RejectedExecutionHandler rejectedExecutionHandler) {
try {
// if it is already cached, reuse the cache; otherwise create one
final ExecutorService executor = this.refCache.get(name, new Callable<ThrottledThreadPoolExecutor>() {
return this.cache.get(name, new Callable<ThrottledThreadPoolExecutor>() {
@Override
public ThrottledThreadPoolExecutor call() {
logger.info("Creating throttled thread pool: {}", name);
ThrottledThreadPoolExecutor executor = null;
if (rejectedExecutionHandler == null) {
executor = new ThrottledThreadPoolExecutor(coreThreadPoolSize, maximumThreadPoolSize, maximumQueueSize,
@@ -103,14 +100,9 @@ public class ExecutorManager implements InitializingBean, DisposableBean, Remova
rejectedExecutionHandler);
}
logger.debug("Created throttled thread pool: {}; threads: {}; queue: {}", name, maximumThreadPoolSize, maximumQueueSize);
executor.prestartAllCoreThreads();
return executor;
}
});
return (ThrottledThreadPoolExecutor) this.expiringCache.get(name, new Callable<ExecutorService>() {
@Override
public ExecutorService call() throws Exception {
logger.trace("Started {} core threads in thread pool: {}", coreThreadPoolSize, name);
return executor;
}
});
@@ -120,19 +112,7 @@ public class ExecutorManager implements InitializingBean, DisposableBean, Remova
}
public ExecutorService get(String name) {
// grab from the expiring cache first, so we can
ExecutorService executor = this.expiringCache.getIfPresent(name);
if (executor != null)
return executor;
executor = this.refCache.getIfPresent(name);
if (executor == null)
return null;
// the executor expired, but it was still referenced by the caller
// re-cache it
this.expiringCache.put(name, executor);
return executor;
return this.cache.getIfPresent(name);
}
}
@@ -1,14 +1,36 @@
package com.inteligr8.alfresco.asie.service;
import java.nio.charset.Charset;
import java.time.Instant;
import java.time.temporal.TemporalAccessor;
import java.util.Date;
import java.util.regex.Matcher;
import org.alfresco.error.AlfrescoRuntimeException;
import org.alfresco.service.cmr.repository.NodeRef;
import org.alfresco.service.cmr.repository.NodeService;
import org.alfresco.service.namespace.NamespaceService;
import org.alfresco.service.namespace.QName;
import org.apache.commons.codec.digest.MurmurHash3;
import org.joda.time.ReadablePartial;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
import com.inteligr8.alfresco.asie.model.ShardSet;
@Component
public class SolrShardHashService {
private final Charset charset = Charset.forName("utf-8");
private final Logger logger = LoggerFactory.getLogger(this.getClass());
@Autowired
private NamespaceService namespaceService;
@Autowired
private NodeService nodeService;
public int hash(Object obj, int shardCount) {
String str = obj.toString();
@@ -23,5 +45,67 @@ public class SolrShardHashService {
hash.add(bytes, 0, bytes.length);
return Math.abs(hash.end()) % shardCount;
}
/**
* @param shardset A ShardSet.
* @param nodeDbId The numeric database ID of an ACS node.
* @return A shard instance from 0 to 1 less than the number of shards; -1 if unable to hash.
*/
public int computeShardInstanceId(ShardSet shardset, long nodeDbId) {
NodeRef nodeRef = this.nodeService.getNodeRef(nodeDbId);
if (nodeRef == null)
throw new AlfrescoRuntimeException("The node " + nodeDbId + " does not exist");
return this.computeShardInstanceId(shardset, nodeRef);
}
/**
* @param shardset A ShardSet.
* @param nodeRef A reference to an ACS node.
* @return A shard instance from 0 to 1 less than the number of shards; -1 if unable to hash.
*/
public int computeShardInstanceId(ShardSet shardset, NodeRef nodeRef) {
switch (shardset.getMethod()) {
case PROPERTY:
QName hashableProperty = QName.createQName(shardset.getPrefixedProperty(), this.namespaceService);
Object fullPropertyValue = this.nodeService.getProperty(nodeRef, hashableProperty);
if (fullPropertyValue == null) {
this.logger.debug("Unable to determine shard instance ID because property does not exist on node: {}: {}", nodeRef, hashableProperty);
return -1;
}
this.logger.trace("Discovered node property for sharding: {} => {}", nodeRef, fullPropertyValue);
String hashableValue = null;
if (fullPropertyValue instanceof TemporalAccessor) {
Instant instant = Instant.from((TemporalAccessor) fullPropertyValue);
hashableValue = instant.toString();
} else if (fullPropertyValue instanceof ReadablePartial) {
hashableValue = ((ReadablePartial) fullPropertyValue).toString();
} else if (fullPropertyValue instanceof Date) {
Instant instant = ((Date) fullPropertyValue).toInstant();
hashableValue = instant.toString();
} else {
hashableValue = fullPropertyValue.toString();
}
if (shardset.getRegex() != null) {
Matcher matcher = shardset.getRegex().matcher(hashableValue);
if (!matcher.find()) {
this.logger.debug("Unable to determine shard instance ID because regex pattern doesn't match the hashing property value: {}: {}: {}", nodeRef, shardset.getRegex(), hashableValue);
return -1;
}
hashableValue = matcher.group(1);
this.logger.trace("Extracted shardable value from node: {}: {} => {}", nodeRef, fullPropertyValue, hashableValue);
}
int shardHash = this.hash(hashableValue, shardset.getShards().intValue());
this.logger.debug("Hash shardable value to shard instance ID: {}: {} => {}", nodeRef, hashableValue, shardHash);
return shardHash;
// TODO replicate hash algorithm for other shard methods
default:
this.logger.trace("Unable to determine shard instance ID due to shard method: {}", shardset.getMethod());
return -1;
}
}
}
@@ -125,7 +125,9 @@ public interface ShardDiscoveryService {
public class ShardedNodeShardStateComparator implements Comparator<Pair<SolrHost, ShardInstanceState>> {
@Override
public int compare(Pair<SolrHost, ShardInstanceState> p1, Pair<SolrHost, ShardInstanceState> p2) {
return - Long.compare(p1.getSecond().getLastIndexedTxId(), p2.getSecond().getLastIndexedTxId());
if (p1 == null) return 1;
else if (p2 == null) return -1;
else return - Long.compare(p1.getSecond().getLastIndexedTxId(), p2.getSecond().getLastIndexedTxId());
}
}
@@ -79,12 +79,14 @@ public class CompositeFuture<T> implements Future<T> {
List<T> results = new ArrayList<>(this.futures.size());
for (Future<T> future : this.futures) {
if (future instanceof RunnableFuture<?>) {
this.logger.debug("Waiting {} ms since the start of the exectuion of the future to complete", unit.toMillis(timeout));
this.logger.trace("Waiting {} ms since the start of the exectuion of the future to complete", unit.toMillis(timeout));
results.add(((RunnableFuture<T>) future).get(timeout, unit));
this.logger.trace("Exectuion completed", unit.toMillis(timeout));
} else {
long remainingTimeMillis = expireTimeMillis - System.currentTimeMillis();
this.logger.debug("Waiting {} ms for the future to complete", remainingTimeMillis);
this.logger.trace("Waiting {} ms for the future to complete", remainingTimeMillis);
results.add(future.get(remainingTimeMillis, TimeUnit.MILLISECONDS));
this.logger.trace("Exectuion completed", unit.toMillis(timeout));
}
}
@@ -124,24 +126,39 @@ public class CompositeFuture<T> implements Future<T> {
List<CompositeFuture<?>> cfutures = new LinkedList<>();
int removedCancelled = 0;
int removedDone = 0;
int remain = 0;
Iterator<Future<T>> i = this.futures.iterator();
while (i.hasNext()) {
Future<T> future = i.next();
if (future.isCancelled()) {
if (includeCancelled) {
this.logger.trace("Removing cancelled future");
removedCancelled++;
i.remove();
} else {
remain++;
}
} else if (future.isDone()) {
removedDone++;
i.remove();
try {
future.get();
} catch (InterruptedException ie) {
this.logger.trace("Future completed because it was interrupted");
} catch (ExecutionException ee) {
this.logger.error(ee.getMessage(), ee);
} finally {
this.logger.trace("Removing completed future");
removedDone++;
i.remove();
}
} else if (future instanceof CompositeFuture<?>) {
cfutures.add((CompositeFuture<?>) future);
} else {
remain++;
}
}
this.logger.debug("Purged {} cancelled and {} completed futures", removedCancelled, removedDone);
this.logger.debug("Purged {} cancelled and {} completed futures; {} remain", removedCancelled, removedDone, remain);
for (CompositeFuture<?> cfuture : cfutures)
cfuture.purge(includeCancelled);
@@ -80,14 +80,11 @@ public class ThrottledThreadPoolExecutor extends ThreadPoolExecutor {
}
private WaitableRunnable submit(WaitableRunnable runnable, long throttlingBlockTimeout, TimeUnit throttlingBlockUnit) throws InterruptedException, TimeoutException {
// if no core threads are running, the queue won't be monitored for runnables
this.prestartAllCoreThreads();
if (throttlingBlockTimeout < 0L) {
this.getQueue().put(runnable);
} else {
if (!this.getQueue().offer(runnable, throttlingBlockTimeout, throttlingBlockUnit))
throw new TimeoutException();
throw new TimeoutException("Timeout waiting for queue space for runnable");
}
return runnable;
@@ -1,8 +1,27 @@
# defaulting to 3 days = 60 * 24 * 3 = 4320
# once the node is selected, no other node for the shard will be used for the backup
inteligr8.asie.backup.persistTimeMinutes=4320
# what authorities (users or groups) may use the REST services provided by this module?
inteligr8.asie.allowedAuthorities=GROUP_ALFRESCO_ADMINISTRATORS
# same as solr.baseUrl, but that property is private to the Search subsystem
inteligr8.asie.basePath=/solr
# How long should idle executors remain before being shutdown?
# They will re-initialize if needed again
inteligr8.asie.executors.expireTimeInMinutes=30
# Reconciliation configuration; each node will be processed in its own thread
inteligr8.asie.reconciliation.nodesChunkSize=250
inteligr8.asie.reconciliation.nodeTimeoutSeconds=10
inteligr8.asie.reconciliation.concurrentQueueSize=32
inteligr8.asie.reconciliation.concurrency=2
inteligr8.asie.reconciliation.waitAfterSolrNodeActionMillis=0
inteligr8.asie.reconciliation.waitAfterSolrNodeReconcileMillis=0
# Action (like indexing and re-indexing) configuration
inteligr8.asie.default.concurrentQueueSize=32
inteligr8.asie.default.concurrency=2
+4 -4
View File
@@ -12,25 +12,25 @@
<groupId>com.inteligr8</groupId>
<artifactId>solr-api</artifactId>
<version>1.1-SNAPSHOT-solr6</version>
<version>1.2-SNAPSHOT-solr6</version>
<packaging>jar</packaging>
<name>Apache Solr Jakarta RS API</name>
<properties>
<jackson.version>2.18.0</jackson.version>
<jackson.version>2.22.1</jackson.version>
</properties>
<dependencies>
<dependency>
<groupId>jakarta.annotation</groupId>
<artifactId>jakarta.annotation-api</artifactId>
<version>2.1.1</version>
<version>3.0.0</version>
</dependency>
<dependency>
<groupId>jakarta.ws.rs</groupId>
<artifactId>jakarta.ws.rs-api</artifactId>
<version>3.1.0</version>
<version>4.0.0</version>
</dependency>
<dependency>
<groupId>com.fasterxml.jackson.module</groupId>