Files
memind/user-model-service/signals.mjs
T
john 212ff3ff80 Add User Model Service and Temporal Recall for MeMind V0.1.
Introduce UMS ingest/snapshot pipeline, Context Planner with multi-source recall, runtime context injection, canonical user mapping, and session snapshot loading on auth/me.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-09-03 23:25:18 +08:00

155 lines
5.2 KiB
JavaScript

import crypto from 'node:crypto';
const STOPWORDS = new Set([
'的', '了', '在', '是', '我', '你', '他', '她', '它', '我们', '你们', '他们',
'这', '那', '有', '和', '与', '或', '就', '也', '都', '还', '要', '会', '能',
'一个', '什么', '怎么', '可以', '没有', '不是', '如果', '因为', '所以', '但是',
'然后', '已经', '还是', '自己', '现在', '今天', '明天', '这个', '那个', '一下',
'the', 'and', 'for', 'with', 'this', 'that', 'from', 'have', 'are', 'was', 'not',
]);
function newId() {
return crypto.randomUUID();
}
function nowMs() {
return Date.now();
}
function extractTerms(text) {
const terms = [];
const cjk = String(text).match(/[\u4e00-\u9fff]{2,12}/g) ?? [];
terms.push(...cjk);
const en = String(text).match(/[a-zA-Z][a-zA-Z0-9]{2,}/g) ?? [];
terms.push(...en.map((w) => w.toLowerCase()));
return terms.filter((t) => !STOPWORDS.has(t));
}
function dayStartUtc(date) {
const d = new Date(date);
return new Date(Date.UTC(d.getUTCFullYear(), d.getUTCMonth(), d.getUTCDate()));
}
function dayEndUtc(date) {
const start = dayStartUtc(date);
return new Date(start.getTime() + 24 * 60 * 60 * 1000 - 1);
}
function hashSignal(userId, signalType, dimensionKey, windowStart, windowEnd, valueJson) {
const raw = `${userId}|${signalType}|${dimensionKey}|${windowStart}|${windowEnd}|${JSON.stringify(valueJson)}`;
return crypto.createHash('sha256').update(raw).digest('hex');
}
/**
* @param {object} envelope Evidence Envelope v1
* @returns {Array<object>} signal drafts
*/
export function extractSignalsFromEnvelope(envelope) {
if (envelope.evidence_type !== 'expression_segment') return [];
const payload = envelope.payload ?? {};
const text = payload.text ?? '';
const occurredAt = envelope.occurred_at;
const windowStart = dayStartUtc(occurredAt).toISOString().slice(0, 23).replace('T', ' ');
const windowEnd = dayEndUtc(occurredAt).toISOString().slice(0, 23).replace('T', ' ');
const signals = [];
for (const term of extractTerms(text)) {
const key = /^[a-z]/.test(term) ? `term:${term}` : `term:${term}`;
signals.push({
signal_type: 'term_frequency',
dimension_key: key,
window_start: windowStart,
window_end: windowEnd,
value_json: { count: 1, chars: text.length },
evidence_ids: [envelope.evidence_id],
});
}
const appId = payload.context?.app_bundle_id;
if (appId) {
signals.push({
signal_type: 'app_usage',
dimension_key: `app:${appId}`,
window_start: windowStart,
window_end: windowEnd,
value_json: { count: 1, app_name: payload.context?.app ?? null },
evidence_ids: [envelope.evidence_id],
});
}
const hour = new Date(occurredAt).getUTCHours();
signals.push({
signal_type: 'segment_count',
dimension_key: `window:daily:${windowStart.slice(0, 10)}`,
window_start: windowStart,
window_end: windowEnd,
value_json: { segments: 1, hour },
evidence_ids: [envelope.evidence_id],
});
return signals;
}
export async function upsertSignals(pool, userId, signalDrafts) {
let touched = 0;
const ts = nowMs();
for (const draft of signalDrafts) {
const contentHash = hashSignal(
userId,
draft.signal_type,
draft.dimension_key,
draft.window_start,
draft.window_end,
draft.value_json,
);
const [existing] = await pool.query(
`SELECT signal_id, value_json, evidence_ids FROM um_signals
WHERE user_id = ? AND signal_type = ? AND dimension_key = ?
AND window_start = ? AND window_end = ?
LIMIT 1`,
[userId, draft.signal_type, draft.dimension_key, draft.window_start, draft.window_end],
);
if (existing[0]) {
const prev = existing[0];
const prevValue = typeof prev.value_json === 'string' ? JSON.parse(prev.value_json) : prev.value_json;
const prevEvidence =
typeof prev.evidence_ids === 'string' ? JSON.parse(prev.evidence_ids) : prev.evidence_ids;
const mergedEvidence = [...new Set([...(prevEvidence ?? []), ...draft.evidence_ids])];
const mergedValue = {
...prevValue,
count: Number(prevValue.count ?? 0) + Number(draft.value_json.count ?? 1),
segments: Number(prevValue.segments ?? 0) + Number(draft.value_json.segments ?? 0),
chars: Number(prevValue.chars ?? 0) + Number(draft.value_json.chars ?? 0),
};
await pool.query(
`UPDATE um_signals SET value_json = ?, evidence_ids = ?, computed_at = ?, content_hash = ?
WHERE signal_id = ?`,
[JSON.stringify(mergedValue), JSON.stringify(mergedEvidence), ts, contentHash, prev.signal_id],
);
} else {
await pool.query(
`INSERT INTO um_signals
(signal_id, user_id, signal_type, dimension_key, window_start, window_end,
value_json, evidence_ids, computed_at, content_hash)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`,
[
newId(),
userId,
draft.signal_type,
draft.dimension_key,
draft.window_start,
draft.window_end,
JSON.stringify(draft.value_json),
JSON.stringify(draft.evidence_ids),
ts,
contentHash,
],
);
}
touched += 1;
}
return touched;
}
export { extractTerms, STOPWORDS };