Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
30 commits
Select commit Hold shift + click to select a range
0d49d12
Add source revision provenance contract
jimador Jul 26, 2026
ed7c188
Add shared provenance evidence key codec
jimador Jul 26, 2026
3c8f8b8
Add source-scoped provenance queries
jimador Jul 26, 2026
09b6557
Carry source revision through pipeline provenance
jimador Jul 26, 2026
84949ec
Carry source revision on analysis events
jimador Jul 26, 2026
2052b30
Verify portable source provenance queries
jimador Jul 26, 2026
21d1233
Record revision-aware collector fold evidence
jimador Jul 26, 2026
f9a0734
Add provenance-aware extraction entrypoints
jimador Jul 26, 2026
a816832
feat(storage): map source revision provenance
jimador Jul 26, 2026
72af52e
fix: make fold evidence undo revision-aware
jimador Jul 26, 2026
83d2c95
chore: correct fold undo evidence ledger
jimador Jul 26, 2026
606230d
Add source revision to REST extraction
jimador Jul 26, 2026
ac8cf5e
test(provenance): pin JSON revision compatibility
jimador Jul 26, 2026
4c85734
feat(discovery): expose lineage provenance
jimador Jul 26, 2026
f51b193
fix(discovery): honor lineage provenance
jimador Jul 26, 2026
937248d
feat(storage): preserve revisioned provenance cardinality
jimador Jul 26, 2026
777b251
fix(storage): isolate provenance race recovery
jimador Jul 26, 2026
6582af7
Prove source revision compatibility boundary
jimador Jul 26, 2026
401740a
Record compatibility gate evidence correction
jimador Jul 26, 2026
5ecea33
Restore unambiguous Kotlin extraction calls
jimador Jul 26, 2026
eb2e7fa
Preserve source revision compatibility
jimador Jul 28, 2026
b057eff
Address source revision review findings
jimador Jul 29, 2026
b32c0eb
Prove REST replay provenance cardinality
jimador Jul 29, 2026
f999cc0
Validate source identity before dedup reuse
jimador Jul 29, 2026
6b08206
Prevent partial writes when collapse participants are missing
jimador Jul 29, 2026
c3fd0c0
test: pin source revision binary compatibility
jimador Jul 29, 2026
90cc842
test: remove synthetic locator dependency
jimador Jul 29, 2026
6a2efb2
test(storage): couple source plans to live queries
jimador Jul 29, 2026
5c3ec75
test(storage): keep test seams out of production
jimador Jul 29, 2026
bb5e5eb
fix(provenance): make collector undo authoritative through store SPI
jimador Aug 1, 2026
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 @@ -117,6 +117,7 @@ object CollectorDecisionRowMapper {
"foldedGrounding" to retired.foldedGrounding,
"foldedProvenanceRefs" to retired.foldedProvenanceRefs,
"foldedSourceIds" to retired.foldedSourceIds,
"foldedProvenanceEvidenceKeys" to retired.foldedProvenanceEvidenceKeys,
)

fun fromRow(row: Map<*, *>, retired: List<RetiredProposition>): CollectorDecision = CollectorDecision(
Expand All @@ -133,6 +134,7 @@ object CollectorDecisionRowMapper {
foldedGrounding = row.stringList("foldedGrounding"),
foldedProvenanceRefs = row.stringList("foldedProvenanceRefs"),
foldedSourceIds = row.stringList("foldedSourceIds"),
foldedProvenanceEvidenceKeys = row.stringList("foldedProvenanceEvidenceKeys"),
)
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -44,7 +44,8 @@ import java.time.Instant
* - `(:CollectorDecision {id, runId, contextId, componentId, survivorId, action, createdAt})`
* with `id = "runId|componentId"`, plus one child
* `(:CollectorRetired {id, runId, contextId, propositionId, priorStatus, foldedGrounding,
* foldedProvenanceRefs, foldedSourceIds})-[:RETIRED_IN]->(:CollectorDecision)` per retired
* foldedProvenanceRefs, foldedSourceIds, foldedProvenanceEvidenceKeys})
* -[:RETIRED_IN]->(:CollectorDecision)` per retired
* proposition, so a reversal has everything a merging sweep folded onto the survivor.
*
* Every write is a single `UNWIND $rows AS r ...` round trip (see [com.embabel.dice.storage.CollectorTraceRowMappers] for
Expand Down Expand Up @@ -227,7 +228,8 @@ class DrivineCollectorTraceStore(
retired: [r IN retiredNodes WHERE r IS NOT NULL | {
propositionId: r.propositionId, priorStatus: r.priorStatus,
foldedGrounding: r.foldedGrounding, foldedProvenanceRefs: r.foldedProvenanceRefs,
foldedSourceIds: r.foldedSourceIds
foldedSourceIds: r.foldedSourceIds,
foldedProvenanceEvidenceKeys: r.foldedProvenanceEvidenceKeys
}]
} AS row
""".trimIndent(),
Expand Down

Large diffs are not rendered by default.

Original file line number Diff line number Diff line change
Expand Up @@ -168,22 +168,25 @@ object PropositionGraphMapper {

// ---- Provenance: ProvenanceEntry <-> (DERIVED_FROM edge + shared Source node) ----

private fun toDerivedFrom(e: ProvenanceEntry): DerivedFrom =
internal fun toDerivedFrom(e: ProvenanceEntry): DerivedFrom =
DerivedFrom(
chunkId = e.chunkId,
startOffset = e.startOffset,
endOffset = e.endOffset,
contentHash = e.contentHash,
sourceRevision = e.sourceRevision,
entryKey = provenanceStorageEntryKey(e),
source = toSourceNode(e.locator),
)

private fun toProvenanceEntry(df: DerivedFrom): ProvenanceEntry =
internal fun toProvenanceEntry(df: DerivedFrom): ProvenanceEntry =
ProvenanceEntry(
locator = toLocator(df.source),
chunkId = df.chunkId,
startOffset = df.startOffset,
endOffset = df.endOffset,
contentHash = df.contentHash,
sourceRevision = df.sourceRevision,
)

private fun toSourceNode(loc: SourceLocator): SourceNode =
Expand Down Expand Up @@ -228,3 +231,21 @@ object PropositionGraphMapper {
)
}
}

/**
* Canonical graph relationship identity for provenance.
*
* Length framing keeps null, empty, and delimiter-containing values distinct without coupling
* graph persistence to the domain's public evidence-key format.
*/
internal fun provenanceStorageEntryKey(entry: ProvenanceEntry): String =
listOf(
entry.locator.key(),
entry.sourceRevision,
entry.chunkId,
entry.startOffset?.toString(),
entry.endOffset?.toString(),
entry.contentHash,
).joinToString(separator = "") { value ->
if (value == null) "-1:" else "${value.length}:$value"
}
Original file line number Diff line number Diff line change
Expand Up @@ -23,10 +23,12 @@ import org.drivine.annotation.RelationshipFragment
* the shared node. Together they reconstitute a dice `ProvenanceEntry`.
*/
@RelationshipFragment
data class DerivedFrom(
data class DerivedFrom @JvmOverloads constructor(
val chunkId: String? = null,
val startOffset: Int? = null,
val endOffset: Int? = null,
val contentHash: String? = null,
val source: SourceNode,
val sourceRevision: String? = null,
val entryKey: String? = null,
)
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,7 @@ class CollectorTraceRowMapperTest {
foldedGrounding = listOf("g1", "g2"),
foldedProvenanceRefs = listOf("prov-1"),
foldedSourceIds = listOf("src-1", "src-2"),
foldedProvenanceEvidenceKeys = listOf("evidence-1"),
)

val bindMap = CollectorDecisionRowMapper.retiredBindMap("run-1", mockDecision(), retired)
Expand All @@ -45,6 +46,23 @@ class CollectorTraceRowMapperTest {
assertEquals(listOf("g1", "g2"), roundTripped.foldedGrounding)
assertEquals(listOf("prov-1"), roundTripped.foldedProvenanceRefs)
assertEquals(listOf("src-1", "src-2"), roundTripped.foldedSourceIds)
assertEquals(listOf("evidence-1"), roundTripped.foldedProvenanceEvidenceKeys)
}

@Test
fun `legacy retired row without evidence keys remains readable`() {
val legacyRow = mapOf(
"propositionId" to "prop-1",
"priorStatus" to "ACTIVE",
"foldedGrounding" to emptyList<String>(),
"foldedProvenanceRefs" to listOf("uri:https://example.com/source"),
"foldedSourceIds" to emptyList<String>(),
)

val retired = CollectorDecisionRowMapper.retiredFromRow(legacyRow)

assertEquals(listOf("uri:https://example.com/source"), retired.foldedProvenanceRefs)
assertEquals(emptyList<String>(), retired.foldedProvenanceEvidenceKeys)
}

@Test
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,11 @@ import com.embabel.dice.projection.memory.collector.CollectorRunContext
import com.embabel.dice.projection.memory.collector.CollectorSurvivorPolicy
import com.embabel.dice.projection.memory.collector.MultiSignalCollectorStrategy
import com.embabel.dice.proposition.Proposition
import com.embabel.dice.proposition.PropositionRepository
import com.embabel.dice.proposition.PropositionStore
import com.embabel.dice.proposition.PropositionStatus
import com.embabel.dice.provenance.ProvenanceEntry
import com.embabel.dice.provenance.UriLocator
import com.embabel.dice.spi.CandidatePair
import com.embabel.dice.spi.CandidatePairSource
import com.embabel.dice.spi.CollectorCandidateEdge
Expand Down Expand Up @@ -57,6 +61,33 @@ class DrivineCollectorTraceStoreIntegrationTest {
@Autowired
private lateinit var propositionRepository: DrivinePropositionRepository

private class RecordingPropositionRepository(
private val delegate: PropositionRepository,
) : PropositionRepository by delegate {

var saveCalls = 0
private set
var setProvenanceCalls = 0
private set

override fun save(proposition: Proposition): Proposition {
saveCalls++
return delegate.save(proposition)
}

override fun setProvenance(
propositionId: String,
entries: List<ProvenanceEntry>,
): Proposition? {
setProvenanceCalls++
return delegate.setProvenance(propositionId, entries)
}
}

private class BaseStoreView(
delegate: PropositionRepository,
) : PropositionStore by delegate

@AfterEach
fun cleanUp() {
CollectorTraceSchema.LABELS.forEach { label ->
Expand All @@ -70,6 +101,7 @@ class DrivineCollectorTraceStoreIntegrationTest {
text: String,
status: PropositionStatus = PropositionStatus.ACTIVE,
grounding: List<String> = emptyList(),
provenance: List<ProvenanceEntry> = emptyList(),
) = Proposition(
id = id,
contextId = ContextId("ctx-undo"),
Expand All @@ -78,6 +110,7 @@ class DrivineCollectorTraceStoreIntegrationTest {
confidence = 0.9,
status = status,
grounding = grounding,
provenanceEntries = provenance,
)

private fun edge(anchorId: String, memberId: String, vetoed: Boolean = false, score: Double = 0.9) = CollectorCandidateEdge(
Expand All @@ -104,6 +137,7 @@ class DrivineCollectorTraceStoreIntegrationTest {
foldedGrounding = listOf("g1", "g2"),
foldedProvenanceRefs = listOf("prov-1"),
foldedSourceIds = listOf("src-1", "src-2"),
foldedProvenanceEvidenceKeys = listOf("evidence-1"),
),
),
)
Expand Down Expand Up @@ -139,6 +173,7 @@ class DrivineCollectorTraceStoreIntegrationTest {
assertEquals(listOf("g1", "g2"), retired.foldedGrounding)
assertEquals(listOf("prov-1"), retired.foldedProvenanceRefs)
assertEquals(listOf("src-1", "src-2"), retired.foldedSourceIds)
assertEquals(listOf("evidence-1"), retired.foldedProvenanceEvidenceKeys)
}

@Test
Expand Down Expand Up @@ -219,6 +254,68 @@ class DrivineCollectorTraceStoreIntegrationTest {
assertNull(traceStore.findRetirement("unknown"))
}

@Test
fun `undoSingleCollapse does not write when the retired proposition is missing`() {
val runId = "run-missing-retired"
val survivor = propositionRepository.save(prop("survivor-missing-retired", "Survivor remains"))
traceStore.recordRunContext(runId, survivor.contextId)
traceStore.recordDecision(
runId,
decisionFor(
componentId = "comp-missing-retired",
survivorId = survivor.id,
retiredId = "missing-retired",
),
)
val before = propositionRepository.findAll().sortedBy { it.id }
val recordingRepository = RecordingPropositionRepository(propositionRepository)

val result = undoSingleCollapse(
traceQuery = traceStore,
propositions = recordingRepository,
survivorId = survivor.id,
retiredId = "missing-retired",
)

val after = propositionRepository.findAll().sortedBy { it.id }
assertNull(result)
assertEquals(before, after)
assertEquals(0, recordingRepository.saveCalls)
assertEquals(0, recordingRepository.setProvenanceCalls)
}

@Test
fun `undoSingleCollapse does not write when the survivor proposition is missing`() {
val runId = "run-missing-survivor"
val retired = propositionRepository.save(
prop("retired-missing-survivor", "Retired remains", status = PropositionStatus.STALE),
)
traceStore.recordRunContext(runId, retired.contextId)
traceStore.recordDecision(
runId,
decisionFor(
componentId = "comp-missing-survivor",
survivorId = "missing-survivor",
retiredId = retired.id,
),
)
val before = propositionRepository.findAll().sortedBy { it.id }
val recordingRepository = RecordingPropositionRepository(propositionRepository)

val result = undoSingleCollapse(
traceQuery = traceStore,
propositions = recordingRepository,
survivorId = "missing-survivor",
retiredId = retired.id,
)

val after = propositionRepository.findAll().sortedBy { it.id }
assertNull(result)
assertEquals(before, after)
assertEquals(0, recordingRepository.saveCalls)
assertEquals(0, recordingRepository.setProvenanceCalls)
}

@Test
fun `undoSingleCollapse restores one member and subtracts only its exclusive grounding, leaving the sibling member retired`() {
val runId = "run-5"
Expand Down Expand Up @@ -339,4 +436,115 @@ class DrivineCollectorTraceStoreIntegrationTest {
// (b) the survivor loses "loser-exclusive" — that one really did come from the loser.
assertEquals(setOf("shared"), updatedSurvivor?.grounding?.toSet())
}

@Test
fun `collector undo removes only the folded revision from persistent provenance`() {
val runId = "run-revision-undo"
val contextId = ContextId("ctx-revision-undo")
val locator = UriLocator("https://example.com/revision-undo")
val revisionOne = ProvenanceEntry(locator = locator, sourceRevision = "r1")
val revisionTwo = ProvenanceEntry(locator = locator, sourceRevision = "r2")
val survivor = propositionRepository.save(
prop(
id = "survivor-revision",
text = "Acme signed the agreement",
provenance = listOf(revisionOne),
),
)
val loser = propositionRepository.save(
prop(
id = "loser-revision",
text = "Acme signed an agreement",
provenance = listOf(revisionOne, revisionTwo),
),
)
val strategy = MultiSignalCollectorStrategy(
pairSources = listOf(
CandidatePairSource {
candidates, _ ->
listOf(CandidatePair(anchor = candidates[0], member = candidates[1]))
},
),
scorers = listOf(
CollectorSignalScorer { _, _ -> CollectorSignalScore(signal = "fixed", score = 1.0) },
),
componentsFinder = InMemoryConnectedComponentsFinder(),
traceStore = traceStore,
survivorPolicy = CollectorSurvivorPolicy { members -> members.single { it.id == survivor.id } },
matchThreshold = 0.5,
)

strategy.mark(listOf(survivor, loser), propositionRepository, CollectorRunContext(runId, contextId))
propositionRepository.save(
propositionRepository.findById(survivor.id)!!.absorbEvidence(loser),
)

undoSingleCollapse(
traceQuery = traceStore,
propositions = propositionRepository,
survivorId = survivor.id,
retiredId = loser.id,
)

assertEquals(
listOf(revisionOne),
propositionRepository.findById(survivor.id)?.provenanceEntries,
)
}

@Test
fun `collector undo removes folded provenance through a base store decorator`() {
val runId = "run-base-store-undo"
val contextId = ContextId("ctx-base-store-undo")
val survivorEvidence = ProvenanceEntry(
locator = UriLocator("https://example.com/base-store-undo/keep"),
)
val foldedEvidence = ProvenanceEntry(
locator = UriLocator("https://example.com/base-store-undo/remove"),
)
val survivor = propositionRepository.save(
prop(
id = "survivor-base-store",
text = "Acme signed the agreement",
provenance = listOf(survivorEvidence, foldedEvidence),
),
)
propositionRepository.save(
prop(
id = "retired-base-store",
text = "Acme signed an agreement",
status = PropositionStatus.STALE,
provenance = listOf(foldedEvidence),
),
)
traceStore.recordRunContext(runId, contextId)
traceStore.recordDecision(
runId,
CollectorDecision(
runId = runId,
componentId = "component-base-store",
survivorId = survivor.id,
action = "duplicate-merge",
retired = listOf(
RetiredProposition(
propositionId = "retired-base-store",
priorStatus = PropositionStatus.ACTIVE,
foldedProvenanceRefs = listOf(foldedEvidence.locator.key()),
),
),
),
)

undoSingleCollapse(
traceQuery = traceStore,
propositions = BaseStoreView(propositionRepository),
survivorId = survivor.id,
retiredId = "retired-base-store",
)

assertEquals(
listOf(survivorEvidence),
propositionRepository.findById(survivor.id)?.provenanceEntries,
)
}
}
Loading