ParsingQueue: Save failed parsing-tasks

Let createAndSchedule... methods in EiffelEventParser throw an
exception upon failure instead of logging. This will make
EiffelEventParserQueue aware of whether the queued task succeeded or
not so that it can store failed tasks to retry at a later date.

Solves: Jira GER-1730
Change-Id: I130fbb53712474b6734ced73d0448d0e387c02f4
diff --git a/src/main/java/com/googlesource/gerrit/plugins/eventseiffel/parsing/EiffelEventParser.java b/src/main/java/com/googlesource/gerrit/plugins/eventseiffel/parsing/EiffelEventParser.java
index 2524fd5..f496f6f 100644
--- a/src/main/java/com/googlesource/gerrit/plugins/eventseiffel/parsing/EiffelEventParser.java
+++ b/src/main/java/com/googlesource/gerrit/plugins/eventseiffel/parsing/EiffelEventParser.java
@@ -19,7 +19,7 @@
 
 public interface EiffelEventParser {
 
-  void createAndScheduleSccFromEvent(Event event);
+  void createAndScheduleSccFromEvent(Event event) throws EventParsingException;
 
   /**
    * Creates SCC events for all commits reachable from branchRef. I.e. create SCCs for already
@@ -27,8 +27,10 @@
    *
    * @param repoName - Name of the repository where the branch exists.
    * @param branchRef - Creates missing events for all commits reachable from branch.
+   * @throws EventParsingException - when event-creation fails.
    */
-  void createAndScheduleSccFromBranch(String repoName, String branchRef);
+  void createAndScheduleSccFromBranch(String repoName, String branchRef)
+      throws EventParsingException;
 
   /**
    * Creates SCC events for all commits reachable from commit.
@@ -36,16 +38,20 @@
    * @param repoName - Name of the repository were the commits exists.
    * @param branchRef - Ref of the branch to create events for.
    * @param commit - Creates missing events for all commits reachable from commit.
+   * @throws EventParsingException - when event-creation fails.
    */
-  void createAndScheduleSccFromCommit(String repoName, String branchRef, String commit);
+  void createAndScheduleSccFromCommit(String repoName, String branchRef, String commit)
+      throws EventParsingException;
 
   /**
    * Creates missing SCS events for repoName, branch.
    *
    * @param repoName - Name of the repository where the branch exists.
    * @param branchRef - Creates missing events for all commits reachable from branchRef.
+   * @throws EventParsingException - when event-creation fails.
    */
-  void createAndScheduleMissingScssFromBranch(String repoName, String branchRef);
+  void createAndScheduleMissingScssFromBranch(String repoName, String branchRef)
+      throws EventParsingException;
 
   /**
    * Creates missing SCS events for a submit transaction. submitter and submittedAt is set for all
@@ -55,12 +61,14 @@
    * @param commitSha1TransactionEnd - Tip before submit transaction.
    * @param submitter - The submitter
    * @param submittedAt - When the commit was submitted.
+   * @throws EventParsingException - when event-creation fails.
    */
   void createAndScheduleMissingScss(
       SourceChangeEventKey scs,
       String commitSha1TransactionEnd,
       AccountInfo submitter,
-      Long submittedAt);
+      Long submittedAt)
+      throws EventParsingException;
 
   /**
    * Creates missing ARTC for a tag, together with CD and (if missing) SCS for referenced commit.
@@ -69,6 +77,8 @@
    * @param tagName - The name of the tag.
    * @param creationTime - The time at which the tag was created.
    * @param force - Whether existing events should be replaced or not.
+   * @throws EventParsingException - when event-creation fails.
    */
-  void createAndScheduleArtc(String repoName, String tagName, Long creationTime, boolean force);
+  void createAndScheduleArtc(String repoName, String tagName, Long creationTime, boolean force)
+      throws EventParsingException;
 }
diff --git a/src/main/java/com/googlesource/gerrit/plugins/eventseiffel/parsing/EiffelEventParserImpl.java b/src/main/java/com/googlesource/gerrit/plugins/eventseiffel/parsing/EiffelEventParserImpl.java
index 263c9a7..9a089d4 100644
--- a/src/main/java/com/googlesource/gerrit/plugins/eventseiffel/parsing/EiffelEventParserImpl.java
+++ b/src/main/java/com/googlesource/gerrit/plugins/eventseiffel/parsing/EiffelEventParserImpl.java
@@ -85,7 +85,7 @@
    * @see com.googlesource.gerrit.plugins.eventseiffel.parsing.EiffelEventParser#createAndScheduleSccFromEvent(com.google.gerrit.extensions.events.RevisionCreatedListener.Event)
    */
   @Override
-  public void createAndScheduleSccFromEvent(Event event) {
+  public void createAndScheduleSccFromEvent(Event event) throws EventParsingException {
     CommitInfo commit = event.getRevision().commit;
     SourceChangeEventKey scc =
         SourceChangeEventKey.sccKey(
@@ -116,9 +116,12 @@
         | NoSuchEntityException
         | EiffelEventIdLookupException
         | InterruptedException e) {
-      logger.atSevere().withCause(e).log(
+      throw new EventParsingException(
+          e,
           "Event creation failed for: %s, %s, %s to SCC.",
-          event.getChange().project, event.getChange().branch, event.getRevision().commit.commit);
+          event.getChange().project,
+          event.getChange().branch,
+          event.getRevision().commit.commit);
     }
   }
 
@@ -126,7 +129,8 @@
    * @see com.googlesource.gerrit.plugins.eventseiffel.parsing.EiffelEventParser#createAndScheduleSccFromBranch(java.lang.String, java.lang.String)
    */
   @Override
-  public void createAndScheduleSccFromBranch(String repoName, String branchRef) {
+  public void createAndScheduleSccFromBranch(String repoName, String branchRef)
+      throws EventParsingException {
     ObjectId tip = getTipOf(repoName, branchRef);
     if (tip == null) {
       return;
@@ -138,7 +142,8 @@
    * @see com.googlesource.gerrit.plugins.eventseiffel.parsing.EiffelEventParser#createAndScheduleSccFromCommit(java.lang.String, java.lang.String, java.lang.String)
    */
   @Override
-  public void createAndScheduleSccFromCommit(String repoName, String branchRef, String commit) {
+  public void createAndScheduleSccFromCommit(String repoName, String branchRef, String commit)
+      throws EventParsingException {
 
     SourceChangeEventKey scc = SourceChangeEventKey.sccKey(repoName, branchRef, commit);
     try {
@@ -148,7 +153,7 @@
         | NoSuchEntityException
         | ConfigInvalidException
         | InterruptedException e) {
-      logger.atSevere().withCause(e).log("Event creation failed for: %s", scc);
+      throw new EventParsingException(e, "Event creation failed for: %s", scc);
     }
   }
 
@@ -156,7 +161,8 @@
    * @see com.googlesource.gerrit.plugins.eventseiffel.parsing.EiffelEventParser#createAndScheduleMissingScssFromBranch(java.lang.String, java.lang.String)
    */
   @Override
-  public void createAndScheduleMissingScssFromBranch(String repoName, String branchRef) {
+  public void createAndScheduleMissingScssFromBranch(String repoName, String branchRef)
+      throws EventParsingException {
     ObjectId tip = getTipOf(repoName, branchRef);
     if (tip == null) {
       return;
@@ -173,7 +179,8 @@
       SourceChangeEventKey scs,
       String commitSha1TransactionEnd,
       AccountInfo submitter,
-      Long submittedAt) {
+      Long submittedAt)
+      throws EventParsingException {
     SourceChangeEventKey currentScs = scs;
     SourceChangeEventKey scc = scs.copy(SCC);
     try {
@@ -226,7 +233,7 @@
         | InterruptedException
         | ConfigInvalidException
         | NoSuchEntityException e) {
-      logger.atSevere().withCause(e).log("Failed to create Eiffel event(s) for %s.", currentScs);
+      throw new EventParsingException(e, "Failed to create Eiffel event(s) for %s.", currentScs);
     }
   }
 
@@ -235,7 +242,8 @@
    */
   @Override
   public void createAndScheduleArtc(
-      String repoName, String tagName, Long creationTime, boolean force) {
+      String repoName, String tagName, Long creationTime, boolean force)
+      throws EventParsingException {
     try {
       CompositionDefinedEventKey cd =
           CompositionDefinedEventKey.create(mapper.tagCompositionName(repoName), tagName);
@@ -258,21 +266,19 @@
             "Event Artc has already been created for: %s, %s", repoName, tagName);
       }
     } catch (EiffelEventIdLookupException | InterruptedException e) {
-      logger.atSevere().withCause(e).log(
-          "Event creation failed for: %s, %s to Artc", repoName, tagName);
+      throw new EventParsingException(
+          e, "Event creation failed for: %s, %s to Artc", repoName, tagName);
     }
   }
 
   private void createAndScheduleCd(
-      String repoName, String tagName, Long creationTime, boolean force) {
+      String repoName, String tagName, Long creationTime, boolean force)
+      throws EventParsingException {
     SourceChangeEventKey scs = null;
     Optional<UUID> scsId = Optional.empty();
 
     try {
       String commitId = peelTag(repoName, tagName);
-      if (commitId == null) {
-        return;
-      }
       scs = SourceChangeEventKey.scsKey(repoName, RefNames.REFS_HEADS + "master", commitId);
       scsId = eventHub.getExistingId(scs);
 
@@ -286,34 +292,32 @@
         try {
           scsId = retryer.call(() -> findSourceChangeEventKey(repoName, commitId));
         } catch (RetryException | ExecutionException e) {
-          logger.atSevere().withCause(e).log("Failed to find SCS for %s in %s", commitId, repoName);
-          return;
+          throw new EventParsingException(e, "Failed to find SCS for %s in %s", commitId, repoName);
         }
       }
       if (scsId.isEmpty()) {
-        logger.atSevere().log("Could not find SCS for: %s in %s", commitId, repoName);
-        return;
+        throw new EventParsingException("Could not find SCS for: %s in %s", commitId, repoName);
       }
       pushToHub(mapper.toCd(repoName, tagName, creationTime, scsId.get()), force);
     } catch (EiffelEventIdLookupException | InterruptedException e) {
-      logger.atSevere().withCause(e).log(
+      throw new EventParsingException(
+          e,
           "Event creation failed for: %s",
           CompositionDefinedEventKey.create(mapper.tagCompositionName(repoName), tagName));
     }
   }
 
-  private String peelTag(String repoName, String tagName) {
+  private String peelTag(String repoName, String tagName) throws EventParsingException {
     try (Repository repo = repoManager.openRepository(Project.nameKey(repoName))) {
       Ref tagRef = repo.getRefDatabase().exactRef(Constants.R_TAGS + tagName);
       if (tagRef != null) {
         ObjectId peeled = repo.getRefDatabase().peel(tagRef).getPeeledObjectId();
         return peeled != null ? peeled.getName() : tagRef.getObjectId().getName();
       }
-      logger.atSevere().log("Cannot find tag: %s:%s", repoName, tagName);
+      throw new EventParsingException("Cannot find tag: %s:%s", repoName, tagName);
     } catch (IOException e) {
-      logger.atSevere().withCause(e).log("Unable to peel tag: %s:%s", repoName, tagName);
+      throw new EventParsingException(e, "Unable to peel tag: %s:%s", repoName, tagName);
     }
-    return null;
   }
 
   private Optional<UUID> findSourceChangeEventKey(String repoName, String commitId)
diff --git a/src/main/java/com/googlesource/gerrit/plugins/eventseiffel/parsing/EiffelEventParsingQueue.java b/src/main/java/com/googlesource/gerrit/plugins/eventseiffel/parsing/EiffelEventParsingQueue.java
index 4a378c7..3849289 100644
--- a/src/main/java/com/googlesource/gerrit/plugins/eventseiffel/parsing/EiffelEventParsingQueue.java
+++ b/src/main/java/com/googlesource/gerrit/plugins/eventseiffel/parsing/EiffelEventParsingQueue.java
@@ -18,6 +18,7 @@
 import static com.googlesource.gerrit.plugins.eventseiffel.eiffel.dto.EiffelEventType.SCC;
 import static com.googlesource.gerrit.plugins.eventseiffel.eiffel.dto.EiffelEventType.SCS;
 
+import com.google.common.annotations.VisibleForTesting;
 import com.google.common.collect.Sets;
 import com.google.common.flogger.FluentLogger;
 import com.google.gerrit.entities.Project;
@@ -29,6 +30,7 @@
 import com.google.gerrit.server.util.time.TimeUtil;
 import com.google.inject.Inject;
 import com.googlesource.gerrit.plugins.eventseiffel.eiffel.SourceChangeEventKey;
+import java.util.HashSet;
 import java.util.Map;
 import java.util.Set;
 import java.util.concurrent.ConcurrentHashMap;
@@ -42,6 +44,8 @@
   private final ParsingQueuePersistence persistence;
   private final ConcurrentHashMap<EventParsingWorker, ScheduledFuture<?>> pending =
       new ConcurrentHashMap<>();
+  private final Set<ParsingQueueTask> failedTasks = new HashSet<>();
+  private final Object failedTasksLock = new Object();
 
   @Inject
   public EiffelEventParsingQueue(
@@ -121,11 +125,34 @@
     }
   }
 
-  protected void markAsCompleted(EventParsingWorker worker) {
-    pending.remove(worker);
+  /**
+   * Reschedules all failed tasks.
+   *
+   * @return the number of failed tasks that were rescheduled.
+   */
+  public int rescheduleFailed() {
+    int failedCnt = 0;
+    synchronized (failedTasksLock) {
+      failedCnt = failedTasks.size();
+      for (ParsingQueueTask task : failedTasks) {
+        schedule(task);
+      }
+      failedTasks.clear();
+    }
+    return failedCnt;
   }
 
-  private void schedule(ParsingQueueTask.Builder taskBuilder) {
+  protected void markAsCompleted(EventParsingWorker worker) {
+    pending.remove(worker);
+    if (!worker.successful) {
+      synchronized (failedTasksLock) {
+        failedTasks.add(worker.task);
+      }
+    }
+  }
+
+  @VisibleForTesting
+  void schedule(ParsingQueueTask.Builder taskBuilder) {
     schedule(taskBuilder.build());
   }
 
@@ -140,41 +167,53 @@
   class EventParsingWorker implements ProjectRunnable {
     private final ParsingQueueTask task;
     boolean running;
+    boolean successful = true;
 
     public EventParsingWorker(ParsingQueueTask task) {
       this.task = task;
     }
 
     public void doRun() {
-      switch (task.type) {
-        case SCC:
-          if (task.patchsetCreatedEvent != null) {
-            eventParser.createAndScheduleSccFromEvent(task.patchsetCreatedEvent);
-          } else if (task.commitId != null) {
-            eventParser.createAndScheduleSccFromCommit(
-                task.repoName, task.branchRefOrTag, task.commitId);
-          } else {
-            eventParser.createAndScheduleSccFromBranch(task.repoName, task.branchRefOrTag);
-          }
-          break;
-        case SCS:
-          if (task.commitId != null) {
-            eventParser.createAndScheduleMissingScss(
-                SourceChangeEventKey.scsKey(task.repoName, task.branchRefOrTag, task.commitId),
-                task.previousTip,
-                task.updater,
-                task.updateTime);
-          } else {
-            eventParser.createAndScheduleMissingScssFromBranch(task.repoName, task.branchRefOrTag);
-          }
-          break;
-        case ARTC:
-          eventParser.createAndScheduleArtc(
-              task.repoName, task.branchRefOrTag, task.updateTime, task.force);
-          break;
-        case CD:
-        default:
-          logger.atWarning().log("Cannot schedule events for type: %s.", task.type);
+      String append = null;
+      try {
+        switch (task.type) {
+          case SCC:
+            if (task.patchsetCreatedEvent != null) {
+              append = ", from event, ";
+              eventParser.createAndScheduleSccFromEvent(task.patchsetCreatedEvent);
+            } else if (task.commitId != null) {
+              eventParser.createAndScheduleSccFromCommit(
+                  task.repoName, task.branchRefOrTag, task.commitId);
+            } else {
+              eventParser.createAndScheduleSccFromBranch(task.repoName, task.branchRefOrTag);
+            }
+            break;
+          case SCS:
+            if (task.commitId != null) {
+              append = ", from submit, ";
+              eventParser.createAndScheduleMissingScss(
+                  SourceChangeEventKey.scsKey(task.repoName, task.branchRefOrTag, task.commitId),
+                  task.previousTip,
+                  task.updater,
+                  task.updateTime);
+            } else {
+              eventParser.createAndScheduleMissingScssFromBranch(
+                  task.repoName, task.branchRefOrTag);
+            }
+            break;
+          case ARTC:
+            eventParser.createAndScheduleArtc(
+                task.repoName, task.branchRefOrTag, task.updateTime, task.force);
+            break;
+          case CD:
+          default:
+            logger.atWarning().log("Cannot schedule events for type: %s.", task.type);
+        }
+        successful = true;
+      } catch (EventParsingException e) {
+        logger.atSevere().withCause(e).log(
+            "Failed to create events%sfrom %s", append != null ? append : " ", task);
+        successful = false;
       }
     }
 
diff --git a/src/main/java/com/googlesource/gerrit/plugins/eventseiffel/parsing/EventParsingException.java b/src/main/java/com/googlesource/gerrit/plugins/eventseiffel/parsing/EventParsingException.java
new file mode 100644
index 0000000..7503b63
--- /dev/null
+++ b/src/main/java/com/googlesource/gerrit/plugins/eventseiffel/parsing/EventParsingException.java
@@ -0,0 +1,38 @@
+// Copyright (C) 2022 The Android Open Source Project
+//
+// Licensed under the Apache License, Version 2.0 (the "License");
+// you may not use this file except in compliance with the License.
+// You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing, software
+// distributed under the License is distributed on an "AS IS" BASIS,
+// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+// See the License for the specific language governing permissions and
+// limitations under the License.
+package com.googlesource.gerrit.plugins.eventseiffel.parsing;
+
+/** Thrown when EventParsing fails. */
+public class EventParsingException extends Exception {
+
+  /** */
+  private static final long serialVersionUID = 1L;
+
+  /**
+   * @param cause - The cause of this exception.
+   * @param formatString - as if passed to String.format(formatString, objects)
+   * @param objects - as if passed to String.format(formatString, objects)
+   */
+  public EventParsingException(Throwable cause, String formatString, Object... objects) {
+    super(String.format(formatString, objects), cause);
+  }
+
+  /**
+   * @param formatString - as if passed to String.format(formatString, objects)
+   * @param objects - as if passed to String.format(formatString, objects)
+   */
+  public EventParsingException(String formatString, Object... objects) {
+    super(String.format(formatString, objects));
+  }
+}
diff --git a/src/test/java/com/googlesource/gerrit/plugins/eventseiffel/parsing/EiffelEventParserIT.java b/src/test/java/com/googlesource/gerrit/plugins/eventseiffel/parsing/EiffelEventParserIT.java
index 77693b7..668bb8e 100644
--- a/src/test/java/com/googlesource/gerrit/plugins/eventseiffel/parsing/EiffelEventParserIT.java
+++ b/src/test/java/com/googlesource/gerrit/plugins/eventseiffel/parsing/EiffelEventParserIT.java
@@ -14,6 +14,8 @@
 
 package com.googlesource.gerrit.plugins.eventseiffel.parsing;
 
+import static com.google.common.truth.Truth.assertThat;
+import static com.google.gerrit.testing.GerritJUnit.assertThrows;
 import static com.googlesource.gerrit.plugins.eventseiffel.eiffel.dto.EiffelEventType.SCC;
 import static com.googlesource.gerrit.plugins.eventseiffel.eiffel.dto.EiffelEventType.SCS;
 import static org.junit.Assert.assertEquals;
@@ -226,8 +228,11 @@
   public void artcNotCreatedWhenMissingScs() throws Exception {
     String tag =
         createTagRef(getHead(repo(), "HEAD").getName(), true).substring(Constants.R_TAGS.length());
-    eventParser.createAndScheduleArtc(project.get(), tag, EPOCH_MILLIS, false);
-    assertEquals(0, TestEventHub.EVENTS.size());
+    EventParsingException thrown =
+        assertThrows(
+            EventParsingException.class,
+            () -> eventParser.createAndScheduleArtc(project.get(), tag, EPOCH_MILLIS, false));
+    assertThat(thrown).hasMessageThat().contains("Failed to find SCS for");
   }
 
   @Test
diff --git a/src/test/java/com/googlesource/gerrit/plugins/eventseiffel/parsing/EiffelEventParsingQueueIT.java b/src/test/java/com/googlesource/gerrit/plugins/eventseiffel/parsing/EiffelEventParsingQueueIT.java
index dc39eff..c07a8bb 100644
--- a/src/test/java/com/googlesource/gerrit/plugins/eventseiffel/parsing/EiffelEventParsingQueueIT.java
+++ b/src/test/java/com/googlesource/gerrit/plugins/eventseiffel/parsing/EiffelEventParsingQueueIT.java
@@ -114,6 +114,29 @@
     assertCorrectEvent(1, artc);
   }
 
+  @Test
+  public void failedTaskIsRescheduled() throws Exception {
+    setScsHandled();
+    String tagName = "a-tag";
+    ArtifactEventKey artc = ArtifactEventKey.create(tagPURL(project.get(), tagName, "localhost"));
+    CompositionDefinedEventKey cd =
+        CompositionDefinedEventKey.create(tagCompositionName(project.get(), "localhost"), tagName);
+
+    /* Scheduling ArtC creation for a non-existing tag fails. */
+    parsingQueue.schedule(
+        ParsingQueueTask.builder(EiffelEventType.ARTC, project.get(), tagName)
+            .updateTime(EPOCH_MILLIS)
+            .force(false));
+    assertEquals(0, TestEventHub.EVENTS.size());
+
+    /* Creating the tag should mean that rescheduling the failed task should work. */
+    createTagRef(tagName, getHead(repo(), "HEAD").getName(), true);
+    parsingQueue.rescheduleFailed();
+    assertEquals(2, TestEventHub.EVENTS.size());
+    assertCorrectEvent(0, cd);
+    assertCorrectEvent(1, artc);
+  }
+
   public static class TestModule extends ParsingTestModule {
 
     @Override