Author: kwright
Date: Thu Jun  5 18:24:48 2014
New Revision: 1600718

URL: http://svn.apache.org/r1600718
Log:
Hook up ProcessActivity class

Modified:
    
manifoldcf/branches/CONNECTORS-946/framework/pull-agent/src/main/java/org/apache/manifoldcf/crawler/system/WorkerThread.java

Modified: 
manifoldcf/branches/CONNECTORS-946/framework/pull-agent/src/main/java/org/apache/manifoldcf/crawler/system/WorkerThread.java
URL: 
http://svn.apache.org/viewvc/manifoldcf/branches/CONNECTORS-946/framework/pull-agent/src/main/java/org/apache/manifoldcf/crawler/system/WorkerThread.java?rev=1600718&r1=1600717&r2=1600718&view=diff
==============================================================================
--- 
manifoldcf/branches/CONNECTORS-946/framework/pull-agent/src/main/java/org/apache/manifoldcf/crawler/system/WorkerThread.java
 (original)
+++ 
manifoldcf/branches/CONNECTORS-946/framework/pull-agent/src/main/java/org/apache/manifoldcf/crawler/system/WorkerThread.java
 Thu Jun  5 18:24:48 2014
@@ -537,9 +537,17 @@ public class WorkerThread extends Thread
                         }
 
                         // First, make the things we will need for all 
subsequent steps.
-                        ProcessActivity activity = new 
ProcessActivity(processID,
+                        ProcessActivity activity = new 
ProcessActivity(job.getID(),processID,
                           threadContext,rt,jobManager,ingester,
-                          
currentTime,job,connection,connector,connMgr,legalLinkTypes,ingestLogger,abortSet,outputVersion,newParameterVersion);
+                          connectionName,outputName,transformationNames,
+                          currentTime,
+                          job.getExpiration(),
+                          job.getForcedMetadata(),
+                          job.getInterval(),
+                          job.getMaxInterval(),
+                          job.getHopcountMode(),
+                          
connection,connector,connMgr,legalLinkTypes,ingestLogger,abortSet,
+                          
outputVersion,newParameterVersion,transformationVersions);
                         try
                         {
 
@@ -1472,12 +1480,20 @@ public class WorkerThread extends Thread
   protected static class ProcessActivity implements IProcessActivity
   {
     // Member variables
+    protected final Long jobID;
     protected final String processID;
     protected final IThreadContext threadContext;
     protected final IJobManager jobManager;
     protected final IIncrementalIngester ingester;
+    protected final String connectionName;
+    protected final String outputName;
+    protected final String[] transformationNames;
     protected final long currentTime;
-    protected final IJobDescription job;
+    protected final Long expireInterval;
+    protected final Map<String,Set<String>> forcedMetadata;
+    protected final Long recrawlInterval;
+    protected final Long maxInterval;
+    protected final int hopcountMode;
     protected final IRepositoryConnection connection;
     protected final IRepositoryConnector connector;
     protected final IRepositoryConnectionManager connMgr;
@@ -1487,6 +1503,7 @@ public class WorkerThread extends Thread
     protected final Set<String> abortSet;
     protected final String outputVersion;
     protected final String parameterVersion;
+    protected final String[] transformationVersions;
     
     // We submit references in bulk, because that's way more efficient.
     protected final Map<DocumentReference,DocumentReference> referenceList = 
new HashMap<DocumentReference,DocumentReference>();
@@ -1504,20 +1521,39 @@ public class WorkerThread extends Thread
     *@param jobManager is the job manager
     *@param ingester is the ingester
     */
-    public ProcessActivity(String processID, IThreadContext threadContext,
+    public ProcessActivity(Long jobID, String processID,
+      IThreadContext threadContext,
       IReprioritizationTracker rt, IJobManager jobManager,
-      IIncrementalIngester ingester, long currentTime,
-      IJobDescription job, IRepositoryConnection connection, 
IRepositoryConnector connector,
+      IIncrementalIngester ingester,
+      String connectionName, String outputName, String[] transformationNames,
+      long currentTime,
+      Long expireInterval,
+      Map<String,Set<String>> forcedMetadata,
+      Long recrawlInterval,
+      Long maxInterval,
+      int hopcountMode,
+      IRepositoryConnection connection, IRepositoryConnector connector,
       IRepositoryConnectionManager connMgr, String[] legalLinkTypes, 
OutputActivity ingestLogger,
-      Set<String> abortSet, String outputVersion, String parameterVersion)
+      Set<String> abortSet,
+      String outputVersion, String parameterVersion, String[] 
transformationVersions)
     {
+      this.jobID = jobID;
       this.processID = processID;
       this.threadContext = threadContext;
       this.rt = rt;
       this.jobManager = jobManager;
       this.ingester = ingester;
+      this.connectionName = connectionName;
+      this.outputName = outputName;
+      this.transformationNames = transformationNames;
       this.currentTime = currentTime;
-      this.job = job;
+      this.expireInterval = expireInterval;
+      this.forcedMetadata = forcedMetadata;
+      this.recrawlInterval = recrawlInterval;
+      this.maxInterval = maxInterval;
+      this.hopcountMode = hopcountMode;
+      
+      //this.job = job;
       this.connection = connection;
       this.connector = connector;
       this.connMgr = connMgr;
@@ -1526,6 +1562,7 @@ public class WorkerThread extends Thread
       this.abortSet = abortSet;
       this.outputVersion = outputVersion;
       this.parameterVersion = parameterVersion;
+      this.transformationVersions = transformationVersions;
     }
 
     /** Clean up any dangling information, before abandoning this process 
activity object */
@@ -1568,7 +1605,7 @@ public class WorkerThread extends Thread
       +" to '"+localIdentifier+"', relationship type 
"+((relationshipType==null)?"null":"'"+relationshipType+"'")
       +", with "+((dataNames==null)?"no":Integer.toString(dataNames.length))+" 
data types, origination 
time="+((originationTime==null)?"unknown":originationTime.toString()));
 
-      Long expireInterval = job.getExpiration();
+      //Long expireInterval = job.getExpiration();
       if (expireInterval != null)
       {
         // We should not queue documents that have already expired; it wastes 
time
@@ -1716,7 +1753,7 @@ public class WorkerThread extends Thread
     public String[] retrieveParentData(String localIdentifier, String dataName)
       throws ManifoldCFException
     {
-      return 
jobManager.retrieveParentData(job.getID(),ManifoldCF.hash(localIdentifier),dataName);
+      return 
jobManager.retrieveParentData(jobID,ManifoldCF.hash(localIdentifier),dataName);
     }
 
     /** Retrieve data passed from parents to a specified child document.
@@ -1728,7 +1765,7 @@ public class WorkerThread extends Thread
     public CharacterInput[] retrieveParentDataAsFiles(String localIdentifier, 
String dataName)
       throws ManifoldCFException
     {
-      return 
jobManager.retrieveParentDataAsFiles(job.getID(),ManifoldCF.hash(localIdentifier),dataName);
+      return 
jobManager.retrieveParentDataAsFiles(jobID,ManifoldCF.hash(localIdentifier),dataName);
     }
 
     /** Record a document version, but don't ingest it.
@@ -1741,7 +1778,7 @@ public class WorkerThread extends Thread
       throws ManifoldCFException, ServiceInterruption
     {
       String documentIdentifierHash = ManifoldCF.hash(documentIdentifier);
-      
ingester.documentRecord(job.getOutputConnectionName(),job.getConnectionName(),documentIdentifierHash,version,currentTime,ingestLogger);
+      
ingester.documentRecord(outputName,connectionName,documentIdentifierHash,version,currentTime,ingestLogger);
     }
 
     /** Ingest the current document.
@@ -1764,7 +1801,7 @@ public class WorkerThread extends Thread
 
       if (data != null)
       {
-        Map<String,Set<String>> forcedMetadata = job.getForcedMetadata();
+        //Map<String,Set<String>> forcedMetadata = job.getForcedMetadata();
         
         // Modify the repository document with forced parameters.
         for (String paramName : forcedMetadata.keySet())
@@ -1781,9 +1818,10 @@ public class WorkerThread extends Thread
       }
         
       // First, we need to add into the metadata the stuff from the job 
description.
-      ingester.documentIngest(job.getOutputConnectionName(),
-        job.getConnectionName(),documentIdentifierHash,
-        version,outputVersion,parameterVersion,
+      ingester.documentIngest(transformationNames,
+        outputName,
+        connectionName,documentIdentifierHash,
+        version,transformationVersions,outputVersion,parameterVersion,
         connection.getACLAuthority(),
         data,currentTime,
         documentURI,
@@ -1817,8 +1855,8 @@ public class WorkerThread extends Thread
       throws ManifoldCFException, ServiceInterruption
     {
       String documentIdentifierHash = ManifoldCF.hash(documentIdentifier);
-      ingester.documentDelete(job.getOutputConnectionName(),
-        job.getConnectionName(),documentIdentifierHash,
+      ingester.documentDelete(outputName,
+        connectionName,documentIdentifierHash,
         ingestLogger);
     }
 
@@ -1905,10 +1943,10 @@ public class WorkerThread extends Thread
     public Long calculateDocumentRescheduleTime(long currentTime, long 
timeAmt, String localIdentifier)
     {
       Long recrawlTime = null;
-      Long recrawlInterval = job.getInterval();
+      //Long recrawlInterval = job.getInterval();
       if (recrawlInterval != null)
       {
-        Long maxInterval = job.getMaxInterval();
+        //Long maxInterval = job.getMaxInterval();
         long actualInterval = recrawlInterval.longValue() + timeAmt;
         if (maxInterval != null && actualInterval > maxInterval.longValue())
           actualInterval = maxInterval.longValue();
@@ -1945,7 +1983,7 @@ public class WorkerThread extends Thread
       Long originationTime = getDocumentOriginationTime(localIdentifier);
       if (originationTime == null)
         originationTime = new Long(currentTime);
-      Long expireInterval = job.getExpiration();
+      //Long expireInterval = job.getExpiration();
       Long expireTime = null;
       if (expireInterval != null)
         expireTime = new Long(originationTime.longValue() + 
expireInterval.longValue());
@@ -2061,7 +2099,7 @@ public class WorkerThread extends Thread
         rt.preloadBinValues();
 
         jobManager.addDocuments(processID,
-          
job.getID(),legalLinkTypes,docidHashes,docids,db.getParentIdentifierHash(),db.getLinkType(),job.getHopcountMode(),
+          
jobID,legalLinkTypes,docidHashes,docids,db.getParentIdentifierHash(),db.getLinkType(),hopcountMode,
           dataNames,dataValues,currentTime,priorities,eventNames);
         
         rt.clearPreloadedValues();
@@ -2079,7 +2117,7 @@ public class WorkerThread extends Thread
     public void checkJobStillActive()
       throws ManifoldCFException, ServiceInterruption
     {
-      if (jobManager.checkJobActive(job.getID()) == false)
+      if (jobManager.checkJobActive(jobID) == false)
         throw new ServiceInterruption("Job no longer 
active",System.currentTimeMillis());
     }
 
@@ -2133,7 +2171,8 @@ public class WorkerThread extends Thread
     public boolean checkMimeTypeIndexable(String mimeType)
       throws ManifoldCFException, ServiceInterruption
     {
-      return 
ingester.checkMimeTypeIndexable(job.getOutputConnectionName(),outputVersion,mimeType);
+      // MHL
+      return 
ingester.checkMimeTypeIndexable(outputName,outputVersion,mimeType);
     }
 
     /** Check whether a document is indexable by the currently specified 
output connector.
@@ -2144,7 +2183,8 @@ public class WorkerThread extends Thread
     public boolean checkDocumentIndexable(File localFile)
       throws ManifoldCFException, ServiceInterruption
     {
-      return 
ingester.checkDocumentIndexable(job.getOutputConnectionName(),outputVersion,localFile);
+      // MHL
+      return 
ingester.checkDocumentIndexable(outputName,outputVersion,localFile);
     }
 
     /** Check whether a document of a specified length is indexable by the 
currently specified output connector.
@@ -2155,7 +2195,8 @@ public class WorkerThread extends Thread
     public boolean checkLengthIndexable(long length)
       throws ManifoldCFException, ServiceInterruption
     {
-      return 
ingester.checkLengthIndexable(job.getOutputConnectionName(),outputVersion,length);
+      // MHL
+      return ingester.checkLengthIndexable(outputName,outputVersion,length);
     }
 
     /** Pre-determine whether a document's URL is indexable by this connector. 
 This method is used by participating repository connectors
@@ -2167,7 +2208,7 @@ public class WorkerThread extends Thread
     public boolean checkURLIndexable(String url)
       throws ManifoldCFException, ServiceInterruption
     {
-      return 
ingester.checkURLIndexable(job.getOutputConnectionName(),outputVersion,url);
+      return ingester.checkURLIndexable(outputName,outputVersion,url);
     }
 
     /** Create a global string from a simple string.
@@ -2197,7 +2238,7 @@ public class WorkerThread extends Thread
     @Override
     public String createJobSpecificString(String simpleString)
     {
-      return ManifoldCF.createJobSpecificString(job.getID(),simpleString);
+      return ManifoldCF.createJobSpecificString(jobID,simpleString);
     }
 
   }


Reply via email to