diff --git a/CHANGELOG.md b/CHANGELOG.md
index 6614123a..defbb0ba 100644
--- a/CHANGELOG.md
+++ b/CHANGELOG.md
@@ -1624,3 +1624,99 @@ and the consumer PRs that deliver it).
default; every existing call site, Kotlin or Java, keeps compiling and keeps its prior behavior
short of the fixes themselves. Design note:
[docs/design/extraction-runs.md](docs/design/extraction-runs.md).
+
+- **EXPERIMENTAL.** `DrivineExtractionRunStore` — the durable Neo4j implementation of
+ `ExtractionRunStore`, completing DICE #67's storage half. It writes three node labels:
+ `(:ExtractionRun)` keyed `(contextId, runId)` for the header,
+ `(:ExtractionRun)-[:RECORDED]->(:ExtractionRunInvocation)` keyed
+ `(contextId, runId, invocationIndex, attempt)` for one attempt at one planned model call, and
+ `(:ExtractionRun)-[:ENDED_BY]->(:ExtractionRunTerminalWrite)` keyed `(contextId, runId)` for the
+ fingerprint of the write that ended the run. Every key is tenant-qualified, and unlike
+ `DrivineDriftReportStore` the tenant needs no `ctx:`-prefixed stand-in: that store's scope is
+ nullable and a Cypher MERGE cannot key on a null, while a run's tenant never is, so the plain value
+ is already an injective key.
+ **Compare-and-set is one Cypher statement, and it holds across processes rather than only across
+ threads.** The statement takes an exclusive node lock with `SET n.casLock = $lockToken` before it
+ reads the run's status, so a second transaction blocks and then reads at read-committed — after the
+ first committed — and takes the no-op branch. The lock token is a fresh UUID on every call so the
+ write is always a real change rather than a no-op a database may optimize away before locking, and
+ no index carries `status` or the terminal fingerprint, so the planner cannot serve the post-lock
+ read from an index entry it read at MATCH time. Underneath that argument sits a fact: the terminal
+ write is a `CREATE` of the terminal-write node under a uniqueness constraint, so two writers that
+ both read a run as `RUNNING` cannot both commit. The loser's transaction rolls back whole, the
+ store catches the violation, re-reads the recorded fingerprint in a fresh transaction, and answers
+ replayed or conflict. The constraint's sufficiency is measured: removing the lock leaves every
+ race test green, and removing both produces six racing writers all reporting `APPLIED` and three
+ contradictory endings recorded for one run. The lock's rests on the documented isolation argument,
+ which is why the constraint exists.
+ **The fingerprint is stored verbatim and compared verbatim, never re-derived.** The node carries the
+ exact string `ExtractionRunTransition.fingerprint` computed, so a correct retry that happened after
+ another attempt was recorded still replays; a store deriving a digest from the stored run would
+ reject it. That string is the `xrun-terminal:v2` digest of the transition's identity, so the counts
+ and failures a terminal write carries ride into the header beside it and reach none of the compared
+ bytes; the store keeps no comparison of its own that could fall out of step.
+ **A run that ends announces itself once, when the write is durable.** `transition` hands an
+ `ExtractionRunTransitioned` to the listener the store was constructed with, defaulting to
+ `DiceEventListener.DEV_NULL` the way the in-memory reference's does. Exactly one call per run
+ reaches that branch, for the schema reason above: reaching it means having created the
+ terminal-write node. A replay, a rejected write, a `save`, a `recordInvocation`, and the writer that
+ lost the race all announce nothing. When the store owns the transaction the listener runs once the
+ template has committed; when a caller's transaction is active the announcement is registered against
+ that caller's commit, so a rollback drops it and no consumer hears about a run nothing can read back. The
+ six event cases the contract suite added run against Neo4j unmodified.
+ **Failures are stored in the closed vocabulary and nothing else fits.** A run's failures are one
+ JSON array on the header node, each element carrying `code`, `stage`, `providerStatus`, a measure as
+ a quantity and a value, `at`, and the attempt it names as an index and an attempt — eight fields,
+ written flat, with optional ones stored as nulls so every failure stores the same keys. No property
+ holds free text, because `ExtractionFailure` has no text-shaped field to write from. A round-trip
+ test reads the properties Neo4j actually holds and checks the key set it finds matches an allowlist
+ the test states itself, so an added `detail` column fails the build before it reaches a graph, and a stored
+ failure holding half a measure or half an invocation id is refused on read.
+ **A header write cannot touch a child row.** Invocation records are their own nodes, so `save` has
+ no way to delete one — the contract's merge-don't-replace rule falls out of the graph model instead
+ of being implemented. One consequence: a durable store keeps identified rows rather than the order a
+ caller listed them in, so it returns attempts in plan order where the in-memory reference returns
+ the caller's order. `invocationsOf` is plan order in both.
+ **Every page scopes in the query ahead of its `LIMIT`**, excludes rows with no sort key (Neo4j sorts
+ null largest, so one would sort to the front of a `DESC` order, spend a slot, and then be dropped by
+ the mapper), and skips corrupt rows with a warning rather than failing the whole read. `runsOfRoot`
+ is one indexed lookup on the denormalized root. The chain walk is client-side and bounded to `limit`
+ keyed lookups in one read transaction, cycle-safe and tenant-scoped at every hop: a parent is a
+ property rather than a relationship, because a run can name a parent not yet stored, and this module
+ takes no APOC dependency. The store passes the whole `AbstractExtractionRunStoreContractTest` suite
+ alongside the in-memory reference, plus Drivine-specific integration tests for multi-writer
+ compare-and-set, cross-tenant fail-closed with identical run ids in two tenants, full-fidelity row
+ round-trip, corrupt-row skip, and the privacy contract asserted over the properties actually
+ written. Design note:
+ [docs/design/extraction-runs.md](docs/design/extraction-runs.md).
+ **Compatibility: additive, and hosts must declare new schema.** No existing class loses a member and
+ no behaviour changes; `DrivineExtractionRunStore`, `ExtractionRunSchema`, `ExtractionRunRowMapper`
+ and `ExtractionInvocationRowMapper` are new types in `com.embabel.dice.storage`. Nothing is
+ auto-configured yet, so a host opts in by declaring the store bean and a `SchemaCatalog` carrying
+ `ExtractionRunSchema.specs()`. That catalog is **eight new schema items**, and the three constraints
+ are required rather than advisory — a MERGE on a natural key is race-free only under one, and the
+ third is what makes the compare-and-set a schema fact:
+ - `UniquenessConstraintSpec("ExtractionRun", ["contextId", "runId"])`
+ - `UniquenessConstraintSpec("ExtractionRunInvocation", ["contextId", "runId", "invocationIndex", "attempt"])`
+ - `UniquenessConstraintSpec("ExtractionRunTerminalWrite", ["contextId", "runId"])`
+ - `RangeIndexSpec("ExtractionRun", "contextId")` — the tenant page; a composite index cannot stand
+ in, because Neo4j will not use one for a predicate on only its leading property
+ - `RangeIndexSpec("ExtractionRun", ["contextId", "rootRunId"])` — the whole-lineage read
+ - `RangeIndexSpec("ExtractionRun", ["contextId", "parentRunId"])` — one hop down the parent axis
+ - `RangeIndexSpec("ExtractionRun", ["contextId", "startedAtEpochSecond"])` — the paging sort key
+ - `RangeIndexSpec("ExtractionRunInvocation", ["contextId", "runId"])` — one run's attempts
+
+ No stored data changes and no migration is required: no released DICE ever wrote these labels. One
+ new bound a host should know about — `ContextId` accepts any non-blank string and the tenant is the
+ leading property of every key here, so the store rejects a tenant id longer than 1024 characters on
+ the write path rather than letting Neo4j fail the write mid-extraction with an index-key-size error.
+ Reads are uncapped, because a read for a longer tenant matches nothing by construction. Every new
+ type carries `@ApiStatus.Experimental` and the shapes may still move while the remaining #67 slices
+ land.
+ **Race detection keys on the driver's status code now, and message text no longer matters.**
+ `Neo4jErrors.isUniquenessViolation` walks the cause chain for a `Neo4jException` whose `code()` is
+ `Neo.ClientError.Schema.ConstraintValidationFailed`, and both `DrivineExtractionRunStore` and
+ `DrivinePropositionRepository` call it to tell a lost compare-and-set race apart from a real
+ failure. `dice-storage` now declares the `neo4j-java-driver` dependency it already ran with
+ through Drivine, so the exception class it checks is visible at compile time too.
+ **Compatibility: additive.** No public signature changes.
diff --git a/dice-storage/pom.xml b/dice-storage/pom.xml
index b029ad5e..e15e68e4 100644
--- a/dice-storage/pom.xml
+++ b/dice-storage/pom.xml
@@ -60,6 +60,15 @@
spring-tx
+
+
+ org.neo4j.driver
+ neo4j-java-driver
+
+
org.jetbrains.kotlin
kotlin-stdlib
diff --git a/dice-storage/src/main/kotlin/com/embabel/dice/storage/DrivineExtractionRunStore.kt b/dice-storage/src/main/kotlin/com/embabel/dice/storage/DrivineExtractionRunStore.kt
new file mode 100644
index 00000000..3a4aac50
--- /dev/null
+++ b/dice-storage/src/main/kotlin/com/embabel/dice/storage/DrivineExtractionRunStore.kt
@@ -0,0 +1,1007 @@
+/*
+ * Copyright 2024-2026 Embabel Pty Ltd.
+ *
+ * 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.embabel.dice.storage
+
+import com.embabel.dice.common.DiceEventListener
+import com.embabel.dice.common.ExtractionRunTransitioned
+import com.embabel.dice.proposition.extraction.ExtractionInvocationOutcome
+import com.embabel.dice.proposition.extraction.ExtractionInvocationRecord
+import com.embabel.dice.proposition.extraction.ExtractionRun
+import com.embabel.dice.proposition.extraction.ExtractionRunConflictException
+import com.embabel.dice.proposition.extraction.ExtractionRunKey
+import com.embabel.dice.proposition.extraction.ExtractionRunNotFoundException
+import com.embabel.dice.proposition.extraction.ExtractionRunRef
+import com.embabel.dice.proposition.extraction.ExtractionRunStatus
+import com.embabel.dice.proposition.extraction.ExtractionRunStore
+import com.embabel.dice.proposition.extraction.ExtractionRunTransition
+import com.embabel.dice.proposition.extraction.ExtractionRunTransitionOutcome
+import com.embabel.dice.proposition.extraction.ExtractionRunTransitionResult
+import org.drivine.manager.PersistenceManager
+import org.drivine.query.QuerySpecification
+import org.slf4j.LoggerFactory
+import org.springframework.transaction.PlatformTransactionManager
+import org.springframework.transaction.annotation.Propagation
+import org.springframework.transaction.annotation.Transactional
+import org.springframework.transaction.support.TransactionSynchronization
+import org.springframework.transaction.support.TransactionSynchronizationManager
+import org.springframework.transaction.support.TransactionTemplate
+import java.time.Clock
+import java.time.Instant
+import java.util.UUID
+
+/**
+ * Drivine / Neo4j implementation of [ExtractionRunStore].
+ *
+ * ## The graph
+ *
+ * - `(:ExtractionRun {contextId, runId, ...})` — the run header, one node per run per tenant.
+ * - `(:ExtractionRun)-[:RECORDED]->(:ExtractionRunInvocation {contextId, runId, invocationIndex, attempt, ...})`
+ * — one child node per attempt at one planned model call.
+ * - `(:ExtractionRun)-[:ENDED_BY]->(:ExtractionRunTerminalWrite {contextId, runId, fingerprint, ...})`
+ * — at most one per run, ever, and the thing that makes that "at most one" true.
+ *
+ * Every key is tenant-qualified because a run id is host-minted and DICE never assumes it is
+ * globally unique. Unlike `DrivineDriftReportStore`, the tenant needs no `ctx:`-prefixed stand-in:
+ * that store's scope is nullable, and a Cypher MERGE cannot key on a null, so a global report needed
+ * a non-null encoding that no real context id could collide with. A run's tenant is never null, so
+ * the plain value is already an injective key.
+ *
+ * Every statement is parameterized; nothing caller-derived is interpolated into Cypher. The
+ * statements are assembled only from this class's own literals.
+ *
+ * See [ExtractionRunSchema] for the constraints and indexes this depends on. They are not tuning:
+ * a MERGE on a natural key is race-free only under a uniqueness constraint on that key.
+ *
+ * ## Compare-and-set: why a terminal write cannot be applied twice
+ *
+ * The whole of [transition] is one Cypher statement, so it is one transaction, and it works two
+ * ways at once — a lock that makes the race rare and a constraint that makes the wrong answer
+ * impossible.
+ *
+ * **The lock.** The statement's first act after matching the run is `SET n.casLock = $lockToken`,
+ * before it reads anything. A property write takes an exclusive lock on the node and holds it until
+ * commit, so a second transaction reaching that line blocks until the first commits. Neo4j reads at
+ * read-committed, and the read of `n.status` is downstream of the `SET` in the same statement, so
+ * the second transaction reads the status *after* it acquired the lock and therefore *after* the
+ * first transaction committed. It sees a terminal run and takes the no-op branch. This is the same
+ * write-lock-before-read idiom `DrivineMetamodelVersionStore` uses on its counter.
+ *
+ * Two details make that argument hold rather than nearly hold. The lock token is a fresh UUID on
+ * every call, so the write is always a real change and never a no-op a future Neo4j might optimize
+ * away before taking the lock. And no index carries `status` or the terminal fingerprint — see
+ * [ExtractionRunSchema] — so the planner cannot serve the post-lock read from an index entry it
+ * read at MATCH time, which is the one way a lock-then-read can quietly read stale.
+ *
+ * **The constraint.** The argument above is about Neo4j's behaviour, and behaviour is a thing to be
+ * wrong about. So the terminal write is also a `CREATE` of an `(:ExtractionRunTerminalWrite)` node
+ * on the run's own key, under a uniqueness constraint on `(contextId, runId)`. If two transactions
+ * ever did both read a run as `RUNNING` — a Neo4j build that locks differently, a cluster, a
+ * planner that reorders in a way this KDoc did not anticipate — they would both try to create that
+ * node and the database would refuse the second one at commit. The loser's whole transaction rolls
+ * back, header included, and this store catches the violation, re-reads the recorded fingerprint in
+ * a fresh transaction, and returns the replay-or-conflict answer the contract asks for. Exactly one
+ * terminal write per run is a schema fact, not an inference.
+ *
+ * That is the multi-process half of the guarantee. Threads in one JVM could be serialized by a
+ * monitor; two processes cannot be, and nothing in this class holds mutable state to serialize on.
+ * Both halves come entirely from the database.
+ *
+ * **The fingerprint is stored, never re-derived.** The node carries the exact string
+ * [ExtractionRunTransition.fingerprint] computed, and a repeated terminal write is decided by
+ * comparing its fingerprint against that string. Re-deriving a digest from the stored run would
+ * reject a correct retry whenever an attempt had been recorded in between, because the run it
+ * derived from would have changed while the terminal write did not.
+ *
+ * ## A header write never touches a child row
+ *
+ * Invocation records are their own nodes on their own key, and [save]'s Cypher names none of them:
+ * it reads and writes `(:ExtractionRun)` alone. [recordInvocation] is the only method that ever
+ * names an `(:ExtractionRunInvocation)` node, so it is the only one that can create, update, or
+ * lock one. Whatever `run.invocations` a caller hands [save] plays no part in what it accepts,
+ * rejects, or persists.
+ *
+ * ## Validation happens before any node is created
+ *
+ * A `MERGE` on a node pattern creates the node the moment it runs, whether or not the write that
+ * follows turns out to be valid — a naive `MERGE` immediately followed by a Kotlin check that
+ * throws leaves the created node behind if whatever catches that exception does not also abort the
+ * transaction. `SAVE_RUN` and `RECORD_INVOCATION` both avoid this: an `OPTIONAL MATCH` reads
+ * whatever is already there, a `CASE` decides in Cypher whether the write should happen, and the
+ * `MERGE` that can create a node sits inside a `FOREACH` gated on that decision. A save naming
+ * version `0` for a run already stored at a later version, or a [recordInvocation] call against a
+ * run that has already ended, stops before reaching any `MERGE` at all, so a caller that catches
+ * the resulting [ExtractionRunConflictException] without rolling back its own transaction finds
+ * the graph exactly as it was.
+ *
+ * ## Scope is applied in the query, ahead of the limit
+ *
+ * Each page puts its tenant in the MATCH pattern and its `LIMIT` after the `ORDER BY`. Reading a
+ * limited page and filtering it in Kotlin would apply the limit first, so a tenant whose neighbour
+ * owns the head of the index would report no runs while the store held plenty. This is
+ * `DrivineDriftReportStore`'s rule carried over, and the cross-backend suite pins it.
+ *
+ * Every read returns each run with its invocation records attached, pages included, because the
+ * in-memory reference does and the two are held to one suite. The cost is bounded — at most `limit`
+ * runs times the run model's cap of 1024 attempts — but it is a real cost on a wide page.
+ *
+ * ## A run that ends announces itself once
+ *
+ * [transition] hands an [ExtractionRunTransitioned] to [listener] for the call that ended the run,
+ * and for no other call: a replay, a rejected write, a [save] and a [recordInvocation] all announce
+ * nothing. Exactly one call per run reaches the applied branch, because reaching it means having
+ * created the run's terminal-write node, and the uniqueness constraint lets one transaction do that.
+ *
+ * The announcement waits for the write to be durable. When this store owns the transaction, that
+ * point is where the template returns, and the listener runs there. When a caller's transaction is
+ * active the write is durable when that caller commits, so the announcement is registered against
+ * the commit and never happens at all if the caller rolls back — a listener told a run ended by a
+ * transaction that was thrown away would be reporting a run nothing can read.
+ *
+ * @param persistenceManager Drivine's handle on the `neo` datasource.
+ * @param transactionManager used to own a transaction when no caller has one, so a lost
+ * compare-and-set race can be recovered in a transaction the race did not already poison.
+ * @param clock supplies the instant a terminal write is recorded at. Informational, and injectable
+ * so a test can pin it; nothing sorts or compares on it.
+ * @param listener Notified when a run ends. Defaults to [DiceEventListener.DEV_NULL], so a host
+ * that has nothing listening constructs the store the way it always did. Handlers run inline on
+ * whichever thread the announcement happens on, and throw isolation belongs to the listener —
+ * wrap it in `SafeDiceEventListener` for graceful degradation.
+ */
+@Transactional
+open class DrivineExtractionRunStore @JvmOverloads constructor(
+ private val persistenceManager: PersistenceManager,
+ transactionManager: PlatformTransactionManager,
+ private val clock: Clock = Clock.systemUTC(),
+ private val listener: DiceEventListener = DiceEventListener.DEV_NULL,
+) : ExtractionRunStore {
+
+ private val logger = LoggerFactory.getLogger(DrivineExtractionRunStore::class.java)
+
+ private val txTemplate = TransactionTemplate(transactionManager)
+
+ // ---- writes ----
+
+ /**
+ * Records a running run, inserting it or updating the one already there — in one statement, so
+ * the check against the stored run and the write that follows it cannot be interleaved by a
+ * concurrent terminal write, and the header's compare-and-set on
+ * [com.embabel.dice.proposition.extraction.ExtractionRun.version] cannot be interleaved by a
+ * concurrent header save either.
+ *
+ * The statement returns what it found so this method can name which rule was broken. Cypher
+ * decides whether to write and Kotlin decides what to say about it, and the conditions are the
+ * same condition, written twice: a save writes when the run is new and names version 0; when it
+ * is still running, agrees with the stored lineage and start time, and its header content
+ * already matches what is stored — a no-op that leaves the version untouched, whatever version
+ * the save named; or when it is still running, agrees with lineage and start time, its content
+ * genuinely differs, and it names the version currently stored. Every other running-and-agreeing
+ * case is a stale write and is rejected without touching the row — see this class's KDoc for why
+ * that rejection never leaves a node behind. `run.invocations` is not part of the statement at
+ * all: this method's Cypher has no clause that reads or writes an `(:ExtractionRunInvocation)`.
+ *
+ * Runs and retries under [ownedTransaction] — see that method for the two shapes.
+ */
+ @Transactional(propagation = Propagation.SUPPORTS)
+ override fun save(run: ExtractionRun): ExtractionRun {
+ require(run.status == ExtractionRunStatus.RUNNING) {
+ "save records a running run; ${run.status} is terminal and belongs to transition()"
+ }
+ require(run.finishedAt == null) {
+ "a running run has not finished, so it carries no finishedAt"
+ }
+ val key = run.key()
+ ExtractionRunSchema.requireStorableTenant(key.contextId)
+
+ val lineageKey = ExtractionRunRowMapper.lineageKeyOf(run.lineage)
+ val startedAt = run.startedAt.toString()
+ val header = ExtractionRunRowMapper.headerBindMap(run)
+ val headerFingerprint = ExtractionRunRowMapper.headerFingerprint(header)
+
+ return ownedTransaction {
+ val row = singleRow(
+ SAVE_RUN,
+ keyBindings(key) + mapOf(
+ "lockToken" to newLockToken(),
+ "running" to RUNNING,
+ "lineageKey" to lineageKey,
+ "startedAt" to startedAt,
+ "header" to header,
+ "headerFingerprint" to headerFingerprint,
+ "version" to run.version,
+ ),
+ ) ?: throw IllegalStateException("SAVE_RUN returned no row for ${key.runRef.runId}")
+
+ when (val priorStatus = row["priorStatus"]?.toString()) {
+ null -> require(run.version == 0L) {
+ "a run's first save must name version 0, the version a run nobody has saved " +
+ "yet carries; this one names ${run.version}"
+ }
+
+ RUNNING -> {
+ if (row["priorLineageKey"]?.toString() != lineageKey) {
+ throw ExtractionRunConflictException(
+ key,
+ "stored lineage differs from the one being saved; lineage is fixed at insert",
+ )
+ }
+ if (row["priorStartedAt"]?.toString() != startedAt) {
+ throw ExtractionRunConflictException(
+ key,
+ "stored startedAt differs from the one being saved; a run starts once",
+ )
+ }
+ if (row["priorHeaderFingerprint"]?.toString() != headerFingerprint) {
+ val priorVersion = (row["priorVersion"] as? Number)?.toLong()
+ ?: throw IllegalStateException(
+ "run ${key.runRef.runId} in context ${key.contextId.value} is RUNNING " +
+ "with no stored version",
+ )
+ if (run.version != priorVersion) {
+ throw ExtractionRunConflictException(
+ key,
+ "the header was read at version ${run.version}; the store is now at " +
+ "$priorVersion. Read the run again with findRun and rebuild this " +
+ "save on what it holds now",
+ )
+ }
+ }
+ // else: content already matches what is stored, so the statement above left the
+ // row exactly as it was — a no-op, whatever version this save named.
+ }
+
+ else -> throw ExtractionRunConflictException(
+ key,
+ "already ended as $priorStatus and cannot be re-opened by a save",
+ )
+ }
+
+ requireStoredRun(key)
+ }
+ }
+
+ /**
+ * Records one attempt against a running run, on its own child node.
+ *
+ * One statement again, and the same shape as [save]: read the run's status and the attempt's
+ * prior outcome and fingerprint, write only if the run is still running and the attempt is not
+ * locked as terminal under a different payload — and, per this class's KDoc, only create the
+ * node at all once that decision comes back yes. A caller that holds only the attempt never has
+ * to hold the header, and this can never overwrite one.
+ *
+ * **Once the attempt is terminal, only an identical write reaches the row.** The incoming
+ * record's fingerprint — [ExtractionInvocationRowMapper.fingerprint] — is compared straight
+ * against the string already stored on the node, always computed fresh from the record a
+ * caller handed this method; the run's own terminal write follows the identical stored-string,
+ * compared-verbatim rule, one level down.
+ *
+ * Runs and retries under [ownedTransaction] — see that method for the two shapes.
+ */
+ @Transactional(propagation = Propagation.SUPPORTS)
+ override fun recordInvocation(
+ key: ExtractionRunKey,
+ record: ExtractionInvocationRecord,
+ ): ExtractionRun {
+ ExtractionRunSchema.requireStorableTenant(key.contextId)
+ val fingerprint = ExtractionInvocationRowMapper.fingerprint(record)
+
+ return ownedTransaction {
+ val row = singleRow(
+ RECORD_INVOCATION,
+ keyBindings(key) + mapOf(
+ "lockToken" to newLockToken(),
+ "running" to RUNNING,
+ "inFlight" to IN_FLIGHT,
+ "invocationIndex" to record.invocationIndex,
+ "attempt" to record.attempt,
+ "recordFingerprint" to fingerprint,
+ "record" to (ExtractionInvocationRowMapper.bindMap(record) + mapOf("recordFingerprint" to fingerprint)),
+ ),
+ ) ?: throw ExtractionRunNotFoundException(key)
+
+ val priorStatus = row["priorStatus"]?.toString()
+ if (priorStatus != RUNNING) {
+ throw ExtractionRunConflictException(
+ key,
+ "already ended as $priorStatus; its invocation records are part of how it ended",
+ )
+ }
+ if (row["locked"] == true) {
+ throw ExtractionRunConflictException(
+ key,
+ "${record.id} already ended as ${row["priorOutcome"]}; once an attempt is terminal " +
+ "only an identical write replays, and this one differs",
+ )
+ }
+ requireStoredRun(key)
+ }
+ }
+
+ /**
+ * Ends a run under compare-and-set. See this class's KDoc for why the answer is a schema fact
+ * rather than a claim about timing.
+ *
+ * Two transaction shapes, for the same reason `DrivinePropositionRepository.save` has two:
+ * - **A caller's transaction is active.** This joins it. A lost race has already ended that
+ * transaction by the time this method could react, so no recovery of ours could commit in it;
+ * the violation propagates and retrying is the caller's. Opening a nested transaction to sneak
+ * a read out would break the atomicity the caller asked for.
+ * - **No caller transaction.** `SUPPORTS` keeps the proxy from opening one, so the template owns
+ * the attempt, and the recovery read runs in a fresh transaction the failed attempt cannot
+ * roll back.
+ */
+ @Transactional(propagation = Propagation.SUPPORTS)
+ override fun transition(
+ key: ExtractionRunKey,
+ transition: ExtractionRunTransition,
+ ): ExtractionRunTransitionResult {
+ ExtractionRunSchema.requireStorableTenant(key.contextId)
+ if (TransactionSynchronizationManager.isActualTransactionActive()) {
+ return announce(endRun(key, transition))
+ }
+
+ var lostRace: RuntimeException? = null
+ val result = try {
+ txTemplate.execute { endRun(key, transition) }
+ } catch (e: RuntimeException) {
+ // Neo4jErrors tells a lost compare-and-set race apart from a real failure by the
+ // driver's own status code, not by the exception's message.
+ if (!Neo4jErrors.isUniquenessViolation(e)) throw e
+ lostRace = e
+ null
+ }
+ if (result != null) return announce(result)
+
+ logger.debug(
+ "Terminal write for run {} in context {} lost the race; re-reading the recorded write",
+ key.runRef.runId,
+ key.contextId.value,
+ )
+ return txTemplate.execute { resolveLostRace(key, transition, lostRace!!) }!!
+ }
+
+ /**
+ * Runs the compare-and-set statement and turns what it found into the contract's answer.
+ *
+ * `priorStatus` is what the run was when this transaction took its lock. `RUNNING` means the
+ * statement applied the terminal fields and created the terminal-write node; anything else means
+ * it wrote nothing, and the decision is a comparison of the stored fingerprint against this one.
+ */
+ private fun endRun(
+ key: ExtractionRunKey,
+ transition: ExtractionRunTransition,
+ ): ExtractionRunTransitionResult {
+ val row = singleRow(
+ END_RUN,
+ keyBindings(key) + mapOf(
+ "lockToken" to newLockToken(),
+ "running" to RUNNING,
+ "status" to transition.status.name,
+ // The string the transition computed, stored verbatim. Nothing here derives it.
+ "fingerprint" to transition.fingerprint,
+ "recordedAt" to clock.instant().toString(),
+ "terminal" to ExtractionRunRowMapper.terminalBindMap(transition),
+ ),
+ ) ?: throw ExtractionRunNotFoundException(key)
+
+ val priorStatus = row["priorStatus"]?.toString()
+ if (priorStatus == RUNNING) {
+ logger.debug(
+ "Ended run {} in context {} as {}",
+ key.runRef.runId,
+ key.contextId.value,
+ transition.status,
+ )
+ return ExtractionRunTransitionResult(
+ requireStoredRun(key),
+ ExtractionRunTransitionOutcome.APPLIED,
+ )
+ }
+ return replayOrConflict(key, transition, row["priorFingerprint"]?.toString(), priorStatus)
+ }
+
+ /**
+ * Decides what a terminal write against an already-ended run means, by comparing fingerprints.
+ *
+ * A missing recorded fingerprint is a conflict, not a replay: a terminal run this store cannot
+ * show a terminal write for is one it cannot prove agrees with the incoming one, and agreeing is
+ * the only thing that makes overwriting safe.
+ */
+ private fun replayOrConflict(
+ key: ExtractionRunKey,
+ transition: ExtractionRunTransition,
+ recordedFingerprint: String?,
+ priorStatus: String?,
+ ): ExtractionRunTransitionResult {
+ if (recordedFingerprint != null && recordedFingerprint == transition.fingerprint) {
+ return ExtractionRunTransitionResult(
+ requireStoredRun(key),
+ ExtractionRunTransitionOutcome.REPLAYED,
+ )
+ }
+ throw ExtractionRunConflictException(
+ key,
+ "already ended as $priorStatus under a different terminal write; " +
+ "this one claims ${transition.status}",
+ )
+ }
+
+ /**
+ * Recovers from losing the create of the terminal-write node: read what the winner recorded and
+ * decide against it.
+ *
+ * If there is no terminal write to find, the violation was not the one this method knows how to
+ * interpret, so the original exception is rethrown rather than guessed about.
+ */
+ private fun resolveLostRace(
+ key: ExtractionRunKey,
+ transition: ExtractionRunTransition,
+ violation: RuntimeException,
+ ): ExtractionRunTransitionResult {
+ val row = singleRow(TERMINAL_WRITE_BY_KEY, keyBindings(key)) ?: throw violation
+ // Nothing announced here: the writer that lost the race created no terminal-write node, so
+ // this path answers replayed or conflict and never applied.
+ return replayOrConflict(
+ key,
+ transition,
+ row["fingerprint"]?.toString(),
+ row["status"]?.toString(),
+ )
+ }
+
+ /**
+ * Tells [listener] about a run this call ended, once the write is durable, and returns [result]
+ * so a caller can announce and return in one line.
+ *
+ * A replay and a rejected write reach neither branch below — a rejected write throws before
+ * getting here, and a replay is not applied. Inside a caller's transaction the write becomes
+ * durable at that caller's commit, so the announcement rides on it and is dropped with the
+ * transaction if the caller rolls back. Otherwise the commit has already happened by the time
+ * this runs, and the listener is called on the spot.
+ */
+ private fun announce(result: ExtractionRunTransitionResult): ExtractionRunTransitionResult {
+ if (!result.isApplied) return result
+ val event = ExtractionRunTransitioned(result.run)
+ if (TransactionSynchronizationManager.isSynchronizationActive()) {
+ TransactionSynchronizationManager.registerSynchronization(
+ object : TransactionSynchronization {
+ override fun afterCommit() {
+ listener.onEvent(event)
+ }
+ },
+ )
+ return result
+ }
+ listener.onEvent(event)
+ return result
+ }
+
+ // ---- reads ----
+
+ @Transactional(readOnly = true)
+ override fun findRun(key: ExtractionRunKey): ExtractionRun? =
+ singleRow(RUN_BY_KEY, keyBindings(key))?.let(::mapRun)
+
+ @Transactional(readOnly = true)
+ override fun invocationsOf(key: ExtractionRunKey): List =
+ queryRows(INVOCATIONS_OF, keyBindings(key)).mapNotNull(::mapInvocation)
+
+ @Transactional(readOnly = true)
+ override fun runsInContext(
+ contextIdValue: String,
+ limit: Int,
+ since: Instant?,
+ ): List = readPage(
+ scope = null,
+ contextIdValue = contextIdValue,
+ limit = limit,
+ since = since,
+ )
+
+ @Transactional(readOnly = true)
+ override fun childrenOf(
+ contextIdValue: String,
+ parentRunId: String,
+ limit: Int,
+ ): List = readPage(
+ scope = ONLY_PARENT,
+ contextIdValue = contextIdValue,
+ limit = limit,
+ since = null,
+ extraBindings = mapOf("parentRunId" to parentRunId),
+ )
+
+ @Transactional(readOnly = true)
+ override fun runsOfRoot(
+ contextIdValue: String,
+ rootRunId: String,
+ limit: Int,
+ since: Instant?,
+ ): List = readPage(
+ scope = ONLY_ROOT,
+ contextIdValue = contextIdValue,
+ limit = limit,
+ since = since,
+ extraBindings = mapOf("rootRunId" to rootRunId),
+ )
+
+ /**
+ * Walks the parent chain upward, one keyed lookup per hop, at most [limit] hops.
+ *
+ * **Why the walk is here and not in Cypher.** A parent is a property, not a relationship: a run
+ * can name a parent that has not been stored yet, and an edge cannot point at a node that does
+ * not exist. Materializing the edge later would mean a second write nothing triggers. Without an
+ * edge there is no variable-length pattern to walk, and the APOC procedures that would do it in
+ * one round trip are not a dependency this module takes.
+ *
+ * So the walk is client-side and bounded by construction: at most `limit` hops, each a lookup on
+ * the uniqueness-constraint index, all inside one read transaction. It stops on a run it has
+ * already seen, which is what makes it safe on a corrupt store holding a cycle — [ExtractionRunLineage]
+ * can reject a run that is its own parent, but a two-hop cycle needs the other runs to see. And
+ * it resolves every hop inside the starting run's tenant, so a parent id that exists only in a
+ * neighbour's tenant resolves to nothing and the walk ends.
+ */
+ @Transactional(readOnly = true)
+ override fun ancestorsOf(key: ExtractionRunKey, limit: Int): List {
+ requirePositiveLimit(limit)
+ val start = findRun(key) ?: return emptyList()
+ val walked = mutableListOf()
+ val seen = mutableSetOf(start.ref)
+ var parentRef: ExtractionRunRef? = start.parentRef
+ while (parentRef != null && walked.size < limit && seen.add(parentRef)) {
+ val parent = findRun(ExtractionRunKey(key.contextId, parentRef)) ?: break
+ walked += parent
+ parentRef = parent.parentRef
+ }
+ return walked
+ }
+
+ /**
+ * Assembles and runs one of the three pages.
+ *
+ * The statement is built from this class's own literals and nothing else; every value travels as
+ * a bound parameter. `scope` narrows inside the `WHERE`, ahead of the `LIMIT`, which is the whole
+ * point.
+ *
+ * A row that will not map has already spent one of the caller's `limit` slots, so a page can come
+ * back shorter than asked for. Reading further to backfill would break the bound the contract
+ * keeps.
+ */
+ private fun readPage(
+ scope: String?,
+ contextIdValue: String,
+ limit: Int,
+ since: Instant?,
+ extraBindings: Map = emptyMap(),
+ ): List {
+ requirePositiveLimit(limit)
+ val statement = buildString {
+ append(PAGE_MATCH)
+ scope?.let { append("\n").append(it) }
+ if (since != null) append("\n").append(SINCE_BOUND)
+ append("\n").append(NEWEST_FIRST_PAGE)
+ }
+ val bindings = buildMap {
+ put("contextId", contextIdValue)
+ put("limit", limit)
+ putAll(extraBindings)
+ if (since != null) {
+ put("sinceEpochSecond", since.epochSecond)
+ put("sinceNano", since.nano)
+ }
+ }
+ return queryRows(statement, bindings).mapNotNull(::mapRun)
+ }
+
+ /**
+ * Turns one `{run, invocations}` row into a run, or logs it and returns null.
+ *
+ * One corrupt node should not fail a whole audit read, and the row mappers throw rather than
+ * inventing defaults so that this can happen. The warning names the property that was missing or
+ * the check that failed, which is what an operator needs to go find the node. A corrupt child row
+ * is dropped on its own, so a run with one unreadable attempt still reads.
+ */
+ private fun mapRun(row: Map<*, *>): ExtractionRun? {
+ val header = row["run"] as? Map<*, *> ?: run {
+ logger.warn("Skipping ExtractionRun row with no header properties")
+ return null
+ }
+ val invocations = (row["invocations"] as? List<*>).orEmpty()
+ .filterIsInstance