Author: kwright
Date: Fri Jun 13 15:06:50 2014
New Revision: 1602456

URL: http://svn.apache.org/r1602456
Log:
More work to get things to build again

Modified:
    
manifoldcf/branches/CONNECTORS-962/framework/agents/src/main/java/org/apache/manifoldcf/agents/incrementalingest/IncrementalIngester.java
    
manifoldcf/branches/CONNECTORS-962/framework/pull-agent/src/main/java/org/apache/manifoldcf/crawler/system/WorkerThread.java

Modified: 
manifoldcf/branches/CONNECTORS-962/framework/agents/src/main/java/org/apache/manifoldcf/agents/incrementalingest/IncrementalIngester.java
URL: 
http://svn.apache.org/viewvc/manifoldcf/branches/CONNECTORS-962/framework/agents/src/main/java/org/apache/manifoldcf/agents/incrementalingest/IncrementalIngester.java?rev=1602456&r1=1602455&r2=1602456&view=diff
==============================================================================
--- 
manifoldcf/branches/CONNECTORS-962/framework/agents/src/main/java/org/apache/manifoldcf/agents/incrementalingest/IncrementalIngester.java
 (original)
+++ 
manifoldcf/branches/CONNECTORS-962/framework/agents/src/main/java/org/apache/manifoldcf/agents/incrementalingest/IncrementalIngester.java
 Fri Jun 13 15:06:50 2014
@@ -565,10 +565,10 @@ public class IncrementalIngester extends
     for (int i = 0; i < pipelineSpecificationWithVersions.getOutputCount(); 
i++)
     {
       int stage = pipelineSpecificationWithVersions.getOutputStage(i);
-      String oldDocumentVersion = 
pipelineSpecificationWithVersions.getOutputDocumentVersionString();
-      String oldParameterVersion = 
pipelineSpecificationWithVersions.getOutputParameterVersionString();
-      String oldOutputVersion = 
pipelineSpecificationWithVersions.getOutputVersionString();
-      String oldAuthorityName = 
pipelineSpecificationWithVersions.getAuthorityNameString();
+      String oldDocumentVersion = 
pipelineSpecificationWithVersions.getOutputDocumentVersionString(i);
+      String oldParameterVersion = 
pipelineSpecificationWithVersions.getOutputParameterVersionString(i);
+      String oldOutputVersion = 
pipelineSpecificationWithVersions.getOutputVersionString(i);
+      String oldAuthorityName = 
pipelineSpecificationWithVersions.getAuthorityNameString(i);
       // If it looks like we never indexed this output before, we need to do 
it now.
       if (oldDocumentVersion == null)
         return true;
@@ -581,7 +581,7 @@ public class IncrementalIngester extends
       
       // Everything matches so far.  Next step is to compute a transformation 
path an corresponding version string.
       String newTransformationVersion = 
computePackedTransformationVersion(pipelineSpecificationWithVersions,stage);
-      if 
(!pipelineSpecificationWithVersions.getOutputTransformationVersionString().equals(newTransformationVersion))
+      if 
(!pipelineSpecificationWithVersions.getOutputTransformationVersionString(i).equals(newTransformationVersion))
         return true;
     }
     // Everything matches, so no reindexing is needed.
@@ -600,7 +600,7 @@ public class IncrementalIngester extends
     int currentStage = stage;
     while (true)
     {
-      int newStage = stage.getStageParent(currentStage);
+      int newStage = pipelineSpecification.getStageParent(currentStage);
       if (newStage == -1)
         break;
       stageCount++;
@@ -613,7 +613,7 @@ public class IncrementalIngester extends
     currentStage = stage;
     while (true)
     {
-      int newStage = stage.getStageParent(currentStage);
+      int newStage = pipelineSpecification.getStageParent(currentStage);
       if (newStage == -1)
         break;
       stageNames[stageCount] = 
pipelineSpecification.getStageConnectionName(newStage);
@@ -804,8 +804,6 @@ public class IncrementalIngester extends
           "","","","",""),
         identifierClass, identifierHash,
         documentVersion,
-        "",
-        outputVersion,
         parameterVersion,
         authorityName,
         data,
@@ -955,7 +953,7 @@ public class IncrementalIngester extends
       throw new ServiceInterruption("Pipeline connector not installed",0L);
     try
     {
-      return 
pipeline.addOrReplaceDocumentWithException(documentURI,document,documentVersion,parameterVersion,authorityNameString,finalActivities,ingestTime);
+      return 
pipeline.addOrReplaceDocumentWithException(docKey,documentURI,document,documentVersion,parameterVersion,authorityName,activities,ingestTime)
 == IPipelineConnector.DOCUMENTSTATUS_ACCEPTED;
     }
     finally
     {
@@ -963,8 +961,6 @@ public class IncrementalIngester extends
     }
   }
 
-
-
   /** Note the fact that we checked a document (and found that it did not need 
to be ingested, because the
   * versions agreed).
   *@param outputConnectionName is the name of the output connection associated 
with this action.
@@ -1056,7 +1052,7 @@ public class IncrementalIngester extends
           list2.clear();
           j = 0;
         }
-        list.add(iter2.next());
+        list2.add(iter2.next());
         j++;
       }
 
@@ -1284,7 +1280,7 @@ public class IncrementalIngester extends
       try
       {
         // Fetch the document URIs for the listed documents
-        for (int i = 0; j < uris.length; i++)
+        for (int i = 0; i < uris.length; i++)
         {
           if (uris[i] != null && uris[i].getURI() != null)
             
removeDocument(connection,uris[i].getURI(),uris[i].getOutputVersion(),activities);
@@ -1352,7 +1348,7 @@ public class IncrementalIngester extends
               list.clear();
               j = 0;
             }
-            list.add(iter.next());
+            list.add(iter2.next());
             j++;
           }
 
@@ -1476,6 +1472,14 @@ public class IncrementalIngester extends
     return findConjunctionClauseMax(new ClauseDescription[]{
       new UnitaryClause(outputConnNameField,outputConnectionName)});
   }
+
+  /** Calculate the maximum number of doc ids we should use.
+  */
+  protected int maxClausesRowIdsForDocIds(String[] outputConnectionNames)
+  {
+    return findConjunctionClauseMax(new ClauseDescription[]{
+      new MultiClause(outputConnNameField,outputConnectionNames)});
+  }
   
   /** Given values and parameters corresponding to a set of hash values, add 
corresponding
   * table row id's to the output map.
@@ -1499,6 +1503,28 @@ public class IncrementalIngester extends
     }
   }
 
+  /** Given values and parameters corresponding to a set of hash values, add 
corresponding
+  * table row id's to the output map.
+  */
+  protected void findRowIdsForDocIds(String[] outputConnectionNames, Set<Long> 
rowIDSet, List<String> paramValues)
+    throws ManifoldCFException
+  {
+    ArrayList list = new ArrayList();
+    String query = buildConjunctionClause(list,new ClauseDescription[]{
+      new MultiClause(docKeyField,paramValues),
+      new MultiClause(outputConnNameField,outputConnectionNames)});
+      
+    IResultSet set = performQuery("SELECT "+idField+" FROM "+
+      getTableName()+" WHERE "+query,list,null,null);
+    
+    for (int i = 0; i < set.getRowCount(); i++)
+    {
+      IResultRow row = set.getRow(i);
+      Long rowID = (Long)row.getValue(idField);
+      rowIDSet.add(rowID);
+    }
+  }
+
   /** Calculate the maximum number of clauses.
   */
   protected int maxClausesDeleteRowIds()
@@ -1508,7 +1534,7 @@ public class IncrementalIngester extends
     
   /** Delete a chunk of row ids.
   */
-  protected void deleteRowIds(ArrayList list)
+  protected void deleteRowIds(List<Long> list)
     throws ManifoldCFException
   {
     ArrayList newList = new ArrayList();
@@ -2304,24 +2330,24 @@ public class IncrementalIngester extends
 
   protected static class OutputRemoveActivitiesWrapper extends 
OutputRecordingActivity implements IOutputRemoveActivity
   {
-    protected final IOutputRemoveActivity activities;
+    protected final IOutputRemoveActivity removeActivities;
     
-    public OutputRemoveActivitiesWrapper(IOutputRemoveActivity activities, 
String outputConnectionName)
+    public OutputRemoveActivitiesWrapper(IOutputRemoveActivity 
removeActivities, String outputConnectionName)
     {
-      super(activities,outputConnectionName);
-      this.activities = activities;
+      super(removeActivities,outputConnectionName);
+      this.removeActivities = removeActivities;
     }
 
   }
   
   protected static class OutputAddActivitiesWrapper extends 
OutputRecordingActivity implements IOutputAddActivity
   {
-    protected final IOutputAddActivity activities;
+    protected final IOutputAddActivity addActivities;
     
-    public OutputAddActivitiesWrapper(IOutputAddActivity activities, String 
outputConnectionName)
+    public OutputAddActivitiesWrapper(IOutputAddActivity addActivities, String 
outputConnectionName)
     {
-      super(activities,outputConnectionName);
-      this.activities = activities;
+      super(addActivities,outputConnectionName);
+      this.addActivities = addActivities;
     }
     
     /** Qualify an access token appropriately, to match access tokens as 
returned by mod_aa.  This method
@@ -2334,7 +2360,7 @@ public class IncrementalIngester extends
     public String qualifyAccessToken(String authorityNameString, String 
accessToken)
       throws ManifoldCFException
     {
-      return activities.qualifyAccessToken(authorityNameString,accessToken);
+      return addActivities.qualifyAccessToken(authorityNameString,accessToken);
     }
 
     /** Send a document via the pipeline to the next output connection.
@@ -2347,7 +2373,7 @@ public class IncrementalIngester extends
     public int sendDocument(String documentURI, RepositoryDocument document, 
String authorityNameString)
       throws ManifoldCFException, ServiceInterruption, IOException
     {
-      return activities.sendDocument(documentURI,document,authorityNameString);
+      return 
addActivities.sendDocument(documentURI,document,authorityNameString);
     }
 
     /** Detect if a mime type is acceptable downstream or not.  This method is 
used to determine whether it makes sense to fetch a document
@@ -2359,7 +2385,7 @@ public class IncrementalIngester extends
     public boolean checkMimeTypeIndexable(String mimeType)
       throws ManifoldCFException, ServiceInterruption
     {
-      return activities.checkMimeTypeIndexable(mimeType);
+      return addActivities.checkMimeTypeIndexable(mimeType);
     }
 
     /** Pre-determine whether a document (passed here as a File object) is 
acceptable downstream.  This method is
@@ -2372,7 +2398,7 @@ public class IncrementalIngester extends
     public boolean checkDocumentIndexable(File localFile)
       throws ManifoldCFException, ServiceInterruption
     {
-      return activities.checkDocumentIndexable(localFile);
+      return addActivities.checkDocumentIndexable(localFile);
     }
 
     /** Pre-determine whether a document's length is acceptable downstream.  
This method is used
@@ -2384,7 +2410,7 @@ public class IncrementalIngester extends
     public boolean checkLengthIndexable(long length)
       throws ManifoldCFException, ServiceInterruption
     {
-      return activities.checkLengthIndexable(length);
+      return addActivities.checkLengthIndexable(length);
     }
 
     /** Pre-determine whether a document's URL is acceptable downstream.  This 
method is used
@@ -2396,7 +2422,7 @@ public class IncrementalIngester extends
     public boolean checkURLIndexable(String url)
       throws ManifoldCFException, ServiceInterruption
     {
-      return activities.checkURLIndexable(url);
+      return addActivities.checkURLIndexable(url);
     }
 
   }
@@ -2488,7 +2514,7 @@ public class IncrementalIngester extends
       // Cycle through the "current set"
       while (true)
       {
-        int parent;
+        int parent = -1;
         int[] siblings = null;
         for (Integer outputStage : currentSet.keySet())
         {
@@ -2499,7 +2525,7 @@ public class IncrementalIngester extends
           boolean skipToNext = false;
           for (int sibling : siblings)
           {
-            if (!currentSet.contains(new Integer(sibling)))
+            if (currentSet.get(new Integer(sibling)) == null)
             {
               skipToNext = true;
               break;
@@ -2518,8 +2544,8 @@ public class IncrementalIngester extends
         if (siblings == null)
           throw new IllegalStateException("Not at root but can't progress");
         
-        PipelineCheckEntryPoint[] siblingEntryPoints = new 
PipelineCheckEntryPoint[siblings.size()];
-        for (int j = 0; j < siblings.size(); j++)
+        PipelineCheckEntryPoint[] siblingEntryPoints = new 
PipelineCheckEntryPoint[siblings.length];
+        for (int j = 0; j < siblings.length; j++)
         {
           siblingEntryPoints[j] = currentSet.remove(new Integer(siblings[j]));
         }
@@ -2544,21 +2570,20 @@ public class IncrementalIngester extends
       ITransformationConnector[] transformationConnectors,
       IOutputConnector[] outputConnectors)
     {
-      
super(pipelineConnectionsWithVersions,transformationConnectors,outputConnectors,
-        transformationDescriptions,outputDescription);
+      
super(pipelineConnectionsWithVersions,transformationConnectors,outputConnectors);
       this.pipelineConnectionsWithVersions = pipelineConnectionsWithVersions;
     }
 
-    public int addOrReplaceDocumentWithException(String documentURI, 
RepositoryDocument document, String newDocumentVersion, String 
newParameterVersion, String authorityNameString, IOutputAddActivity 
finalActivity, long ingestTime)
+    public int addOrReplaceDocumentWithException(String docKey, String 
documentURI, RepositoryDocument document, String newDocumentVersion, String 
newParameterVersion, String authorityNameString, IOutputActivity finalActivity, 
long ingestTime)
       throws ManifoldCFException, ServiceInterruption, IOException
     {
-      PipelineAddFanout entryPoint = 
buildAddPipeline(finalActivity,newDocumentVersion,newParameterVersion,authorityNameString,ingestTime);
-      return 
entryPoint.addOrReplaceDocumentWithException(documentURI,document,authorityNameString);
+      PipelineAddFanout entryPoint = 
buildAddPipeline(finalActivity,newDocumentVersion,newParameterVersion,authorityNameString,ingestTime,docKey);
+      return entryPoint.sendDocument(documentURI,document,authorityNameString);
     }
     
-    protected PipelineAddFanout buildAddPipeline(IOutputAddActivity 
finalActivity,
+    protected PipelineAddFanout buildAddPipeline(IOutputActivity finalActivity,
       String newDocumentVersion, String newParameterVersion, String 
newAuthorityNameString,
-      long ingestTime)
+      long ingestTime, String docKey)
     {
       // Algorithm for building a pipeline:
       // (1) We start with the set of final output connection stages, and 
build an entry point for each one.  That's our "current set".
@@ -2604,17 +2629,20 @@ public class IncrementalIngester extends
         PipelineAddEntryPoint outputStageEntryPoint = new OutputAddEntryPoint(
           outputConnectors[connectionIndex],
           spec.getStageDescriptionString(outputStage),
-          new 
OutputAddActivitiesWrapper(finalActivity,spec.getStageConnectionName(outputStage)),
+          new 
OutputActivitiesWrapper(finalActivity,spec.getStageConnectionName(outputStage)),
           needToReindex,
           spec.getStageConnectionName(outputStage),
           newTransformationVersion,
-          ingestTime);
+          ingestTime,
+          newDocumentVersion,
+          newParameterVersion,
+          docKey);
         currentSet.put(new Integer(outputStage), outputStageEntryPoint);
       }
       // Cycle through the "current set"
       while (true)
       {
-        int parent;
+        int parent = -1;
         int[] siblings = null;
         for (Integer outputStage : currentSet.keySet())
         {
@@ -2625,7 +2653,7 @@ public class IncrementalIngester extends
           boolean skipToNext = false;
           for (int sibling : siblings)
           {
-            if (!currentSet.contains(new Integer(sibling)))
+            if (currentSet.get(new Integer(sibling)) == null)
             {
               skipToNext = true;
               break;
@@ -2644,8 +2672,8 @@ public class IncrementalIngester extends
         if (siblings == null)
           throw new IllegalStateException("Not at root but can't progress");
         
-        PipelineAddEntryPoint[] siblingEntryPoints = new 
PipelineAddEntryPoint[siblings.size()];
-        for (int j = 0; j < siblings.size(); j++)
+        PipelineAddEntryPoint[] siblingEntryPoints = new 
PipelineAddEntryPoint[siblings.length];
+        for (int j = 0; j < siblings.length; j++)
         {
           siblingEntryPoints[j] = currentSet.remove(new Integer(siblings[j]));
         }
@@ -2864,16 +2892,16 @@ public class IncrementalIngester extends
       throws ManifoldCFException, ServiceInterruption, IOException
     {
       // MHL to clone document if there are multiple targets!!
-      if (entryPoints.size() > 1)
+      if (entryPoints.length > 1)
         throw new RuntimeException("Cannot handle fanouts yet");
       // If any of them accept the document, we return "accept".
-      int rval = IPipelineConnector.DOCUMENTSTATUS_REJECT;
+      int rval = IPipelineConnector.DOCUMENTSTATUS_REJECTED;
       for (PipelineAddEntryPoint p : entryPoints)
       {
         if (!p.isActive())
           continue;
-        if 
(p.addOrReplaceDocumentWithException(documentURI,document,authorityNameString) 
== IPipelineConnector.DOCUMENTSTATUS_ACCEPT)
-          rval = IPipelineConnector.DOCUMENTSTATUS_ACCEPT;
+        if 
(p.addOrReplaceDocumentWithException(documentURI,document,authorityNameString) 
== IPipelineConnector.DOCUMENTSTATUS_ACCEPTED)
+          rval = IPipelineConnector.DOCUMENTSTATUS_ACCEPTED;
       }
       return rval;
     }
@@ -2982,20 +3010,31 @@ public class IncrementalIngester extends
     protected final String outputConnectionName;
     protected final String transformationVersion;
     protected final long ingestTime;
+    protected final String documentVersion;
+    protected final String parameterVersion;
+    protected final String docKey;
+    protected final IOutputActivity activity;
     
     public OutputAddEntryPoint(IOutputConnector outputConnector,
       String outputDescriptionString,
-      IOutputAddActivity addActivity,
+      IOutputActivity activity,
       boolean isActive,
       String outputConnectionName,
       String transformationVersion,
-      long ingestTime)
+      long ingestTime,
+      String documentVersion,
+      String parameterVersion,
+      String docKey)
     {
-      super(outputConnector,outputDescriptionString,addActivity,isActive);
+      super(outputConnector,outputDescriptionString,activity,isActive);
       this.outputConnector = outputConnector;
       this.outputConnectionName = outputConnectionName;
       this.transformationVersion = transformationVersion;
       this.ingestTime = ingestTime;
+      this.documentVersion = documentVersion;
+      this.parameterVersion = parameterVersion;
+      this.docKey = docKey;
+      this.activity = activity;
     }
     
     @Override
@@ -3078,7 +3117,7 @@ public class IncrementalIngester extends
             new UnitaryClause(outputConnNameField,outputConnectionName)});
           list.add(docKey);
           performDelete("WHERE "+query+" AND "+docKeyField+"!=?",list,null);
-          outputConnector.removeDocument(oldURI,oldOutputVersion,activities);
+          outputConnector.removeDocument(oldURI,oldOutputVersion,activity);
         }
 
         if (documentURI != null)
@@ -3111,15 +3150,15 @@ public class IncrementalIngester extends
           // that we don't know anything about it.  That means it will be 
reingested when the
           // next version comes along, and will be deleted if called for also.
           
noteDocumentIngest(outputConnectionName,docKey,null,null,null,null,null,ingestTime,documentURI,documentURIHash);
-          int result = super.addOrReplaceDocument(documentURI, document, 
authorityNameString);
-          
noteDocumentIngest(outputConnectionName,docKey,documentVersion,transformationVersion,outputDescriptionString,parameterVersion,authorityNameString,ingestTime,documentURI,documentURIHash);
+          int result = super.addOrReplaceDocumentWithException(documentURI, 
document, authorityNameString);
+          
noteDocumentIngest(outputConnectionName,docKey,documentVersion,transformationVersion,pipelineDescriptionString,parameterVersion,authorityNameString,ingestTime,documentURI,documentURIHash);
           return result;
         }
 
         // If we get here, it means we are noting that the document was 
examined, but that no change was required.  This is signaled
         // to noteDocumentIngest by having the null documentURI.
-        
noteDocumentIngest(outputConnectionName,docKey,documentVersion,transformationVersion,outputDescriptionString,parameterVersion,authorityNameString,ingestTime,null,null);
-        return DOCUMENTSTATUS_ACCEPTED;
+        
noteDocumentIngest(outputConnectionName,docKey,documentVersion,transformationVersion,pipelineDescriptionString,parameterVersion,authorityNameString,ingestTime,null,null);
+        return IPipelineConnector.DOCUMENTSTATUS_ACCEPTED;
       }
       finally
       {
@@ -3149,7 +3188,7 @@ public class IncrementalIngester extends
   {
     protected final String outputConnectionName;
     
-    public RuntPipelineSpecification(String outputConnectionName)
+    public RuntPipelineSpecificationBasic(String outputConnectionName)
     {
       this.outputConnectionName = outputConnectionName;
     }
@@ -3357,9 +3396,9 @@ public class IncrementalIngester extends
       for (int i = 0; i < count; i++)
       {
         if (spec.checkStageOutputConnection(i))
-          outputs.add(spec.getStageConnectionName());
+          outputs.add(spec.getStageConnectionName(i));
         else
-          transformations.add(spec.getStageConnectionName());
+          transformations.add(spec.getStageConnectionName(i));
       }
       
       Map<String,Integer> transformationNameMap = new 
HashMap<String,Integer>();
@@ -3387,11 +3426,11 @@ public class IncrementalIngester extends
         Integer k;
         if (spec.checkStageOutputConnection(i))
         {
-          outputConnectionLookupMap.put(new 
Integer(i),outputNameMap.get(spec.getStageConnectionName()));
+          outputConnectionLookupMap.put(new 
Integer(i),outputNameMap.get(spec.getStageConnectionName(i)));
         }
         else
         {
-          transformationConnectionLookupMap.put(new 
Integer(i),transformationNameMap.get(spec.getStageConnectionName()));
+          transformationConnectionLookupMap.put(new 
Integer(i),transformationNameMap.get(spec.getStageConnectionName(i)));
         }
       }
     }

Modified: 
manifoldcf/branches/CONNECTORS-962/framework/pull-agent/src/main/java/org/apache/manifoldcf/crawler/system/WorkerThread.java
URL: 
http://svn.apache.org/viewvc/manifoldcf/branches/CONNECTORS-962/framework/pull-agent/src/main/java/org/apache/manifoldcf/crawler/system/WorkerThread.java?rev=1602456&r1=1602455&r2=1602456&view=diff
==============================================================================
--- 
manifoldcf/branches/CONNECTORS-962/framework/pull-agent/src/main/java/org/apache/manifoldcf/crawler/system/WorkerThread.java
 (original)
+++ 
manifoldcf/branches/CONNECTORS-962/framework/pull-agent/src/main/java/org/apache/manifoldcf/crawler/system/WorkerThread.java
 Fri Jun 13 15:06:50 2014
@@ -2879,7 +2879,7 @@ public class WorkerThread extends Thread
     protected final String[] transformationDescriptionStrings;
     protected final String outputDescriptionString;
     
-    public PipelineSpecificationBasic(IJobDescription job, String[] 
transformationDescriptionStrings,
+    public PipelineSpecification(IJobDescription job, String[] 
transformationDescriptionStrings,
       String outputDescriptionString)
     {
       super(job);


Reply via email to