Author: kwright
Date: Sat Nov 23 13:34:09 2013
New Revision: 1544790
URL: http://svn.apache.org/r1544790
Log:
Agents process revisions.
Modified:
manifoldcf/branches/CONNECTORS-781/framework/agents/src/main/java/org/apache/manifoldcf/agents/AgentRun.java
manifoldcf/branches/CONNECTORS-781/framework/agents/src/main/java/org/apache/manifoldcf/agents/AgentStop.java
manifoldcf/branches/CONNECTORS-781/framework/agents/src/main/java/org/apache/manifoldcf/agents/interfaces/IAgent.java
manifoldcf/branches/CONNECTORS-781/framework/agents/src/main/java/org/apache/manifoldcf/agents/system/ManifoldCF.java
manifoldcf/branches/CONNECTORS-781/framework/combined-service/src/main/java/org/apache/manifoldcf/combinedservice/ServletListener.java
manifoldcf/branches/CONNECTORS-781/framework/core/src/main/java/org/apache/manifoldcf/core/interfaces/ILockManager.java
manifoldcf/branches/CONNECTORS-781/framework/core/src/main/java/org/apache/manifoldcf/core/lockmanager/BaseLockManager.java
manifoldcf/branches/CONNECTORS-781/framework/core/src/main/java/org/apache/manifoldcf/core/lockmanager/LockManager.java
manifoldcf/branches/CONNECTORS-781/framework/jetty-runner/src/main/java/org/apache/manifoldcf/jettyrunner/ManifoldCFJettyRunner.java
manifoldcf/branches/CONNECTORS-781/framework/pull-agent/src/main/java/org/apache/manifoldcf/crawler/system/CrawlerAgent.java
manifoldcf/branches/CONNECTORS-781/framework/pull-agent/src/main/java/org/apache/manifoldcf/crawler/system/ManifoldCF.java
manifoldcf/branches/CONNECTORS-781/framework/pull-agent/src/test/java/org/apache/manifoldcf/crawler/tests/ManifoldCFInstance.java
Modified:
manifoldcf/branches/CONNECTORS-781/framework/agents/src/main/java/org/apache/manifoldcf/agents/AgentRun.java
URL:
http://svn.apache.org/viewvc/manifoldcf/branches/CONNECTORS-781/framework/agents/src/main/java/org/apache/manifoldcf/agents/AgentRun.java?rev=1544790&r1=1544789&r2=1544790&view=diff
==============================================================================
---
manifoldcf/branches/CONNECTORS-781/framework/agents/src/main/java/org/apache/manifoldcf/agents/AgentRun.java
(original)
+++
manifoldcf/branches/CONNECTORS-781/framework/agents/src/main/java/org/apache/manifoldcf/agents/AgentRun.java
Sat Nov 23 13:34:09 2013
@@ -29,7 +29,6 @@ public class AgentRun extends BaseAgents
public static final String _rcsid = "@(#)$Id: AgentRun.java 988245
2010-08-23 18:39:35Z kwright $";
public static final String agentServiceType = "AGENT";
- public static final String agentShutdownSignal = "_AGENTRUN_";
public AgentRun()
{
@@ -50,35 +49,18 @@ public class AgentRun extends BaseAgents
//
// Note well that the agents shutdown signal is NEVER modified by this
code; it will be set/cleared by
// AgentStop only, and AgentStop will wait until all services become
inactive before exiting.
-
+ String processID = ManifoldCF.getProcessID();
ILockManager lockManager = LockManagerFactory.make(tc);
- // Don't come up at all if shutdown signal in force
- if (lockManager.checkGlobalFlag(agentShutdownSignal))
- return;
- lockManager.registerServiceBeginServiceActivity(agentServiceType,
ManifoldCF.getProcessID());
+ lockManager.registerServiceBeginServiceActivity(agentServiceType,
processID);
try
{
- ManifoldCF.addShutdownHook(new AgentRunShutdownRunner());
+ // Register a shutdown hook to make sure we signal that the main agents
process is going inactive.
+ ManifoldCF.addShutdownHook(new AgentRunShutdownRunner(processID));
Logging.root.info("Running...");
- while (true)
- {
- // Any shutdown signal yet?
- if (lockManager.checkGlobalFlag(agentShutdownSignal))
- break;
-
- // Start whatever agents need to be started
- ManifoldCF.startAgents(tc);
-
- try
- {
- ManifoldCF.sleep(5000L);
- }
- catch (InterruptedException e)
- {
- break;
- }
- }
+ // Register hook first so stopAgents() not required
+ ManifoldCF.registerAgentsShutdownHook(tc, processID);
+ ManifoldCF.runAgents(tc, processID);
Logging.root.info("Shutting down...");
}
catch (ManifoldCFException e)
@@ -90,7 +72,7 @@ public class AgentRun extends BaseAgents
{
// Exit service
// This is a courtesy; some lock managers (i.e. ZooKeeper) manage to do
this anyway
- lockManager.endServiceActivity(agentServiceType,
ManifoldCF.getProcessID());
+ lockManager.endServiceActivity(agentServiceType, processID);
}
}
@@ -120,8 +102,11 @@ public class AgentRun extends BaseAgents
protected static class AgentRunShutdownRunner implements IShutdownHook
{
- public AgentRunShutdownRunner()
+ protected final String processID;
+
+ public AgentRunShutdownRunner(String processID)
{
+ this.processID = processID;
}
public void doCleanup()
@@ -131,15 +116,9 @@ public class AgentRun extends BaseAgents
ILockManager lockManager = LockManagerFactory.make(tc);
// We can blast the active flag off here; we may have already exited
though and an exception will
// therefore be thrown.
- try
- {
- lockManager.endServiceActivity(agentServiceType,
ManifoldCF.getProcessID());
- }
- catch (ManifoldCFException e)
+ if (lockManager.checkServiceActive(agentServiceType, processID))
{
- if (e.getErrorCode() == ManifoldCFException.INTERRUPTED)
- throw e;
- // Otherwise eat the exception
+ lockManager.endServiceActivity(agentServiceType, processID);
}
}
Modified:
manifoldcf/branches/CONNECTORS-781/framework/agents/src/main/java/org/apache/manifoldcf/agents/AgentStop.java
URL:
http://svn.apache.org/viewvc/manifoldcf/branches/CONNECTORS-781/framework/agents/src/main/java/org/apache/manifoldcf/agents/AgentStop.java?rev=1544790&r1=1544789&r2=1544790&view=diff
==============================================================================
---
manifoldcf/branches/CONNECTORS-781/framework/agents/src/main/java/org/apache/manifoldcf/agents/AgentStop.java
(original)
+++
manifoldcf/branches/CONNECTORS-781/framework/agents/src/main/java/org/apache/manifoldcf/agents/AgentStop.java
Sat Nov 23 13:34:09 2013
@@ -37,7 +37,7 @@ public class AgentStop extends BaseAgent
// As part of the work for CONNECTORS-781, this method is now synchronous.
// We assert the shutdown signal, and then wait until all active services
have shut down.
ILockManager lockManager = LockManagerFactory.make(tc);
- lockManager.setGlobalFlag(AgentRun.agentShutdownSignal);
+ ManifoldCF.assertAgentsShutdownSignal(tc);
try
{
Logging.root.info("Shutdown signal sent");
@@ -70,7 +70,7 @@ public class AgentStop extends BaseAgent
finally
{
// Clear shutdown signal
- lockManager.clearGlobalFlag(AgentRun.agentShutdownSignal);
+ ManifoldCF.clearAgentsShutdownSignal(tc);
}
}
Modified:
manifoldcf/branches/CONNECTORS-781/framework/agents/src/main/java/org/apache/manifoldcf/agents/interfaces/IAgent.java
URL:
http://svn.apache.org/viewvc/manifoldcf/branches/CONNECTORS-781/framework/agents/src/main/java/org/apache/manifoldcf/agents/interfaces/IAgent.java?rev=1544790&r1=1544789&r2=1544790&view=diff
==============================================================================
---
manifoldcf/branches/CONNECTORS-781/framework/agents/src/main/java/org/apache/manifoldcf/agents/interfaces/IAgent.java
(original)
+++
manifoldcf/branches/CONNECTORS-781/framework/agents/src/main/java/org/apache/manifoldcf/agents/interfaces/IAgent.java
Sat Nov 23 13:34:09 2013
@@ -75,13 +75,15 @@ public interface IAgent
/** Start the agent. This method should spin up the agent threads, and
* then return.
+ *@param processID is the process ID to start up an agent for.
*/
- public void startAgent()
+ public void startAgent(String processID)
throws ManifoldCFException;
/** Stop the agent. This should shut down the agent threads.
+ *@param processID is the process ID to stop an agent for.
*/
- public void stopAgent()
+ public void stopAgent(String processID)
throws ManifoldCFException;
/** Request permission from agent to delete an output connection.
Modified:
manifoldcf/branches/CONNECTORS-781/framework/agents/src/main/java/org/apache/manifoldcf/agents/system/ManifoldCF.java
URL:
http://svn.apache.org/viewvc/manifoldcf/branches/CONNECTORS-781/framework/agents/src/main/java/org/apache/manifoldcf/agents/system/ManifoldCF.java?rev=1544790&r1=1544789&r2=1544790&view=diff
==============================================================================
---
manifoldcf/branches/CONNECTORS-781/framework/agents/src/main/java/org/apache/manifoldcf/agents/system/ManifoldCF.java
(original)
+++
manifoldcf/branches/CONNECTORS-781/framework/agents/src/main/java/org/apache/manifoldcf/agents/system/ManifoldCF.java
Sat Nov 23 13:34:09 2013
@@ -27,6 +27,8 @@ public class ManifoldCF extends org.apac
{
public static final String _rcsid = "@(#)$Id: ManifoldCF.java 988245
2010-08-23 18:39:35Z kwright $";
+ public static final String agentShutdownSignal = "_AGENTRUN_";
+
// Agents initialized flag
protected static boolean agentsInitialized = false;
@@ -68,9 +70,6 @@ public class ManifoldCF extends org.apac
if (agentsInitialized)
return;
- // Create the shutdown hook for agents. All activity will be keyed off
of runningHash, so it is safe to do this under all conditions.
- org.apache.manifoldcf.core.system.ManifoldCF.addShutdownHook(new
AgentsShutdownHook());
-
// Initialize the local loggers
Logging.initializeLoggers();
Logging.setLogLevels(threadContext);
@@ -129,12 +128,95 @@ public class ManifoldCF extends org.apac
mgr.deinstall();
}
+ // There are a number of different ways of running the agents framework.
+ // (1) Repeatedly call checkAgents(), and when all done make sure to call
stopAgents().
+ // (2) Call registerAgentsShutdownHook(), then repeatedly run checkAgents(),
Agent shutdown happens on JVM exit.
+ // (3) Call runAgents(), which will wait for someone else to call
assertAgentsShutdownSignal(). Before exit, stopAgents() must be called.
+ // (4) Call registerAgentsShutdownHook(), then call runAgents(), which will
wait for someone else to call assertAgentsShutdownSignal(). Shutdown happens
on JVM exit.
+
+ /** Assert shutdown signal.
+ */
+ public static void assertAgentsShutdownSignal(IThreadContext threadContext)
+ throws ManifoldCFException
+ {
+ ILockManager lockManager = LockManagerFactory.make(threadContext);
+ lockManager.setGlobalFlag(agentShutdownSignal);
+ }
+
+ /** Clear shutdown signal.
+ */
+ public static void clearAgentsShutdownSignal(IThreadContext threadContext)
+ throws ManifoldCFException
+ {
+ ILockManager lockManager = LockManagerFactory.make(threadContext);
+ lockManager.clearGlobalFlag(agentShutdownSignal);
+ }
+
+
+ /** Register agents shutdown hook.
+ * Call this ONCE before calling startAgents or checkAgents the first time,
if you want automatic cleanup of agents on JVM stop.
+ */
+ public static void registerAgentsShutdownHook(IThreadContext threadContext,
String processID)
+ throws ManifoldCFException
+ {
+ // Create the shutdown hook for agents. All activity will be keyed off of
runningHash, so it is safe to do this under all conditions.
+ org.apache.manifoldcf.core.system.ManifoldCF.addShutdownHook(new
AgentsShutdownHook(processID));
+ }
+
+ /** Agent service name prefix (followed by agent class name) */
+ public static final String agentServicePrefix = "AGENT_";
+
+ /** Run agents process.
+ * This method will not return until a shutdown signal is sent.
+ */
+ public static void runAgents(IThreadContext threadContext, String processID)
+ throws ManifoldCFException
+ {
+ ILockManager lockManager = LockManagerFactory.make(threadContext);
+
+ // Don't come up at all if shutdown signal in force
+ if (lockManager.checkGlobalFlag(agentShutdownSignal))
+ return;
+
+ while (true)
+ {
+ // Any shutdown signal yet?
+ if (lockManager.checkGlobalFlag(agentShutdownSignal))
+ break;
+
+ // Start whatever agents need to be started
+ checkAgents(threadContext, processID);
+
+ try
+ {
+ ManifoldCF.sleep(5000);
+ }
+ catch (InterruptedException e)
+ {
+ break;
+ }
+ }
+ }
+
+ protected final static String agentClassLockPrefix = "_AGENTCLASSLOCK_";
+
+ protected static String getAgentsClassLockName(String agentClassName)
+ {
+ return agentClassLockPrefix + agentClassName;
+ }
+
+ protected static String getAgentsClassServiceType(String agentClassName)
+ {
+ return agentServicePrefix + agentClassName;
+ }
+
/** Start all not-running agents.
*@param threadContext is the thread context.
*/
- public static void startAgents(IThreadContext threadContext)
+ public static void checkAgents(IThreadContext threadContext, String
processID)
throws ManifoldCFException
{
+ ILockManager lockManager = LockManagerFactory.make(threadContext);
// Get agent manager
IAgentManager manager = AgentManagerFactory.make(threadContext);
ManifoldCFException problem = null;
@@ -145,6 +227,8 @@ public class ManifoldCF extends org.apac
if (stopAgentsRun)
return;
String[] classes = manager.getAllAgents();
+ Set<String> currentAgentClasses = new HashSet<String>();
+
int i = 0;
while (i < classes.length)
{
@@ -153,12 +237,29 @@ public class ManifoldCF extends org.apac
{
// Start this agent
IAgent agent = AgentFactory.make(threadContext,className);
- agent.initialize();
try
{
+ // Throw a lock, so that cleanup processes and startup processes
don't collide.
+ String lockName = getAgentsClassLockName(className);
+ boolean firstTime;
+ lockManager.enterWriteLock(lockName);
+ try
+ {
+ firstTime =
lockManager.registerServiceBeginServiceActivity(getAgentsClassServiceType(className),
processID);
+ }
+ finally
+ {
+ lockManager.leaveWriteLock(lockName);
+ }
+ // Now initialize agent, being sure to clean up data from previous
incarnations
+ agent.initialize();
+ if (firstTime)
+ agent.cleanUpAgentData();
+ else
+ agent.cleanUpAgentData(processID);
// There is a potential race condition where the agent has been
started but hasn't yet appeared in runningHash.
// But having runningHash be the synchronizer for this activity
will prevent any problems.
- agent.startAgent();
+ agent.startAgent(processID);
// Successful!
runningHash.put(className,agent);
}
@@ -168,7 +269,65 @@ public class ManifoldCF extends org.apac
agent.cleanUp();
}
}
+ currentAgentClasses.add(className);
+ }
+
+ // Go through running hash and look for agents processes that have left
+ Iterator<String> runningAgentsIterator = runningHash.keySet().iterator();
+ while (runningAgentsIterator.hasNext())
+ {
+ String runningAgentClass = runningAgentsIterator.next();
+ if (!currentAgentClasses.contains(runningAgentClass))
+ {
+ // Shut down this one agent.
+ IAgent agent = runningHash.get(runningAgentClass);
+ try
+ {
+ // Stop it
+ agent.stopAgent(processID);
+
lockManager.endServiceActivity(getAgentsClassServiceType(runningAgentClass),
processID);
+ runningAgentsIterator.remove();
+ agent.cleanUp();
+ }
+ catch (ManifoldCFException e)
+ {
+ problem = e;
+ }
+ }
}
+
+ // For every class we're supposed to be running, find registered but
no-longer-active instances and clean
+ // up after them.
+ for (String agentsClass : runningHash.keySet())
+ {
+ IAgent agent = runningHash.get(agentsClass);
+ // Look for dead service instances for this class.
+ // This cannot happen at the same time as other processes doing this
check, or at the same
+ // time as a service of that class starting up, so we need a lock to
prevent those situations.
+ String lockName = getAgentsClassLockName(agentsClass);
+ lockManager.enterWriteLock(lockName);
+ try
+ {
+ // Find the derelict agents of this class, clean them up, and
deregister them.
+ String agentsClassServiceType =
getAgentsClassServiceType(agentsClass);
+ String[] inactiveAgents =
lockManager.getInactiveServices(agentsClassServiceType);
+ for (String inactiveAgentProcessID : inactiveAgents)
+ {
+ agent.cleanUpAgentData(inactiveAgentProcessID);
+ // Deregister
+ lockManager.unregisterService(agentsClassServiceType,
inactiveAgentProcessID);
+ }
+ }
+ catch (ManifoldCFException e)
+ {
+ problem = e;
+ }
+ finally
+ {
+ lockManager.leaveWriteLock(lockName);
+ }
+ }
+
}
if (problem != null)
throw problem;
@@ -177,9 +336,10 @@ public class ManifoldCF extends org.apac
/** Stop all started agents.
*/
- public static void stopAgents(IThreadContext threadContext)
+ public static void stopAgents(IThreadContext threadContext, String processID)
throws ManifoldCFException
{
+ ILockManager lockManager = LockManagerFactory.make(threadContext);
synchronized (runningHash)
{
// This is supposedly safe; iterator remove is used
@@ -189,7 +349,8 @@ public class ManifoldCF extends org.apac
String className = iter.next();
IAgent agent = runningHash.get(className);
// Stop it
- agent.stopAgent();
+ agent.stopAgent(processID);
+ lockManager.endServiceActivity(getAgentsClassServiceType(className),
processID);
iter.remove();
agent.cleanUp();
}
@@ -217,9 +378,11 @@ public class ManifoldCF extends org.apac
/** Agents shutdown hook class */
protected static class AgentsShutdownHook implements IShutdownHook
{
-
- public AgentsShutdownHook()
+ protected final String processID;
+
+ public AgentsShutdownHook(String processID)
{
+ this.processID = processID;
}
public void doCleanup()
@@ -231,7 +394,7 @@ public class ManifoldCF extends org.apac
stopAgentsRun = true;
}
IThreadContext tc = ThreadContextFactory.make();
- stopAgents(tc);
+ stopAgents(tc,processID);
}
}
Modified:
manifoldcf/branches/CONNECTORS-781/framework/combined-service/src/main/java/org/apache/manifoldcf/combinedservice/ServletListener.java
URL:
http://svn.apache.org/viewvc/manifoldcf/branches/CONNECTORS-781/framework/combined-service/src/main/java/org/apache/manifoldcf/combinedservice/ServletListener.java?rev=1544790&r1=1544789&r2=1544790&view=diff
==============================================================================
---
manifoldcf/branches/CONNECTORS-781/framework/combined-service/src/main/java/org/apache/manifoldcf/combinedservice/ServletListener.java
(original)
+++
manifoldcf/branches/CONNECTORS-781/framework/combined-service/src/main/java/org/apache/manifoldcf/combinedservice/ServletListener.java
Sat Nov 23 13:34:09 2013
@@ -29,8 +29,7 @@ public class ServletListener implements
{
public static final String _rcsid = "@(#)$Id$";
- public static final String agentShutdownSignal =
org.apache.manifoldcf.agents.AgentRun.agentShutdownSignal;
- private Thread jobsThread = null;
+ protected static AgentsThread agentsThread = null;
public void contextInitialized(ServletContextEvent sce)
{
@@ -44,7 +43,8 @@ public class ServletListener implements
ManifoldCF.registerThisAgent(tc);
ManifoldCF.reregisterAllConnectors(tc);
- ManifoldCF.startAgents(tc);
+ agentsThread = new AgentsThread(ManifoldCF.getProcessID());
+ agentsThread.start();
}
catch (ManifoldCFException e)
{
@@ -57,13 +57,74 @@ public class ServletListener implements
IThreadContext tc = ThreadContextFactory.make();
try
{
- ManifoldCF.stopAgents(tc);
+ if (agentsThread != null)
+ {
+ ManifoldCF.assertAgentsShutdownSignal(tc);
+ agentsThread.finishUp();
+ agentsThread = null;
+ ManifoldCF.clearAgentsShutdownSignal(tc);
+ }
+ }
+ catch (InterruptedException e)
+ {
}
catch (ManifoldCFException e)
{
- throw new RuntimeException("Cannot shutdown servlet cleanly;
"+e.getMessage(),e);
+ if (e.getErrorCode() != ManifoldCFException.INTERRUPTED)
+ throw new RuntimeException("Cannot shutdown servlet cleanly;
"+e.getMessage(),e);
}
ManifoldCF.cleanUpEnvironment(tc);
}
+ protected static class AgentsThread extends Thread
+ {
+
+ protected final String processID;
+
+ protected Throwable exception = null;
+
+ public AgentsThread(String processID)
+ {
+ setName("Agents");
+ this.processID = processID;
+ }
+
+ public void run()
+ {
+ IThreadContext tc = ThreadContextFactory.make();
+ try
+ {
+ ManifoldCF.clearAgentsShutdownSignal(tc);
+ try
+ {
+ ManifoldCF.runAgents(tc, processID);
+ }
+ finally
+ {
+ ManifoldCF.stopAgents(tc, processID);
+ }
+ }
+ catch (Throwable e)
+ {
+ exception = e;
+ }
+ }
+
+ public void finishUp()
+ throws ManifoldCFException, InterruptedException
+ {
+ join();
+ if (exception != null)
+ {
+ if (exception instanceof RuntimeException)
+ throw (RuntimeException)exception;
+ if (exception instanceof Error)
+ throw (Error)exception;
+ if (exception instanceof ManifoldCFException)
+ throw (ManifoldCFException)exception;
+ throw new RuntimeException("Unknown exception type thrown:
"+exception.getClass().getName()+": "+exception.getMessage(),exception);
+ }
+ }
+ }
+
}
Modified:
manifoldcf/branches/CONNECTORS-781/framework/core/src/main/java/org/apache/manifoldcf/core/interfaces/ILockManager.java
URL:
http://svn.apache.org/viewvc/manifoldcf/branches/CONNECTORS-781/framework/core/src/main/java/org/apache/manifoldcf/core/interfaces/ILockManager.java?rev=1544790&r1=1544789&r2=1544790&view=diff
==============================================================================
---
manifoldcf/branches/CONNECTORS-781/framework/core/src/main/java/org/apache/manifoldcf/core/interfaces/ILockManager.java
(original)
+++
manifoldcf/branches/CONNECTORS-781/framework/core/src/main/java/org/apache/manifoldcf/core/interfaces/ILockManager.java
Sat Nov 23 13:34:09 2013
@@ -45,8 +45,9 @@ public interface ILockManager
* If the transient registration already exists, it is treated as an error
and an exception will be thrown.
*@param serviceType is the type of service.
*@param serviceName is the name of the service to register.
+ *@return true if this is the only active service of this type at this time.
*/
- public void registerServiceBeginServiceActivity(String serviceType, String
serviceName)
+ public boolean registerServiceBeginServiceActivity(String serviceType,
String serviceName)
throws ManifoldCFException;
/** Un-register a service.
@@ -64,7 +65,14 @@ public interface ILockManager
*/
public String[] getRegisteredServices(String serviceType)
throws ManifoldCFException;
-
+
+ /** List services that are registered but not active.
+ *@param serviceType is the service type.
+ *@return the list of service names.
+ */
+ public String[] getInactiveServices(String serviceType)
+ throws ManifoldCFException;
+
/** End service activity.
* This operation exits the "active" zone for the service. This must take
place using the same ILockManager
* object that was used to registerServiceBeginServiceActivity() - which
implies that it is the same thread.
Modified:
manifoldcf/branches/CONNECTORS-781/framework/core/src/main/java/org/apache/manifoldcf/core/lockmanager/BaseLockManager.java
URL:
http://svn.apache.org/viewvc/manifoldcf/branches/CONNECTORS-781/framework/core/src/main/java/org/apache/manifoldcf/core/lockmanager/BaseLockManager.java?rev=1544790&r1=1544789&r2=1544790&view=diff
==============================================================================
---
manifoldcf/branches/CONNECTORS-781/framework/core/src/main/java/org/apache/manifoldcf/core/lockmanager/BaseLockManager.java
(original)
+++
manifoldcf/branches/CONNECTORS-781/framework/core/src/main/java/org/apache/manifoldcf/core/lockmanager/BaseLockManager.java
Sat Nov 23 13:34:09 2013
@@ -100,9 +100,10 @@ public class BaseLockManager implements
* If the transient registration already exists, it is treated as an error
and an exception will be thrown.
*@param serviceType is the type of service.
*@param serviceName is the name of the service to register.
+ *@return true if this is the only active service of this type at this time.
*/
@Override
- public void registerServiceBeginServiceActivity(String serviceType, String
serviceName)
+ public boolean registerServiceBeginServiceActivity(String serviceType,
String serviceName)
throws ManifoldCFException
{
enterWriteLock(serviceLock);
@@ -112,7 +113,9 @@ public class BaseLockManager implements
String serviceActiveFlag = makeActiveServiceFlagName(serviceType,
serviceName);
if (checkGlobalFlag(serviceActiveFlag))
throw new ManifoldCFException("Service '"+serviceName+"' of type
'"+serviceType+"' is already active");
- // First, register the service
+ // First, register the service and find out how many such services are
active.
+ boolean foundService = false;
+ boolean foundActiveService = false;
int i = 0;
while (true)
{
@@ -120,31 +123,37 @@ public class BaseLockManager implements
String x = readServiceName(resourceName);
if (x == null)
{
- writeServiceName(resourceName, serviceName);
- try
- {
- setGlobalFlag(makeRegisteredServiceFlagName(serviceType,
serviceName));
- }
- catch (Throwable e)
+ if (!foundService)
{
- writeServiceName(resourceName, null);
- if (e instanceof Error)
- throw (Error)e;
- if (e instanceof RuntimeException)
- throw (RuntimeException)e;
- if (e instanceof ManifoldCFException)
- throw (ManifoldCFException)e;
- else
- throw new RuntimeException("Unknown exception of type:
"+e.getClass().getName()+": "+e.getMessage(),e);
+ writeServiceName(resourceName, serviceName);
+ try
+ {
+ setGlobalFlag(makeRegisteredServiceFlagName(serviceType,
serviceName));
+ }
+ catch (Throwable e)
+ {
+ writeServiceName(resourceName, null);
+ if (e instanceof Error)
+ throw (Error)e;
+ if (e instanceof RuntimeException)
+ throw (RuntimeException)e;
+ if (e instanceof ManifoldCFException)
+ throw (ManifoldCFException)e;
+ else
+ throw new RuntimeException("Unknown exception of type:
"+e.getClass().getName()+": "+e.getMessage(),e);
+ }
}
break;
}
if (x.equals(serviceName))
- break;
+ foundService = true;
+ else if (checkGlobalFlag(makeActiveServiceFlagName(serviceType, x)))
+ foundActiveService = true;
i++;
}
// Now, set the appropriate active flag
setGlobalFlag(serviceActiveFlag);
+ return !foundActiveService;
}
finally
{
@@ -250,7 +259,43 @@ public class BaseLockManager implements
}
}
-
+ /** List services that are registered but not active.
+ *@param serviceType is the service type.
+ *@return the list of service names.
+ */
+ @Override
+ public String[] getInactiveServices(String serviceType)
+ throws ManifoldCFException
+ {
+ enterWriteLock(serviceLock);
+ try
+ {
+ int i = 0;
+ List<String> inactiveServices = new ArrayList<String>();
+ while (true)
+ {
+ String resourceName = buildServiceListEntry(serviceType, i);
+ String x = readServiceName(resourceName);
+ if (x == null)
+ break;
+ if (!checkGlobalFlag(makeActiveServiceFlagName(serviceType, x)))
+ inactiveServices.add(x);
+ i++;
+ }
+ String[] rval = new String[inactiveServices.size()];
+ i = 0;
+ for (String x : inactiveServices)
+ {
+ rval[i++] = x;
+ }
+ return rval;
+ }
+ finally
+ {
+ leaveWriteLock(serviceLock);
+ }
+ }
+
/** End service activity.
* This operation exits the "active" zone for the service. This must take
place using the same ILockManager
* object that was used to registerServiceBeginServiceActivity() - which
implies that it is the same thread.
Modified:
manifoldcf/branches/CONNECTORS-781/framework/core/src/main/java/org/apache/manifoldcf/core/lockmanager/LockManager.java
URL:
http://svn.apache.org/viewvc/manifoldcf/branches/CONNECTORS-781/framework/core/src/main/java/org/apache/manifoldcf/core/lockmanager/LockManager.java?rev=1544790&r1=1544789&r2=1544790&view=diff
==============================================================================
---
manifoldcf/branches/CONNECTORS-781/framework/core/src/main/java/org/apache/manifoldcf/core/lockmanager/LockManager.java
(original)
+++
manifoldcf/branches/CONNECTORS-781/framework/core/src/main/java/org/apache/manifoldcf/core/lockmanager/LockManager.java
Sat Nov 23 13:34:09 2013
@@ -54,12 +54,13 @@ public class LockManager implements ILoc
* If the transient registration already exists, it is treated as an error
and an exception will be thrown.
*@param serviceType is the type of service.
*@param serviceName is the name of the service to register.
+ *@return true if this is the only active service of this type at this time.
*/
@Override
- public void registerServiceBeginServiceActivity(String serviceType, String
serviceName)
+ public boolean registerServiceBeginServiceActivity(String serviceType,
String serviceName)
throws ManifoldCFException
{
- lockManager.registerServiceBeginServiceActivity(serviceType, serviceName);
+ return lockManager.registerServiceBeginServiceActivity(serviceType,
serviceName);
}
/** Un-register a service.
@@ -85,7 +86,17 @@ public class LockManager implements ILoc
{
return lockManager.getRegisteredServices(serviceType);
}
-
+
+ /** List services that are registered but not active.
+ *@param serviceType is the service type.
+ *@return the list of service names.
+ */
+ public String[] getInactiveServices(String serviceType)
+ throws ManifoldCFException
+ {
+ return lockManager.getInactiveServices(serviceType);
+ }
+
/** End service activity.
* This operation exits the "active" zone for the service. This must take
place using the same ILockManager
* object that was used to registerServiceBeginServiceActivity() - which
implies that it is the same thread.
Modified:
manifoldcf/branches/CONNECTORS-781/framework/jetty-runner/src/main/java/org/apache/manifoldcf/jettyrunner/ManifoldCFJettyRunner.java
URL:
http://svn.apache.org/viewvc/manifoldcf/branches/CONNECTORS-781/framework/jetty-runner/src/main/java/org/apache/manifoldcf/jettyrunner/ManifoldCFJettyRunner.java?rev=1544790&r1=1544789&r2=1544790&view=diff
==============================================================================
---
manifoldcf/branches/CONNECTORS-781/framework/jetty-runner/src/main/java/org/apache/manifoldcf/jettyrunner/ManifoldCFJettyRunner.java
(original)
+++
manifoldcf/branches/CONNECTORS-781/framework/jetty-runner/src/main/java/org/apache/manifoldcf/jettyrunner/ManifoldCFJettyRunner.java
Sat Nov 23 13:34:09 2013
@@ -50,8 +50,6 @@ public class ManifoldCFJettyRunner
public static final String useJettyParentClassLoaderProperty =
"org.apache.manifoldcf.usejettyparentclassloader";
public static final String jettyPortProperty =
"org.apache.manifoldcf.jettyport";
- public static final String agentShutdownSignal =
org.apache.manifoldcf.agents.AgentRun.agentShutdownSignal;
-
protected Server server;
public ManifoldCFJettyRunner( int port, String crawlerWarPath, String
authorityServiceWarPath, String apiWarPath, boolean useParentLoader )
@@ -138,26 +136,10 @@ public class ManifoldCFJettyRunner
public static void runAgents(IThreadContext tc)
throws ManifoldCFException
{
- ILockManager lockManager = LockManagerFactory.make(tc);
-
- while (true)
- {
- // Any shutdown signal yet?
- if (lockManager.checkGlobalFlag(agentShutdownSignal))
- break;
-
- // Start whatever agents need to be started
- ManifoldCF.startAgents(tc);
-
- try
- {
- ManifoldCF.sleep(5000);
- }
- catch (InterruptedException e)
- {
- break;
- }
- }
+ String processID = ManifoldCF.getProcessID();
+ // Do this so we don't have to call stopAgents() ourselves.
+ ManifoldCF.registerAgentsShutdownHook(tc, processID);
+ ManifoldCF.runAgents(tc, processID);
}
/**
@@ -216,8 +198,7 @@ public class ManifoldCFJettyRunner
if (useParentClassLoader)
{
// Clear the agents shutdown signal.
- ILockManager lockManager = LockManagerFactory.make(tc);
- lockManager.clearGlobalFlag(agentShutdownSignal);
+ ManifoldCF.clearAgentsShutdownSignal(tc);
// Do the basic initialization of the database and its schema
ManifoldCF.createSystemDatabase(tc);
Modified:
manifoldcf/branches/CONNECTORS-781/framework/pull-agent/src/main/java/org/apache/manifoldcf/crawler/system/CrawlerAgent.java
URL:
http://svn.apache.org/viewvc/manifoldcf/branches/CONNECTORS-781/framework/pull-agent/src/main/java/org/apache/manifoldcf/crawler/system/CrawlerAgent.java?rev=1544790&r1=1544789&r2=1544790&view=diff
==============================================================================
---
manifoldcf/branches/CONNECTORS-781/framework/pull-agent/src/main/java/org/apache/manifoldcf/crawler/system/CrawlerAgent.java
(original)
+++
manifoldcf/branches/CONNECTORS-781/framework/pull-agent/src/main/java/org/apache/manifoldcf/crawler/system/CrawlerAgent.java
Sat Nov 23 13:34:09 2013
@@ -110,19 +110,19 @@ public class CrawlerAgent implements IAg
* then return.
*/
@Override
- public void startAgent()
+ public void startAgent(String processID)
throws ManifoldCFException
{
- ManifoldCF.startSystem(threadContext);
+ ManifoldCF.startSystem(threadContext, processID);
}
/** Stop the agent. This should shut down the agent threads.
*/
@Override
- public void stopAgent()
+ public void stopAgent(String processID)
throws ManifoldCFException
{
- ManifoldCF.stopSystem(threadContext);
+ ManifoldCF.stopSystem(threadContext, processID);
}
/** Request permission from agent to delete an output connection.
Modified:
manifoldcf/branches/CONNECTORS-781/framework/pull-agent/src/main/java/org/apache/manifoldcf/crawler/system/ManifoldCF.java
URL:
http://svn.apache.org/viewvc/manifoldcf/branches/CONNECTORS-781/framework/pull-agent/src/main/java/org/apache/manifoldcf/crawler/system/ManifoldCF.java?rev=1544790&r1=1544789&r2=1544790&view=diff
==============================================================================
---
manifoldcf/branches/CONNECTORS-781/framework/pull-agent/src/main/java/org/apache/manifoldcf/crawler/system/ManifoldCF.java
(original)
+++
manifoldcf/branches/CONNECTORS-781/framework/pull-agent/src/main/java/org/apache/manifoldcf/crawler/system/ManifoldCF.java
Sat Nov 23 13:34:09 2013
@@ -592,6 +592,7 @@ public class ManifoldCF extends org.apac
throws ManifoldCFException
{
IJobManager jobManager = JobManagerFactory.make(threadContext);
+ // This replaces prepareForStart(), and is always called before system
starts
jobManager.cleanupProcessData(processID);
}
@@ -606,7 +607,7 @@ public class ManifoldCF extends org.apac
/** Start everything.
*/
- public static void startSystem(IThreadContext threadContext)
+ public static void startSystem(IThreadContext threadContext, String
processID)
throws ManifoldCFException
{
Logging.root.info("Starting up pull-agent...");
@@ -722,7 +723,6 @@ public class ManifoldCF extends org.apac
{
int i;
- // Initialize the database
try
{
IThreadContext threadContext = ThreadContextFactory.make();
@@ -731,11 +731,14 @@ public class ManifoldCF extends org.apac
IJobManager jobManager = JobManagerFactory.make(threadContext);
IRepositoryConnectionManager mgr =
RepositoryConnectionManagerFactory.make(threadContext);
+ /* No longer needed, because IAgents specifically initializes/cleans
up.
+
Logging.threads.debug("Agents process starting initialization...");
// Call the database to get it ready
jobManager.prepareForStart();
-
+ */
+
Logging.threads.debug("Agents process reprioritizing documents...");
Map<String,IRepositoryConnection> connectionMap = new
HashMap<String,IRepositoryConnection>();
@@ -830,7 +833,7 @@ public class ManifoldCF extends org.apac
/** Stop the system.
*/
- public static void stopSystem(IThreadContext threadContext)
+ public static void stopSystem(IThreadContext threadContext, String processID)
throws ManifoldCFException
{
Logging.root.info("Shutting down pull-agent...");
Modified:
manifoldcf/branches/CONNECTORS-781/framework/pull-agent/src/test/java/org/apache/manifoldcf/crawler/tests/ManifoldCFInstance.java
URL:
http://svn.apache.org/viewvc/manifoldcf/branches/CONNECTORS-781/framework/pull-agent/src/test/java/org/apache/manifoldcf/crawler/tests/ManifoldCFInstance.java?rev=1544790&r1=1544789&r2=1544790&view=diff
==============================================================================
---
manifoldcf/branches/CONNECTORS-781/framework/pull-agent/src/test/java/org/apache/manifoldcf/crawler/tests/ManifoldCFInstance.java
(original)
+++
manifoldcf/branches/CONNECTORS-781/framework/pull-agent/src/test/java/org/apache/manifoldcf/crawler/tests/ManifoldCFInstance.java
Sat Nov 23 13:34:09 2013
@@ -56,8 +56,6 @@ import org.apache.http.entity.ContentTyp
/** Tests that run the "agents daemon" should be derived from this */
public class ManifoldCFInstance
{
- public static final String agentShutdownSignal = "agent-process";
-
protected boolean singleWar = false;
protected int testPort = 8346;
@@ -548,8 +546,7 @@ public class ManifoldCFInstance
// If all worked, then we can start the daemon.
// Clear the agents shutdown signal.
IThreadContext tc = ThreadContextFactory.make();
- ILockManager lockManager = LockManagerFactory.make(tc);
- lockManager.clearGlobalFlag(agentShutdownSignal);
+ ManifoldCF.clearAgentsShutdownSignal(tc);
daemonThread = new DaemonThread();
daemonThread.start();
@@ -644,8 +641,7 @@ public class ManifoldCFInstance
if (!singleWar)
{
// Shut down daemon
- ILockManager lockManager = LockManagerFactory.make(tc);
- lockManager.setGlobalFlag(agentShutdownSignal);
+ ManifoldCF.assertAgentsShutdownSignal(tc);
// Wait for daemon thread to exit.
while (true)
@@ -692,30 +688,13 @@ public class ManifoldCFInstance
public void run()
{
+ String processID = ManifoldCF.getProcessID();
IThreadContext tc = ThreadContextFactory.make();
// Now, start the server, and then wait for the shutdown signal. On
shutdown, we have to actually do the cleanup,
// because the JVM isn't going away.
try
{
- ILockManager lockManager = LockManagerFactory.make(tc);
- while (true)
- {
- // Any shutdown signal yet?
- if (lockManager.checkGlobalFlag(agentShutdownSignal))
- break;
-
- // Start whatever agents need to be started
- ManifoldCF.startAgents(tc);
-
- try
- {
- ManifoldCF.sleep(5000);
- }
- catch (InterruptedException e)
- {
- break;
- }
- }
+ ManifoldCF.runAgents(tc, processID);
}
catch (ManifoldCFException e)
{
@@ -725,7 +704,7 @@ public class ManifoldCFInstance
{
try
{
- ManifoldCF.stopAgents(tc);
+ ManifoldCF.stopAgents(tc, processID);
}
catch (ManifoldCFException e)
{