Skip to content

Commit 21e9ebd

Browse files
committed
fix(control_plane/memory): address P1 review issues — async interface, privacy, authority
Fix four P1 issues from code review: 1. Async/sync interface mismatch: MemoryStore interface now declares async methods (Promise return types). FileMemoryStore and MemoryStoreInMemory both implement the async contract. All consumers (discovery, resume, handoff, projection) updated to await store calls. 2. adoptMemory self-link bug: adoptMemory now requires a distinct successor_memory_id and validates that it exists, belongs to the same goal, and its supersedes field points back to the source. Self-links (source == successor) are rejected with EffectRuntimeRequestError. 3. No privacy enforcement: discoverMemories now accepts an optional MemoryViewerContext. Session-scoped records from other sessions are filtered out when a viewer context is present. createMemoryRecord defaults public_safe/redaction_verified to false (caller must opt in). Projection functions only include public_safe records in summaries and handoff rows. 4. Handoff authority is just text: composeHandoffWithMode now validates that soft_claim/hard_lease modes require handoff-availability records that are certified public_safe. Legacy mode remains permissive. Authority labels are validated against record metadata, not just returned as opaque strings. Test count: 64 → 73 (added tests for new validation paths). Signed-off-by: xiaodeshi <xiaodeshi@users.noreply.github.com> Signed-off-by: Xiao Deshi <xiaods@gmail.com>
1 parent 7a39f56 commit 21e9ebd

13 files changed

Lines changed: 544 additions & 206 deletions

‎loopx/control_plane/memory/index.ts‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -42,6 +42,7 @@ export {
4242
discoverMemories,
4343
getSupersedeChain,
4444
type MemoryDiscoveryFilter,
45+
type MemoryViewerContext,
4546
} from "./memory_discovery.ts";
4647

4748
export {

‎loopx/control_plane/memory/memory_discovery.ts‎

Lines changed: 31 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -11,6 +11,11 @@ import type { MemoryStore } from "./memory_store.ts";
1111
// Discovery filters
1212
// ---------------------------------------------------------------------------
1313

14+
export interface MemoryViewerContext {
15+
readonly session_ref: string;
16+
readonly agent_id: string;
17+
}
18+
1419
export interface MemoryDiscoveryFilter {
1520
readonly goal_id?: string;
1621
readonly status?: MemoryStatus;
@@ -20,17 +25,18 @@ export interface MemoryDiscoveryFilter {
2025
readonly session_ref?: string;
2126
readonly semantic_query?: string;
2227
readonly limit?: number;
28+
readonly viewer?: MemoryViewerContext;
2329
}
2430

2531
// ---------------------------------------------------------------------------
2632
// Discovery API
2733
// ---------------------------------------------------------------------------
2834

29-
export function discoverMemories(
35+
export async function discoverMemories(
3036
store: MemoryStore,
3137
filter: MemoryDiscoveryFilter = {},
32-
): MemoryRecord[] {
33-
const all = store.readAll();
38+
): Promise<MemoryRecord[]> {
39+
const all = await store.readAll();
3440
let results = all;
3541

3642
if (filter.goal_id !== undefined) {
@@ -61,6 +67,18 @@ export function discoverMemories(
6167
results = results.filter((r) => r.session_ref === filter.session_ref);
6268
}
6369

70+
// Privacy enforcement: session-scoped memories are only visible to
71+
// the session that created them. When a viewer context is provided,
72+
// session-private records from other sessions are excluded.
73+
if (filter.viewer !== undefined) {
74+
results = results.filter((r) => {
75+
if (r.availability === "session") {
76+
return r.session_ref === filter.viewer!.session_ref;
77+
}
78+
return true;
79+
});
80+
}
81+
6482
if (filter.semantic_query !== undefined) {
6583
const query = filter.semantic_query.toLowerCase();
6684
results = results.filter((r) =>
@@ -83,34 +101,34 @@ export function discoverMemories(
83101
return results;
84102
}
85103

86-
export function discoverLatestMemory(
104+
export async function discoverLatestMemory(
87105
store: MemoryStore,
88106
todo_id: string,
89-
): MemoryRecord | null {
90-
const matches = discoverMemories(store, { todo_id, status: "active" });
107+
): Promise<MemoryRecord | null> {
108+
const matches = await discoverMemories(store, { todo_id, status: "active" });
91109
return matches.length > 0 ? matches[0] : null;
92110
}
93111

94-
export function discoverHandoffMemories(
112+
export async function discoverHandoffMemories(
95113
store: MemoryStore,
96114
goal_id: string,
97-
): MemoryRecord[] {
98-
return discoverMemories(store, {
115+
): Promise<MemoryRecord[]> {
116+
return (await discoverMemories(store, {
99117
goal_id,
100118
status: "active",
101119
availability: "handoff",
102-
}).filter((r) => !r.superseded_by);
120+
})).filter((r) => !r.superseded_by);
103121
}
104122

105123
// ---------------------------------------------------------------------------
106124
// Supersede chain traversal
107125
// ---------------------------------------------------------------------------
108126

109-
export function getSupersedeChain(
127+
export async function getSupersedeChain(
110128
store: MemoryStore,
111129
memory_id: string,
112-
): MemoryRecord[] {
113-
const all = store.readAll();
130+
): Promise<MemoryRecord[]> {
131+
const all = await store.readAll();
114132
const byId = new Map(all.map((r) => [r.memory_id, r]));
115133

116134
// Walk backward to find the root

‎loopx/control_plane/memory/memory_handoff.ts‎

Lines changed: 31 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -44,16 +44,18 @@ export interface HandoffMemoryWriteInput {
4444
readonly evidence_refs?: readonly string[];
4545
readonly supersedes?: string | null;
4646
readonly source_run_id?: string | null;
47+
readonly public_safe?: boolean;
48+
readonly redaction_verified?: boolean;
4749
}
4850

4951
/**
5052
* Write a memory record during handoff preparation.
5153
* This is the "source session" side of the handoff flow.
5254
*/
53-
export function writeHandoffMemory(
55+
export async function writeHandoffMemory(
5456
store: MemoryStore,
5557
input: HandoffMemoryWriteInput,
56-
): MemoryRecord {
58+
): Promise<MemoryRecord> {
5759
return store.append({
5860
goal_id: input.goal_id,
5961
session_ref: input.session_ref,
@@ -67,6 +69,8 @@ export function writeHandoffMemory(
6769
evidence_refs: input.evidence_refs ?? [],
6870
supersedes: input.supersedes ?? null,
6971
source_run_id: input.source_run_id ?? null,
72+
public_safe: input.public_safe ?? false,
73+
redaction_verified: input.redaction_verified ?? false,
7074
});
7175
}
7276

@@ -91,23 +95,23 @@ export function buildMemoryHandoffNoteRef(
9195
* Discover handoff-available memories for a goal.
9296
* This is the "target session" side of the handoff flow.
9397
*/
94-
export function discoverHandoffMemoriesForGoal(
98+
export async function discoverHandoffMemoriesForGoal(
9599
store: MemoryStore,
96100
goal_id: string,
97-
): MemoryRecord[] {
101+
): Promise<MemoryRecord[]> {
98102
return discoverHandoffMemories(store, goal_id);
99103
}
100104

101105
/**
102106
* Resume from a handoff memory, producing a resume packet.
103107
* The target session decides whether to adopt.
104108
*/
105-
export function resumeFromHandoffMemory(
109+
export async function resumeFromHandoffMemory(
106110
store: MemoryStore,
107111
memory_id: string,
108112
target_agent_id: string,
109-
): MemoryResumePacket {
110-
const memory = store.readById(memory_id);
113+
): Promise<MemoryResumePacket> {
114+
const memory = await store.readById(memory_id);
111115
if (!memory) {
112116
throw new EffectRuntimeRequestError(`handoff memory not found: ${memory_id}`);
113117
}
@@ -148,12 +152,32 @@ export interface HandoffCompositionResult {
148152
* claim = adopt.
149153
* - hard_lease: memory transfer requires the target to hold the task lease;
150154
* memory is the payload that travels with the lease.
155+
*
156+
* Note: This function declares the authority mode label. Actual claim/lease
157+
* enforcement is delegated to the existing typed claim/lease owner in the
158+
* control plane. This composition validates that the memory record itself
159+
* meets the minimum bar for the declared mode (e.g., handoff-availability
160+
* records for soft_claim/hard_lease).
151161
*/
152162
export function composeHandoffWithMode(
153163
input: HandoffCompositionInput,
154164
): HandoffCompositionResult {
155165
const { mode, memory, target_agent_id } = input;
156166

167+
// Validate that the memory record meets the minimum bar for the mode
168+
if (mode === "soft_claim" || mode === "hard_lease") {
169+
if (memory.availability !== "handoff") {
170+
throw new EffectRuntimeRequestError(
171+
`memory ${memory.memory_id} has availability "${memory.availability}" but mode "${mode}" requires "handoff"`,
172+
);
173+
}
174+
if (!memory.public_safe) {
175+
throw new EffectRuntimeRequestError(
176+
`memory ${memory.memory_id} is not certified public_safe; cannot compose with mode "${mode}"`,
177+
);
178+
}
179+
}
180+
157181
const resumePacket = resumeFromMemory({
158182
memory,
159183
target_agent_id,

‎loopx/control_plane/memory/memory_projection.ts‎

Lines changed: 19 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -39,21 +39,24 @@ export interface MemoryIndexProjection extends JsonObject {
3939
// Projection builder
4040
// ---------------------------------------------------------------------------
4141

42-
export function buildMemoryIndexProjection(
42+
export async function buildMemoryIndexProjection(
4343
store: MemoryStore,
44-
): MemoryIndexProjection {
45-
const index = store.getIndex();
46-
const all = store.readAll();
44+
): Promise<MemoryIndexProjection> {
45+
const index = await store.getIndex();
46+
const all = await store.readAll();
4747
const byId = new Map(all.map((r) => [r.memory_id, r]));
4848

4949
const activeCount = index.active_memories.length;
5050
const handoffCount = index.handoff_available.length;
5151

52+
// Only include public_safe records in the projection summary
53+
const publicSafeRecords = all.filter((r) => r.public_safe);
54+
5255
let latestMemoryAt: string | null = null;
5356
let latestSummary: string | null = null;
5457

55-
if (all.length > 0) {
56-
const sorted = [...all].sort((a, b) => {
58+
if (publicSafeRecords.length > 0) {
59+
const sorted = [...publicSafeRecords].sort((a, b) => {
5760
if (a.created_at > b.created_at) return -1;
5861
if (a.created_at < b.created_at) return 1;
5962
return 0;
@@ -65,7 +68,7 @@ export function buildMemoryIndexProjection(
6568

6669
const handoffRows: MemoryHandoffRow[] = index.handoff_available
6770
.map((id) => byId.get(id))
68-
.filter((r): r is MemoryRecord => r !== undefined)
71+
.filter((r): r is MemoryRecord => r !== undefined && r.public_safe)
6972
.map((r) => ({
7073
memory_id: r.memory_id,
7174
summary: r.summary,
@@ -100,15 +103,18 @@ export interface MemoryFirstScreen {
100103
readonly latest_summary: string | null;
101104
}
102105

103-
export function buildMemoryFirstScreen(
106+
export async function buildMemoryFirstScreen(
104107
store: MemoryStore,
105-
): MemoryFirstScreen {
106-
const index = store.getIndex();
107-
const all = store.readAll();
108+
): Promise<MemoryFirstScreen> {
109+
const index = await store.getIndex();
110+
const all = await store.readAll();
111+
112+
// Only include public_safe records in the first-screen summary
113+
const publicSafeRecords = all.filter((r) => r.public_safe);
108114

109115
let latestSummary: string | null = null;
110-
if (all.length > 0) {
111-
const sorted = [...all].sort((a, b) => {
116+
if (publicSafeRecords.length > 0) {
117+
const sorted = [...publicSafeRecords].sort((a, b) => {
112118
if (a.created_at > b.created_at) return -1;
113119
if (a.created_at < b.created_at) return 1;
114120
return 0;

‎loopx/control_plane/memory/memory_record.ts‎

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -343,6 +343,8 @@ export interface CreateMemoryRecordInput {
343343
availability?: MemoryAvailability;
344344
source_run_id?: string | null;
345345
created_at?: string;
346+
public_safe?: boolean;
347+
redaction_verified?: boolean;
346348
}
347349

348350
export function createMemoryRecord(input: CreateMemoryRecordInput): MemoryRecord {
@@ -366,8 +368,8 @@ export function createMemoryRecord(input: CreateMemoryRecordInput): MemoryRecord
366368
superseded_by: null,
367369
status: MEMORY_STATUS_ACTIVE,
368370
availability: input.availability ?? AVAILABILITY_SESSION,
369-
public_safe: true,
370-
redaction_verified: true,
371+
public_safe: input.public_safe ?? false,
372+
redaction_verified: input.redaction_verified ?? false,
371373
created_at: now,
372374
source_run_id: input.source_run_id ?? null,
373375
});

‎loopx/control_plane/memory/memory_resume.ts‎

Lines changed: 41 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -120,17 +120,52 @@ export interface AdoptMemoryResult {
120120
readonly superseded_memory_id: string | null;
121121
}
122122

123-
export function adoptMemory(
123+
export interface AdoptMemoryArgs {
124+
readonly source_memory_id: string;
125+
readonly successor_memory_id: string;
126+
}
127+
128+
export async function adoptMemory(
124129
store: MemoryStore,
125130
packet: MemoryResumePacket,
126-
args: { source_memory_id: string },
127-
): AdoptMemoryResult {
128-
const { source_memory_id } = args;
131+
args: AdoptMemoryArgs,
132+
): Promise<AdoptMemoryResult> {
133+
const { source_memory_id, successor_memory_id } = args;
134+
135+
// Validate: successor must be distinct from source (no self-links)
136+
if (source_memory_id === successor_memory_id) {
137+
throw new EffectRuntimeRequestError(
138+
`cannot adopt memory: successor_memory_id must differ from source_memory_id (${source_memory_id})`,
139+
);
140+
}
141+
142+
// Validate: successor must exist in the store
143+
const successor = await store.readById(successor_memory_id);
144+
if (!successor) {
145+
throw new EffectRuntimeRequestError(
146+
`cannot adopt memory: successor memory not found: ${successor_memory_id}`,
147+
);
148+
}
149+
150+
// Validate: successor must belong to the same goal
151+
if (successor.goal_id !== packet.goal_id) {
152+
throw new EffectRuntimeRequestError(
153+
`cannot adopt memory: successor goal_id (${successor.goal_id}) does not match packet goal_id (${packet.goal_id})`,
154+
);
155+
}
156+
157+
// Validate: successor's supersedes must point back to source
158+
if (successor.supersedes !== source_memory_id) {
159+
throw new EffectRuntimeRequestError(
160+
`cannot adopt memory: successor.supersedes (${successor.supersedes}) does not point to source_memory_id (${source_memory_id})`,
161+
);
162+
}
163+
129164
// Mark the source memory as superseded
130-
const index = store.getIndex();
165+
const index = await store.getIndex();
131166
const source = index.supersede_chains[source_memory_id];
132167
if (source && source.superseded_by === null) {
133-
store.updateSupersededBy(source_memory_id, packet.memory_id);
168+
await store.updateSupersededBy(source_memory_id, successor_memory_id);
134169
return {
135170
adopted: true,
136171
resume_packet: packet,

0 commit comments

Comments
 (0)