University project for the Concurrent Programming course. The task was to write a parallel version of a Datalog engine's deriver: the part that decides which queries can be derived from the rules. The parser, the data structures and the single-threaded SimpleDeriver were given. My part is ParallelDeriver.java (in src/main/java/cp2025/engine/) plus some tests.
A program has constants, rules like reach(X) :- edge(X, Y), reach(Y). and queries like reach(a). For every query we have to say whether it can be derived. Some predicates are not defined by rules but by an oracle, an external object whose calculate(atom) can block for a long time.
Every query is a separate task in a fixed thread pool. Each task does a depth-first search like SimpleDeriver. It tries every rule whose head matches the goal and every assignment of constants to the remaining variables, then derives the atoms in the body one by one.
The threads share two concurrent maps:
resultCache: atoms whose answer is already known. A thread always checks it before deriving anything.computingThread: which thread is currently computing a given atom.
Each thread also keeps its own inProgress set with the atoms on its current recursion path. If a goal is already there, we hit a cycle (e.g. p :- q. and q :- p.), so that path returns false.
Caching follows the same rules as SimpleDeriver. A true result is stored right away. A false result is stored only when the thread is back at the root of its query, because a false found deeper can depend on some atom higher up being treated as false to cut a cycle.
Say thread T1 is deriving a and needs b, and T2 is already working on b. T1 does not wait for T2. If b is not in the cache yet, T1 starts deriving b itself and registers as the thread computing it.
When a thread gets the answer for b, it puts it in the cache and interrupts the thread registered in computingThread for b. The interrupted thread checks the cache for its goal. If the answer is there, it returns it and drops its own work. This matters most for oracle calls: a thread stuck in a slow calculate(b) gets an InterruptedException, finds the result in the cache and moves on immediately.
The same interrupt is used for real cancellation. When derive() finishes or the caller interrupts it, it sets the globalInterrupted flag and calls shutdownNow() on the pool. A worker interrupted while this flag is set throws InterruptedException and ends. So one mechanism covers both "the thing you're computing is ready" and "stop everything".
No waiting, so no deadlocks. The obvious alternative is for T1 to block until T2 finishes b. But Datalog rules can be cyclic. If T1 computes a and needs b while T2 computes b and needs a, they would wait for each other forever. With a fixed pool, threads blocked on each other would also just occupy workers. Here a thread never waits for another thread and never holds a lock while computing. The worst case is that two threads compute the same atom at the same time.
Interrupts instead of a flag or wait/notify. A thread can be blocked inside the oracle, e.g. on I/O. A boolean flag would only be noticed after that call returns, but an interrupt breaks it out right away. It's also the standard way to cancel work in Java, and the deriver has to support it anyway because derive() must throw InterruptedException when the caller interrupts it.
Correctness doesn't depend on the signal. The cache is the source of truth and every thread checks it before deriving an atom. If a signal comes late or not at all, a thread just does some extra work and still gets the same answer.
Needs Java 21.
make # compile
make test # run tests
make examples # run examples/ and compare with the expected results
Run a single program:
java -cp target/classes:lib/antlr-4.13.2-complete.jar cp2025.engine.Main < examples/14_BasicTransitivity.d