ParsingQueue: Move scheduling logic into worker This allows for scheduling of tasks instead of workers which in turn aligns the scheduling flow with the rescheduling of persisted tasks. Solves: Jira GER-1715 Change-Id: Ica261badbcb3e127293f77ce89ed682ec01dca74
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 d7bfd68..6f23236 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
@@ -22,8 +22,6 @@ import com.google.common.flogger.FluentLogger; import com.google.gerrit.entities.Project; import com.google.gerrit.entities.Project.NameKey; -import com.google.gerrit.entities.RefNames; -import com.google.gerrit.extensions.common.AccountInfo; import com.google.gerrit.extensions.events.GitReferenceUpdatedListener; import com.google.gerrit.extensions.events.RevisionCreatedListener; import com.google.gerrit.server.git.ProjectRunnable; @@ -56,22 +54,7 @@ } public void scheduleSccCreation(RevisionCreatedListener.Event event) { - schedule( - new EventParsingWorker( - ParsingQueueTask.builder(SCC, event.getChange().project, event.getChange().branch) - .commit(event.getRevision().commit.commit)) { - - @Override - public void doRun() { - try { - eventParser.createAndScheduleSccFromEvent(event); - } catch (Exception e) { - logger.atSevere().withCause(e).log( - "Failed to create SCC from event for %s:%s:%s.", - event.getChange().project, event.getChange().branch, event.getRevision().commit); - } - } - }); + schedule(ParsingQueueTask.sccBuilder(event)); } /** @@ -81,19 +64,7 @@ * @param branchRef - create events for commits reachable from tip of branchRef. */ public void scheduleSccCreation(String repoName, String branchRef) { - schedule( - new EventParsingWorker(ParsingQueueTask.builder(SCS, repoName, branchRef)) { - - @Override - public void doRun() { - try { - eventParser.createAndScheduleSccFromBranch(repoName, branchRef); - } catch (Exception e) { - logger.atSevere().withCause(e).log( - "Failed to create SCC for %s:%s.", repoName, branchRef); - } - } - }); + schedule(ParsingQueueTask.builder(SCC, repoName, branchRef)); } /** @@ -104,56 +75,34 @@ * @param commit - create events from commits reachable from commit. */ public void scheduleSccCreation(String repoName, String branchRef, String commit) { - schedule( - new EventParsingWorker(ParsingQueueTask.builder(SCS, repoName, branchRef).commit(commit)) { - - @Override - public void doRun() { - try { - eventParser.createAndScheduleSccFromCommit(repoName, branchRef, commit); - } catch (Exception e) { - logger.atSevere().withCause(e).log( - "Failed to create SCC for %s:%s.", repoName, branchRef); - } - } - }); + schedule(ParsingQueueTask.builder(SCC, repoName, branchRef).commit(commit)); } public void scheduleScsCreation(GitReferenceUpdatedListener.Event event) { - - scheduleScsCreation( - event.getProjectName(), - event.getRefName(), - event.getNewObjectId(), - event.getUpdater(), - TimeUtil.nowMs(), - event.getOldObjectId()); + schedule( + ParsingQueueTask.builder(SCS, event.getProjectName(), event.getRefName()) + .commit(event.getNewObjectId()) + .updater(event.getUpdater()) + .updateTime(TimeUtil.nowMs()) + .previousTip(event.getOldObjectId())); } public void scheduleScsCreation(String repoName, String branchRef) { - schedule( - new EventParsingWorker(ParsingQueueTask.builder(SCS, repoName, branchRef)) { - - @Override - public void doRun() { - try { - eventParser.createAndScheduleMissingScssFromBranch(repoName, branchRef); - } catch (Exception e) { - logger.atSevere().withCause(e).log( - "Failed to create SCS for %s:%s", repoName, branchRef); - } - } - }); + schedule(ParsingQueueTask.builder(SCS, repoName, branchRef)); } public void scheduleArtcCreation(GitReferenceUpdatedListener.Event event) { - scheduleArtcCreation(event.getProjectName(), event.getRefName(), TimeUtil.nowMs(), false); + schedule( + ParsingQueueTask.builder(ARTC, event.getProjectName(), event.getRefName()) + .updateTime(TimeUtil.nowMs()) + .force(false)); } public void scheduleArtcCreation(TagResource resource, boolean force) { - String tagRef = resource.getRef(); - scheduleArtcCreation( - resource.getName(), tagRef, resource.getTagInfo().created.getTime(), force); + schedule( + ParsingQueueTask.builder(ARTC, resource.getName(), resource.getRef()) + .updateTime(resource.getTagInfo().created.getTime()) + .force(force)); } public void shutDown() { @@ -168,7 +117,7 @@ public void init() { for (ParsingQueueTask task : persistance.getPersistedTasks()) { - requeueTask(task); + schedule(task); } } @@ -176,108 +125,58 @@ pending.remove(worker); } - private void scheduleScsCreation( - String repoName, - String branchRef, - String commit, - AccountInfo updater, - Long submitTime, - String previousTip) { - schedule( - new EventParsingWorker( - ParsingQueueTask.builder(SCS, repoName, branchRef) - .commit(commit) - .updater(updater) - .updateTime(submitTime) - .previousTip(previousTip)) { - - @Override - public void doRun() { - try { - eventParser.createAndScheduleMissingScss( - SourceChangeEventKey.scsKey(repoName, branchRef, commit), - previousTip, - updater, - submitTime); - } catch (Exception e) { - logger.atSevere().withCause(e).log( - "Failed to create SCS from event for %s:%s:%s.", repoName, branchRef, commit); - } - } - }); + private void schedule(ParsingQueueTask.Builder taskBuilder) { + schedule(taskBuilder.build()); } - private void scheduleArtcCreation( - String projectName, String tagRefOrName, Long creationTime, boolean force) { - String tagName = - tagRefOrName.startsWith(RefNames.REFS_TAGS) - ? tagRefOrName.substring(RefNames.REFS_TAGS.length()) - : tagRefOrName; - schedule( - new EventParsingWorker( - ParsingQueueTask.builder(ARTC, projectName, tagName) - .updateTime(creationTime) - .force(force)) { - - @Override - public void doRun() { - try { - eventParser.createAndScheduleArtc(projectName, tagName, creationTime, force); - } catch (Exception e) { - logger.atSevere().withCause(e).log( - "Failed to create ARTC for %s:%s", projectName, tagName); - } - } - }); - } - - private void schedule(EventParsingWorker worker) { + private void schedule(ParsingQueueTask task) { + EventParsingWorker worker = new EventParsingWorker(task); ScheduledFuture<?> future = pool.schedule(worker); if (future != null) { pending.put(worker, future); } } - private void requeueTask(ParsingQueueTask task) { - switch (task.type) { - case SCC: - if (task.commitId != null) { - scheduleSccCreation(task.repoName, task.branchRefOrTag, task.commitId); - } else { - scheduleSccCreation(task.repoName, task.branchRefOrTag); - } - break; - case SCS: - if (task.commitId != null) { - scheduleScsCreation( - task.repoName, - task.branchRefOrTag, - task.commitId, - task.updater, - task.updateTime, - task.previousTip); - } else { - scheduleScsCreation(task.repoName, task.branchRefOrTag); - } - break; - case ARTC: - scheduleArtcCreation(task.repoName, task.branchRefOrTag, task.updateTime, task.force); - break; - case CD: // CD is never scheduled for creation directly. - default: - break; - } - } - - abstract class EventParsingWorker implements ProjectRunnable { + class EventParsingWorker implements ProjectRunnable { private final ParsingQueueTask task; boolean running; - public EventParsingWorker(ParsingQueueTask.Builder taskBuilder) { - this.task = taskBuilder.build(); + public EventParsingWorker(ParsingQueueTask task) { + this.task = task; } - public abstract void doRun(); + 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); + } + } @Override public void run() {
diff --git a/src/main/java/com/googlesource/gerrit/plugins/eventseiffel/parsing/ParsingQueueTask.java b/src/main/java/com/googlesource/gerrit/plugins/eventseiffel/parsing/ParsingQueueTask.java index 4513e24..245c4b3 100644 --- a/src/main/java/com/googlesource/gerrit/plugins/eventseiffel/parsing/ParsingQueueTask.java +++ b/src/main/java/com/googlesource/gerrit/plugins/eventseiffel/parsing/ParsingQueueTask.java
@@ -15,6 +15,7 @@ import com.google.gerrit.entities.RefNames; import com.google.gerrit.extensions.common.AccountInfo; +import com.google.gerrit.extensions.events.RevisionCreatedListener; import com.googlesource.gerrit.plugins.eventseiffel.eiffel.dto.EiffelEventType; import java.util.Objects; @@ -27,15 +28,25 @@ return new Builder(type, repoName, branchRefOrTag); } + /** + * The event contains all information necessary to create a SCC without querying the index. + * Setting the event also automatically sets repoName, branchRefOrTag, commitId from the event and + * type to SCC and ignores all other fields. + */ + static Builder sccBuilder(RevisionCreatedListener.Event event) { + return new Builder(event); + } + public static class Builder { - private final EiffelEventType _type; - private final String _repoName; - private final String _branchRefOrTag; + private EiffelEventType _type; + private String _repoName; + private String _branchRefOrTag; private String _commitId = null; private AccountInfo _updater = null; private Long _updateTime = null; private Boolean _force = null; private String _previousTip = null; + private RevisionCreatedListener.Event _patchsetCreatedEvent; private Builder(EiffelEventType type, String repoName, String branchRefOrTag) { this._type = type; @@ -43,6 +54,10 @@ this._branchRefOrTag = branchRefOrTag; } + private Builder(RevisionCreatedListener.Event event) { + this._patchsetCreatedEvent = event; + } + public Builder commit(String commitId) { this._commitId = commitId; return this; @@ -69,6 +84,9 @@ } public ParsingQueueTask build() { + if (this._patchsetCreatedEvent != null) { + return new ParsingQueueTask(_patchsetCreatedEvent); + } return new ParsingQueueTask( this._type, this._repoName, @@ -89,6 +107,7 @@ final Long updateTime; final Boolean force; final String previousTip; + final RevisionCreatedListener.Event patchsetCreatedEvent; private ParsingQueueTask( EiffelEventType type, @@ -101,12 +120,28 @@ Boolean force) { this.type = type; this.repoName = repoName; - this.branchRefOrTag = branchRefOrTag; + this.branchRefOrTag = + branchRefOrTag.startsWith(RefNames.REFS_TAGS) + ? branchRefOrTag.substring(RefNames.REFS_TAGS.length()) + : branchRefOrTag; this.commitId = commitId; this.updater = updater; this.updateTime = updateTime; this.previousTip = previousTip; this.force = force; + this.patchsetCreatedEvent = null; + } + + private ParsingQueueTask(RevisionCreatedListener.Event patchsetCreatedEvent) { + this.type = EiffelEventType.SCC; + this.repoName = patchsetCreatedEvent.getChange().project; + this.branchRefOrTag = patchsetCreatedEvent.getChange().branch; + this.commitId = patchsetCreatedEvent.getRevision().commit.commit; + this.updater = null; + this.updateTime = null; + this.previousTip = null; + this.force = null; + this.patchsetCreatedEvent = patchsetCreatedEvent; } @Override