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
237 changes: 114 additions & 123 deletions infrastructure/evault-core/src/core/db/db.service.ts
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
import neo4j, { type Driver } from "neo4j-driver";
import { W3IDBuilder } from "w3id";
import { timed } from "../utils/timing";
import { deserializeValue, serializeValue } from "./schema";
import type {
AppendEnvelopeOperationLogParams,
Expand DownExpand Up@@ -41,12 +42,15 @@ export class DbService {
* @returns The result of the query execution
*/
private async runQueryInternal(query: string, params: Record<string, any>) {
const session = this.driver.session();
try {
return await session.run(query, params);
} finally {
await session.close();
}
const firstLine = query.trim().split("\n")[0].slice(0, 80);
return timed(`db.query "${firstLine}"`, async () => {
const session = this.driver.session();
try {
return await session.run(query, params);
} finally {
await session.close();
}
});
}

/**
Expand DownExpand Up@@ -74,11 +78,14 @@ export class DbService {
acl: string[],
eName: string,
): Promise<StoreMetaEnvelopeResult<T>> {
return timed("db.storeMetaEnvelope", async () => {
if (!eName) {
throw new Error("eName is required for storing meta-envelopes");
}

const w3id = await new W3IDBuilder().build();
const w3id = await timed("db.storeMetaEnvelope.buildMetaId", () =>
new W3IDBuilder().build(),
);

const cypher: string[] = [
`CREATE (m:MetaEnvelope { id: $metaId, ontology: $ontology, acl: $acl, eName: $eName })`,
Expand DownExpand Up@@ -128,7 +135,9 @@ export class DbService {
counter++;
}

await this.runQueryInternal(cypher.join("\n"), envelopeParams);
await timed("db.storeMetaEnvelope.runQuery", () =>
this.runQueryInternal(cypher.join("\n"), envelopeParams),
);

return {
metaEnvelope: {
Expand All@@ -138,6 +147,7 @@ export class DbService {
},
envelopes: createdEnvelopes,
};
});
}

/**
Expand DownExpand Up@@ -528,91 +538,81 @@ export class DbService {
acl: string[],
eName: string,
): Promise<StoreMetaEnvelopeResult<T>> {
return timed("db.updateMetaEnvelopeById", async () => {
if (!eName) {
throw new Error("eName is required for updating meta-envelopes");
}

// The whole read-modify-write cycle runs inside a single Neo4j write
// transaction. The opening MERGE+SET acquires a write lock on the
// MetaEnvelope node, so concurrent updates to the same id serialize
// here — without this, request B's "delete stale envelopes" step
// could clobber fields that request A just wrote.
const session = this.driver.session();
try {
let existing = await this.findMetaEnvelopeById<T>(id, eName);
if (!existing) {
const metaW3id = await new W3IDBuilder().build();
await this.runQueryInternal(
return await session.executeWrite(async (tx) => {
const findResult = await tx.run(
`
CREATE (m:MetaEnvelope {
id: $id,
ontology: $ontology,
acl: $acl,
eName: $eName
})
MERGE (m:MetaEnvelope { id: $id, eName: $eName })
ON CREATE SET m.ontology = $ontology, m.acl = $acl
ON MATCH SET m.ontology = $ontology, m.acl = $acl
WITH m
OPTIONAL MATCH (m)-[:LINKS_TO]->(e:Envelope)
RETURN collect(e) AS envelopes
`,
{ id, ontology: meta.ontology, acl, eName },
{ id, eName, ontology: meta.ontology, acl },
);
existing = {
id,
ontology: meta.ontology,
acl,
parsed: meta.payload,
envelopes: [],
};
}

// Update the meta-envelope properties (ensure eName matches)
await this.runQueryInternal(
`
MATCH (m:MetaEnvelope { id: $id, eName: $eName })
SET m.ontology = $ontology, m.acl = $acl
`,
{ id, ontology: meta.ontology, acl, eName },
);
const envelopeNodes: any[] = (
findResult.records[0]?.get("envelopes") ?? []
).filter((n: any) => n !== null && n !== undefined);

// Deduplicate envelopes — if multiple Envelope nodes share the
// same ontology (field name), keep the first and delete the rest.
// This prevents non-deterministic reads where collect(e) returns
// duplicates in undefined order and reduce picks the wrong one.
const seen = new Map<string, string>(); // ontology → kept envelope id
const dupsToDelete: string[] = [];
for (const env of existing.envelopes) {
if (seen.has(env.ontology)) {
dupsToDelete.push(env.id);
} else {
seen.set(env.ontology, env.id);
let workingEnvelopes: Envelope<T[keyof T]>[] =
envelopeNodes.map((node: any) => ({
id: node.properties.id,
ontology: node.properties.ontology,
value: deserializeValue(
node.properties.value,
node.properties.valueType,
) as T[keyof T],
valueType: node.properties.valueType,
}));

// Deduplicate envelopes — if multiple Envelope nodes share the
// same ontology, keep the first and delete the rest.
const seen = new Map<string, string>();
const dupsToDelete: string[] = [];
for (const env of workingEnvelopes) {
if (seen.has(env.ontology)) {
dupsToDelete.push(env.id);
} else {
seen.set(env.ontology, env.id);
}
}
}
if (dupsToDelete.length > 0) {
console.warn(
`[eVault] Cleaning ${dupsToDelete.length} duplicate envelope(s) for MetaEnvelope ${id}`,
);
for (const dupId of dupsToDelete) {
await this.runQueryInternal(
`MATCH (e:Envelope { id: $envelopeId }) DETACH DELETE e`,
{ envelopeId: dupId },
if (dupsToDelete.length > 0) {
console.warn(
`[eVault] Cleaning ${dupsToDelete.length} duplicate envelope(s) for MetaEnvelope ${id}`,
);
await tx.run(
`MATCH (e:Envelope) WHERE e.id IN $ids DETACH DELETE e`,
{ ids: dupsToDelete },
);
workingEnvelopes = workingEnvelopes.filter(
(e) => !dupsToDelete.includes(e.id),
);
}
// Remove deleted dupes from the existing list so the update
// loop below doesn't try to reference them.
existing.envelopes = existing.envelopes.filter(
(e) => !dupsToDelete.includes(e.id),
);
}

const createdEnvelopes: Envelope<T[keyof T]>[] = [];
let counter = 0;
const createdEnvelopes: Envelope<T[keyof T]>[] = [];

// For each field in the new payload
for (const [key, value] of Object.entries(meta.payload)) {
try {
for (const [key, value] of Object.entries(meta.payload)) {
const { value: storedValue, type: valueType } =
serializeValue(value);
const alias = `e${counter}`;

// Check if an envelope with this ontology already exists
const existingEnvelope = existing.envelopes.find(
const existingEnvelope = workingEnvelopes.find(
(e) => e.ontology === key,
);

if (existingEnvelope) {
// Update existing envelope
await this.runQueryInternal(
await tx.run(
`
MATCH (e:Envelope { id: $envelopeId })
SET e.value = $newValue, e.valueType = $valueType
Expand All@@ -623,88 +623,79 @@ export class DbService {
valueType,
},
);

createdEnvelopes.push({
id: existingEnvelope.id,
ontology: key,
value: value as T[keyof T],
valueType,
});
} else {
// Create new envelope — use MERGE on the relationship
// + ontology to prevent duplicate Envelopes if two
// concurrent updates race.
const envW3id = await new W3IDBuilder().build();
const envelopeId = envW3id.id;

await this.runQueryInternal(
await tx.run(
`
MATCH (m:MetaEnvelope { id: $metaId, eName: $eName })
MERGE (m)-[:LINKS_TO]->(${alias}:Envelope { ontology: $${alias}_ontology })
ON CREATE SET ${alias}.id = $${alias}_id, ${alias}.value = $${alias}_value, ${alias}.valueType = $${alias}_type
ON MATCH SET ${alias}.value = $${alias}_value, ${alias}.valueType = $${alias}_type
MERGE (m)-[:LINKS_TO]->(e:Envelope { ontology: $ontology })
ON CREATE SET e.id = $envelopeId, e.value = $newValue, e.valueType = $valueType
ON MATCH SET e.value = $newValue, e.valueType = $valueType
`,
{
metaId: id,
eName: eName,
[`${alias}_id`]: envelopeId,
[`${alias}_ontology`]: key,
[`${alias}_value`]: storedValue,
[`${alias}_type`]: valueType,
eName,
envelopeId,
ontology: key,
newValue: storedValue,
valueType,
},
);

createdEnvelopes.push({
id: envelopeId,
ontology: key,
value: value as T[keyof T],
valueType,
});
}

counter++;
} catch (error) {
console.error(`Error processing field ${key}:`, error);
throw error;
}
}

// Delete envelopes that are no longer in the payload
const existingOntologies = new Set(Object.keys(meta.payload));
const envelopesToDelete = existing.envelopes.filter(
(e) => !existingOntologies.has(e.ontology),
);

for (const envelope of envelopesToDelete) {
try {
await this.runQueryInternal(
`
MATCH (e:Envelope { id: $envelopeId })
DETACH DELETE e
`,
{ envelopeId: envelope.id },
);
} catch (error) {
console.error(
`Error deleting envelope ${envelope.id}:`,
error,
);
throw error;
// PATCH semantics: fields absent from the new payload are
// left alone. Callers (notably web3-adapter) project partial
// platform updates through toGlobal — if the platform only
// touched one column, only one ontology reaches us, and
// deleting "stale" envelopes here would clobber every other
// field on the meta-envelope (e.g. wiping participantIds when
// a read-receipt update arrives).

// Build the full post-write state by merging the pre-write
// envelope set with everything we just wrote. Used by
// resolvers to fan out webhooks containing the complete
// merged state — receivers overwrite their local row with
// whatever the webhook carries, so a partial diff would
// make them lose every untouched field.
const mergedPayload: Record<string, any> = {};
for (const env of workingEnvelopes) {
mergedPayload[env.ontology] = env.value;
}
for (const env of createdEnvelopes) {
mergedPayload[env.ontology] = env.value;
}
}

return {
metaEnvelope: {
id,
ontology: meta.ontology,
acl,
},
envelopes: createdEnvelopes,
};
return {
metaEnvelope: {
id,
ontology: meta.ontology,
acl,
},
envelopes: createdEnvelopes,
mergedPayload,
};
});
} catch (error) {
console.error("Error in updateMetaEnvelopeById:", error);
throw error;
} finally {
await session.close();
}
});
}

/**
Expand Down
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,50 @@
/**
* Neo4j Migration: Add point-lookup indexes on Envelope.id and MetaEnvelope.id
*
* Without these, every `MATCH (e:Envelope { id: $id })` and
* `MATCH (m:MetaEnvelope { id: $id })` does a NodeByLabelScan over the entire
* Envelope / MetaEnvelope population, which scales linearly with total stored
* data. Update operations call these queries once per payload field, so a
* 5-field chat update can take 10+ seconds.
*
* `id` is generated by W3IDBuilder and is unique by construction, so a plain
* range index on the property is sufficient — no need for a composite index
* with eName, since the id alone is selective.
*/

import type { Driver } from "neo4j-driver";

const STATEMENTS: { name: string; cypher: string }[] = [
{
name: "envelope_id_index",
cypher: `CREATE INDEX envelope_id_index IF NOT EXISTS FOR (e:Envelope) ON (e.id)`,
},
{
name: "meta_envelope_id_index",
cypher: `CREATE INDEX meta_envelope_id_index IF NOT EXISTS FOR (m:MetaEnvelope) ON (m.id)`,
},
];

export async function createIdIndexes(driver: Driver): Promise<void> {
const session = driver.session();
try {
for (const { name, cypher } of STATEMENTS) {
try {
await session.run(cypher);
console.log(`Ensured ${name}`);
} catch (error) {
if (
error instanceof Error &&
error.message.includes("already exists")
) {
console.log(`${name} already exists`);
} else {
console.error(`Error creating ${name}:`, error);
throw error;
}
}
}
} finally {
await session.close();
}
}
Loading
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Add copy buttons to all
 blocks\n(function() {\n function addCopyButtons() {\n document.querySelectorAll('pre code').forEach(function(codeBlock) {\n if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;\n codeBlock.parentElement.setAttribute('data-copy-added', 'true');\n \n var btn = document.createElement('button');\n btn.textContent = 'Copy';\n btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';\n btn.onmouseover = function() { this.style.opacity = '1'; };\n btn.onmouseout = function() { this.style.opacity = '0.7'; };\n btn.onclick = function() {\n navigator.clipboard.writeText(codeBlock.textContent).then(function() {\n btn.textContent = 'Copied!';\n setTimeout(function() { btn.textContent = 'Copy'; }, 1500);\n });\n };\n codeBlock.parentElement.style.position = 'relative';\n codeBlock.parentElement.appendChild(btn);\n });\n }\n \n addCopyButtons();\n \n // Re-run on dynamic content\n var observer = new MutationObserver(addCopyButtons);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Add Copy Buttons to Code Blocks");
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
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
237 changes: 114 additions & 123 deletions infrastructure/evault-core/src/core/db/db.service.ts
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
import neo4j, { type Driver } from "neo4j-driver";
import { W3IDBuilder } from "w3id";
import { timed } from "../utils/timing";
import { deserializeValue, serializeValue } from "./schema";
import type {
AppendEnvelopeOperationLogParams,
Expand DownExpand Up@@ -41,12 +42,15 @@ export class DbService {
* @returns The result of the query execution
*/
private async runQueryInternal(query: string, params: Record<string, any>) {
const session = this.driver.session();
try {
return await session.run(query, params);
} finally {
await session.close();
}
const firstLine = query.trim().split("\n")[0].slice(0, 80);
return timed(`db.query "${firstLine}"`, async () => {
const session = this.driver.session();
try {
return await session.run(query, params);
} finally {
await session.close();
}
});
}

/**
Expand DownExpand Up@@ -74,11 +78,14 @@ export class DbService {
acl: string[],
eName: string,
): Promise<StoreMetaEnvelopeResult<T>> {
return timed("db.storeMetaEnvelope", async () => {
if (!eName) {
throw new Error("eName is required for storing meta-envelopes");
}

const w3id = await new W3IDBuilder().build();
const w3id = await timed("db.storeMetaEnvelope.buildMetaId", () =>
new W3IDBuilder().build(),
);

const cypher: string[] = [
`CREATE (m:MetaEnvelope { id: $metaId, ontology: $ontology, acl: $acl, eName: $eName })`,
Expand DownExpand Up@@ -128,7 +135,9 @@ export class DbService {
counter++;
}

await this.runQueryInternal(cypher.join("\n"), envelopeParams);
await timed("db.storeMetaEnvelope.runQuery", () =>
this.runQueryInternal(cypher.join("\n"), envelopeParams),
);

return {
metaEnvelope: {
Expand All@@ -138,6 +147,7 @@ export class DbService {
},
envelopes: createdEnvelopes,
};
});
}

/**
Expand DownExpand Up@@ -528,91 +538,81 @@ export class DbService {
acl: string[],
eName: string,
): Promise<StoreMetaEnvelopeResult<T>> {
return timed("db.updateMetaEnvelopeById", async () => {
if (!eName) {
throw new Error("eName is required for updating meta-envelopes");
}

// The whole read-modify-write cycle runs inside a single Neo4j write
// transaction. The opening MERGE+SET acquires a write lock on the
// MetaEnvelope node, so concurrent updates to the same id serialize
// here — without this, request B's "delete stale envelopes" step
// could clobber fields that request A just wrote.
const session = this.driver.session();
try {
let existing = await this.findMetaEnvelopeById<T>(id, eName);
if (!existing) {
const metaW3id = await new W3IDBuilder().build();
await this.runQueryInternal(
return await session.executeWrite(async (tx) => {
const findResult = await tx.run(
`
CREATE (m:MetaEnvelope {
id: $id,
ontology: $ontology,
acl: $acl,
eName: $eName
})
MERGE (m:MetaEnvelope { id: $id, eName: $eName })
ON CREATE SET m.ontology = $ontology, m.acl = $acl
ON MATCH SET m.ontology = $ontology, m.acl = $acl
WITH m
OPTIONAL MATCH (m)-[:LINKS_TO]->(e:Envelope)
RETURN collect(e) AS envelopes
`,
{ id, ontology: meta.ontology, acl, eName },
{ id, eName, ontology: meta.ontology, acl },
);
existing = {
id,
ontology: meta.ontology,
acl,
parsed: meta.payload,
envelopes: [],
};
}

// Update the meta-envelope properties (ensure eName matches)
await this.runQueryInternal(
`
MATCH (m:MetaEnvelope { id: $id, eName: $eName })
SET m.ontology = $ontology, m.acl = $acl
`,
{ id, ontology: meta.ontology, acl, eName },
);
const envelopeNodes: any[] = (
findResult.records[0]?.get("envelopes") ?? []
).filter((n: any) => n !== null && n !== undefined);

// Deduplicate envelopes — if multiple Envelope nodes share the
// same ontology (field name), keep the first and delete the rest.
// This prevents non-deterministic reads where collect(e) returns
// duplicates in undefined order and reduce picks the wrong one.
const seen = new Map<string, string>(); // ontology → kept envelope id
const dupsToDelete: string[] = [];
for (const env of existing.envelopes) {
if (seen.has(env.ontology)) {
dupsToDelete.push(env.id);
} else {
seen.set(env.ontology, env.id);
let workingEnvelopes: Envelope<T[keyof T]>[] =
envelopeNodes.map((node: any) => ({
id: node.properties.id,
ontology: node.properties.ontology,
value: deserializeValue(
node.properties.value,
node.properties.valueType,
) as T[keyof T],
valueType: node.properties.valueType,
}));

// Deduplicate envelopes — if multiple Envelope nodes share the
// same ontology, keep the first and delete the rest.
const seen = new Map<string, string>();
const dupsToDelete: string[] = [];
for (const env of workingEnvelopes) {
if (seen.has(env.ontology)) {
dupsToDelete.push(env.id);
} else {
seen.set(env.ontology, env.id);
}
}
}
if (dupsToDelete.length > 0) {
console.warn(
`[eVault] Cleaning ${dupsToDelete.length} duplicate envelope(s) for MetaEnvelope ${id}`,
);
for (const dupId of dupsToDelete) {
await this.runQueryInternal(
`MATCH (e:Envelope { id: $envelopeId }) DETACH DELETE e`,
{ envelopeId: dupId },
if (dupsToDelete.length > 0) {
console.warn(
`[eVault] Cleaning ${dupsToDelete.length} duplicate envelope(s) for MetaEnvelope ${id}`,
);
await tx.run(
`MATCH (e:Envelope) WHERE e.id IN $ids DETACH DELETE e`,
{ ids: dupsToDelete },
);
workingEnvelopes = workingEnvelopes.filter(
(e) => !dupsToDelete.includes(e.id),
);
}
// Remove deleted dupes from the existing list so the update
// loop below doesn't try to reference them.
existing.envelopes = existing.envelopes.filter(
(e) => !dupsToDelete.includes(e.id),
);
}

const createdEnvelopes: Envelope<T[keyof T]>[] = [];
let counter = 0;
const createdEnvelopes: Envelope<T[keyof T]>[] = [];

// For each field in the new payload
for (const [key, value] of Object.entries(meta.payload)) {
try {
for (const [key, value] of Object.entries(meta.payload)) {
const { value: storedValue, type: valueType } =
serializeValue(value);
const alias = `e${counter}`;

// Check if an envelope with this ontology already exists
const existingEnvelope = existing.envelopes.find(
const existingEnvelope = workingEnvelopes.find(
(e) => e.ontology === key,
);

if (existingEnvelope) {
// Update existing envelope
await this.runQueryInternal(
await tx.run(
`
MATCH (e:Envelope { id: $envelopeId })
SET e.value = $newValue, e.valueType = $valueType
Expand All@@ -623,88 +623,79 @@ export class DbService {
valueType,
},
);

createdEnvelopes.push({
id: existingEnvelope.id,
ontology: key,
value: value as T[keyof T],
valueType,
});
} else {
// Create new envelope — use MERGE on the relationship
// + ontology to prevent duplicate Envelopes if two
// concurrent updates race.
const envW3id = await new W3IDBuilder().build();
const envelopeId = envW3id.id;

await this.runQueryInternal(
await tx.run(
`
MATCH (m:MetaEnvelope { id: $metaId, eName: $eName })
MERGE (m)-[:LINKS_TO]->(${alias}:Envelope { ontology: $${alias}_ontology })
ON CREATE SET ${alias}.id = $${alias}_id, ${alias}.value = $${alias}_value, ${alias}.valueType = $${alias}_type
ON MATCH SET ${alias}.value = $${alias}_value, ${alias}.valueType = $${alias}_type
MERGE (m)-[:LINKS_TO]->(e:Envelope { ontology: $ontology })
ON CREATE SET e.id = $envelopeId, e.value = $newValue, e.valueType = $valueType
ON MATCH SET e.value = $newValue, e.valueType = $valueType
`,
{
metaId: id,
eName: eName,
[`${alias}_id`]: envelopeId,
[`${alias}_ontology`]: key,
[`${alias}_value`]: storedValue,
[`${alias}_type`]: valueType,
eName,
envelopeId,
ontology: key,
newValue: storedValue,
valueType,
},
);

createdEnvelopes.push({
id: envelopeId,
ontology: key,
value: value as T[keyof T],
valueType,
});
}

counter++;
} catch (error) {
console.error(`Error processing field ${key}:`, error);
throw error;
}
}

// Delete envelopes that are no longer in the payload
const existingOntologies = new Set(Object.keys(meta.payload));
const envelopesToDelete = existing.envelopes.filter(
(e) => !existingOntologies.has(e.ontology),
);

for (const envelope of envelopesToDelete) {
try {
await this.runQueryInternal(
`
MATCH (e:Envelope { id: $envelopeId })
DETACH DELETE e
`,
{ envelopeId: envelope.id },
);
} catch (error) {
console.error(
`Error deleting envelope ${envelope.id}:`,
error,
);
throw error;
// PATCH semantics: fields absent from the new payload are
// left alone. Callers (notably web3-adapter) project partial
// platform updates through toGlobal — if the platform only
// touched one column, only one ontology reaches us, and
// deleting "stale" envelopes here would clobber every other
// field on the meta-envelope (e.g. wiping participantIds when
// a read-receipt update arrives).

// Build the full post-write state by merging the pre-write
// envelope set with everything we just wrote. Used by
// resolvers to fan out webhooks containing the complete
// merged state — receivers overwrite their local row with
// whatever the webhook carries, so a partial diff would
// make them lose every untouched field.
const mergedPayload: Record<string, any> = {};
for (const env of workingEnvelopes) {
mergedPayload[env.ontology] = env.value;
}
for (const env of createdEnvelopes) {
mergedPayload[env.ontology] = env.value;
}
}

return {
metaEnvelope: {
id,
ontology: meta.ontology,
acl,
},
envelopes: createdEnvelopes,
};
return {
metaEnvelope: {
id,
ontology: meta.ontology,
acl,
},
envelopes: createdEnvelopes,
mergedPayload,
};
});
} catch (error) {
console.error("Error in updateMetaEnvelopeById:", error);
throw error;
} finally {
await session.close();
}
});
}

/**
Expand Down
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,50 @@
/**
* Neo4j Migration: Add point-lookup indexes on Envelope.id and MetaEnvelope.id
*
* Without these, every `MATCH (e:Envelope { id: $id })` and
* `MATCH (m:MetaEnvelope { id: $id })` does a NodeByLabelScan over the entire
* Envelope / MetaEnvelope population, which scales linearly with total stored
* data. Update operations call these queries once per payload field, so a
* 5-field chat update can take 10+ seconds.
*
* `id` is generated by W3IDBuilder and is unique by construction, so a plain
* range index on the property is sufficient — no need for a composite index
* with eName, since the id alone is selective.
*/

import type { Driver } from "neo4j-driver";

const STATEMENTS: { name: string; cypher: string }[] = [
{
name: "envelope_id_index",
cypher: `CREATE INDEX envelope_id_index IF NOT EXISTS FOR (e:Envelope) ON (e.id)`,
},
{
name: "meta_envelope_id_index",
cypher: `CREATE INDEX meta_envelope_id_index IF NOT EXISTS FOR (m:MetaEnvelope) ON (m.id)`,
},
];

export async function createIdIndexes(driver: Driver): Promise<void> {
const session = driver.session();
try {
for (const { name, cypher } of STATEMENTS) {
try {
await session.run(cypher);
console.log(`Ensured ${name}`);
} catch (error) {
if (
error instanceof Error &&
error.message.includes("already exists")
) {
console.log(`${name} already exists`);
} else {
console.error(`Error creating ${name}:`, error);
throw error;
}
}
}
} finally {
await session.close();
}
}
Loading
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Force GitHub README to respect dark mode\n(function() {\n var style = document.createElement('style');\n style.textContent = '\n .markdown-body {\n color-scheme: dark light;\n }\n .markdown-body pre { background: #161b22 !important; }\n .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; }\n .markdown-body table th, .markdown-body table td { border-color: #30363d !important; }\n .markdown-body img { background: #0d1117; }\n .markdown-body blockquote { border-left-color: #8b949e; }\n .markdown-body hr { border-color: #30363d; }\n ';\n document.head.appendChild(style);\n})();", "GitHub Dark Mode README Fix"); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
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
237 changes: 114 additions & 123 deletions infrastructure/evault-core/src/core/db/db.service.ts
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
import neo4j, { type Driver } from "neo4j-driver";
import { W3IDBuilder } from "w3id";
import { timed } from "../utils/timing";
import { deserializeValue, serializeValue } from "./schema";
import type {
AppendEnvelopeOperationLogParams,
Expand DownExpand Up@@ -41,12 +42,15 @@ export class DbService {
* @returns The result of the query execution
*/
private async runQueryInternal(query: string, params: Record<string, any>) {
const session = this.driver.session();
try {
return await session.run(query, params);
} finally {
await session.close();
}
const firstLine = query.trim().split("\n")[0].slice(0, 80);
return timed(`db.query "${firstLine}"`, async () => {
const session = this.driver.session();
try {
return await session.run(query, params);
} finally {
await session.close();
}
});
}

/**
Expand DownExpand Up@@ -74,11 +78,14 @@ export class DbService {
acl: string[],
eName: string,
): Promise<StoreMetaEnvelopeResult<T>> {
return timed("db.storeMetaEnvelope", async () => {
if (!eName) {
throw new Error("eName is required for storing meta-envelopes");
}

const w3id = await new W3IDBuilder().build();
const w3id = await timed("db.storeMetaEnvelope.buildMetaId", () =>
new W3IDBuilder().build(),
);

const cypher: string[] = [
`CREATE (m:MetaEnvelope { id: $metaId, ontology: $ontology, acl: $acl, eName: $eName })`,
Expand DownExpand Up@@ -128,7 +135,9 @@ export class DbService {
counter++;
}

await this.runQueryInternal(cypher.join("\n"), envelopeParams);
await timed("db.storeMetaEnvelope.runQuery", () =>
this.runQueryInternal(cypher.join("\n"), envelopeParams),
);

return {
metaEnvelope: {
Expand All@@ -138,6 +147,7 @@ export class DbService {
},
envelopes: createdEnvelopes,
};
});
}

/**
Expand DownExpand Up@@ -528,91 +538,81 @@ export class DbService {
acl: string[],
eName: string,
): Promise<StoreMetaEnvelopeResult<T>> {
return timed("db.updateMetaEnvelopeById", async () => {
if (!eName) {
throw new Error("eName is required for updating meta-envelopes");
}

// The whole read-modify-write cycle runs inside a single Neo4j write
// transaction. The opening MERGE+SET acquires a write lock on the
// MetaEnvelope node, so concurrent updates to the same id serialize
// here — without this, request B's "delete stale envelopes" step
// could clobber fields that request A just wrote.
const session = this.driver.session();
try {
let existing = await this.findMetaEnvelopeById<T>(id, eName);
if (!existing) {
const metaW3id = await new W3IDBuilder().build();
await this.runQueryInternal(
return await session.executeWrite(async (tx) => {
const findResult = await tx.run(
`
CREATE (m:MetaEnvelope {
id: $id,
ontology: $ontology,
acl: $acl,
eName: $eName
})
MERGE (m:MetaEnvelope { id: $id, eName: $eName })
ON CREATE SET m.ontology = $ontology, m.acl = $acl
ON MATCH SET m.ontology = $ontology, m.acl = $acl
WITH m
OPTIONAL MATCH (m)-[:LINKS_TO]->(e:Envelope)
RETURN collect(e) AS envelopes
`,
{ id, ontology: meta.ontology, acl, eName },
{ id, eName, ontology: meta.ontology, acl },
);
existing = {
id,
ontology: meta.ontology,
acl,
parsed: meta.payload,
envelopes: [],
};
}

// Update the meta-envelope properties (ensure eName matches)
await this.runQueryInternal(
`
MATCH (m:MetaEnvelope { id: $id, eName: $eName })
SET m.ontology = $ontology, m.acl = $acl
`,
{ id, ontology: meta.ontology, acl, eName },
);
const envelopeNodes: any[] = (
findResult.records[0]?.get("envelopes") ?? []
).filter((n: any) => n !== null && n !== undefined);

// Deduplicate envelopes — if multiple Envelope nodes share the
// same ontology (field name), keep the first and delete the rest.
// This prevents non-deterministic reads where collect(e) returns
// duplicates in undefined order and reduce picks the wrong one.
const seen = new Map<string, string>(); // ontology → kept envelope id
const dupsToDelete: string[] = [];
for (const env of existing.envelopes) {
if (seen.has(env.ontology)) {
dupsToDelete.push(env.id);
} else {
seen.set(env.ontology, env.id);
let workingEnvelopes: Envelope<T[keyof T]>[] =
envelopeNodes.map((node: any) => ({
id: node.properties.id,
ontology: node.properties.ontology,
value: deserializeValue(
node.properties.value,
node.properties.valueType,
) as T[keyof T],
valueType: node.properties.valueType,
}));

// Deduplicate envelopes — if multiple Envelope nodes share the
// same ontology, keep the first and delete the rest.
const seen = new Map<string, string>();
const dupsToDelete: string[] = [];
for (const env of workingEnvelopes) {
if (seen.has(env.ontology)) {
dupsToDelete.push(env.id);
} else {
seen.set(env.ontology, env.id);
}
}
}
if (dupsToDelete.length > 0) {
console.warn(
`[eVault] Cleaning ${dupsToDelete.length} duplicate envelope(s) for MetaEnvelope ${id}`,
);
for (const dupId of dupsToDelete) {
await this.runQueryInternal(
`MATCH (e:Envelope { id: $envelopeId }) DETACH DELETE e`,
{ envelopeId: dupId },
if (dupsToDelete.length > 0) {
console.warn(
`[eVault] Cleaning ${dupsToDelete.length} duplicate envelope(s) for MetaEnvelope ${id}`,
);
await tx.run(
`MATCH (e:Envelope) WHERE e.id IN $ids DETACH DELETE e`,
{ ids: dupsToDelete },
);
workingEnvelopes = workingEnvelopes.filter(
(e) => !dupsToDelete.includes(e.id),
);
}
// Remove deleted dupes from the existing list so the update
// loop below doesn't try to reference them.
existing.envelopes = existing.envelopes.filter(
(e) => !dupsToDelete.includes(e.id),
);
}

const createdEnvelopes: Envelope<T[keyof T]>[] = [];
let counter = 0;
const createdEnvelopes: Envelope<T[keyof T]>[] = [];

// For each field in the new payload
for (const [key, value] of Object.entries(meta.payload)) {
try {
for (const [key, value] of Object.entries(meta.payload)) {
const { value: storedValue, type: valueType } =
serializeValue(value);
const alias = `e${counter}`;

// Check if an envelope with this ontology already exists
const existingEnvelope = existing.envelopes.find(
const existingEnvelope = workingEnvelopes.find(
(e) => e.ontology === key,
);

if (existingEnvelope) {
// Update existing envelope
await this.runQueryInternal(
await tx.run(
`
MATCH (e:Envelope { id: $envelopeId })
SET e.value = $newValue, e.valueType = $valueType
Expand All@@ -623,88 +623,79 @@ export class DbService {
valueType,
},
);

createdEnvelopes.push({
id: existingEnvelope.id,
ontology: key,
value: value as T[keyof T],
valueType,
});
} else {
// Create new envelope — use MERGE on the relationship
// + ontology to prevent duplicate Envelopes if two
// concurrent updates race.
const envW3id = await new W3IDBuilder().build();
const envelopeId = envW3id.id;

await this.runQueryInternal(
await tx.run(
`
MATCH (m:MetaEnvelope { id: $metaId, eName: $eName })
MERGE (m)-[:LINKS_TO]->(${alias}:Envelope { ontology: $${alias}_ontology })
ON CREATE SET ${alias}.id = $${alias}_id, ${alias}.value = $${alias}_value, ${alias}.valueType = $${alias}_type
ON MATCH SET ${alias}.value = $${alias}_value, ${alias}.valueType = $${alias}_type
MERGE (m)-[:LINKS_TO]->(e:Envelope { ontology: $ontology })
ON CREATE SET e.id = $envelopeId, e.value = $newValue, e.valueType = $valueType
ON MATCH SET e.value = $newValue, e.valueType = $valueType
`,
{
metaId: id,
eName: eName,
[`${alias}_id`]: envelopeId,
[`${alias}_ontology`]: key,
[`${alias}_value`]: storedValue,
[`${alias}_type`]: valueType,
eName,
envelopeId,
ontology: key,
newValue: storedValue,
valueType,
},
);

createdEnvelopes.push({
id: envelopeId,
ontology: key,
value: value as T[keyof T],
valueType,
});
}

counter++;
} catch (error) {
console.error(`Error processing field ${key}:`, error);
throw error;
}
}

// Delete envelopes that are no longer in the payload
const existingOntologies = new Set(Object.keys(meta.payload));
const envelopesToDelete = existing.envelopes.filter(
(e) => !existingOntologies.has(e.ontology),
);

for (const envelope of envelopesToDelete) {
try {
await this.runQueryInternal(
`
MATCH (e:Envelope { id: $envelopeId })
DETACH DELETE e
`,
{ envelopeId: envelope.id },
);
} catch (error) {
console.error(
`Error deleting envelope ${envelope.id}:`,
error,
);
throw error;
// PATCH semantics: fields absent from the new payload are
// left alone. Callers (notably web3-adapter) project partial
// platform updates through toGlobal — if the platform only
// touched one column, only one ontology reaches us, and
// deleting "stale" envelopes here would clobber every other
// field on the meta-envelope (e.g. wiping participantIds when
// a read-receipt update arrives).

// Build the full post-write state by merging the pre-write
// envelope set with everything we just wrote. Used by
// resolvers to fan out webhooks containing the complete
// merged state — receivers overwrite their local row with
// whatever the webhook carries, so a partial diff would
// make them lose every untouched field.
const mergedPayload: Record<string, any> = {};
for (const env of workingEnvelopes) {
mergedPayload[env.ontology] = env.value;
}
for (const env of createdEnvelopes) {
mergedPayload[env.ontology] = env.value;
}
}

return {
metaEnvelope: {
id,
ontology: meta.ontology,
acl,
},
envelopes: createdEnvelopes,
};
return {
metaEnvelope: {
id,
ontology: meta.ontology,
acl,
},
envelopes: createdEnvelopes,
mergedPayload,
};
});
} catch (error) {
console.error("Error in updateMetaEnvelopeById:", error);
throw error;
} finally {
await session.close();
}
});
}

/**
Expand Down
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,50 @@
/**
* Neo4j Migration: Add point-lookup indexes on Envelope.id and MetaEnvelope.id
*
* Without these, every `MATCH (e:Envelope { id: $id })` and
* `MATCH (m:MetaEnvelope { id: $id })` does a NodeByLabelScan over the entire
* Envelope / MetaEnvelope population, which scales linearly with total stored
* data. Update operations call these queries once per payload field, so a
* 5-field chat update can take 10+ seconds.
*
* `id` is generated by W3IDBuilder and is unique by construction, so a plain
* range index on the property is sufficient — no need for a composite index
* with eName, since the id alone is selective.
*/

import type { Driver } from "neo4j-driver";

const STATEMENTS: { name: string; cypher: string }[] = [
{
name: "envelope_id_index",
cypher: `CREATE INDEX envelope_id_index IF NOT EXISTS FOR (e:Envelope) ON (e.id)`,
},
{
name: "meta_envelope_id_index",
cypher: `CREATE INDEX meta_envelope_id_index IF NOT EXISTS FOR (m:MetaEnvelope) ON (m.id)`,
},
];

export async function createIdIndexes(driver: Driver): Promise<void> {
const session = driver.session();
try {
for (const { name, cypher } of STATEMENTS) {
try {
await session.run(cypher);
console.log(`Ensured ${name}`);
} catch (error) {
if (
error instanceof Error &&
error.message.includes("already exists")
) {
console.log(`${name} already exists`);
} else {
console.error(`Error creating ${name}:`, error);
throw error;
}
}
}
} finally {
await session.close();
}
}
Loading
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Highlight search terms from Google/DuckDuckGo/Bing referrer\n(function() {\n var ref = document.referrer;\n var terms = [];\n \n if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) {\n var url = new URL(ref);\n var q = url.searchParams.get('q') || url.searchParams.get('p');\n if (q) {\n terms = q.split(/\\s+/).filter(function(t) { return t.length > 2; });\n }\n }\n \n if (terms.length === 0) return;\n \n var style = document.createElement('style');\n style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }';\n document.head.appendChild(style);\n \n function highlight(node) {\n if (node.nodeType === 3) { // text node\n var text = node.textContent;\n var found = false;\n terms.forEach(function(term) {\n var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\\]\\\\]/g, '\\\\') + ')', 'gi');\n if (regex.test(text)) {\n found = true;\n var frag = document.createDocumentFragment();\n var parts = text.split(regex);\n parts.forEach(function(part, i) {\n if (i % 2 === 0) {\n frag.appendChild(document.createTextNode(part));\n } else {\n var span = document.createElement('span');\n span.className = 'userscript-highlight';\n span.textContent = part;\n frag.appendChild(span);\n }\n });\n node.parentNode.replaceChild(frag, node);\n }\n });\n } else if (node.nodeType === 1 && node.childNodes) { // element\n var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT'];\n if (!skipTags.includes(node.tagName)) {\n Array.from(node.childNodes).forEach(highlight);\n }\n }\n }\n \n highlight(document.body);\n \n // Re-highlight on dynamic content\n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1 || node.nodeType === 3) highlight(node);\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Highlight Search Terms"); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
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
237 changes: 114 additions & 123 deletions infrastructure/evault-core/src/core/db/db.service.ts
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
import neo4j, { type Driver } from "neo4j-driver";
import { W3IDBuilder } from "w3id";
import { timed } from "../utils/timing";
import { deserializeValue, serializeValue } from "./schema";
import type {
AppendEnvelopeOperationLogParams,
Expand DownExpand Up@@ -41,12 +42,15 @@ export class DbService {
* @returns The result of the query execution
*/
private async runQueryInternal(query: string, params: Record<string, any>) {
const session = this.driver.session();
try {
return await session.run(query, params);
} finally {
await session.close();
}
const firstLine = query.trim().split("\n")[0].slice(0, 80);
return timed(`db.query "${firstLine}"`, async () => {
const session = this.driver.session();
try {
return await session.run(query, params);
} finally {
await session.close();
}
});
}

/**
Expand DownExpand Up@@ -74,11 +78,14 @@ export class DbService {
acl: string[],
eName: string,
): Promise<StoreMetaEnvelopeResult<T>> {
return timed("db.storeMetaEnvelope", async () => {
if (!eName) {
throw new Error("eName is required for storing meta-envelopes");
}

const w3id = await new W3IDBuilder().build();
const w3id = await timed("db.storeMetaEnvelope.buildMetaId", () =>
new W3IDBuilder().build(),
);

const cypher: string[] = [
`CREATE (m:MetaEnvelope { id: $metaId, ontology: $ontology, acl: $acl, eName: $eName })`,
Expand DownExpand Up@@ -128,7 +135,9 @@ export class DbService {
counter++;
}

await this.runQueryInternal(cypher.join("\n"), envelopeParams);
await timed("db.storeMetaEnvelope.runQuery", () =>
this.runQueryInternal(cypher.join("\n"), envelopeParams),
);

return {
metaEnvelope: {
Expand All@@ -138,6 +147,7 @@ export class DbService {
},
envelopes: createdEnvelopes,
};
});
}

/**
Expand DownExpand Up@@ -528,91 +538,81 @@ export class DbService {
acl: string[],
eName: string,
): Promise<StoreMetaEnvelopeResult<T>> {
return timed("db.updateMetaEnvelopeById", async () => {
if (!eName) {
throw new Error("eName is required for updating meta-envelopes");
}

// The whole read-modify-write cycle runs inside a single Neo4j write
// transaction. The opening MERGE+SET acquires a write lock on the
// MetaEnvelope node, so concurrent updates to the same id serialize
// here — without this, request B's "delete stale envelopes" step
// could clobber fields that request A just wrote.
const session = this.driver.session();
try {
let existing = await this.findMetaEnvelopeById<T>(id, eName);
if (!existing) {
const metaW3id = await new W3IDBuilder().build();
await this.runQueryInternal(
return await session.executeWrite(async (tx) => {
const findResult = await tx.run(
`
CREATE (m:MetaEnvelope {
id: $id,
ontology: $ontology,
acl: $acl,
eName: $eName
})
MERGE (m:MetaEnvelope { id: $id, eName: $eName })
ON CREATE SET m.ontology = $ontology, m.acl = $acl
ON MATCH SET m.ontology = $ontology, m.acl = $acl
WITH m
OPTIONAL MATCH (m)-[:LINKS_TO]->(e:Envelope)
RETURN collect(e) AS envelopes
`,
{ id, ontology: meta.ontology, acl, eName },
{ id, eName, ontology: meta.ontology, acl },
);
existing = {
id,
ontology: meta.ontology,
acl,
parsed: meta.payload,
envelopes: [],
};
}

// Update the meta-envelope properties (ensure eName matches)
await this.runQueryInternal(
`
MATCH (m:MetaEnvelope { id: $id, eName: $eName })
SET m.ontology = $ontology, m.acl = $acl
`,
{ id, ontology: meta.ontology, acl, eName },
);
const envelopeNodes: any[] = (
findResult.records[0]?.get("envelopes") ?? []
).filter((n: any) => n !== null && n !== undefined);

// Deduplicate envelopes — if multiple Envelope nodes share the
// same ontology (field name), keep the first and delete the rest.
// This prevents non-deterministic reads where collect(e) returns
// duplicates in undefined order and reduce picks the wrong one.
const seen = new Map<string, string>(); // ontology → kept envelope id
const dupsToDelete: string[] = [];
for (const env of existing.envelopes) {
if (seen.has(env.ontology)) {
dupsToDelete.push(env.id);
} else {
seen.set(env.ontology, env.id);
let workingEnvelopes: Envelope<T[keyof T]>[] =
envelopeNodes.map((node: any) => ({
id: node.properties.id,
ontology: node.properties.ontology,
value: deserializeValue(
node.properties.value,
node.properties.valueType,
) as T[keyof T],
valueType: node.properties.valueType,
}));

// Deduplicate envelopes — if multiple Envelope nodes share the
// same ontology, keep the first and delete the rest.
const seen = new Map<string, string>();
const dupsToDelete: string[] = [];
for (const env of workingEnvelopes) {
if (seen.has(env.ontology)) {
dupsToDelete.push(env.id);
} else {
seen.set(env.ontology, env.id);
}
}
}
if (dupsToDelete.length > 0) {
console.warn(
`[eVault] Cleaning ${dupsToDelete.length} duplicate envelope(s) for MetaEnvelope ${id}`,
);
for (const dupId of dupsToDelete) {
await this.runQueryInternal(
`MATCH (e:Envelope { id: $envelopeId }) DETACH DELETE e`,
{ envelopeId: dupId },
if (dupsToDelete.length > 0) {
console.warn(
`[eVault] Cleaning ${dupsToDelete.length} duplicate envelope(s) for MetaEnvelope ${id}`,
);
await tx.run(
`MATCH (e:Envelope) WHERE e.id IN $ids DETACH DELETE e`,
{ ids: dupsToDelete },
);
workingEnvelopes = workingEnvelopes.filter(
(e) => !dupsToDelete.includes(e.id),
);
}
// Remove deleted dupes from the existing list so the update
// loop below doesn't try to reference them.
existing.envelopes = existing.envelopes.filter(
(e) => !dupsToDelete.includes(e.id),
);
}

const createdEnvelopes: Envelope<T[keyof T]>[] = [];
let counter = 0;
const createdEnvelopes: Envelope<T[keyof T]>[] = [];

// For each field in the new payload
for (const [key, value] of Object.entries(meta.payload)) {
try {
for (const [key, value] of Object.entries(meta.payload)) {
const { value: storedValue, type: valueType } =
serializeValue(value);
const alias = `e${counter}`;

// Check if an envelope with this ontology already exists
const existingEnvelope = existing.envelopes.find(
const existingEnvelope = workingEnvelopes.find(
(e) => e.ontology === key,
);

if (existingEnvelope) {
// Update existing envelope
await this.runQueryInternal(
await tx.run(
`
MATCH (e:Envelope { id: $envelopeId })
SET e.value = $newValue, e.valueType = $valueType
Expand All@@ -623,88 +623,79 @@ export class DbService {
valueType,
},
);

createdEnvelopes.push({
id: existingEnvelope.id,
ontology: key,
value: value as T[keyof T],
valueType,
});
} else {
// Create new envelope — use MERGE on the relationship
// + ontology to prevent duplicate Envelopes if two
// concurrent updates race.
const envW3id = await new W3IDBuilder().build();
const envelopeId = envW3id.id;

await this.runQueryInternal(
await tx.run(
`
MATCH (m:MetaEnvelope { id: $metaId, eName: $eName })
MERGE (m)-[:LINKS_TO]->(${alias}:Envelope { ontology: $${alias}_ontology })
ON CREATE SET ${alias}.id = $${alias}_id, ${alias}.value = $${alias}_value, ${alias}.valueType = $${alias}_type
ON MATCH SET ${alias}.value = $${alias}_value, ${alias}.valueType = $${alias}_type
MERGE (m)-[:LINKS_TO]->(e:Envelope { ontology: $ontology })
ON CREATE SET e.id = $envelopeId, e.value = $newValue, e.valueType = $valueType
ON MATCH SET e.value = $newValue, e.valueType = $valueType
`,
{
metaId: id,
eName: eName,
[`${alias}_id`]: envelopeId,
[`${alias}_ontology`]: key,
[`${alias}_value`]: storedValue,
[`${alias}_type`]: valueType,
eName,
envelopeId,
ontology: key,
newValue: storedValue,
valueType,
},
);

createdEnvelopes.push({
id: envelopeId,
ontology: key,
value: value as T[keyof T],
valueType,
});
}

counter++;
} catch (error) {
console.error(`Error processing field ${key}:`, error);
throw error;
}
}

// Delete envelopes that are no longer in the payload
const existingOntologies = new Set(Object.keys(meta.payload));
const envelopesToDelete = existing.envelopes.filter(
(e) => !existingOntologies.has(e.ontology),
);

for (const envelope of envelopesToDelete) {
try {
await this.runQueryInternal(
`
MATCH (e:Envelope { id: $envelopeId })
DETACH DELETE e
`,
{ envelopeId: envelope.id },
);
} catch (error) {
console.error(
`Error deleting envelope ${envelope.id}:`,
error,
);
throw error;
// PATCH semantics: fields absent from the new payload are
// left alone. Callers (notably web3-adapter) project partial
// platform updates through toGlobal — if the platform only
// touched one column, only one ontology reaches us, and
// deleting "stale" envelopes here would clobber every other
// field on the meta-envelope (e.g. wiping participantIds when
// a read-receipt update arrives).

// Build the full post-write state by merging the pre-write
// envelope set with everything we just wrote. Used by
// resolvers to fan out webhooks containing the complete
// merged state — receivers overwrite their local row with
// whatever the webhook carries, so a partial diff would
// make them lose every untouched field.
const mergedPayload: Record<string, any> = {};
for (const env of workingEnvelopes) {
mergedPayload[env.ontology] = env.value;
}
for (const env of createdEnvelopes) {
mergedPayload[env.ontology] = env.value;
}
}

return {
metaEnvelope: {
id,
ontology: meta.ontology,
acl,
},
envelopes: createdEnvelopes,
};
return {
metaEnvelope: {
id,
ontology: meta.ontology,
acl,
},
envelopes: createdEnvelopes,
mergedPayload,
};
});
} catch (error) {
console.error("Error in updateMetaEnvelopeById:", error);
throw error;
} finally {
await session.close();
}
});
}

/**
Expand Down
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,50 @@
/**
* Neo4j Migration: Add point-lookup indexes on Envelope.id and MetaEnvelope.id
*
* Without these, every `MATCH (e:Envelope { id: $id })` and
* `MATCH (m:MetaEnvelope { id: $id })` does a NodeByLabelScan over the entire
* Envelope / MetaEnvelope population, which scales linearly with total stored
* data. Update operations call these queries once per payload field, so a
* 5-field chat update can take 10+ seconds.
*
* `id` is generated by W3IDBuilder and is unique by construction, so a plain
* range index on the property is sufficient — no need for a composite index
* with eName, since the id alone is selective.
*/

import type { Driver } from "neo4j-driver";

const STATEMENTS: { name: string; cypher: string }[] = [
{
name: "envelope_id_index",
cypher: `CREATE INDEX envelope_id_index IF NOT EXISTS FOR (e:Envelope) ON (e.id)`,
},
{
name: "meta_envelope_id_index",
cypher: `CREATE INDEX meta_envelope_id_index IF NOT EXISTS FOR (m:MetaEnvelope) ON (m.id)`,
},
];

export async function createIdIndexes(driver: Driver): Promise<void> {
const session = driver.session();
try {
for (const { name, cypher } of STATEMENTS) {
try {
await session.run(cypher);
console.log(`Ensured ${name}`);
} catch (error) {
if (
error instanceof Error &&
error.message.includes("already exists")
) {
console.log(`${name} already exists`);
} else {
console.error(`Error creating ${name}:`, error);
throw error;
}
}
}
} finally {
await session.close();
}
}
Loading
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Strip utm_, fbclid, gclid, etc. from all links on page\n(function() {\n var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content',\n 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid',\n 'ref', 'ref_src', 'source', 'medium', 'campaign'];\n \n function cleanUrl(url) {\n try {\n var u = new URL(url, window.location.origin);\n var changed = false;\n trackingParams.forEach(function(p) {\n if (u.searchParams.has(p)) {\n u.searchParams.delete(p);\n changed = true;\n }\n });\n return changed ? u.toString() : url;\n } catch (e) {\n return url;\n }\n }\n \n function cleanLinks() {\n document.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n \n cleanLinks();\n \n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1) {\n if (node.tagName === 'A') cleanLinks();\n node.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Remove Tracking Parameters from Links"); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + '
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
237 changes: 114 additions & 123 deletions infrastructure/evault-core/src/core/db/db.service.ts
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
import neo4j, { type Driver } from "neo4j-driver";
import { W3IDBuilder } from "w3id";
import { timed } from "../utils/timing";
import { deserializeValue, serializeValue } from "./schema";
import type {
AppendEnvelopeOperationLogParams,
Expand DownExpand Up@@ -41,12 +42,15 @@ export class DbService {
* @returns The result of the query execution
*/
private async runQueryInternal(query: string, params: Record<string, any>) {
const session = this.driver.session();
try {
return await session.run(query, params);
} finally {
await session.close();
}
const firstLine = query.trim().split("\n")[0].slice(0, 80);
return timed(`db.query "${firstLine}"`, async () => {
const session = this.driver.session();
try {
return await session.run(query, params);
} finally {
await session.close();
}
});
}

/**
Expand DownExpand Up@@ -74,11 +78,14 @@ export class DbService {
acl: string[],
eName: string,
): Promise<StoreMetaEnvelopeResult<T>> {
return timed("db.storeMetaEnvelope", async () => {
if (!eName) {
throw new Error("eName is required for storing meta-envelopes");
}

const w3id = await new W3IDBuilder().build();
const w3id = await timed("db.storeMetaEnvelope.buildMetaId", () =>
new W3IDBuilder().build(),
);

const cypher: string[] = [
`CREATE (m:MetaEnvelope { id: $metaId, ontology: $ontology, acl: $acl, eName: $eName })`,
Expand DownExpand Up@@ -128,7 +135,9 @@ export class DbService {
counter++;
}

await this.runQueryInternal(cypher.join("\n"), envelopeParams);
await timed("db.storeMetaEnvelope.runQuery", () =>
this.runQueryInternal(cypher.join("\n"), envelopeParams),
);

return {
metaEnvelope: {
Expand All@@ -138,6 +147,7 @@ export class DbService {
},
envelopes: createdEnvelopes,
};
});
}

/**
Expand DownExpand Up@@ -528,91 +538,81 @@ export class DbService {
acl: string[],
eName: string,
): Promise<StoreMetaEnvelopeResult<T>> {
return timed("db.updateMetaEnvelopeById", async () => {
if (!eName) {
throw new Error("eName is required for updating meta-envelopes");
}

// The whole read-modify-write cycle runs inside a single Neo4j write
// transaction. The opening MERGE+SET acquires a write lock on the
// MetaEnvelope node, so concurrent updates to the same id serialize
// here — without this, request B's "delete stale envelopes" step
// could clobber fields that request A just wrote.
const session = this.driver.session();
try {
let existing = await this.findMetaEnvelopeById<T>(id, eName);
if (!existing) {
const metaW3id = await new W3IDBuilder().build();
await this.runQueryInternal(
return await session.executeWrite(async (tx) => {
const findResult = await tx.run(
`
CREATE (m:MetaEnvelope {
id: $id,
ontology: $ontology,
acl: $acl,
eName: $eName
})
MERGE (m:MetaEnvelope { id: $id, eName: $eName })
ON CREATE SET m.ontology = $ontology, m.acl = $acl
ON MATCH SET m.ontology = $ontology, m.acl = $acl
WITH m
OPTIONAL MATCH (m)-[:LINKS_TO]->(e:Envelope)
RETURN collect(e) AS envelopes
`,
{ id, ontology: meta.ontology, acl, eName },
{ id, eName, ontology: meta.ontology, acl },
);
existing = {
id,
ontology: meta.ontology,
acl,
parsed: meta.payload,
envelopes: [],
};
}

// Update the meta-envelope properties (ensure eName matches)
await this.runQueryInternal(
`
MATCH (m:MetaEnvelope { id: $id, eName: $eName })
SET m.ontology = $ontology, m.acl = $acl
`,
{ id, ontology: meta.ontology, acl, eName },
);
const envelopeNodes: any[] = (
findResult.records[0]?.get("envelopes") ?? []
).filter((n: any) => n !== null && n !== undefined);

// Deduplicate envelopes — if multiple Envelope nodes share the
// same ontology (field name), keep the first and delete the rest.
// This prevents non-deterministic reads where collect(e) returns
// duplicates in undefined order and reduce picks the wrong one.
const seen = new Map<string, string>(); // ontology → kept envelope id
const dupsToDelete: string[] = [];
for (const env of existing.envelopes) {
if (seen.has(env.ontology)) {
dupsToDelete.push(env.id);
} else {
seen.set(env.ontology, env.id);
let workingEnvelopes: Envelope<T[keyof T]>[] =
envelopeNodes.map((node: any) => ({
id: node.properties.id,
ontology: node.properties.ontology,
value: deserializeValue(
node.properties.value,
node.properties.valueType,
) as T[keyof T],
valueType: node.properties.valueType,
}));

// Deduplicate envelopes — if multiple Envelope nodes share the
// same ontology, keep the first and delete the rest.
const seen = new Map<string, string>();
const dupsToDelete: string[] = [];
for (const env of workingEnvelopes) {
if (seen.has(env.ontology)) {
dupsToDelete.push(env.id);
} else {
seen.set(env.ontology, env.id);
}
}
}
if (dupsToDelete.length > 0) {
console.warn(
`[eVault] Cleaning ${dupsToDelete.length} duplicate envelope(s) for MetaEnvelope ${id}`,
);
for (const dupId of dupsToDelete) {
await this.runQueryInternal(
`MATCH (e:Envelope { id: $envelopeId }) DETACH DELETE e`,
{ envelopeId: dupId },
if (dupsToDelete.length > 0) {
console.warn(
`[eVault] Cleaning ${dupsToDelete.length} duplicate envelope(s) for MetaEnvelope ${id}`,
);
await tx.run(
`MATCH (e:Envelope) WHERE e.id IN $ids DETACH DELETE e`,
{ ids: dupsToDelete },
);
workingEnvelopes = workingEnvelopes.filter(
(e) => !dupsToDelete.includes(e.id),
);
}
// Remove deleted dupes from the existing list so the update
// loop below doesn't try to reference them.
existing.envelopes = existing.envelopes.filter(
(e) => !dupsToDelete.includes(e.id),
);
}

const createdEnvelopes: Envelope<T[keyof T]>[] = [];
let counter = 0;
const createdEnvelopes: Envelope<T[keyof T]>[] = [];

// For each field in the new payload
for (const [key, value] of Object.entries(meta.payload)) {
try {
for (const [key, value] of Object.entries(meta.payload)) {
const { value: storedValue, type: valueType } =
serializeValue(value);
const alias = `e${counter}`;

// Check if an envelope with this ontology already exists
const existingEnvelope = existing.envelopes.find(
const existingEnvelope = workingEnvelopes.find(
(e) => e.ontology === key,
);

if (existingEnvelope) {
// Update existing envelope
await this.runQueryInternal(
await tx.run(
`
MATCH (e:Envelope { id: $envelopeId })
SET e.value = $newValue, e.valueType = $valueType
Expand All@@ -623,88 +623,79 @@ export class DbService {
valueType,
},
);

createdEnvelopes.push({
id: existingEnvelope.id,
ontology: key,
value: value as T[keyof T],
valueType,
});
} else {
// Create new envelope — use MERGE on the relationship
// + ontology to prevent duplicate Envelopes if two
// concurrent updates race.
const envW3id = await new W3IDBuilder().build();
const envelopeId = envW3id.id;

await this.runQueryInternal(
await tx.run(
`
MATCH (m:MetaEnvelope { id: $metaId, eName: $eName })
MERGE (m)-[:LINKS_TO]->(${alias}:Envelope { ontology: $${alias}_ontology })
ON CREATE SET ${alias}.id = $${alias}_id, ${alias}.value = $${alias}_value, ${alias}.valueType = $${alias}_type
ON MATCH SET ${alias}.value = $${alias}_value, ${alias}.valueType = $${alias}_type
MERGE (m)-[:LINKS_TO]->(e:Envelope { ontology: $ontology })
ON CREATE SET e.id = $envelopeId, e.value = $newValue, e.valueType = $valueType
ON MATCH SET e.value = $newValue, e.valueType = $valueType
`,
{
metaId: id,
eName: eName,
[`${alias}_id`]: envelopeId,
[`${alias}_ontology`]: key,
[`${alias}_value`]: storedValue,
[`${alias}_type`]: valueType,
eName,
envelopeId,
ontology: key,
newValue: storedValue,
valueType,
},
);

createdEnvelopes.push({
id: envelopeId,
ontology: key,
value: value as T[keyof T],
valueType,
});
}

counter++;
} catch (error) {
console.error(`Error processing field ${key}:`, error);
throw error;
}
}

// Delete envelopes that are no longer in the payload
const existingOntologies = new Set(Object.keys(meta.payload));
const envelopesToDelete = existing.envelopes.filter(
(e) => !existingOntologies.has(e.ontology),
);

for (const envelope of envelopesToDelete) {
try {
await this.runQueryInternal(
`
MATCH (e:Envelope { id: $envelopeId })
DETACH DELETE e
`,
{ envelopeId: envelope.id },
);
} catch (error) {
console.error(
`Error deleting envelope ${envelope.id}:`,
error,
);
throw error;
// PATCH semantics: fields absent from the new payload are
// left alone. Callers (notably web3-adapter) project partial
// platform updates through toGlobal — if the platform only
// touched one column, only one ontology reaches us, and
// deleting "stale" envelopes here would clobber every other
// field on the meta-envelope (e.g. wiping participantIds when
// a read-receipt update arrives).

// Build the full post-write state by merging the pre-write
// envelope set with everything we just wrote. Used by
// resolvers to fan out webhooks containing the complete
// merged state — receivers overwrite their local row with
// whatever the webhook carries, so a partial diff would
// make them lose every untouched field.
const mergedPayload: Record<string, any> = {};
for (const env of workingEnvelopes) {
mergedPayload[env.ontology] = env.value;
}
for (const env of createdEnvelopes) {
mergedPayload[env.ontology] = env.value;
}
}

return {
metaEnvelope: {
id,
ontology: meta.ontology,
acl,
},
envelopes: createdEnvelopes,
};
return {
metaEnvelope: {
id,
ontology: meta.ontology,
acl,
},
envelopes: createdEnvelopes,
mergedPayload,
};
});
} catch (error) {
console.error("Error in updateMetaEnvelopeById:", error);
throw error;
} finally {
await session.close();
}
});
}

/**
Expand Down
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,50 @@
/**
* Neo4j Migration: Add point-lookup indexes on Envelope.id and MetaEnvelope.id
*
* Without these, every `MATCH (e:Envelope { id: $id })` and
* `MATCH (m:MetaEnvelope { id: $id })` does a NodeByLabelScan over the entire
* Envelope / MetaEnvelope population, which scales linearly with total stored
* data. Update operations call these queries once per payload field, so a
* 5-field chat update can take 10+ seconds.
*
* `id` is generated by W3IDBuilder and is unique by construction, so a plain
* range index on the property is sufficient — no need for a composite index
* with eName, since the id alone is selective.
*/

import type { Driver } from "neo4j-driver";

const STATEMENTS: { name: string; cypher: string }[] = [
{
name: "envelope_id_index",
cypher: `CREATE INDEX envelope_id_index IF NOT EXISTS FOR (e:Envelope) ON (e.id)`,
},
{
name: "meta_envelope_id_index",
cypher: `CREATE INDEX meta_envelope_id_index IF NOT EXISTS FOR (m:MetaEnvelope) ON (m.id)`,
},
];

export async function createIdIndexes(driver: Driver): Promise<void> {
const session = driver.session();
try {
for (const { name, cypher } of STATEMENTS) {
try {
await session.run(cypher);
console.log(`Ensured ${name}`);
} catch (error) {
if (
error instanceof Error &&
error.message.includes("already exists")
) {
console.log(`${name} already exists`);
} else {
console.error(`Error creating ${name}:`, error);
throw error;
}
}
}
} finally {
await session.close();
}
}
Loading
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Auto-enable theater mode on YouTube\n(function() {\n function tryTheater() {\n var btn = document.querySelector('button[aria-label=\"Theater mode\"], ytd-player #player button[title=\"Theater mode\"]');\n if (btn && !btn.classList.contains('activated')) {\n btn.click();\n }\n }\n \n // Try immediately\n tryTheater();\n \n // Try after navigation (SPA)\n var lastUrl = location.href;\n setInterval(function() {\n if (location.href !== lastUrl) {\n lastUrl = location.href;\n setTimeout(tryTheater, 500);\n }\n }, 1000);\n \n // Also try on player load\n var observer = new MutationObserver(tryTheater);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "YouTube Theater Mode Default"); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
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
237 changes: 114 additions & 123 deletions infrastructure/evault-core/src/core/db/db.service.ts
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
import neo4j, { type Driver } from "neo4j-driver";
import { W3IDBuilder } from "w3id";
import { timed } from "../utils/timing";
import { deserializeValue, serializeValue } from "./schema";
import type {
AppendEnvelopeOperationLogParams,
Expand DownExpand Up@@ -41,12 +42,15 @@ export class DbService {
* @returns The result of the query execution
*/
private async runQueryInternal(query: string, params: Record<string, any>) {
const session = this.driver.session();
try {
return await session.run(query, params);
} finally {
await session.close();
}
const firstLine = query.trim().split("\n")[0].slice(0, 80);
return timed(`db.query "${firstLine}"`, async () => {
const session = this.driver.session();
try {
return await session.run(query, params);
} finally {
await session.close();
}
});
}

/**
Expand DownExpand Up@@ -74,11 +78,14 @@ export class DbService {
acl: string[],
eName: string,
): Promise<StoreMetaEnvelopeResult<T>> {
return timed("db.storeMetaEnvelope", async () => {
if (!eName) {
throw new Error("eName is required for storing meta-envelopes");
}

const w3id = await new W3IDBuilder().build();
const w3id = await timed("db.storeMetaEnvelope.buildMetaId", () =>
new W3IDBuilder().build(),
);

const cypher: string[] = [
`CREATE (m:MetaEnvelope { id: $metaId, ontology: $ontology, acl: $acl, eName: $eName })`,
Expand DownExpand Up@@ -128,7 +135,9 @@ export class DbService {
counter++;
}

await this.runQueryInternal(cypher.join("\n"), envelopeParams);
await timed("db.storeMetaEnvelope.runQuery", () =>
this.runQueryInternal(cypher.join("\n"), envelopeParams),
);

return {
metaEnvelope: {
Expand All@@ -138,6 +147,7 @@ export class DbService {
},
envelopes: createdEnvelopes,
};
});
}

/**
Expand DownExpand Up@@ -528,91 +538,81 @@ export class DbService {
acl: string[],
eName: string,
): Promise<StoreMetaEnvelopeResult<T>> {
return timed("db.updateMetaEnvelopeById", async () => {
if (!eName) {
throw new Error("eName is required for updating meta-envelopes");
}

// The whole read-modify-write cycle runs inside a single Neo4j write
// transaction. The opening MERGE+SET acquires a write lock on the
// MetaEnvelope node, so concurrent updates to the same id serialize
// here — without this, request B's "delete stale envelopes" step
// could clobber fields that request A just wrote.
const session = this.driver.session();
try {
let existing = await this.findMetaEnvelopeById<T>(id, eName);
if (!existing) {
const metaW3id = await new W3IDBuilder().build();
await this.runQueryInternal(
return await session.executeWrite(async (tx) => {
const findResult = await tx.run(
`
CREATE (m:MetaEnvelope {
id: $id,
ontology: $ontology,
acl: $acl,
eName: $eName
})
MERGE (m:MetaEnvelope { id: $id, eName: $eName })
ON CREATE SET m.ontology = $ontology, m.acl = $acl
ON MATCH SET m.ontology = $ontology, m.acl = $acl
WITH m
OPTIONAL MATCH (m)-[:LINKS_TO]->(e:Envelope)
RETURN collect(e) AS envelopes
`,
{ id, ontology: meta.ontology, acl, eName },
{ id, eName, ontology: meta.ontology, acl },
);
existing = {
id,
ontology: meta.ontology,
acl,
parsed: meta.payload,
envelopes: [],
};
}

// Update the meta-envelope properties (ensure eName matches)
await this.runQueryInternal(
`
MATCH (m:MetaEnvelope { id: $id, eName: $eName })
SET m.ontology = $ontology, m.acl = $acl
`,
{ id, ontology: meta.ontology, acl, eName },
);
const envelopeNodes: any[] = (
findResult.records[0]?.get("envelopes") ?? []
).filter((n: any) => n !== null && n !== undefined);

// Deduplicate envelopes — if multiple Envelope nodes share the
// same ontology (field name), keep the first and delete the rest.
// This prevents non-deterministic reads where collect(e) returns
// duplicates in undefined order and reduce picks the wrong one.
const seen = new Map<string, string>(); // ontology → kept envelope id
const dupsToDelete: string[] = [];
for (const env of existing.envelopes) {
if (seen.has(env.ontology)) {
dupsToDelete.push(env.id);
} else {
seen.set(env.ontology, env.id);
let workingEnvelopes: Envelope<T[keyof T]>[] =
envelopeNodes.map((node: any) => ({
id: node.properties.id,
ontology: node.properties.ontology,
value: deserializeValue(
node.properties.value,
node.properties.valueType,
) as T[keyof T],
valueType: node.properties.valueType,
}));

// Deduplicate envelopes — if multiple Envelope nodes share the
// same ontology, keep the first and delete the rest.
const seen = new Map<string, string>();
const dupsToDelete: string[] = [];
for (const env of workingEnvelopes) {
if (seen.has(env.ontology)) {
dupsToDelete.push(env.id);
} else {
seen.set(env.ontology, env.id);
}
}
}
if (dupsToDelete.length > 0) {
console.warn(
`[eVault] Cleaning ${dupsToDelete.length} duplicate envelope(s) for MetaEnvelope ${id}`,
);
for (const dupId of dupsToDelete) {
await this.runQueryInternal(
`MATCH (e:Envelope { id: $envelopeId }) DETACH DELETE e`,
{ envelopeId: dupId },
if (dupsToDelete.length > 0) {
console.warn(
`[eVault] Cleaning ${dupsToDelete.length} duplicate envelope(s) for MetaEnvelope ${id}`,
);
await tx.run(
`MATCH (e:Envelope) WHERE e.id IN $ids DETACH DELETE e`,
{ ids: dupsToDelete },
);
workingEnvelopes = workingEnvelopes.filter(
(e) => !dupsToDelete.includes(e.id),
);
}
// Remove deleted dupes from the existing list so the update
// loop below doesn't try to reference them.
existing.envelopes = existing.envelopes.filter(
(e) => !dupsToDelete.includes(e.id),
);
}

const createdEnvelopes: Envelope<T[keyof T]>[] = [];
let counter = 0;
const createdEnvelopes: Envelope<T[keyof T]>[] = [];

// For each field in the new payload
for (const [key, value] of Object.entries(meta.payload)) {
try {
for (const [key, value] of Object.entries(meta.payload)) {
const { value: storedValue, type: valueType } =
serializeValue(value);
const alias = `e${counter}`;

// Check if an envelope with this ontology already exists
const existingEnvelope = existing.envelopes.find(
const existingEnvelope = workingEnvelopes.find(
(e) => e.ontology === key,
);

if (existingEnvelope) {
// Update existing envelope
await this.runQueryInternal(
await tx.run(
`
MATCH (e:Envelope { id: $envelopeId })
SET e.value = $newValue, e.valueType = $valueType
Expand All@@ -623,88 +623,79 @@ export class DbService {
valueType,
},
);

createdEnvelopes.push({
id: existingEnvelope.id,
ontology: key,
value: value as T[keyof T],
valueType,
});
} else {
// Create new envelope — use MERGE on the relationship
// + ontology to prevent duplicate Envelopes if two
// concurrent updates race.
const envW3id = await new W3IDBuilder().build();
const envelopeId = envW3id.id;

await this.runQueryInternal(
await tx.run(
`
MATCH (m:MetaEnvelope { id: $metaId, eName: $eName })
MERGE (m)-[:LINKS_TO]->(${alias}:Envelope { ontology: $${alias}_ontology })
ON CREATE SET ${alias}.id = $${alias}_id, ${alias}.value = $${alias}_value, ${alias}.valueType = $${alias}_type
ON MATCH SET ${alias}.value = $${alias}_value, ${alias}.valueType = $${alias}_type
MERGE (m)-[:LINKS_TO]->(e:Envelope { ontology: $ontology })
ON CREATE SET e.id = $envelopeId, e.value = $newValue, e.valueType = $valueType
ON MATCH SET e.value = $newValue, e.valueType = $valueType
`,
{
metaId: id,
eName: eName,
[`${alias}_id`]: envelopeId,
[`${alias}_ontology`]: key,
[`${alias}_value`]: storedValue,
[`${alias}_type`]: valueType,
eName,
envelopeId,
ontology: key,
newValue: storedValue,
valueType,
},
);

createdEnvelopes.push({
id: envelopeId,
ontology: key,
value: value as T[keyof T],
valueType,
});
}

counter++;
} catch (error) {
console.error(`Error processing field ${key}:`, error);
throw error;
}
}

// Delete envelopes that are no longer in the payload
const existingOntologies = new Set(Object.keys(meta.payload));
const envelopesToDelete = existing.envelopes.filter(
(e) => !existingOntologies.has(e.ontology),
);

for (const envelope of envelopesToDelete) {
try {
await this.runQueryInternal(
`
MATCH (e:Envelope { id: $envelopeId })
DETACH DELETE e
`,
{ envelopeId: envelope.id },
);
} catch (error) {
console.error(
`Error deleting envelope ${envelope.id}:`,
error,
);
throw error;
// PATCH semantics: fields absent from the new payload are
// left alone. Callers (notably web3-adapter) project partial
// platform updates through toGlobal — if the platform only
// touched one column, only one ontology reaches us, and
// deleting "stale" envelopes here would clobber every other
// field on the meta-envelope (e.g. wiping participantIds when
// a read-receipt update arrives).

// Build the full post-write state by merging the pre-write
// envelope set with everything we just wrote. Used by
// resolvers to fan out webhooks containing the complete
// merged state — receivers overwrite their local row with
// whatever the webhook carries, so a partial diff would
// make them lose every untouched field.
const mergedPayload: Record<string, any> = {};
for (const env of workingEnvelopes) {
mergedPayload[env.ontology] = env.value;
}
for (const env of createdEnvelopes) {
mergedPayload[env.ontology] = env.value;
}
}

return {
metaEnvelope: {
id,
ontology: meta.ontology,
acl,
},
envelopes: createdEnvelopes,
};
return {
metaEnvelope: {
id,
ontology: meta.ontology,
acl,
},
envelopes: createdEnvelopes,
mergedPayload,
};
});
} catch (error) {
console.error("Error in updateMetaEnvelopeById:", error);
throw error;
} finally {
await session.close();
}
});
}

/**
Expand Down
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,50 @@
/**
* Neo4j Migration: Add point-lookup indexes on Envelope.id and MetaEnvelope.id
*
* Without these, every `MATCH (e:Envelope { id: $id })` and
* `MATCH (m:MetaEnvelope { id: $id })` does a NodeByLabelScan over the entire
* Envelope / MetaEnvelope population, which scales linearly with total stored
* data. Update operations call these queries once per payload field, so a
* 5-field chat update can take 10+ seconds.
*
* `id` is generated by W3IDBuilder and is unique by construction, so a plain
* range index on the property is sufficient — no need for a composite index
* with eName, since the id alone is selective.
*/

import type { Driver } from "neo4j-driver";

const STATEMENTS: { name: string; cypher: string }[] = [
{
name: "envelope_id_index",
cypher: `CREATE INDEX envelope_id_index IF NOT EXISTS FOR (e:Envelope) ON (e.id)`,
},
{
name: "meta_envelope_id_index",
cypher: `CREATE INDEX meta_envelope_id_index IF NOT EXISTS FOR (m:MetaEnvelope) ON (m.id)`,
},
];

export async function createIdIndexes(driver: Driver): Promise<void> {
const session = driver.session();
try {
for (const { name, cypher } of STATEMENTS) {
try {
await session.run(cypher);
console.log(`Ensured ${name}`);
} catch (error) {
if (
error instanceof Error &&
error.message.includes("already exists")
) {
console.log(`${name} already exists`);
} else {
console.error(`Error creating ${name}:`, error);
throw error;
}
}
}
} finally {
await session.close();
}
}
Loading
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Remove or un-stick sticky/fixed headers that block content\n(function() {\n function unstick() {\n document.querySelectorAll('header, nav, [role=\"banner\"], .header, .navbar, .sticky, .fixed-top, [style*=\"position: fixed\"], [style*=\"position:sticky\"]').forEach(function(el) {\n if (el.style.position === 'fixed' || el.style.position === 'sticky' || \n getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') {\n el.style.position = 'static';\n el.style.top = 'auto';\n el.style.zIndex = 'auto';\n }\n });\n }\n \n unstick();\n \n var observer = new MutationObserver(unstick);\n observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] });\n})();", "Kill Sticky Headers"); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
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
237 changes: 114 additions & 123 deletions infrastructure/evault-core/src/core/db/db.service.ts
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
import neo4j, { type Driver } from "neo4j-driver";
import { W3IDBuilder } from "w3id";
import { timed } from "../utils/timing";
import { deserializeValue, serializeValue } from "./schema";
import type {
AppendEnvelopeOperationLogParams,
Expand DownExpand Up@@ -41,12 +42,15 @@ export class DbService {
* @returns The result of the query execution
*/
private async runQueryInternal(query: string, params: Record<string, any>) {
const session = this.driver.session();
try {
return await session.run(query, params);
} finally {
await session.close();
}
const firstLine = query.trim().split("\n")[0].slice(0, 80);
return timed(`db.query "${firstLine}"`, async () => {
const session = this.driver.session();
try {
return await session.run(query, params);
} finally {
await session.close();
}
});
}

/**
Expand DownExpand Up@@ -74,11 +78,14 @@ export class DbService {
acl: string[],
eName: string,
): Promise<StoreMetaEnvelopeResult<T>> {
return timed("db.storeMetaEnvelope", async () => {
if (!eName) {
throw new Error("eName is required for storing meta-envelopes");
}

const w3id = await new W3IDBuilder().build();
const w3id = await timed("db.storeMetaEnvelope.buildMetaId", () =>
new W3IDBuilder().build(),
);

const cypher: string[] = [
`CREATE (m:MetaEnvelope { id: $metaId, ontology: $ontology, acl: $acl, eName: $eName })`,
Expand DownExpand Up@@ -128,7 +135,9 @@ export class DbService {
counter++;
}

await this.runQueryInternal(cypher.join("\n"), envelopeParams);
await timed("db.storeMetaEnvelope.runQuery", () =>
this.runQueryInternal(cypher.join("\n"), envelopeParams),
);

return {
metaEnvelope: {
Expand All@@ -138,6 +147,7 @@ export class DbService {
},
envelopes: createdEnvelopes,
};
});
}

/**
Expand DownExpand Up@@ -528,91 +538,81 @@ export class DbService {
acl: string[],
eName: string,
): Promise<StoreMetaEnvelopeResult<T>> {
return timed("db.updateMetaEnvelopeById", async () => {
if (!eName) {
throw new Error("eName is required for updating meta-envelopes");
}

// The whole read-modify-write cycle runs inside a single Neo4j write
// transaction. The opening MERGE+SET acquires a write lock on the
// MetaEnvelope node, so concurrent updates to the same id serialize
// here — without this, request B's "delete stale envelopes" step
// could clobber fields that request A just wrote.
const session = this.driver.session();
try {
let existing = await this.findMetaEnvelopeById<T>(id, eName);
if (!existing) {
const metaW3id = await new W3IDBuilder().build();
await this.runQueryInternal(
return await session.executeWrite(async (tx) => {
const findResult = await tx.run(
`
CREATE (m:MetaEnvelope {
id: $id,
ontology: $ontology,
acl: $acl,
eName: $eName
})
MERGE (m:MetaEnvelope { id: $id, eName: $eName })
ON CREATE SET m.ontology = $ontology, m.acl = $acl
ON MATCH SET m.ontology = $ontology, m.acl = $acl
WITH m
OPTIONAL MATCH (m)-[:LINKS_TO]->(e:Envelope)
RETURN collect(e) AS envelopes
`,
{ id, ontology: meta.ontology, acl, eName },
{ id, eName, ontology: meta.ontology, acl },
);
existing = {
id,
ontology: meta.ontology,
acl,
parsed: meta.payload,
envelopes: [],
};
}

// Update the meta-envelope properties (ensure eName matches)
await this.runQueryInternal(
`
MATCH (m:MetaEnvelope { id: $id, eName: $eName })
SET m.ontology = $ontology, m.acl = $acl
`,
{ id, ontology: meta.ontology, acl, eName },
);
const envelopeNodes: any[] = (
findResult.records[0]?.get("envelopes") ?? []
).filter((n: any) => n !== null && n !== undefined);

// Deduplicate envelopes — if multiple Envelope nodes share the
// same ontology (field name), keep the first and delete the rest.
// This prevents non-deterministic reads where collect(e) returns
// duplicates in undefined order and reduce picks the wrong one.
const seen = new Map<string, string>(); // ontology → kept envelope id
const dupsToDelete: string[] = [];
for (const env of existing.envelopes) {
if (seen.has(env.ontology)) {
dupsToDelete.push(env.id);
} else {
seen.set(env.ontology, env.id);
let workingEnvelopes: Envelope<T[keyof T]>[] =
envelopeNodes.map((node: any) => ({
id: node.properties.id,
ontology: node.properties.ontology,
value: deserializeValue(
node.properties.value,
node.properties.valueType,
) as T[keyof T],
valueType: node.properties.valueType,
}));

// Deduplicate envelopes — if multiple Envelope nodes share the
// same ontology, keep the first and delete the rest.
const seen = new Map<string, string>();
const dupsToDelete: string[] = [];
for (const env of workingEnvelopes) {
if (seen.has(env.ontology)) {
dupsToDelete.push(env.id);
} else {
seen.set(env.ontology, env.id);
}
}
}
if (dupsToDelete.length > 0) {
console.warn(
`[eVault] Cleaning ${dupsToDelete.length} duplicate envelope(s) for MetaEnvelope ${id}`,
);
for (const dupId of dupsToDelete) {
await this.runQueryInternal(
`MATCH (e:Envelope { id: $envelopeId }) DETACH DELETE e`,
{ envelopeId: dupId },
if (dupsToDelete.length > 0) {
console.warn(
`[eVault] Cleaning ${dupsToDelete.length} duplicate envelope(s) for MetaEnvelope ${id}`,
);
await tx.run(
`MATCH (e:Envelope) WHERE e.id IN $ids DETACH DELETE e`,
{ ids: dupsToDelete },
);
workingEnvelopes = workingEnvelopes.filter(
(e) => !dupsToDelete.includes(e.id),
);
}
// Remove deleted dupes from the existing list so the update
// loop below doesn't try to reference them.
existing.envelopes = existing.envelopes.filter(
(e) => !dupsToDelete.includes(e.id),
);
}

const createdEnvelopes: Envelope<T[keyof T]>[] = [];
let counter = 0;
const createdEnvelopes: Envelope<T[keyof T]>[] = [];

// For each field in the new payload
for (const [key, value] of Object.entries(meta.payload)) {
try {
for (const [key, value] of Object.entries(meta.payload)) {
const { value: storedValue, type: valueType } =
serializeValue(value);
const alias = `e${counter}`;

// Check if an envelope with this ontology already exists
const existingEnvelope = existing.envelopes.find(
const existingEnvelope = workingEnvelopes.find(
(e) => e.ontology === key,
);

if (existingEnvelope) {
// Update existing envelope
await this.runQueryInternal(
await tx.run(
`
MATCH (e:Envelope { id: $envelopeId })
SET e.value = $newValue, e.valueType = $valueType
Expand All@@ -623,88 +623,79 @@ export class DbService {
valueType,
},
);

createdEnvelopes.push({
id: existingEnvelope.id,
ontology: key,
value: value as T[keyof T],
valueType,
});
} else {
// Create new envelope — use MERGE on the relationship
// + ontology to prevent duplicate Envelopes if two
// concurrent updates race.
const envW3id = await new W3IDBuilder().build();
const envelopeId = envW3id.id;

await this.runQueryInternal(
await tx.run(
`
MATCH (m:MetaEnvelope { id: $metaId, eName: $eName })
MERGE (m)-[:LINKS_TO]->(${alias}:Envelope { ontology: $${alias}_ontology })
ON CREATE SET ${alias}.id = $${alias}_id, ${alias}.value = $${alias}_value, ${alias}.valueType = $${alias}_type
ON MATCH SET ${alias}.value = $${alias}_value, ${alias}.valueType = $${alias}_type
MERGE (m)-[:LINKS_TO]->(e:Envelope { ontology: $ontology })
ON CREATE SET e.id = $envelopeId, e.value = $newValue, e.valueType = $valueType
ON MATCH SET e.value = $newValue, e.valueType = $valueType
`,
{
metaId: id,
eName: eName,
[`${alias}_id`]: envelopeId,
[`${alias}_ontology`]: key,
[`${alias}_value`]: storedValue,
[`${alias}_type`]: valueType,
eName,
envelopeId,
ontology: key,
newValue: storedValue,
valueType,
},
);

createdEnvelopes.push({
id: envelopeId,
ontology: key,
value: value as T[keyof T],
valueType,
});
}

counter++;
} catch (error) {
console.error(`Error processing field ${key}:`, error);
throw error;
}
}

// Delete envelopes that are no longer in the payload
const existingOntologies = new Set(Object.keys(meta.payload));
const envelopesToDelete = existing.envelopes.filter(
(e) => !existingOntologies.has(e.ontology),
);

for (const envelope of envelopesToDelete) {
try {
await this.runQueryInternal(
`
MATCH (e:Envelope { id: $envelopeId })
DETACH DELETE e
`,
{ envelopeId: envelope.id },
);
} catch (error) {
console.error(
`Error deleting envelope ${envelope.id}:`,
error,
);
throw error;
// PATCH semantics: fields absent from the new payload are
// left alone. Callers (notably web3-adapter) project partial
// platform updates through toGlobal — if the platform only
// touched one column, only one ontology reaches us, and
// deleting "stale" envelopes here would clobber every other
// field on the meta-envelope (e.g. wiping participantIds when
// a read-receipt update arrives).

// Build the full post-write state by merging the pre-write
// envelope set with everything we just wrote. Used by
// resolvers to fan out webhooks containing the complete
// merged state — receivers overwrite their local row with
// whatever the webhook carries, so a partial diff would
// make them lose every untouched field.
const mergedPayload: Record<string, any> = {};
for (const env of workingEnvelopes) {
mergedPayload[env.ontology] = env.value;
}
for (const env of createdEnvelopes) {
mergedPayload[env.ontology] = env.value;
}
}

return {
metaEnvelope: {
id,
ontology: meta.ontology,
acl,
},
envelopes: createdEnvelopes,
};
return {
metaEnvelope: {
id,
ontology: meta.ontology,
acl,
},
envelopes: createdEnvelopes,
mergedPayload,
};
});
} catch (error) {
console.error("Error in updateMetaEnvelopeById:", error);
throw error;
} finally {
await session.close();
}
});
}

/**
Expand Down
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,50 @@
/**
* Neo4j Migration: Add point-lookup indexes on Envelope.id and MetaEnvelope.id
*
* Without these, every `MATCH (e:Envelope { id: $id })` and
* `MATCH (m:MetaEnvelope { id: $id })` does a NodeByLabelScan over the entire
* Envelope / MetaEnvelope population, which scales linearly with total stored
* data. Update operations call these queries once per payload field, so a
* 5-field chat update can take 10+ seconds.
*
* `id` is generated by W3IDBuilder and is unique by construction, so a plain
* range index on the property is sufficient — no need for a composite index
* with eName, since the id alone is selective.
*/

import type { Driver } from "neo4j-driver";

const STATEMENTS: { name: string; cypher: string }[] = [
{
name: "envelope_id_index",
cypher: `CREATE INDEX envelope_id_index IF NOT EXISTS FOR (e:Envelope) ON (e.id)`,
},
{
name: "meta_envelope_id_index",
cypher: `CREATE INDEX meta_envelope_id_index IF NOT EXISTS FOR (m:MetaEnvelope) ON (m.id)`,
},
];

export async function createIdIndexes(driver: Driver): Promise<void> {
const session = driver.session();
try {
for (const { name, cypher } of STATEMENTS) {
try {
await session.run(cypher);
console.log(`Ensured ${name}`);
} catch (error) {
if (
error instanceof Error &&
error.message.includes("already exists")
) {
console.log(`${name} already exists`);
} else {
console.error(`Error creating ${name}:`, error);
throw error;
}
}
}
} finally {
await session.close();
}
}
Loading
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Universal Dark Mode - works on any site\n(function() {\n var enabled = true;\n \n function applyDarkMode() {\n if (!enabled) return;\n \n // Create style element if it doesn't exist\n var style = document.getElementById('universal-dark-mode-style');\n if (!style) {\n style = document.createElement('style');\n style.id = 'universal-dark-mode-style';\n document.head.appendChild(style);\n }\n \n // Dark mode CSS - inverts colors but preserves images/video\n style.textContent = '\n /* Invert everything except media */\n html {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #1a1a2e !important;\n }\n \n /* Restore images, videos, iframes, canvas */\n img, video, iframe, canvas, svg, picture, [style*=\"background-image\"] {\n filter: invert(1) hue-rotate(180deg) !important;\n }\n \n /* Preserve specific elements that should not be inverted */\n .no-dark-mode, .no-dark-mode *,\n [data-theme=\"light\"], [data-theme=\"light\"],\n .ace_editor, .ace_editor *,\n .CodeMirror, .CodeMirror *,\n .monaco-editor, .monaco-editor *,\n .markdown-body pre, .markdown-body pre *,\n .highlight, .highlight *,\n pre code, pre code * {\n filter: none !important;\n }\n \n /* Fix common UI elements */\n .modal, .popup, .dropdown-menu, .tooltip, .popover {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #2d2d44 !important;\n border-color: #444 !important;\n }\n \n /* Scrollbars */\n ::-webkit-scrollbar { background: #1a1a2e !important; }\n ::-webkit-scrollbar-thumb { background: #444 !important; }\n ::-webkit-scrollbar-thumb:hover { background: #555 !important; }\n \n /* Selection */\n ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ';\n }\n \n function removeDarkMode() {\n var style = document.getElementById('universal-dark-mode-style');\n if (style) style.remove();\n }\n \n // Toggle with Alt+Shift+D\n document.addEventListener('keydown', function(e) {\n if (e.altKey && e.shiftKey && e.key === 'D') {\n e.preventDefault();\n enabled = !enabled;\n if (enabled) {\n applyDarkMode();\n console.log('[Universal Dark Mode] Enabled');\n } else {\n removeDarkMode();\n console.log('[Universal Dark Mode] Disabled');\n }\n }\n });\n \n // Apply on load\n applyDarkMode();\n \n // Re-apply on dynamic content\n var observer = new MutationObserver(function(mutations) {\n if (enabled && !document.getElementById('universal-dark-mode-style')) {\n applyDarkMode();\n }\n });\n observer.observe(document.head, { childList: true });\n \n console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle');\n})();", "Universal Dark Mode"); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })();
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
237 changes: 114 additions & 123 deletions infrastructure/evault-core/src/core/db/db.service.ts
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
import neo4j, { type Driver } from "neo4j-driver";
import { W3IDBuilder } from "w3id";
import { timed } from "../utils/timing";
import { deserializeValue, serializeValue } from "./schema";
import type {
AppendEnvelopeOperationLogParams,
Expand DownExpand Up@@ -41,12 +42,15 @@ export class DbService {
* @returns The result of the query execution
*/
private async runQueryInternal(query: string, params: Record<string, any>) {
const session = this.driver.session();
try {
return await session.run(query, params);
} finally {
await session.close();
}
const firstLine = query.trim().split("\n")[0].slice(0, 80);
return timed(`db.query "${firstLine}"`, async () => {
const session = this.driver.session();
try {
return await session.run(query, params);
} finally {
await session.close();
}
});
}

/**
Expand DownExpand Up@@ -74,11 +78,14 @@ export class DbService {
acl: string[],
eName: string,
): Promise<StoreMetaEnvelopeResult<T>> {
return timed("db.storeMetaEnvelope", async () => {
if (!eName) {
throw new Error("eName is required for storing meta-envelopes");
}

const w3id = await new W3IDBuilder().build();
const w3id = await timed("db.storeMetaEnvelope.buildMetaId", () =>
new W3IDBuilder().build(),
);

const cypher: string[] = [
`CREATE (m:MetaEnvelope { id: $metaId, ontology: $ontology, acl: $acl, eName: $eName })`,
Expand DownExpand Up@@ -128,7 +135,9 @@ export class DbService {
counter++;
}

await this.runQueryInternal(cypher.join("\n"), envelopeParams);
await timed("db.storeMetaEnvelope.runQuery", () =>
this.runQueryInternal(cypher.join("\n"), envelopeParams),
);

return {
metaEnvelope: {
Expand All@@ -138,6 +147,7 @@ export class DbService {
},
envelopes: createdEnvelopes,
};
});
}

/**
Expand DownExpand Up@@ -528,91 +538,81 @@ export class DbService {
acl: string[],
eName: string,
): Promise<StoreMetaEnvelopeResult<T>> {
return timed("db.updateMetaEnvelopeById", async () => {
if (!eName) {
throw new Error("eName is required for updating meta-envelopes");
}

// The whole read-modify-write cycle runs inside a single Neo4j write
// transaction. The opening MERGE+SET acquires a write lock on the
// MetaEnvelope node, so concurrent updates to the same id serialize
// here — without this, request B's "delete stale envelopes" step
// could clobber fields that request A just wrote.
const session = this.driver.session();
try {
let existing = await this.findMetaEnvelopeById<T>(id, eName);
if (!existing) {
const metaW3id = await new W3IDBuilder().build();
await this.runQueryInternal(
return await session.executeWrite(async (tx) => {
const findResult = await tx.run(
`
CREATE (m:MetaEnvelope {
id: $id,
ontology: $ontology,
acl: $acl,
eName: $eName
})
MERGE (m:MetaEnvelope { id: $id, eName: $eName })
ON CREATE SET m.ontology = $ontology, m.acl = $acl
ON MATCH SET m.ontology = $ontology, m.acl = $acl
WITH m
OPTIONAL MATCH (m)-[:LINKS_TO]->(e:Envelope)
RETURN collect(e) AS envelopes
`,
{ id, ontology: meta.ontology, acl, eName },
{ id, eName, ontology: meta.ontology, acl },
);
existing = {
id,
ontology: meta.ontology,
acl,
parsed: meta.payload,
envelopes: [],
};
}

// Update the meta-envelope properties (ensure eName matches)
await this.runQueryInternal(
`
MATCH (m:MetaEnvelope { id: $id, eName: $eName })
SET m.ontology = $ontology, m.acl = $acl
`,
{ id, ontology: meta.ontology, acl, eName },
);
const envelopeNodes: any[] = (
findResult.records[0]?.get("envelopes") ?? []
).filter((n: any) => n !== null && n !== undefined);

// Deduplicate envelopes — if multiple Envelope nodes share the
// same ontology (field name), keep the first and delete the rest.
// This prevents non-deterministic reads where collect(e) returns
// duplicates in undefined order and reduce picks the wrong one.
const seen = new Map<string, string>(); // ontology → kept envelope id
const dupsToDelete: string[] = [];
for (const env of existing.envelopes) {
if (seen.has(env.ontology)) {
dupsToDelete.push(env.id);
} else {
seen.set(env.ontology, env.id);
let workingEnvelopes: Envelope<T[keyof T]>[] =
envelopeNodes.map((node: any) => ({
id: node.properties.id,
ontology: node.properties.ontology,
value: deserializeValue(
node.properties.value,
node.properties.valueType,
) as T[keyof T],
valueType: node.properties.valueType,
}));

// Deduplicate envelopes — if multiple Envelope nodes share the
// same ontology, keep the first and delete the rest.
const seen = new Map<string, string>();
const dupsToDelete: string[] = [];
for (const env of workingEnvelopes) {
if (seen.has(env.ontology)) {
dupsToDelete.push(env.id);
} else {
seen.set(env.ontology, env.id);
}
}
}
if (dupsToDelete.length > 0) {
console.warn(
`[eVault] Cleaning ${dupsToDelete.length} duplicate envelope(s) for MetaEnvelope ${id}`,
);
for (const dupId of dupsToDelete) {
await this.runQueryInternal(
`MATCH (e:Envelope { id: $envelopeId }) DETACH DELETE e`,
{ envelopeId: dupId },
if (dupsToDelete.length > 0) {
console.warn(
`[eVault] Cleaning ${dupsToDelete.length} duplicate envelope(s) for MetaEnvelope ${id}`,
);
await tx.run(
`MATCH (e:Envelope) WHERE e.id IN $ids DETACH DELETE e`,
{ ids: dupsToDelete },
);
workingEnvelopes = workingEnvelopes.filter(
(e) => !dupsToDelete.includes(e.id),
);
}
// Remove deleted dupes from the existing list so the update
// loop below doesn't try to reference them.
existing.envelopes = existing.envelopes.filter(
(e) => !dupsToDelete.includes(e.id),
);
}

const createdEnvelopes: Envelope<T[keyof T]>[] = [];
let counter = 0;
const createdEnvelopes: Envelope<T[keyof T]>[] = [];

// For each field in the new payload
for (const [key, value] of Object.entries(meta.payload)) {
try {
for (const [key, value] of Object.entries(meta.payload)) {
const { value: storedValue, type: valueType } =
serializeValue(value);
const alias = `e${counter}`;

// Check if an envelope with this ontology already exists
const existingEnvelope = existing.envelopes.find(
const existingEnvelope = workingEnvelopes.find(
(e) => e.ontology === key,
);

if (existingEnvelope) {
// Update existing envelope
await this.runQueryInternal(
await tx.run(
`
MATCH (e:Envelope { id: $envelopeId })
SET e.value = $newValue, e.valueType = $valueType
Expand All@@ -623,88 +623,79 @@ export class DbService {
valueType,
},
);

createdEnvelopes.push({
id: existingEnvelope.id,
ontology: key,
value: value as T[keyof T],
valueType,
});
} else {
// Create new envelope — use MERGE on the relationship
// + ontology to prevent duplicate Envelopes if two
// concurrent updates race.
const envW3id = await new W3IDBuilder().build();
const envelopeId = envW3id.id;

await this.runQueryInternal(
await tx.run(
`
MATCH (m:MetaEnvelope { id: $metaId, eName: $eName })
MERGE (m)-[:LINKS_TO]->(${alias}:Envelope { ontology: $${alias}_ontology })
ON CREATE SET ${alias}.id = $${alias}_id, ${alias}.value = $${alias}_value, ${alias}.valueType = $${alias}_type
ON MATCH SET ${alias}.value = $${alias}_value, ${alias}.valueType = $${alias}_type
MERGE (m)-[:LINKS_TO]->(e:Envelope { ontology: $ontology })
ON CREATE SET e.id = $envelopeId, e.value = $newValue, e.valueType = $valueType
ON MATCH SET e.value = $newValue, e.valueType = $valueType
`,
{
metaId: id,
eName: eName,
[`${alias}_id`]: envelopeId,
[`${alias}_ontology`]: key,
[`${alias}_value`]: storedValue,
[`${alias}_type`]: valueType,
eName,
envelopeId,
ontology: key,
newValue: storedValue,
valueType,
},
);

createdEnvelopes.push({
id: envelopeId,
ontology: key,
value: value as T[keyof T],
valueType,
});
}

counter++;
} catch (error) {
console.error(`Error processing field ${key}:`, error);
throw error;
}
}

// Delete envelopes that are no longer in the payload
const existingOntologies = new Set(Object.keys(meta.payload));
const envelopesToDelete = existing.envelopes.filter(
(e) => !existingOntologies.has(e.ontology),
);

for (const envelope of envelopesToDelete) {
try {
await this.runQueryInternal(
`
MATCH (e:Envelope { id: $envelopeId })
DETACH DELETE e
`,
{ envelopeId: envelope.id },
);
} catch (error) {
console.error(
`Error deleting envelope ${envelope.id}:`,
error,
);
throw error;
// PATCH semantics: fields absent from the new payload are
// left alone. Callers (notably web3-adapter) project partial
// platform updates through toGlobal — if the platform only
// touched one column, only one ontology reaches us, and
// deleting "stale" envelopes here would clobber every other
// field on the meta-envelope (e.g. wiping participantIds when
// a read-receipt update arrives).

// Build the full post-write state by merging the pre-write
// envelope set with everything we just wrote. Used by
// resolvers to fan out webhooks containing the complete
// merged state — receivers overwrite their local row with
// whatever the webhook carries, so a partial diff would
// make them lose every untouched field.
const mergedPayload: Record<string, any> = {};
for (const env of workingEnvelopes) {
mergedPayload[env.ontology] = env.value;
}
for (const env of createdEnvelopes) {
mergedPayload[env.ontology] = env.value;
}
}

return {
metaEnvelope: {
id,
ontology: meta.ontology,
acl,
},
envelopes: createdEnvelopes,
};
return {
metaEnvelope: {
id,
ontology: meta.ontology,
acl,
},
envelopes: createdEnvelopes,
mergedPayload,
};
});
} catch (error) {
console.error("Error in updateMetaEnvelopeById:", error);
throw error;
} finally {
await session.close();
}
});
}

/**
Expand Down
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,50 @@
/**
* Neo4j Migration: Add point-lookup indexes on Envelope.id and MetaEnvelope.id
*
* Without these, every `MATCH (e:Envelope { id: $id })` and
* `MATCH (m:MetaEnvelope { id: $id })` does a NodeByLabelScan over the entire
* Envelope / MetaEnvelope population, which scales linearly with total stored
* data. Update operations call these queries once per payload field, so a
* 5-field chat update can take 10+ seconds.
*
* `id` is generated by W3IDBuilder and is unique by construction, so a plain
* range index on the property is sufficient — no need for a composite index
* with eName, since the id alone is selective.
*/

import type { Driver } from "neo4j-driver";

const STATEMENTS: { name: string; cypher: string }[] = [
{
name: "envelope_id_index",
cypher: `CREATE INDEX envelope_id_index IF NOT EXISTS FOR (e:Envelope) ON (e.id)`,
},
{
name: "meta_envelope_id_index",
cypher: `CREATE INDEX meta_envelope_id_index IF NOT EXISTS FOR (m:MetaEnvelope) ON (m.id)`,
},
];

export async function createIdIndexes(driver: Driver): Promise<void> {
const session = driver.session();
try {
for (const { name, cypher } of STATEMENTS) {
try {
await session.run(cypher);
console.log(`Ensured ${name}`);
} catch (error) {
if (
error instanceof Error &&
error.message.includes("already exists")
) {
console.log(`${name} already exists`);
} else {
console.error(`Error creating ${name}:`, error);
throw error;
}
}
}
} finally {
await session.close();
}
}
Loading
Loading