Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@
package gorsat;

import gorsat.Commands.Analysis;
import org.gorpipe.exceptions.GorCancelledException;
import org.gorpipe.exceptions.GorException;
import org.gorpipe.exceptions.GorSystemException;
import org.gorpipe.gor.model.GenomicIterator;
Expand Down Expand Up @@ -218,8 +219,9 @@ public boolean hasNext() {
}
return ret;
} catch (InterruptedException e) {
// Must not look like end of stream: callers would treat the truncated output as complete (ENGKNOW-3979).
Thread.currentThread().interrupt();
return false;
throw new GorCancelledException("Query interrupted", e);
}
}

Expand Down
14 changes: 13 additions & 1 deletion gortools/src/main/java/gorsat/process/ParallelExecutor.java
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@

import java.util.Arrays;
import java.util.concurrent.ConcurrentLinkedQueue;
import java.util.concurrent.atomic.AtomicBoolean;

/**
* This class encapsulates a general execution in parallel of the pgor command when a
Expand All @@ -37,11 +38,21 @@ public class ParallelExecutor {
private Throwable firstException;
private final Thread[] threads;
private final Function0<Unit>[] commands;
private final AtomicBoolean cancelled;

public ParallelExecutor(int workers, Function0<Unit>[] commands) {
this(workers, commands, new AtomicBoolean(false));
}

/**
* @param cancelled set when any command fails, before the other workers are interrupted. Commands can check
* it to avoid committing partial results, since the interrupt flag alone is easily lost.
*/
public ParallelExecutor(int workers, Function0<Unit>[] commands, AtomicBoolean cancelled) {
this.commands = commands;
this.threads = new Thread[workers];
this.firstException = null;
this.cancelled = cancelled;
}

@SuppressWarnings("squid:S00112") // We need to handle Throwable here, sorry
Expand All @@ -50,7 +61,7 @@ public void parallelExecute() throws Throwable {
for( int i = 0; i < threads.length; i++ ) {
Thread t = new Thread(() -> {
Function0<Unit> func = clq.poll();
while( func != null ) {
while( func != null && !cancelled.get() ) {
func.apply();
func = clq.poll();
}
Expand All @@ -70,6 +81,7 @@ public void parallelExecute() throws Throwable {
private synchronized void parallelExcecuteUncaughtExceptionHandler(Thread thread, Throwable throwable) {
if (firstException == null) {
firstException = throwable;
cancelled.set(true);
for (Thread t : threads) {
if (t != thread) {
t.interrupt();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -35,7 +35,7 @@ import gorsat.Utilities.{AnalysisUtilities, MacroUtilities}
import gorsat.process.{GorJavaUtilities, ParallelExecutor}
import org.apache.commons.io.FilenameUtils
import org.gorpipe.client.FileCache
import org.gorpipe.exceptions.{GorException, GorSystemException, GorUserException}
import org.gorpipe.exceptions.{GorCancelledException, GorException, GorSystemException, GorUserException}
import org.gorpipe.gor.binsearch.GorIndexType
import org.gorpipe.gor.driver.meta.DataType
import org.gorpipe.gor.model.{DriverBackedFileReader, FileReader, GorMeta, GorOptions, GorParallelQueryHandler}
Expand All @@ -48,6 +48,8 @@ import org.gorpipe.gor.util.DataUtil
import org.slf4j.LoggerFactory

import java.util.Optional
import java.util.concurrent.atomic.AtomicBoolean
import java.util.function.BooleanSupplier
import scala.jdk.CollectionConverters.IteratorHasAsScala

class GeneralQueryHandler(context: GorContext, header: Boolean) extends GorParallelQueryHandler {
Expand All @@ -67,7 +69,8 @@ class GeneralQueryHandler(context: GorContext, header: Boolean) extends GorParal
(linkCacheFilePath, extension)
}

def runAndStoreLinkFileInCache(nested: GorContext, writeLocationPath: String, fileCache: FileCache, useMd5: Boolean): String = {
def runAndStoreLinkFileInCache(nested: GorContext, writeLocationPath: String, fileCache: FileCache, useMd5: Boolean,
isCancelled: BooleanSupplier = GeneralQueryHandler.NotCancelled): String = {
val startTime = System.currentTimeMillis
val fileReader = nested.getSession.getProjectContext.getFileReader
val commandToExecute = nested.getCommand
Expand All @@ -76,7 +79,7 @@ class GeneralQueryHandler(context: GorContext, header: Boolean) extends GorParal
val noDict = commandToExecute.toLowerCase.contains(" -nodict ")
val writeGord = isGord && !noDict
var cacheRes = writeLocationPath
val resultFileName = runCommand(nested, commandToExecute, if (isGord) writeLocationPath else null, useMd5, theTheDict = true)
val resultFileName = runCommand(nested, commandToExecute, if (isGord) writeLocationPath else null, useMd5, theTheDict = true, isCancelled)
val isCacheDir = fileReader.resolveUrl(writeLocationPath,true).isDirectory()

if(fileCache != null && (!isCacheDir || writeGord)) {
Expand All @@ -89,12 +92,13 @@ class GeneralQueryHandler(context: GorContext, header: Boolean) extends GorParal
cacheRes
}

def runAndStoreInCache(nested: GorContext, fileCache: FileCache, useMd5: Boolean): String = {
def runAndStoreInCache(nested: GorContext, fileCache: FileCache, useMd5: Boolean,
isCancelled: BooleanSupplier = GeneralQueryHandler.NotCancelled): String = {
val startTime = System.currentTimeMillis
val commandToExecute = nested.getCommand
val commandSignature = nested.getSignature
var cacheFile = findCacheFile(commandSignature, commandToExecute, header, fileCache, AnalysisUtilities.theCacheDirectory(context.getSession))
val resultFileName = runCommand(nested, commandToExecute, cacheFile, useMd5, theTheDict = false)
val resultFileName = runCommand(nested, commandToExecute, cacheFile, useMd5, theTheDict = false, isCancelled)
if (fileCache != null) {
val extension = CommandParseUtilities.getExtensionForQuery(commandToExecute, header)
val overheadTime = findOverheadTime(commandToExecute)
Expand Down Expand Up @@ -135,6 +139,9 @@ class GeneralQueryHandler(context: GorContext, header: Boolean) extends GorParal
val fileReader = context.getSession.getProjectContext.getFileReader
var commandList: List[() => Unit] = Nil
val useMd5 = System.getProperty("gor.caching.md5.enabled", "false").toBoolean
// Set by ParallelExecutor when any part fails, so the other parts never commit partial output (ENGKNOW-3979)
val cancelled = new AtomicBoolean(false)
val isCancelled: BooleanSupplier = () => cancelled.get()

for (i <- commandSignatures.indices) {
val executeFunction = block2Function {
Expand All @@ -149,9 +156,9 @@ class GeneralQueryHandler(context: GorContext, header: Boolean) extends GorParal
fileNames(i) = if (cacheFile == null || !fileReader.exists(cacheFile)) {
val writeLocationPath = cacheFiles(i)
if (writeLocationPath != null) {
runAndStoreLinkFileInCache(nested, writeLocationPath, fileCache, useMd5)
runAndStoreLinkFileInCache(nested, writeLocationPath, fileCache, useMd5, isCancelled)
} else {
runAndStoreInCache(nested, fileCache, useMd5)
runAndStoreInCache(nested, fileCache, useMd5, isCancelled)
}
} else {
generateDictionaryFile(commandToExecute, fileReader, useMd5, cacheFile)
Expand All @@ -168,13 +175,13 @@ class GeneralQueryHandler(context: GorContext, header: Boolean) extends GorParal
commandList ::= executeFunction
}

if (commandList != Nil) parallelExecution(commandList.reverse.toArray)
if (commandList != Nil) parallelExecution(commandList.reverse.toArray, cancelled)
fileNames
}


def parallelExecution(commands: Array[() => Unit]): Unit = {
val pe = new ParallelExecutor(context.getSession.getSystemContext.getWorkers, commands)
def parallelExecution(commands: Array[() => Unit], cancelled: AtomicBoolean = new AtomicBoolean(false)): Unit = {
val pe = new ParallelExecutor(context.getSession.getSystemContext.getWorkers, commands, cancelled)
try
pe.parallelExecute()
catch {
Expand Down Expand Up @@ -210,24 +217,50 @@ object GeneralQueryHandler {
CommandParseUtilities.getExtensionForQuery(commandToExecute, header))
}

def runCommand(context: GorContext, commandToExecute: String, outfile: String, useMd5: Boolean, theTheDict: Boolean): String = {
private val NotCancelled: BooleanSupplier = () => false

/**
* True if the output must not be committed: the run was cancelled (e.g. a sibling parallel part failed) or the
* thread was interrupted. A cancelled source may end early without an exception, so its output is partial.
*/
private def mustNotCommit(isCancelled: BooleanSupplier): Boolean =
isCancelled.getAsBoolean || Thread.currentThread().isInterrupted

def runCommand(context: GorContext, commandToExecute: String, outfile: String, useMd5: Boolean, theTheDict: Boolean,
isCancelled: BooleanSupplier = NotCancelled): String = {
context.start(outfile)
// We are using absolute paths here
val fileReader = context.getSession.getProjectContext.getSystemFileReader
val result = if (commandToExecute.toUpperCase().startsWith(CommandParseUtilities.GOR_DICTIONARY_PART) || commandToExecute.toUpperCase().startsWith(CommandParseUtilities.GOR_DICTIONARY_FOLDER_PART)) {
writeOutGorDictionaryPart(commandToExecute, fileReader, outfile, theTheDict)
checkDictionaryCommit(writeOutGorDictionaryPart(commandToExecute, fileReader, outfile, theTheDict), fileReader, commandToExecute, isCancelled)
} else if (commandToExecute.toUpperCase().startsWith(CommandParseUtilities.GOR_DICTIONARY)) {
writeOutGorDictionary(commandToExecute, fileReader, outfile, theTheDict)
checkDictionaryCommit(writeOutGorDictionary(commandToExecute, fileReader, outfile, theTheDict), fileReader, commandToExecute, isCancelled)
} else if (commandToExecute.toUpperCase().startsWith(CommandParseUtilities.NOR_DICTIONARY)) {
writeOutNorDictionaryPart(commandToExecute, fileReader, outfile)
checkDictionaryCommit(writeOutNorDictionaryPart(commandToExecute, fileReader, outfile), fileReader, commandToExecute, isCancelled)
} else {
runCommandInternal(context, commandToExecute, outfile, useMd5)
runCommandInternal(context, commandToExecute, outfile, useMd5, isCancelled)
}
context.end()
result
}

private def runCommandInternal(context: GorContext, commandToExecute: String, outfile: String, useMd5: Boolean): String = {
/**
* Dictionaries are written in place, so on cancel remove a possibly partial dictionary file to keep it from
* being picked up from the cache (ENGKNOW-3979).
*/
private def checkDictionaryCommit(outfile: String, fileReader: FileReader, commandToExecute: String, isCancelled: BooleanSupplier): String = {
if (mustNotCommit(isCancelled)) {
try {
if (outfile != null && fileReader.exists(outfile) && !fileReader.isDirectory(outfile)) fileReader.delete(outfile)
} catch {
case _: Exception => /* do nothing */
}
throw new GorCancelledException(s"Query cancelled, result not stored: $commandToExecute", null)
}
outfile
}

private def runCommandInternal(context: GorContext, commandToExecute: String, outfile: String, useMd5: Boolean, isCancelled: BooleanSupplier): String = {
val theSource = new DynamicRowSource(commandToExecute, context)
val theHeader = theSource.getHeader

Expand Down Expand Up @@ -294,6 +327,12 @@ object GeneralQueryHandler {
}
}

// A cancelled run (e.g. a sibling parallel part failed) may have ended its source early without an
// exception. Never commit that partial output to the cache; the catch below removes the temp file (ENGKNOW-3979).
if (mustNotCommit(isCancelled)) {
throw new GorCancelledException(s"Query cancelled, result not stored: $commandToExecute", null)
}

if(oldName!=null && fileReader.exists(oldName) && !oldName.equals(newName)) {
fileReader.move(oldName, newName)
val oldMetaName = DataUtil.toFile(oldName, DataType.META)
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,92 @@
package gorsat.QueryHandlers;

import gorsat.DynIterator;
import gorsat.process.GorInputSources;
import gorsat.process.GorPipeCommands;
import gorsat.process.PipeInstance;
import gorsat.process.PipeOptions;
import gorsat.process.TestSessionFactory;
import org.gorpipe.exceptions.GorCancelledException;
import org.gorpipe.gor.session.GorSession;
import org.junit.After;
import org.junit.Assert;
import org.junit.Before;
import org.junit.Rule;
import org.junit.Test;
import org.junit.rules.TemporaryFolder;

import java.io.IOException;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.stream.Stream;

/**
* ENGKNOW-3979: output of a cancelled parallel part must never be committed, whether or not the part's
* thread still has its interrupt flag set.
*/
public class UTestGeneralQueryHandler {

@Rule
public TemporaryFolder workDir = new TemporaryFolder();

private Path workPath;
private GorSession session;

@Before
public void setUp() {
GorPipeCommands.register();
GorInputSources.register();
DynIterator.createGorIterator_$eq(PipeInstance::createGorIterator);
workPath = workDir.getRoot().toPath();
var options = new PipeOptions();
options.gorRoot_$eq(workPath.toString());
options.cacheDir_$eq(workPath.resolve("cache").toString());
options.requestId_$eq("test");
session = new TestSessionFactory(options, null, false, null, null).create();
}

@After
public void tearDown() {
if (session != null) session.close();
}

@Test
public void cancelledQueryIsNotCommitted() throws IOException {
var outfile = workPath.resolve("out.gor");

Assert.assertThrows(GorCancelledException.class, () -> GeneralQueryHandler.runCommand(
session.getGorContext(), "gorrows -p chr1:1-100", outfile.toString(), false, false, () -> true));

Assert.assertFalse("Cancelled output must not be committed", Files.exists(outfile));
assertNoFilesLeft();
}

@Test
public void cancelledDictionaryIsNotCommitted() throws IOException {
Files.writeString(workPath.resolve("a.tsv"), "#col\n1\n");
Files.writeString(workPath.resolve("b.tsv"), "#col\n2\n");
var outfile = workPath.resolve("out.nord");

Assert.assertThrows(GorCancelledException.class, () -> GeneralQueryHandler.runCommand(
session.getGorContext(), "NORDICT a.tsv a b.tsv b", outfile.toString(), false, false, () -> true));

Assert.assertFalse("Cancelled dictionary must not be committed", Files.exists(outfile));
}

@Test
public void notCancelledQueryIsCommitted() {
var outfile = workPath.resolve("out.gor");

GeneralQueryHandler.runCommand(
session.getGorContext(), "gorrows -p chr1:1-100", outfile.toString(), false, false, () -> false);

Assert.assertTrue(Files.exists(outfile));
}

private void assertNoFilesLeft() throws IOException {
try (Stream<Path> files = Files.list(workPath)) {
var left = files.filter(p -> p.getFileName().toString().startsWith("out")).toList();
Assert.assertTrue("Temp output must be removed, found: " + left, left.isEmpty());
}
}
}
1 change: 1 addition & 0 deletions gortools/src/test/java/gorsat/Script/UTestSignature.java
Original file line number Diff line number Diff line change
Expand Up @@ -228,6 +228,7 @@ public void testSignatureWithVersionedLinkFile() throws IOException, Interrupted
}

Files.writeString(dataPath1, "chr1\t2\n", StandardOpenOption.APPEND);
Files.setLastModifiedTime(dataPath1, FileTime.fromMillis(Files.getLastModifiedTime(dataPath1).toMillis() + 1000));

try (var session = factory.create()) {
var engine = ScriptEngineFactory.create(session.getGorContext());
Expand Down
39 changes: 39 additions & 0 deletions gortools/src/test/java/gorsat/UTestParallel.java
Original file line number Diff line number Diff line change
Expand Up @@ -22,12 +22,23 @@

package gorsat;

import org.gorpipe.exceptions.GorException;
import org.gorpipe.exceptions.GorParsingException;
import org.junit.Assert;
import org.junit.Rule;
import org.junit.Test;
import org.junit.rules.TemporaryFolder;

import java.io.IOException;
import java.nio.file.Files;
import java.nio.file.Path;
import java.nio.file.attribute.FileTime;

public class UTestParallel {

@Rule
public TemporaryFolder workDir = new TemporaryFolder();

@Test
public void testNorParallelQuery() {
String query = "create pnlist = norrows 100 -offset 100 | calc pn 'PN_'+rownum | signature -timeres 1; parallel -parts [pnlist] <(norrows 100 -offset #{col:rownum} | calc pn '#{col:pn}')";
Expand Down Expand Up @@ -79,4 +90,32 @@ public void testGorParallelQueryExceedingLimit() {
String query = "create pnlist = norrows 100 -offset 100 | calc pn 'PN_'+rownum | signature -timeres 1; parallel -parts [pnlist] -limit 10 <(gorrows -p chr1:0-#{col:rownum} | calc pn '#{col:pn}')";
TestUtils.runGorPipeCount(query);
}

/**
* ENGKNOW-3979: when one part of a parallel create fails, a sibling part that is still running gets
* interrupted. Its partial (header-only) output must not be committed to the result cache, otherwise
* a rerun after the failure is fixed reuses the empty part and silently drops its rows.
*/
@Test
public void testFailedPartDoesNotLeaveSiblingPartInCache() throws IOException {
Path root = workDir.getRoot().toPath();
Path cacheDir = Files.createDirectory(root.resolve("result_cache"));
Path failingPart = root.resolve("a.gor");
Files.writeString(failingPart, "chrom\tpos\tval\tdelay\nchr1\t1\tbad\t0\n");
Files.writeString(root.resolve("b.gor"),
"chrom\tpos\tval\tdelay\nchr2\t1\tok\t500\nchr2\t2\tok\t500\nchr2\t3\tok\t500\n");
Files.writeString(root.resolve("parts.tsv"), "#name\na\nb\n");

String query = "create xx = parallel -parts parts.tsv <(gor #{col:name}.gor | calc s sleep(delay) | throwif val = 'bad'); gor [xx]";

Assert.assertThrows(GorException.class,
() -> TestUtils.runGorPipe(query, root.toString(), cacheDir.toString(), false, null, null));

Files.writeString(failingPart, "chrom\tpos\tval\tdelay\nchr1\t1\tok\t0\n");
Files.setLastModifiedTime(failingPart, FileTime.fromMillis(System.currentTimeMillis() + 10000));

String result = TestUtils.runGorPipe(query, root.toString(), cacheDir.toString(), false, null, null);
Assert.assertEquals("All rows from all parts expected on rerun, got:\n" + result,
5, result.split("\n").length);
}
}
Loading
Loading