Skip to content
Draft
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
11 changes: 11 additions & 0 deletions packages/agent/src/dkg-agent-lifecycle.ts
Original file line number Diff line number Diff line change
Expand Up @@ -104,6 +104,7 @@ import { decodeChangelogRequest, encodeChangelogResponse } from './sync/changelo
import { runChangelogSync, planPageApply } from './sync/requester/changelog-sync.js';
import {
authenticateVerifiedGraphScopedAsset,
isAuthenticatedGraphScopedAssetMaterialized,
materializeVerifiedGraphScopedAsset,
type VerifiedGraphScopedAsset,
type VerifyContextGraphBinding,
Expand Down Expand Up @@ -4429,6 +4430,16 @@ export class LifecycleSyncMethods extends DKGAgentBase {
source: 'agent.durableSync.storeInsert',
}),
storeGraphScopedAsset: async (asset, deadline) => {
if (await isAuthenticatedGraphScopedAssetMaterialized({
store: this.store,
asset,
options: {
priority: 'background',
source: 'agent.durableSync.graphScopedReplayProbe',
},
})) {
return 'stale';
}
const authentication = await authenticateDurableGraphScopedAsset({
chain: this.chain,
asset,
Expand Down
114 changes: 114 additions & 0 deletions packages/agent/src/sync/requester/graph-scoped-materialization.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ import {
type TripleStore,
} from '@origintrail-official/dkg-storage';
import {
computeFlatKCRootV10,
mergeSameVersionGraphKnowledgeAssetMetadataV1,
readGraphKnowledgeAssetConfirmationKindV1,
readLocallyTrustedKnowledgeAssetControls,
Expand All @@ -21,7 +22,9 @@ import {
} from '../oversize-filter.js';

const ASSERTION_VERSION = 'http://dkg.io/ontology/assertionVersion';
const ASSERTION_GRAPH = 'http://dkg.io/ontology/assertionGraph';
const MERKLE_ROOT = 'http://dkg.io/ontology/merkleRoot';
const PRIVATE_MERKLE_ROOT = 'http://dkg.io/ontology/privateMerkleRoot';
const STATUS = 'http://dkg.io/ontology/status';
const TRANSACTION_HASH = 'http://dkg.io/ontology/transactionHash';
const MATERIALIZED_VERSION = 'http://dkg.io/ontology/materializedVersion';
Expand Down Expand Up @@ -52,6 +55,96 @@ export type VerifyContextGraphBinding = (

export type GraphScopedMaterializationOutcome = 'applied' | 'stale' | 'quarantined';

/**
* Detect an exact replay that was already admitted through this requester's
* authenticated materialization path. `materializedVersion` and `status` are
* stripped from peer input before storage, so their local presence can be used
* as a trust marker only when the immutable assertion identity still matches.
*
* This is deliberately an optimization, not an admission fallback: missing,
* malformed, duplicated, or mismatched metadata returns false and forces the
* caller through normal chain authentication again.
*/
export async function isAuthenticatedGraphScopedAssetMaterialized(params: {
store: TripleStore;
asset: VerifiedGraphScopedAsset;
options?: QueryOptions;
}): Promise<boolean> {
const { store, asset, options = {} } = params;
const incomingRoots = asset.metadataQuads
.filter((quad) => quad.predicate === MERKLE_ROOT)
.map((quad) => quad.object);
if (incomingRoots.length !== 1) return false;

let incomingRoot: Uint8Array;
try {
incomingRoot = parseBytes32Literal(incomingRoots[0]!, 'merkleRoot');
} catch {
return false;
}
const privateRoots: Uint8Array[] = [];
try {
for (const quad of asset.metadataQuads) {
if (quad.predicate === PRIVATE_MERKLE_ROOT) {
privateRoots.push(parseBytes32Literal(quad.object, 'privateMerkleRoot'));
}
}
} catch {
return false;
}

const [metadataResult, dataResult] = await Promise.all([
store.query(`
SELECT ?assertionVersion ?assertionGraph ?merkleRoot ?materializedVersion ?status WHERE {
GRAPH <${assertSafeIri(asset.metaGraph)}> {
<${assertSafeIri(asset.ual)}> <${ASSERTION_VERSION}> ?assertionVersion ;
<${ASSERTION_GRAPH}> ?assertionGraph ;
<${MERKLE_ROOT}> ?merkleRoot ;
<${MATERIALIZED_VERSION}> ?materializedVersion ;
<${STATUS}> ?status .
}
}
`, options),
store.query(`
SELECT ?s ?p ?o WHERE {
GRAPH <${assertSafeIri(asset.assertionGraph)}> { ?s ?p ?o }
}
`, options),
]);
if (
metadataResult.type !== 'bindings'
|| metadataResult.bindings.length !== 1
|| dataResult.type !== 'bindings'
) return false;

const row = metadataResult.bindings[0]!;
const version = parseUnsignedIntegerLiteral(row.assertionVersion);
const storedRoot = parseBytes32LiteralOrUndefined(row.merkleRoot);
const materializedVersion = parseRdfLiteral(row.materializedVersion ?? '');
const status = parseRdfLiteral(row.status ?? '');
const storedAssertionGraph = stripIriBinding(row.assertionGraph);
const localData = dataResult.bindings.flatMap((binding): Quad[] => (
binding.s !== undefined && binding.p !== undefined && binding.o !== undefined
? [{
subject: binding.s,
predicate: binding.p,
object: binding.o,
graph: asset.assertionGraph,
}]
: []
));
const localRoot = computeFlatKCRootV10(localData, privateRoots);

return version === asset.assertionVersion
&& storedAssertionGraph === asset.assertionGraph
&& storedRoot !== undefined
&& bytesEqual(storedRoot, incomingRoot)
&& bytesEqual(localRoot, incomingRoot)
&& materializedVersion !== undefined
&& /^\d+:\d+$/.test(materializedVersion)
&& status === 'confirmed';
}

/**
* Bind the peer-verified payload to current chain truth before its structural
* metadata can influence local assertion ordering. No-chain development keeps
Expand Down Expand Up @@ -436,6 +529,27 @@ function parseBytes32Literal(raw: string, field: string): Uint8Array {
return Uint8Array.from(hex.match(/.{2}/g)!.map((pair) => Number.parseInt(pair, 16)));
}

function parseBytes32LiteralOrUndefined(raw: string | undefined): Uint8Array | undefined {
if (raw === undefined) return undefined;
try {
return parseBytes32Literal(raw, 'merkleRoot');
} catch {
return undefined;
}
}

function parseUnsignedIntegerLiteral(raw: string | undefined): bigint | undefined {
if (raw === undefined) return undefined;
const lexical = parseRdfLiteral(raw) ?? raw;
if (!/^\d+$/.test(lexical)) return undefined;
return BigInt(lexical);
}

function stripIriBinding(raw: string | undefined): string | undefined {
if (raw === undefined) return undefined;
return raw.match(/^<([^>]*)>$/)?.[1] ?? raw;
}

function parseTransactionHashLiteral(raw: string): string {
const lexical = parseRdfLiteral(raw) ?? raw;
if (!/^0x[0-9a-f]{64}$/i.test(lexical)) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@ import type { SyncPageResult } from '../src/sync/requester/page-fetch.js';
import { DKGAgent } from '../src/dkg-agent.js';
import {
authenticateVerifiedGraphScopedAsset,
isAuthenticatedGraphScopedAssetMaterialized,
materializeVerifiedGraphScopedAsset,
type GraphScopedMaterializationOutcome,
type VerifyContextGraphBinding,
Expand Down Expand Up @@ -243,6 +244,69 @@ function runGraphScopedDurableSync(options: {
}

describe('durable graph-scoped KA materialization', () => {
it('recognizes only an exact replay with local authenticated materialization markers', async () => {
const store = new OxigraphStore();
const root = computeFlatKCRootV10([dataQuad(2)], []);
const rootHex = toHex(root);
const asset: VerifiedGraphScopedAsset = {
contextGraphId,
ual,
assertionVersion: 2n,
assertionGraph,
metaGraph,
dataQuads: [dataQuad(2)],
metadataQuads: metadata(2, rootHex),
};
await store.insert([
...asset.dataQuads,
...asset.metadataQuads,
{
subject: ual,
predicate: `${DKG}materializedVersion`,
object: '"123:4"',
graph: metaGraph,
},
{
subject: ual,
predicate: `${DKG}status`,
object: '"confirmed"',
graph: metaGraph,
},
]);

await expect(isAuthenticatedGraphScopedAssetMaterialized({
store,
asset,
})).resolves.toBe(true);

const changedRoot = new Uint8Array(32);
changedRoot[31] = 3;
await expect(isAuthenticatedGraphScopedAssetMaterialized({
store,
asset: {
...asset,
metadataQuads: metadata(2, toHex(changedRoot)),
},
})).resolves.toBe(false);

await store.deleteByPattern({ graph: assertionGraph });
await store.insert([dataQuad(1)]);
await expect(isAuthenticatedGraphScopedAssetMaterialized({
store,
asset,
})).resolves.toBe(false);

await store.deleteByPattern({
graph: metaGraph,
subject: ual,
predicate: `${DKG}materializedVersion`,
});
await expect(isAuthenticatedGraphScopedAssetMaterialized({
store,
asset,
})).resolves.toBe(false);
});

it('hands the exact context-graph deadline to graph-scoped storage', async () => {
const deadline = 1_800_000_123_456;
const storeGraphScopedAsset = vi.fn(async (
Expand Down
32 changes: 32 additions & 0 deletions packages/agent/test/durable-sync-lifecycle-binding.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@ vi.mock('../src/sync/requester/graph-scoped-materialization.js', async (importOr
>();
return {
...actual,
isAuthenticatedGraphScopedAssetMaterialized: vi.fn(async () => false),
materializeVerifiedGraphScopedAsset: vi.fn(async () => 'applied' as const),
};
});
Expand All @@ -22,6 +23,7 @@ import { DKGAgent } from '../src/dkg-agent.js';
import { LifecycleSyncMethods } from '../src/dkg-agent-lifecycle.js';
import { runDurableSync } from '../src/sync/requester/durable-sync.js';
import {
isAuthenticatedGraphScopedAssetMaterialized,
materializeVerifiedGraphScopedAsset,
type VerifiedGraphScopedAsset,
} from '../src/sync/requester/graph-scoped-materialization.js';
Expand All @@ -34,6 +36,7 @@ const metaGraph = `did:dkg:context-graph:${contextGraphId}/_meta`;
const ctx = { kind: 'sync', id: 'lifecycle-binding-test', startedAt: 0 } as OperationContext;

const mockedRunDurableSync = vi.mocked(runDurableSync);
const mockedReplayProbe = vi.mocked(isAuthenticatedGraphScopedAssetMaterialized);
const mockedMaterialize = vi.mocked(materializeVerifiedGraphScopedAsset);

function graphScopedAsset(
Expand Down Expand Up @@ -105,9 +108,38 @@ async function captureGraphScopedStore(
describe('durable sync lifecycle chain binding', () => {
beforeEach(() => {
mockedRunDurableSync.mockClear();
mockedReplayProbe.mockReset();
mockedReplayProbe.mockResolvedValue(false);
mockedMaterialize.mockClear();
});

it('skips chain authentication for an exact locally authenticated replay', async () => {
const root = new Uint8Array(32);
root[31] = 2;
const getLatestMerkleRoot = vi.fn(async () => root);
const getMerkleRootCount = vi.fn(async () => 2n);
const getKAContextGraphId = vi.fn(async () => 14n);
const chain = {
chainId: 'otp:2043',
getLatestMerkleRoot,
getMerkleRootCount,
getKAContextGraphId,
} as ChainAdapter;
mockedReplayProbe.mockResolvedValueOnce(true);

const storeGraphScopedAsset = await captureGraphScopedStore(chain);
await expect(storeGraphScopedAsset(
graphScopedAsset(root),
Date.now() + 60_000,
)).resolves.toBe('stale');

expect(mockedReplayProbe).toHaveBeenCalledTimes(1);
expect(getLatestMerkleRoot).not.toHaveBeenCalled();
expect(getMerkleRootCount).not.toHaveBeenCalled();
expect(getKAContextGraphId).not.toHaveBeenCalled();
expect(mockedMaterialize).not.toHaveBeenCalled();
});

it('retries a transient binding read, caches only the successful proof, and persists the CG id', async () => {
const root = new Uint8Array(32);
root[31] = 2;
Expand Down
Loading