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
280 changes: 280 additions & 0 deletions CHANGELOG.md

Large diffs are not rendered by default.

Large diffs are not rendered by default.

Original file line number Diff line number Diff line change
@@ -0,0 +1,30 @@
/*
* 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.proposition.extraction.ExtractionRunStore
import com.embabel.dice.proposition.extraction.InMemoryExtractionRunStore

/**
* Runs the [AbstractExtractionRunStoreContractTest] suite against the in-memory backend. No Docker,
* so it runs in the normal test phase — the always-on half of the cross-backend parity check the
* Drivine run store completes when it lands.
*/
class InMemoryExtractionRunStoreContractTest : AbstractExtractionRunStoreContractTest() {
override fun store(listener: DiceEventListener): ExtractionRunStore =
InMemoryExtractionRunStore(listener)
}
26 changes: 26 additions & 0 deletions dice/src/main/kotlin/com/embabel/dice/common/DiceEvent.kt
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,9 @@ import com.embabel.common.core.types.Timestamped
import com.embabel.dice.pipeline.PropositionExtractionStats
import com.embabel.dice.proposition.Proposition
import com.embabel.dice.proposition.PropositionStatus
import com.embabel.dice.proposition.extraction.ExtractionRun
import com.fasterxml.jackson.annotation.JsonTypeInfo
import org.jetbrains.annotations.ApiStatus
import java.time.Instant

/**
Expand Down Expand Up @@ -160,6 +162,30 @@ data class ExtractionBatchCompleted @JvmOverloads constructor(
override val timestamp: Instant = Instant.now(),
) : DiceEvent

/**
* An extraction run ended, and this is the call that ended it.
*
* It fires once per run, from the store, after the terminal write has landed. A coordinator
* retrying a terminal write it never saw the answer to gets a replay, and a replay fires nothing —
* so a listener counting finished runs, or kicking off work behind one, sees each run once however
* many times its terminal write was sent.
*
* The run carries everything there is to know about how it ended: the tenant, the lineage, the
* terminal status, the finish time, the counts, and the failures in their closed vocabulary. There
* is no source text or provider message anywhere in it, so this event is safe to hand to a listener
* that logs or forwards what it receives.
*
* EXPERIMENTAL. The shape may still change while extraction runs (DICE #67) land.
*
* @property run The run in its terminal state.
* @property timestamp When the event was created.
*/
@ApiStatus.Experimental
data class ExtractionRunTransitioned @JvmOverloads constructor(
Comment thread
jimador marked this conversation as resolved.
val run: ExtractionRun,
override val timestamp: Instant = Instant.now(),
) : DiceEvent

/**
* A proposition moved to a different lifecycle [PropositionStatus] — for example going stale
* during a decay sweep, or coming back to life when it's seen again.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -83,6 +83,11 @@ enum class ExtractionInvocationOutcome {

/** The attempt was stopped before it produced anything. */
CANCELLED,
;

/** True for the three outcomes an attempt does not leave, same rule as [ExtractionRunStatus]. */
val isTerminal: Boolean
get() = this != IN_FLIGHT
}

/**
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,20 @@ data class ExtractionRunKey(

/** The tenant id as a plain string, for Java callers, since `ContextId` is a value class. */
fun getContextIdValue(): String = contextId.value

companion object {

/**
* Java-friendly factory taking both halves as plain strings.
*
* `ContextId` is a Kotlin value class, so the constructor has a mangled JVM name. Every
* read on `ExtractionRunStore` keys on this type, so without a factory the store would be
* unreachable from Java.
*/
@JvmStatic
fun of(contextIdValue: String, runId: String): ExtractionRunKey =
ExtractionRunKey(ContextId(contextIdValue), ExtractionRunRef(runId))
}
}

/**
Expand Down Expand Up @@ -81,7 +95,7 @@ data class ExtractionRunKey(
* This is a plain class rather than a data class on purpose. A data class has to declare its
* collection parameters as properties, which means the field is the caller's list and there is
* nowhere to copy it; its generated `copy` and `componentN` methods would also pin an ABI across
* seventeen fields while #67 is still moving. Equality and hash are written out over every
* eighteen fields while #67 is still moving. Equality and hash are written out over every
* component instead.
*
* EXPERIMENTAL. The shape may still change while extraction runs (DICE #67) land.
Expand All @@ -103,6 +117,15 @@ data class ExtractionRunKey(
* @property counts How much the run got through
* @property invocations One record per attempt at each planned model call
* @property failures Bounded record of what went wrong, in the failure vocabulary
* @property version The compare-and-set generation [ExtractionRunStore.save] checks this header
* against. A run that has never been saved, and the run its first accepted save produces, both
* carry 0 — the first save inserts the row, and there is no earlier generation for it to raise
* past. A store rejects a first save naming any other value. Every later save [save] accepts that
* actually changes the header raises it by one; a save whose content already matches what is
* stored is accepted too, as a no-op replay, and leaves the generation exactly where it stood.
* [ExtractionRunStore.recordInvocation] never changes it, since an invocation record writes a
* child row of its own — and [ExtractionRunTransition.applyTo] carries whatever value the run
* already has, because a terminal run takes no more saves and nothing compares its version again
*/
@ApiStatus.Experimental
class ExtractionRun @JvmOverloads constructor(
Expand All @@ -123,6 +146,7 @@ class ExtractionRun @JvmOverloads constructor(
val counts: ExtractionRunCounts = ExtractionRunCounts(),
invocations: List<ExtractionInvocationRecord> = emptyList(),
failures: List<ExtractionFailure> = emptyList(),
val version: Long = 0,
) {

/** Which revisions of which sources this run read, in order. */
Expand All @@ -138,6 +162,8 @@ class ExtractionRun @JvmOverloads constructor(
Collections.unmodifiableList(ArrayList(failures))

init {
require(version >= 0) { "version must not be negative, was $version" }

require(finishedAt == null || !finishedAt.isBefore(startedAt)) {
"finishedAt must not be before startedAt"
}
Expand Down Expand Up @@ -233,7 +259,8 @@ class ExtractionRun @JvmOverloads constructor(
replayFidelity == other.replayFidelity &&
counts == other.counts &&
invocations == other.invocations &&
failures == other.failures
failures == other.failures &&
version == other.version
}

override fun hashCode(): Int {
Expand All @@ -254,6 +281,7 @@ class ExtractionRun @JvmOverloads constructor(
result = 31 * result + counts.hashCode()
result = 31 * result + invocations.hashCode()
result = 31 * result + failures.hashCode()
result = 31 * result + version.hashCode()
return result
}

Expand All @@ -267,7 +295,8 @@ class ExtractionRun @JvmOverloads constructor(
"ExtractionRun(contextId=${contextId.value}, runId=${ref.runId}, rootRunId=${rootRef.runId}, " +
"parentRunId=${parentRef?.runId}, pass=${lineage.passIndex}, status=$status, " +
"startedAt=$startedAt, finishedAt=$finishedAt, sourceRevisions=${sourceRevisions.size}, " +
"invocations=${invocations.size}, failures=${failures.size}, replayFidelity=$replayFidelity)"
"invocations=${invocations.size}, failures=${failures.size}, replayFidelity=$replayFidelity, " +
"version=$version)"

companion object {

Expand Down Expand Up @@ -297,6 +326,7 @@ class ExtractionRun @JvmOverloads constructor(
counts: ExtractionRunCounts = ExtractionRunCounts(),
invocations: List<ExtractionInvocationRecord> = emptyList(),
failures: List<ExtractionFailure> = emptyList(),
version: Long = 0,
): ExtractionRun = ExtractionRun(
contextId = ContextId(contextIdValue),
lineage = lineage,
Expand All @@ -315,6 +345,7 @@ class ExtractionRun @JvmOverloads constructor(
counts = counts,
invocations = invocations,
failures = failures,
version = version,
)
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,141 @@
/*
* 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.proposition.extraction

import org.jetbrains.annotations.ApiStatus
import java.security.MessageDigest
import java.time.Instant

/**
* The digest a store compares two terminal writes by.
*
* A run store keeps the fingerprint of the write that terminalized a run. A second write under the
* same key either matches it — a retry, replayed as success — or does not, and is rejected. The
* whole mechanism rests on two payloads that mean the same thing producing the same bytes, so the
* encoding is specified here rather than left to whatever a serializer happens to emit.
*
* ## What the digest covers
*
* The transition's identity: the terminal status, and when the run finished. Two writes that agree
* on those two things are the same terminal write, and the second replays.
*
* The counts and failures a transition carries are data it delivers, and they stay outside the
* digest. A run's outcome is written once — the first accepted terminal write wins, and a retry
* carrying different numbers replays against what is already stored. Keeping them out of the digest
* is what lets the outcome payload grow: DICE #69 adds typed product outcomes to what a terminal
* write reports, and no field it adds can move a digest already recorded beside a run.
*
* ## Why not JSON
*
* RFC 8785 exists because naive JSON serialization is not byte-stable. Three failure modes, all of
* which would surface here as a correct retry being rejected:
*
* - **Key order.** Most serializers emit fields in declaration or reflection order, and neither is
* guaranteed stable across versions of the code or the library.
* - **Number rendering.** The same value can serialize as `1`, `1.0`, `1e0` or `1.0E+0` depending
* on the writer, and floating-point round-tripping differs between implementations.
* - **Insignificant text.** Whitespace, escaping choices and Unicode normalization all move the
* bytes without moving the meaning.
*
* This encoding is the one DICE already uses for `MetamodelVersion.contentHash`, applied to a
* different payload: length-prefixed tokens, count-prefixed field sets, SHA-256, lowercase hex.
*
* ## The rules
*
* 1. **Every token is length-prefixed**, `<length>:<token>`. A delimiter-joined encoding lets
* `["a;b"]` and `["a", "b"]` hash the same, which hides a real difference. Length prefixes make
* that collision unreachable whatever characters a token happens to carry.
* 2. **A field set is preceded by how many fields it holds**, so a shorter one can never be a
* prefix of a longer one.
* 3. **Fields are emitted as `(name, value)` pairs sorted by name**, so the bytes do not depend on
* the order the fields happen to be declared in.
* 4. **Instants render as `<epochSecond>.<nanos padded to 9>`.** Fixed width, and independent of
* `java.time`'s own formatting. `Instant.toString()` varies its precision with the value —
* `…:47Z` for a whole second, `…:47.500Z` for half of one — so the encoded length moves with the
* data and a persisted digest would depend on a formatting rule DICE does not own. Number
* rendering is the equivalent hazard in JSON and is most of what RFC 8785 is about.
* 5. **The digest is SHA-256, rendered lowercase hex**, and the encoded input carries a version
* tag. A reader meeting a version it does not know matches nothing rather than guessing.
*
* ## This is a persisted format
*
* The digest is stored beside the run. Changing the encoding makes every recorded fingerprint
* unmatchable, so a correct retry against an old run would be rejected as an incompatible rewrite.
* That is what [TERMINAL_VERSION] is for, and it is why the payload is the transition's identity
* alone: a field added to what a terminal write reports never reaches these bytes.
* `ExtractionRunFingerprintTest` pins the digest of a fixed payload with a literal assertion:
* changing the encoding means changing that literal deliberately.
*
* EXPERIMENTAL. The shape may still change while extraction runs (DICE #67) land.
*/
@ApiStatus.Experimental
object ExtractionRunFingerprint {

/**
* Version tag on the terminal-write encoding.
*
* `v2` narrowed the payload to the transition's identity. `v1` also folded in the counts and
* the failures a transition carried, which made the digest move whenever the outcome payload
* gained a field.
*/
const val TERMINAL_VERSION: String = "xrun-terminal:v2"

/**
* The digest of one terminal write: the status it asserts, and when the run finished.
*
* Two things stay out of it deliberately. The run, because a coordinator that records another
* invocation between a terminal write and its retry made the same terminal write both times,
* and folding the run's invocation list in would turn that correct retry into a rejected
* conflict. And the counts and failures the transition carries, because those are the outcome
* a run reports, while the digest names which terminal write this is — see the class doc.
*/
@JvmStatic
fun ofTerminal(status: ExtractionRunStatus, finishedAt: Instant): String {
val payload = fields(
"finishedAt" to encodeInstant(finishedAt),
"status" to status.name,
)
return digest(TERMINAL_VERSION + "|" + payload)
}

/**
* Fixed-width rendering: seconds since the epoch, a dot, then nanoseconds padded to nine
* digits. Two instants that compare equal always render identically, and no instant renders
* two ways.
*/
private fun encodeInstant(instant: Instant): String =
"${instant.epochSecond}.${instant.nano.toString().padStart(9, '0')}"

/** Emit `(name, value)` pairs sorted by name, each half length-prefixed. */
private fun fields(vararg pairs: Pair<String, String>): String = buildString {
append(pairs.size).append('|')
pairs.sortedBy { it.first }.forEach { (name, value) ->
appendSized(name)
appendSized(value)
}
}

private fun StringBuilder.appendSized(token: String) {
append(token.length).append(':').append(token)
}

private fun digest(input: String): String =
MessageDigest.getInstance("SHA-256")
.digest(input.toByteArray(Charsets.UTF_8))
.joinToString("") { byte -> HEX[(byte.toInt() shr 4) and 0xF].toString() + HEX[byte.toInt() and 0xF] }

private const val HEX = "0123456789abcdef"
}
Loading
Loading