Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
77 changes: 77 additions & 0 deletions packages/convex/convex/models/facts/facts.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ import { afterEach, beforeEach, describe, expect, test } from "vitest";
import { api, internal } from "../../_generated/api";
import schema from "../../schema";
import { modules } from "../../test.setup";
import { listFacts, searchFacts } from "./model";

const issuer = "https://brain.example.test";

Expand Down Expand Up @@ -321,4 +322,80 @@ describe("structured durable facts", () => {
"Jordan — home city: Seattle.",
]);
});

test("fills fact result limits after lifecycle filtering", async () => {
const t = convexTest(schema, modules);
const userId = await t.run((ctx) => ctx.db.insert("users", {}));

await t.run(async (ctx) => {
const subjectEntityId = await ctx.db.insert("entities", {
userId,
key: "person:pagination-test",
kind: "person",
canonicalName: "Pagination Test",
normalizedName: "pagination test",
aliases: [],
normalizedAliases: [],
});
const insertFact = (
index: number,
status: "current" | "retracted",
validFrom?: number,
) =>
ctx.db.insert("facts", {
userId,
subjectEntityId,
predicate: "school",
value: { type: "text", value: `School ${index}` },
statement: `Pagination Test — school: School ${index}.`,
searchText: `pagination sentinel school School ${index}`,
sourceType: "user_stated",
confidence: 1,
isCore: true,
validFrom,
status,
});

// Retrievable rows are deliberately older. The scheduled rows exhaust
// the old current-only `take(limit * 5)` window, while the still-newer
// retractions exhaust its historical window.
for (let index = 0; index < 10; index += 1) {
await insertFact(index, "current");
}
for (let index = 10; index < 70; index += 1) {
await insertFact(index, "current", Date.now() + 86_400_000);
}
for (let index = 70; index < 130; index += 1) {
await insertFact(index, "retracted");
}
});

const recent = await t.run((ctx) => listFacts(ctx, userId, { limit: 10 }));
const core = await t.run((ctx) =>
listFacts(ctx, userId, { limit: 10, coreOnly: true }),
);
const historical = await t.run((ctx) =>
listFacts(ctx, userId, { limit: 10, includeHistorical: true }),
);
const search = await t.run((ctx) =>
searchFacts(ctx, userId, "pagination sentinel school", { limit: 10 }),
);
const historicalSearch = await t.run((ctx) =>
searchFacts(ctx, userId, "pagination sentinel school", {
limit: 10,
includeHistorical: true,
}),
);

for (const results of [
recent,
core,
historical,
search,
historicalSearch,
]) {
expect(results).toHaveLength(10);
expect(results.every((fact) => fact.status !== "retracted")).toBe(true);
}
});
});
100 changes: 76 additions & 24 deletions packages/convex/convex/models/facts/model.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
import type { Infer } from "convex/values";
import type { Expression, FilterBuilder, NamedTableInfo } from "convex/server";

import type { Doc, Id } from "../../_generated/dataModel";
import type { DataModel, Doc, Id } from "../../_generated/dataModel";
import type { MutationCtx, QueryCtx } from "../../_generated/server";
import {
assertValidMemoryValidity,
Expand Down Expand Up @@ -304,6 +305,30 @@ export function isFactRetrievable(
);
}

type FactFilterBuilder = FilterBuilder<NamedTableInfo<DataModel, "facts">>;

/**
* Express current business-time validity inside the Convex query so `take`
* counts retrievable rows rather than candidates later discarded in memory.
*/
function currentFactValidityFilter(
q: FactFilterBuilder,
activeAt: number,
): Expression<boolean> {
const validFrom = q.field("validFrom");
const validTo = q.field("validTo");
return q.and(
q.or(
q.eq(validFrom, undefined),
q.lte(validFrom as Expression<number>, activeAt),
),
q.or(
q.eq(validTo, undefined),
q.gt(validTo as Expression<number>, activeAt),
),
);
}

export type RememberFactArgs = {
subject: EntitySelector;
predicate: string;
Expand Down Expand Up @@ -532,22 +557,43 @@ export async function listFacts(
throw new Error("Fact limit must be a positive integer");
}
const limit = Math.min(requested, MAX_FACT_SEARCH_LIMIT);
const facts = options.coreOnly
? await ctx.db
.query("facts")
.withIndex("by_userId_and_isCore", (q) =>
q.eq("userId", userId).eq("isCore", true),
)
.order("desc")
.take(limit * 5)
: await ctx.db
.query("facts")
.withIndex("by_userId", (q) => q.eq("userId", userId))
.order("desc")
.take(limit * 5);
const selected = facts
.filter((fact) => isFactRetrievable(fact, options.includeHistorical))
.slice(0, limit);
const activeAt = Date.now();
let selected: Doc<"facts">[];
if (options.includeHistorical) {
selected = options.coreOnly
? await ctx.db
.query("facts")
.withIndex("by_userId_and_isCore", (q) =>
q.eq("userId", userId).eq("isCore", true),
)
.order("desc")
.filter((q) => q.neq(q.field("status"), "retracted"))
.take(limit)
: await ctx.db
.query("facts")
.withIndex("by_userId", (q) => q.eq("userId", userId))
.order("desc")
.filter((q) => q.neq(q.field("status"), "retracted"))
.take(limit);
} else {
selected = options.coreOnly
? await ctx.db
.query("facts")
.withIndex("by_userId_isCore_status", (q) =>
q.eq("userId", userId).eq("isCore", true).eq("status", "current"),
)
.order("desc")
.filter((q) => currentFactValidityFilter(q, activeAt))
.take(limit)
: await ctx.db
.query("facts")
.withIndex("by_userId_and_status", (q) =>
q.eq("userId", userId).eq("status", "current"),
)
.order("desc")
.filter((q) => currentFactValidityFilter(q, activeAt))
.take(limit);
}
return await Promise.all(selected.map((fact) => hydrateFact(ctx, fact)));
}

Expand All @@ -563,14 +609,20 @@ export async function searchFacts(
throw new Error("Fact search limit must be a positive integer");
}
const limit = Math.min(requested, MAX_FACT_SEARCH_LIMIT);
const hits = await ctx.db
const activeAt = Date.now();
const selected = await ctx.db
.query("facts")
.withSearchIndex("by_searchText", (q) =>
q.search("searchText", cleanedQuery).eq("userId", userId),
.withSearchIndex("by_searchText", (q) => {
const search = q.search("searchText", cleanedQuery).eq("userId", userId);
return options.includeHistorical
? search
: search.eq("status", "current");
})
.filter((q) =>
options.includeHistorical
? q.neq(q.field("status"), "retracted")
: currentFactValidityFilter(q, activeAt),
)
.take(limit * 5);
const selected = hits
.filter((fact) => isFactRetrievable(fact, options.includeHistorical))
.slice(0, limit);
.take(limit);
return await Promise.all(selected.map((fact) => hydrateFact(ctx, fact)));
}
12 changes: 12 additions & 0 deletions packages/convex/convex/models/thoughts/coreMemory.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -89,6 +89,18 @@ describe("core memories", () => {
isCore: true,
memoryStatus: "current",
});
// These are newer than the retrievable core set. The previous fixed
// 250-candidate window returned nothing once enough history accumulated.
for (let index = 0; index < 260; index += 1) {
await ctx.db.insert("thoughts", {
userId: ownerId,
content: `Owner retracted core memory ${index}`,
embedding,
metadata,
isCore: true,
memoryStatus: "retracted",
});
}
});

const owner = t.withIdentity({ issuer, subject: ownerId });
Expand Down
56 changes: 56 additions & 0 deletions packages/convex/convex/models/thoughts/memoryTransition.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -65,6 +65,62 @@ describe("temporal memory transitions", () => {
]);
});

test("fills memory result limits after lifecycle filtering", async () => {
const t = convexTest(schema, modules);
const userId = await t.run((ctx) => ctx.db.insert("users", {}));

await t.run(async (ctx) => {
const insertMemory = (
index: number,
memoryStatus: "current" | "retracted" | undefined,
validFrom?: number,
) =>
ctx.db.insert("thoughts", {
userId,
content: `Pagination sentinel memory ${index}`,
embedding,
metadata,
memoryStatus,
validFrom,
});

for (let index = 0; index < 10; index += 1) {
await insertMemory(index, undefined);
}
for (let index = 10; index < 70; index += 1) {
await insertMemory(index, "current", Date.now() + 86_400_000);
}
for (let index = 70; index < 130; index += 1) {
await insertMemory(index, "retracted");
}
});

const current = await t.run((ctx) => _listByUser(ctx, userId, 10));
const historical = await t.run((ctx) => _listByUser(ctx, userId, 10, true));
const search = await t.query(
internal.models.thoughts.private.searchByText,
{
userId,
query: "pagination sentinel memory",
limit: 10,
activeAt: Date.now(),
},
);

expect(current).toHaveLength(10);
expect(current.every((memory) => memory.memoryStatus === undefined)).toBe(
true,
);
expect(historical).toHaveLength(10);
expect(
historical.every((memory) => memory.memoryStatus !== "retracted"),
).toBe(true);
expect(search).toHaveLength(10);
expect(search.every((memory) => memory.memoryStatus === undefined)).toBe(
true,
);
});

test("atomically preserves and links a superseded memory", async () => {
const t = convexTest(schema, modules);
const userId = await t.run((ctx) => ctx.db.insert("users", {}));
Expand Down
Loading
Loading