diff --git a/apps/memos-local-plugin/adapters/hermes/README.md b/apps/memos-local-plugin/adapters/hermes/README.md index 9f05755c9..f1b20b896 100644 --- a/apps/memos-local-plugin/adapters/hermes/README.md +++ b/apps/memos-local-plugin/adapters/hermes/README.md @@ -25,6 +25,36 @@ keepalive, reconnect generation, and host callback dispatch. All algorithm logic (L1/L2/L3, skills, retrieval, feedback, decision repair) remains in the shared TypeScript core. +## Desktop and custom Unix installations + +Use the backend source directory and Python environment used by the desktop +app, which may differ from the CLI installation. The Unix installer validates +`hermes_cli` and `plugins.memory.load_memory_provider` before stopping Hermes +or deploying the plugin. It checks literal Python/Bash launchers and the default +backend's `venv` and `.venv` directories. + +If auto-detection fails, use the updated `install.sh` with explicit paths: + +```bash +HERMES_INSTALL_DIR="/actual/path/to/hermes-agent" \ +HERMES_PYTHON="/actual/path/to/hermes-agent/venv/bin/python" \ +bash install.sh --agent hermes +``` + +`HERMES_INSTALL_DIR` is the backend source directory containing `hermes_cli` +and `plugins/memory`, not simply the `.app` bundle. `HERMES_PYTHON` must be the +interpreter that runs that backend. Explicit paths fail with diagnostic output +instead of falling back to another installation. Paths containing spaces work +when quoted. For a non-default data/config directory, also set `HERMES_HOME`; +it defaults to `~/.hermes`. The MemOS package/data location remains +`~/.hermes/memos-plugin` for compatibility with existing installations. + +If the chosen interpreter cannot import the host's memory provider API, repair +or update that Hermes environment. Creating an empty `plugins/memory` directory +does not supply the missing API. Restart the desktop app after installation. +Desktop distributions with unrecognized layouts require explicit paths; these +options do not imply that every desktop release has been tested. + ## Protocol surface The adapter calls the following methods on the bridge: diff --git a/apps/memos-local-plugin/core/config/defaults.ts b/apps/memos-local-plugin/core/config/defaults.ts index 2575a8405..b09307c2a 100644 --- a/apps/memos-local-plugin/core/config/defaults.ts +++ b/apps/memos-local-plugin/core/config/defaults.ts @@ -222,6 +222,8 @@ export const DEFAULT_CONFIG: ResolvedConfig = { // an early-life install can still cluster into a world model; // strict 0.6 starved L3 in real usage. clusterMinSimilarity: 0.3, + maxPoliciesPerCluster: 20, + maxPromptChars: 32_000, policyCharCap: 800, traceCharCap: 500, traceEvidencePerPolicy: 1, @@ -249,9 +251,9 @@ export const DEFAULT_CONFIG: ResolvedConfig = { // real usage; 1 lets the candidate→active transition happen // immediately on first successful invocation. candidateTrials: 1, - // Lowered from 6 hours → 0: no cooldown, skills can re-evolve - // as soon as new evidence arrives. - cooldownMs: 0, + // Verification failures are retried after six hours by default; + // operators may set this to 0 when immediate re-evaluation is desired. + cooldownMs: 6 * 60 * 60 * 1000, traceCharCap: 500, evidenceLimit: 6, useLlm: true, diff --git a/apps/memos-local-plugin/core/config/schema.ts b/apps/memos-local-plugin/core/config/schema.ts index 9cd79b612..a5d394b49 100644 --- a/apps/memos-local-plugin/core/config/schema.ts +++ b/apps/memos-local-plugin/core/config/schema.ts @@ -326,6 +326,10 @@ const AlgorithmSchema = Type.Object({ * are ignored (policies too disparate to share a world model). */ clusterMinSimilarity: NumberInRange(0.6, 0, 1), + /** Maximum policies included in one L3 abstraction prompt. */ + maxPoliciesPerCluster: NumberInRange(20, 1, 100), + /** Hard total character cap for one L3 abstraction prompt. */ + maxPromptChars: NumberInRange(32_000, 4_000, 128_000), /** Chars of L2 body handed to `l3.abstraction`. */ policyCharCap: NumberInRange(800, 200, 4_000), /** Chars of trace body handed per evidence trace. */ diff --git a/apps/memos-local-plugin/core/llm/client.ts b/apps/memos-local-plugin/core/llm/client.ts index ee456ac12..7e19f25ca 100644 --- a/apps/memos-local-plugin/core/llm/client.ts +++ b/apps/memos-local-plugin/core/llm/client.ts @@ -282,6 +282,7 @@ export function createLlmClientWithProvider( maxTokens: opts?.maxTokens ?? config.maxTokens ?? DEFAULT_MAX_TOKENS, jsonMode, stop: opts?.stop, + op: opts?.op, }; } @@ -355,24 +356,25 @@ export function createLlmClientWithProvider( notifyError: true, }); } catch (hostErr) { + const normalizedHostErr = normalizeError(hostErr, ERROR_CODES.LLM_UNAVAILABLE, "host fallback failed"); failures++; - const failAt = markFail(hostErr); + const failAt = markFail(normalizedHostErr); facadeLog.error("host.fallback_failed", { primary: summarizeErr(err), - host: summarizeErr(hostErr), + host: summarizeErr(normalizedHostErr), }); // Primary AND host bridge both failed. Trip on a terminal // primary error (the one the operator typically needs to fix // — host bridge failures are usually transient stdio issues). if (breakerIsTerminal(err)) breakerTrip(err); - notifyOnError(hostErr); + notifyOnError(normalizedHostErr); notifyStatus({ status: "error", provider: provider.name, model: config.model, - message: summarizeErrMessage(hostErr), - code: hostErr instanceof MemosError ? hostErr.code : undefined, - ...extractRetryDiagnostics(hostErr instanceof MemosError ? hostErr.details : undefined), + message: summarizeErrMessage(normalizedHostErr), + code: normalizedHostErr.code, + ...extractRetryDiagnostics(normalizedHostErr.details), at: failAt, durationMs: Date.now() - startedAt, fallbackProvider: "host", @@ -380,12 +382,7 @@ export function createLlmClientWithProvider( episodeId: opts?.episodeId, phase: opts?.phase, }); - throw hostErr instanceof MemosError - ? hostErr - : new MemosError( - ERROR_CODES.LLM_UNAVAILABLE, - `host fallback failed: ${(hostErr as Error).message ?? String(hostErr)}`, - ); + throw normalizedHostErr; } } failures++; @@ -840,3 +837,21 @@ function summarizeErrMessage(e: unknown): string { if (e instanceof Error) return e.message; return String(e); } + +function normalizeError( + err: unknown, + fallbackCode: (typeof ERROR_CODES)[keyof typeof ERROR_CODES], + prefix: string, +): MemosError { + if (err instanceof MemosError) return err; + if (err instanceof Error) return new MemosError(fallbackCode, `${prefix}: ${err.message}`); + if (typeof err === "object" && err !== null) { + const record = err as { code?: unknown; message?: unknown; data?: unknown }; + const code = typeof record.code === "string" ? record.code : fallbackCode; + const message = typeof record.message === "string" ? record.message : String(record.data ?? err); + return new MemosError(code as (typeof ERROR_CODES)[keyof typeof ERROR_CODES], `${prefix}: ${message}`, { + bridgeError: err as Record, + }); + } + return new MemosError(fallbackCode, `${prefix}: ${String(err)}`); +} diff --git a/apps/memos-local-plugin/core/llm/prompts/index.ts b/apps/memos-local-plugin/core/llm/prompts/index.ts index 504398a92..ab5aa68ba 100644 --- a/apps/memos-local-plugin/core/llm/prompts/index.ts +++ b/apps/memos-local-plugin/core/llm/prompts/index.ts @@ -50,8 +50,13 @@ export function languageSteeringLine(lang: PromptLanguage): string { * Heuristic: * - Count CJK Unified Ideographs (U+4E00..U+9FFF) as `zh`. * - Count ASCII letters A-Z/a-z as `en`. - * - If CJK accounts for more than `zhRatioThreshold` of counted - * CJK+ASCII signal, pick `zh`. + * - Treat Japanese kana as an explicit non-Chinese signal. This keeps + * Japanese prompts from being mistaken for Chinese just because they + * contain a few shared Han characters. + * - CJK characters carry a small weight because technical identifiers + * (package names, commands, file paths) can contribute many ASCII + * characters inside an otherwise Chinese sentence. If weighted CJK + * accounts for more than `zhRatioThreshold` of the signal, pick `zh`. * - Otherwise pick `en`. * * This intentionally treats Japanese / Korean prompts with filenames, @@ -66,18 +71,25 @@ export function detectDominantLanguage( samples: ReadonlyArray, opts: { zhRatioThreshold?: number } = {}, ): PromptLanguage { - const zhRatioThreshold = opts.zhRatioThreshold ?? 0.7; + const zhRatioThreshold = opts.zhRatioThreshold ?? 0.6; let zh = 0; let en = 0; + let kana = 0; for (const s of samples) { if (!s) continue; for (let i = 0; i < s.length; i++) { const code = s.charCodeAt(i); if (code >= 0x4e00 && code <= 0x9fff) zh++; + else if ( + (code >= 0x3040 && code <= 0x30ff) || + (code >= 0x31f0 && code <= 0x31ff) + ) kana++; else if ((code >= 0x41 && code <= 0x5a) || (code >= 0x61 && code <= 0x7a)) en++; } } + if (kana > 0) return "en"; const total = zh + en; if (total === 0) return "en"; - return zh / total > zhRatioThreshold ? "zh" : "en"; + const weightedZh = zh * 4; + return weightedZh / (weightedZh + en) > zhRatioThreshold ? "zh" : "en"; } diff --git a/apps/memos-local-plugin/core/llm/types.ts b/apps/memos-local-plugin/core/llm/types.ts index 2dd761f32..ef61c439f 100644 --- a/apps/memos-local-plugin/core/llm/types.ts +++ b/apps/memos-local-plugin/core/llm/types.ts @@ -257,6 +257,14 @@ export interface ProviderCallInput { maxTokens: number; jsonMode: boolean; stop?: string[]; + /** + * Logical call site (e.g. `capture.summarize`, `retrieval.filter`, + * `skill.evolve`). Forwarded from `LlmCallOptions.op` so providers + * can apply per-op behavior (request-body tweaks, routing overrides, + * reasoning kill-switches, per-op budget caps). Optional — providers + * must not assume it is set. + */ + op?: string; } /** What providers return — pre-facade post-processing. */ diff --git a/apps/memos-local-plugin/core/memory/l2/induce.ts b/apps/memos-local-plugin/core/memory/l2/induce.ts index 2ea750cd2..6d752f099 100644 --- a/apps/memos-local-plugin/core/memory/l2/induce.ts +++ b/apps/memos-local-plugin/core/memory/l2/induce.ts @@ -23,6 +23,7 @@ import type { EmbeddingVector, EpisodeId, PolicyId, + PolicyMetadata, PolicyRow, TraceId, TraceRow, @@ -156,6 +157,7 @@ export function buildPolicyRow(args: { inducedBy: string; // prompt id + version now?: number; id?: PolicyId; + sourceSignature?: string; }): PolicyRow { const now = args.now ?? Date.now(); const vec = centroid(args.evidenceTraces.map((t) => t.vecSummary ?? t.vecAction ?? null)); @@ -170,6 +172,7 @@ export function buildPolicyRow(args: { gain: 0, status: "candidate", sourceEpisodeIds: Array.from(new Set(args.episodeIds)), + sourceTraceIds: Array.from(new Set(args.evidenceTraces.map((trace) => trace.id))), inducedBy: args.inducedBy, // Fresh policy starts without learned guidance — populated by the // decision-repair pipeline as user feedback / failure bursts arrive. @@ -177,9 +180,72 @@ export function buildPolicyRow(args: { vec: vec as EmbeddingVector | null, createdAt: now, updatedAt: now, + metadata: derivePolicyMetadata(args.evidenceTraces, args.sourceSignature), }; } +function derivePolicyMetadata( + traces: readonly TraceRow[], + sourceSignature?: string, +): PolicyMetadata { + const domainTags = uniqueStrings(traces.flatMap((t) => t.tags ?? [])); + const toolNames = uniqueStrings( + traces.flatMap((t) => (t.toolCalls ?? []).map((c) => c.name ?? "")), + ); + const errorCodes = uniqueStrings( + traces.flatMap((t) => { + const text = [ + t.agentText, + t.reflection ?? "", + ...(t.toolCalls ?? []).map((c) => + typeof c.output === "string" ? c.output : "", + ), + ].join(" "); + return Array.from( + text.matchAll(/\b[A-Z][A-Z0-9]{2,}_[A-Z0-9_]+\b/g), + (m) => m[0], + ); + }), + ); + let zh = 0; + let en = 0; + for (const t of traces) { + for (const s of [t.userText, t.agentText, t.reflection ?? ""]) { + for (const ch of s) { + const code = ch.charCodeAt(0); + if (code >= 0x4e00 && code <= 0x9fff) zh++; + else if ( + (code >= 0x41 && code <= 0x5a) || + (code >= 0x61 && code <= 0x7a) + ) en++; + } + } + } + const total = zh + en; + const language = + total === 0 + ? "unknown" + : zh / total >= 0.7 + ? "zh" + : en / total >= 0.7 + ? "en" + : "mixed"; + return { + version: 1, + language, + domainTags, + toolNames, + errorCodes, + ...(sourceSignature ? { sourceSignature } : {}), + }; +} + +function uniqueStrings(values: readonly string[]): string[] { + return Array.from( + new Set(values.map((v) => v.trim().toLowerCase()).filter(Boolean)), + ).slice(0, 32); +} + // ─── helpers ──────────────────────────────────────────────────────────────── function packTraces( diff --git a/apps/memos-local-plugin/core/memory/l2/l2.ts b/apps/memos-local-plugin/core/memory/l2/l2.ts index 903502b32..b880df52d 100644 --- a/apps/memos-local-plugin/core/memory/l2/l2.ts +++ b/apps/memos-local-plugin/core/memory/l2/l2.ts @@ -253,6 +253,7 @@ export async function runL2( evidenceTraces: traces, inducedBy: `${L2_INDUCTION_PROMPT.id}.v${L2_INDUCTION_PROMPT.version}`, now: input.now ?? Date.now(), + sourceSignature: bucket.signature, }); const owner = ownerFromTraces(traces); policy.ownerAgentKind = owner.ownerAgentKind; @@ -632,6 +633,7 @@ function mergePolicyEvidence(existing: PolicyRow, incoming: PolicyRow, now: numb ...incoming.sourceEpisodeIds, ]), vec: existing.vec ?? incoming.vec, + metadata: existing.metadata ?? incoming.metadata, updatedAt: now as PolicyRow["updatedAt"], }; } diff --git a/apps/memos-local-plugin/core/memory/l3/ALGORITHMS.md b/apps/memos-local-plugin/core/memory/l3/ALGORITHMS.md index bdccc8290..ae3594d51 100644 --- a/apps/memos-local-plugin/core/memory/l3/ALGORITHMS.md +++ b/apps/memos-local-plugin/core/memory/l3/ALGORITHMS.md @@ -50,13 +50,15 @@ rust|cargo|go|java|maven|gradle|typescript|javascript ``` Matches are ordered by the first-hit position so the sweep is -deterministic, then we pick the top two distinct tokens. We deliberately -**do not** embed free-form LLM tags here — domain keys must be cheap -and stable enough to hash. +deterministic, then we pick the top two distinct tokens. Newly induced L2 +rows also persist trace-derived metadata (language, tags, tool names, error +codes, and the source signature); that structured metadata is preferred over +this legacy prose heuristic. We deliberately **do not** embed free-form LLM +tags here — domain keys must be cheap and stable enough to hash. -Policies with no recognised domain keyword fall into the bucket -`__generic|` and are still candidates for clustering by vector -similarity. +Policies with no recognised domain keyword fall into the generic bucket and +are only clustered when at least two vector-bearing policies pass the cosine +gate. This prevents unrelated no-label policies from becoming an L3 prompt. --- @@ -126,7 +128,11 @@ strict, high-gain clusters surface first. ## 4. Evidence packing -Per cluster we assemble a prompt payload: +Per cluster we assemble one or more prompt payloads. Policies are sorted +deterministically and split into batches of at most +`maxPoliciesPerCluster`; the limit bounds prompt size and never discards +cluster members. Batch drafts are then unioned into one draft before the +single merge/create decision. ``` { @@ -134,8 +140,8 @@ Per cluster we assemble a prompt payload: domain_tags: string[], avg_gain: number, avg_support: number, - policies: PolicyPrompt[], // up to |cluster|, each capped - evidence: TracePrompt[] // at most traceEvidencePerPolicy × |cluster| + policies: PolicyPrompt[], // up to maxPoliciesPerCluster, each capped + evidence: TracePrompt[] // at most traceEvidencePerPolicy × batch size } ``` @@ -144,8 +150,8 @@ Per cluster we assemble a prompt payload: * For each policy we fetch the most recent non-redacted supporting trace (by `episodeId`) and include up to `traceCharCap` characters of `userText + reflection`. Evidence is **read-only**, never mutated. -* Total token budget is bounded by `policyCharCap × |cluster| + - traceCharCap × evidencePerPolicy × |cluster|`, which is deterministic +* Per-call token budget is bounded by `policyCharCap × batchSize + + traceCharCap × evidencePerPolicy × batchSize`, which is deterministic and easy to debug. --- @@ -250,9 +256,15 @@ if (now − kv.get(key)) < cooldownDays × 86_400_000: ## 9. Failure policy -* Storage error → propagate. Partial state remains; next run sees the - same eligible policies and re-drives. +* Storage error → warn for the affected cluster and retain retry state. + Other clusters continue; a later run re-drives the failed cluster. * LLM error → single-cluster skip, reason logged. No cooldown update. + +Failed LLM drafts use a persisted retry key scoped by cluster membership. The +retry delays are 5 minutes, 30 minutes, 2 hours, then 6 hours (capped), so a +repeated provider failure cannot consume one LLM call per episode. A successful +world-model insert/update clears the retry key and only then records the normal +cooldown. Storage failures keep the retry state and do not record success. Other clusters continue. * Invalid draft (missing `environment/inference/constraints`) → treated as LLM error. diff --git a/apps/memos-local-plugin/core/memory/l3/README.md b/apps/memos-local-plugin/core/memory/l3/README.md index 0d71f64de..7db55d3c8 100644 --- a/apps/memos-local-plugin/core/memory/l3/README.md +++ b/apps/memos-local-plugin/core/memory/l3/README.md @@ -37,8 +37,8 @@ l2.policy.induced ── triggers ──▶ attachL3Subscriber 2. cluster by (domainKey, centroid cosine ≥ similarity) 3. cooldown check per primary domain tag 4. for each cluster: - a. pack policies + a small evidence trace slice - b. `l3.abstraction` prompt → draft + a. split policies into prompt-sized batches without dropping members + b. `l3.abstraction` prompt per batch → one combined draft c. gather candidate WMs via findByDomainTag d. chooseMergeTarget(cluster, candidates, draft) ├── update: mergeForUpdate + updateBody + bump confidence @@ -56,11 +56,18 @@ No single step blocks reward/L2. Any LLM failure is captured as a `clusterPolicies` (see [`cluster.ts`](./cluster.ts) and [`ALGORITHMS.md`](./ALGORITHMS.md)) bucket-sorts policies by a compact -**domain key** derived from the policy's trigger/procedure text +**domain key** derived first from trace-derived policy metadata (language, +tags, tools, error codes, and source signature), with a legacy +trigger/procedure-text fallback for pre-migration rows (`docker|pip`, `node|npm`, …) and then splits each bucket by centroid cosine, so policies in the same bucket that are still semantically far apart (different sub-environments) end up in separate clusters. +`maxPoliciesPerCluster` is a prompt-size bound, not a retention bound. +Clusters larger than that value are processed in deterministic policy-id +batches. Their drafts are merged before persistence, so all source policy +and episode ids remain attached to a single world model. + ### Merge vs create Whenever a cluster's centroid cosine-matches an existing WM that shares @@ -146,6 +153,8 @@ See `algorithm.l3Abstraction` in | `traceEvidencePerPolicy` | `1` | Evidence traces per policy in the prompt. | | `useLlm` | `true` | Toggle the LLM abstractor off for tests. | | `cooldownDays` | `1` | Debounce per domain tag. | +| `maxPoliciesPerCluster` | `20` | Batch size for one abstraction prompt; overflow is retained. | +| `maxPromptChars` | `32000` | Hard total prompt cap; oversized legacy batches are skipped and quarantined. | | `confidenceDelta` | `0.05` | Confidence step per merge / feedback. | | `minConfidenceForRetrieval` | `0.2` | Tier-3 hide threshold. | @@ -160,6 +169,11 @@ All L3 work is logged on dedicated channels (see * `core.memory.l3.merge` — merge decisions. * `core.memory.l3.confidence` — confidence bumps. * `core.memory.l3.feedback` — human feedback-driven confidence changes. + +Failed legacy clusters use a bounded retry policy (5m/30m/2h/6h). After the +fourth deterministic failure they are quarantined rather than retried forever; +clear the `l3.retry.*` record through the L3 retry-state helper after fixing +the provider or prompt configuration. * `core.memory.l3.events` — listener dispatch errors. ## Tests diff --git a/apps/memos-local-plugin/core/memory/l3/abstract.ts b/apps/memos-local-plugin/core/memory/l3/abstract.ts index 0bde96fb3..e01c70734 100644 --- a/apps/memos-local-plugin/core/memory/l3/abstract.ts +++ b/apps/memos-local-plugin/core/memory/l3/abstract.ts @@ -49,7 +49,7 @@ export interface AbstractInput { export interface AbstractDeps { llm: LlmClient | null; log: Logger; - config: Pick; + config: Pick; /** Optional extra validation executed after the base validator. */ validate?: (d: L3AbstractionDraft) => void; } @@ -74,19 +74,38 @@ export async function abstractDraft( } const userPayload = packPrompt(input, config); + const maxPromptChars = Math.max(4_000, Math.floor(config.maxPromptChars ?? 32_000)); + if (userPayload.length > maxPromptChars) { + log.warn("l3.abstract.prompt_too_large", { + clusterKey: input.cluster.key, + promptChars: userPayload.length, + maxPromptChars, + policyCount: input.cluster.policies.length, + }); + return { + ok: false, + reason: "prompt_too_large", + detail: `prompt has ${userPayload.length} chars; limit is ${maxPromptChars}`, + }; + } // Pick the world-model's rendering language from the underlying // policies + trace evidence. A Chinese user generating "docker alpine // 依赖" policies should see the environment/inference/constraint bullets // written in Chinese; an English user should see them in English. - const langSamples: Array = []; + const policySamples: Array = []; for (const p of input.cluster.policies) { - langSamples.push(p.title, p.trigger, p.procedure, p.boundary, p.verification); + policySamples.push(p.title, p.trigger, p.procedure, p.boundary, p.verification); } + const evidenceSamples: Array = []; for (const traces of input.evidenceByPolicy.values()) { - for (const t of traces) langSamples.push(t.userText, t.agentText, t.reflection); + for (const t of traces) evidenceSamples.push(t.userText, t.agentText, t.reflection); } - const evidenceLang = detectDominantLanguage(langSamples); + // Evidence language wins over legacy policy prose: old English policies + // must not force a Chinese trace cluster back to English. + const evidenceLang = detectDominantLanguage( + evidenceSamples.some((sample) => sample?.trim()) ? evidenceSamples : policySamples, + ); try { const rsp = await llm.completeJson>( diff --git a/apps/memos-local-plugin/core/memory/l3/cluster.ts b/apps/memos-local-plugin/core/memory/l3/cluster.ts index ab20809e0..4620b3d3a 100644 --- a/apps/memos-local-plugin/core/memory/l3/cluster.ts +++ b/apps/memos-local-plugin/core/memory/l3/cluster.ts @@ -27,7 +27,7 @@ export interface ClusterInput { } export interface ClusterDeps { - config: Pick; + config: Pick; } // ─── Domain key extraction ───────────────────────────────────────────────── @@ -60,6 +60,32 @@ const TOOL_REGEXES: Array<{ re: RegExp; tag: string }> = [ ]; export function domainKeyOf(policy: PolicyRow): { key: PolicyClusterKey; tags: string[] } { + // New policies carry structured provenance from their L1 traces. Prefer it + // over English-only regexes so Chinese/Japanese policies and tool-heavy + // traces do not collapse into the generic `_|_` bucket. + if (policy.metadata) { + const tags = uniqueLower([ + ...(policy.metadata.domainTags ?? []), + ...(policy.metadata.toolNames ?? []), + ...(policy.metadata.errorCodes ?? []), + ]); + const signature = policy.metadata.sourceSignature?.split("|") ?? []; + const primary = firstUseful( + policy.metadata.domainTags, + signature[0] && signature[0] !== "_" ? [signature[0]] : [], + ); + const tool = firstUseful( + policy.metadata.toolNames, + signature[2] && signature[2] !== "_" ? [signature[2]] : [], + ); + // Legacy backfill can only provide language (and an empty tag set) when + // traces were already compacted. Preserve the old text heuristics in that + // case instead of turning every such policy into the generic `_|_` bucket. + if (primary || tool || tags.length > 0) { + return { key: `${primary ?? "_"}|${tool ?? "_"}`, tags }; + } + } + const haystack = [policy.title, policy.trigger, policy.procedure, policy.boundary] .filter(Boolean) .join(" \n "); @@ -87,6 +113,24 @@ export function domainKeyOf(policy: PolicyRow): { key: PolicyClusterKey; tags: s }; } +function uniqueLower(values: readonly string[] | undefined): string[] { + return Array.from( + new Set( + (values ?? []) + .map((value) => value.trim().toLowerCase()) + .filter(Boolean), + ), + ).slice(0, 64); +} + +function firstUseful(...groups: Array): string | undefined { + for (const group of groups) { + const value = group?.find((item) => item.trim() && item.trim() !== "_"); + if (value) return value.trim().toLowerCase(); + } + return undefined; +} + // ─── Clustering ──────────────────────────────────────────────────────────── interface PolicyWithMeta { @@ -123,6 +167,11 @@ export function clusterPolicies( for (const [key, members] of byKey) { if (members.length < config.minPolicies) continue; + if (key === "_|_") { + out.push(...clusterUntagged(members, config)); + continue; + } + const vecs: Array = members.map((m) => m.policy.vec ?? null); const center = centroid(vecs); @@ -173,25 +222,30 @@ export function clusterPolicies( // here is only "strict subset" vs "whole bucket". let cohort: PolicyWithMeta[]; let admission: "strict" | "loose"; - if (strict.length >= config.minPolicies) { + const requiredPolicies = config.minPolicies; + if (strict.length >= requiredPolicies) { cohort = strict; admission = "strict"; - } else if (members.length >= config.minPolicies) { + } else if (key !== "_|_" && members.length >= requiredPolicies) { cohort = members; admission = "loose"; } else { continue; } + const ordered = cohort + .slice() + .sort((a, b) => String(a.policy.id).localeCompare(String(b.policy.id))); + if (ordered.length < requiredPolicies) continue; const tags = new Set(); - for (const m of cohort) for (const t of m.tags) tags.add(t); + for (const m of ordered) for (const t of m.tags) tags.add(t); const avgGain = - cohort.reduce((s, m) => s + m.policy.gain, 0) / Math.max(1, cohort.length); + ordered.reduce((s, m) => s + m.policy.gain, 0) / Math.max(1, ordered.length); out.push({ key, - policies: cohort.map((m) => m.policy), + policies: ordered.map((m) => m.policy), domainTags: Array.from(tags), centroidVec: center, avgGain, @@ -211,3 +265,45 @@ export function clusterPolicies( }); return out; } + +function clusterUntagged( + members: readonly PolicyWithMeta[], + config: ClusterDeps["config"], +): PolicyCluster[] { + const requiredPolicies = Math.max(2, config.minPolicies); + const groups: PolicyWithMeta[][] = []; + + for (const member of members + .filter((m) => m.policy.vec) + .slice() + .sort((a, b) => String(a.policy.id).localeCompare(String(b.policy.id)))) { + let target: PolicyWithMeta[] | undefined; + for (const group of groups) { + const center = centroid(group.map((m) => m.policy.vec ?? null)); + if (center && member.policy.vec && cosine(center, member.policy.vec) >= config.clusterMinSimilarity) { + target = group; + break; + } + } + if (target) target.push(member); + else groups.push([member]); + } + + return groups + .filter((group) => group.length >= requiredPolicies) + .map((group) => { + const center = centroid(group.map((m) => m.policy.vec ?? null)); + const cohesion = center + ? group.reduce((sum, m) => sum + cosine(center, m.policy.vec!), 0) / group.length + : 0; + return { + key: `_|_:vec:${String(group[0]!.policy.id)}`, + policies: group.map((m) => m.policy), + domainTags: [], + centroidVec: center, + avgGain: group.reduce((sum, m) => sum + m.policy.gain, 0) / group.length, + cohesion, + admission: "strict" as const, + }; + }); +} diff --git a/apps/memos-local-plugin/core/memory/l3/index.ts b/apps/memos-local-plugin/core/memory/l3/index.ts index 4443a023a..65cc789ee 100644 --- a/apps/memos-local-plugin/core/memory/l3/index.ts +++ b/apps/memos-local-plugin/core/memory/l3/index.ts @@ -13,7 +13,7 @@ export { abstractDraft, buildWorldModelRow } from "./abstract.js"; export { clusterPolicies, domainKeyOf } from "./cluster.js"; export { createL3EventBus } from "./events.js"; -export { adjustConfidence, runL3 } from "./l3.js"; +export { adjustConfidence, clearL3RetryState, runL3 } from "./l3.js"; export type { RunL3Deps } from "./l3.js"; export { chooseMergeTarget, diff --git a/apps/memos-local-plugin/core/memory/l3/l3.ts b/apps/memos-local-plugin/core/memory/l3/l3.ts index 3a786b08d..bdc131e9c 100644 --- a/apps/memos-local-plugin/core/memory/l3/l3.ts +++ b/apps/memos-local-plugin/core/memory/l3/l3.ts @@ -43,12 +43,16 @@ import { } from "./merge.js"; import type { AbstractionResult, + L3AbstractionDraft, + L3AbstractionDraftEntry, + L3AbstractionDraftResult, L3Config, L3Event, L3EventBus, L3ProcessInput, L3ProcessResult, PolicyCluster, + PolicyClusterKey, } from "./types.js"; // ─── Deps ────────────────────────────────────────────────────────────────── @@ -62,6 +66,18 @@ export interface RunL3Deps { } const KV_COOLDOWN_PREFIX = "l3.lastRun."; +const KV_RETRY_PREFIX = "l3.retry."; +const FAILURE_BACKOFF_MS = [5 * 60_000, 30 * 60_000, 2 * 60 * 60_000, 6 * 60 * 60_000]; +const RETRY_STATE_VERSION = 2; +const MAX_FAILURE_ATTEMPTS = FAILURE_BACKOFF_MS.length; + +interface L3RetryState { + version: number; + nextRetryAt: number; + failures: number; + quarantined: boolean; + reason?: string; +} // ─── Public entry ────────────────────────────────────────────────────────── @@ -102,6 +118,7 @@ export async function runL3( config: { clusterMinSimilarity: config.clusterMinSimilarity, minPolicies: config.minPolicies, + maxPoliciesPerCluster: config.maxPoliciesPerCluster ?? 20, }, }, ); @@ -139,6 +156,12 @@ export async function runL3( if (!cluster.centroidVec) { abstractions.push(skipped(cluster, "no_centroid")); + emit(bus, { + kind: "l3.abstraction.skipped", + clusterKey: cluster.key, + reason: "no_centroid", + policyIds: cluster.policies.map((p) => p.id), + }); continue; } @@ -150,6 +173,21 @@ export async function runL3( abstractions.push(skipped(cluster, "cooldown")); continue; } + const retry = readRetryState(repos.kv.get(retryKey(cluster), null)); + if (retry?.quarantined) { + abstractions.push(skipped(cluster, "quarantined")); + emit(bus, { + kind: "l3.abstraction.skipped", + clusterKey: cluster.key, + reason: "quarantined", + policyIds: cluster.policies.map((p) => p.id), + }); + continue; + } + if (retry && retry.nextRetryAt > now) { + abstractions.push(skipped(cluster, "retry_cooldown")); + continue; + } const evidenceByPolicy = loadEvidence(cluster, repos, config.traceEvidencePerPolicy); const episodeIds = collectEpisodeIds(cluster.policies, evidenceByPolicy); @@ -163,20 +201,58 @@ export async function runL3( const triggerEpisodeId = input.episodeId ?? episodeIds[0]; const t0 = Date.now(); - const draftRes = await abstractDraft( - { cluster, evidenceByPolicy, episodeId: triggerEpisodeId }, - { llm: deps.llm, log: abstractLog, config }, - ); + const batchSize = Math.max(1, config.maxPoliciesPerCluster ?? 20); + const drafts: Array<{ draft: L3AbstractionDraft; policyCount: number }> = []; + let draftRes: L3AbstractionDraftResult | null = null; + for (let offset = 0; offset < cluster.policies.length; offset += batchSize) { + const batchPolicies = cluster.policies.slice(offset, offset + batchSize); + const batchPolicyIds = new Set(batchPolicies.map((policy) => policy.id)); + const batchEvidence = new Map( + Array.from(evidenceByPolicy.entries()).filter(([policyId]) => batchPolicyIds.has(policyId)), + ); + const batchCluster: PolicyCluster = { + ...cluster, + policies: batchPolicies, + }; + const batchResult = await abstractDraft( + { cluster: batchCluster, evidenceByPolicy: batchEvidence, episodeId: triggerEpisodeId }, + { llm: deps.llm, log: abstractLog, config }, + ); + if (!batchResult.ok) { + draftRes = batchResult; + break; + } + drafts.push({ draft: batchResult.draft, policyCount: batchPolicies.length }); + } + draftRes ??= { ok: true, draft: combineBatchDrafts(drafts) }; timings.abstract += Date.now() - t0; if (!draftRes.ok) { - abstractions.push(skipped(cluster, draftRes.reason, { episodeIds, policyIds: cluster.policies.map((p) => p.id) })); - emit(bus, { - kind: "l3.failed", - stage: "abstract", - error: { code: draftRes.reason, message: draftRes.detail ?? "" }, - clusterKey: cluster.key, - }); + if ( + draftRes.reason === "llm_failed" || + draftRes.reason === "draft_invalid" || + draftRes.reason === "prompt_too_large" + ) { + recordFailure(cluster, repos.kv, now, draftRes.reason, draftRes.reason === "prompt_too_large"); + } + const policyIds = cluster.policies.map((p) => p.id); + abstractions.push(skipped(cluster, draftRes.reason, { episodeIds, policyIds })); + if (draftRes.reason === "prompt_too_large") { + emit(bus, { + kind: "l3.abstraction.skipped", + clusterKey: cluster.key, + reason: draftRes.reason, + policyIds, + }); + } else { + emit(bus, { + kind: "l3.failed", + stage: "abstract", + error: { code: draftRes.reason, message: draftRes.detail ?? "" }, + clusterKey: cluster.key, + policyIds, + }); + } continue; } @@ -186,6 +262,7 @@ export async function runL3( lookup: repos.worldModel, config, }); + let persisted = false; if (decision.kind === "update") { const patch = mergeForUpdate({ @@ -259,6 +336,7 @@ export async function runL3( policyIds: patch.policyIds as PolicyId[], confidence: bumped, }); + persisted = true; } catch (err) { warnings.push(stageWarn("merge", err, { clusterKey: cluster.key })); } @@ -308,12 +386,18 @@ export async function runL3( policyIds: wm.policyIds, confidence: wm.confidence, }); + persisted = true; } catch (err) { warnings.push(stageWarn("insert", err, { clusterKey: cluster.key })); } } - markCooldown(cluster, repos.kv, now); + if (persisted) { + markCooldown(cluster, repos.kv, now); + repos.kv.del(retryKey(cluster)); + } else { + recordFailure(cluster, repos.kv, now); + } timings.persist += Date.now() - t1; } @@ -372,6 +456,52 @@ function worldModelVectorText(title: string, body: string): string { return [title.trim(), body.trim()].filter(Boolean).join("\n\n") || "(empty)"; } +function combineBatchDrafts( + drafts: readonly { draft: L3AbstractionDraft; policyCount: number }[], +): L3AbstractionDraft { + const first = drafts[0]!; + const weightedPolicyCount = drafts.reduce((sum, item) => sum + item.policyCount, 0); + return { + title: first.draft.title, + domainTags: dedupeStrings(drafts.flatMap((item) => item.draft.domainTags)), + environment: combineDraftEntries(drafts.flatMap((item) => item.draft.environment)), + inference: combineDraftEntries(drafts.flatMap((item) => item.draft.inference)), + constraints: combineDraftEntries(drafts.flatMap((item) => item.draft.constraints)), + body: dedupeStrings(drafts.map((item) => item.draft.body).filter(Boolean)).join("\n\n---\n\n"), + confidence: + drafts.reduce( + (sum, item) => sum + item.draft.confidence * item.policyCount, + 0, + ) / weightedPolicyCount, + supersedesWorldIds: Array.from( + new Set(drafts.flatMap((item) => item.draft.supersedesWorldIds ?? [])), + ), + }; +} + +function combineDraftEntries( + entries: readonly L3AbstractionDraftEntry[], +): L3AbstractionDraftEntry[] { + const combined = new Map(); + for (const entry of entries) { + const key = `${entry.label.trim().toLowerCase()}\u0000${entry.description.trim().toLowerCase()}`; + const previous = combined.get(key); + if (!previous) { + combined.set(key, { ...entry, evidenceIds: dedupeStrings(entry.evidenceIds ?? []) }); + continue; + } + previous.evidenceIds = dedupeStrings([ + ...(previous.evidenceIds ?? []), + ...(entry.evidenceIds ?? []), + ]); + } + return Array.from(combined.values()); +} + +function dedupeStrings(values: readonly string[]): string[] { + return Array.from(new Set(values.map((value) => value.trim()).filter(Boolean))); +} + function skipped( cluster: PolicyCluster, reason: Exclude, @@ -441,6 +571,62 @@ function cooldownKey(cluster: PolicyCluster): string { return `${KV_COOLDOWN_PREFIX}${primary}`; } +function retryKey(cluster: PolicyCluster): string { + return retryKeyFor(cluster.key, cluster.policies.map((p) => p.id)); +} + +function retryKeyFor(clusterKey: PolicyClusterKey, policyIds: readonly PolicyId[]): string { + const members = policyIds.map((id) => String(id)).sort().join(","); + return `${KV_RETRY_PREFIX}${clusterKey}:${members}`; +} + +/** Clear a retry/quarantine record after a config or prompt fix. */ +export function clearL3RetryState( + clusterKey: PolicyClusterKey, + policyIds: readonly PolicyId[], + kv: Repos["kv"], +): void { + kv.del(retryKeyFor(clusterKey, policyIds)); +} + +function recordFailure( + cluster: PolicyCluster, + kv: Repos["kv"], + now: number, + reason?: string, + quarantine = false, +): void { + const previous = readRetryState(kv.get(retryKey(cluster), null)); + const failures = Math.min((previous?.failures ?? 0) + 1, MAX_FAILURE_ATTEMPTS); + const delay = FAILURE_BACKOFF_MS[failures - 1] ?? FAILURE_BACKOFF_MS[FAILURE_BACKOFF_MS.length - 1]!; + kv.set(retryKey(cluster), { + version: RETRY_STATE_VERSION, + failures, + nextRetryAt: now + delay, + quarantined: quarantine || failures >= MAX_FAILURE_ATTEMPTS, + ...(reason ? { reason } : {}), + }); +} + +function readRetryState(raw: unknown): L3RetryState | null { + if (!raw || typeof raw !== "object") return null; + const row = raw as Record; + const failures = typeof row.failures === "number" && Number.isFinite(row.failures) + ? Math.max(0, Math.floor(row.failures)) + : 0; + const nextRetryAt = typeof row.nextRetryAt === "number" && Number.isFinite(row.nextRetryAt) + ? row.nextRetryAt + : 0; + if (failures <= 0 && nextRetryAt <= 0) return null; + return { + version: typeof row.version === "number" ? row.version : 1, + failures, + nextRetryAt, + quarantined: row.quarantined === true, + reason: typeof row.reason === "string" ? row.reason : undefined, + }; +} + function isInCooldown( cluster: PolicyCluster, kv: Repos["kv"], diff --git a/apps/memos-local-plugin/core/memory/l3/types.ts b/apps/memos-local-plugin/core/memory/l3/types.ts index 69d0751a8..99f081f9f 100644 --- a/apps/memos-local-plugin/core/memory/l3/types.ts +++ b/apps/memos-local-plugin/core/memory/l3/types.ts @@ -39,6 +39,10 @@ export interface L3Config { minPolicySupport: number; /** Cosine floor for two L2s to share a cluster. */ clusterMinSimilarity: number; + /** Maximum policies admitted to one abstraction prompt. */ + maxPoliciesPerCluster?: number; + /** Hard total character cap for one abstraction prompt. */ + maxPromptChars?: number; /** Char cap for each L2 body section handed to the prompt. */ policyCharCap: number; /** Char cap for each L1 evidence trace handed to the prompt. */ @@ -130,7 +134,11 @@ export interface L3AbstractionDraft { export type L3AbstractionDraftResult = | { ok: true; draft: L3AbstractionDraft } - | { ok: false; reason: "llm_disabled" | "llm_failed" | "draft_invalid"; detail?: string }; + | { + ok: false; + reason: "llm_disabled" | "llm_failed" | "draft_invalid" | "prompt_too_large"; + detail?: string; + }; // ─── Abstraction outcomes ────────────────────────────────────────────────── @@ -151,7 +159,10 @@ export interface AbstractionResult { | "llm_disabled" | "llm_failed" | "draft_invalid" + | "prompt_too_large" + | "quarantined" | "cooldown" + | "retry_cooldown" | "no_centroid" | "duplicate_of"; /** When `skippedReason === "duplicate_of"`, the existing WM id. */ @@ -229,6 +240,13 @@ export type L3Event = stage: string; error: { code: string; message: string }; clusterKey?: PolicyClusterKey; + policyIds?: PolicyId[]; + } + | { + kind: "l3.abstraction.skipped"; + clusterKey: PolicyClusterKey; + reason: string; + policyIds: PolicyId[]; }; export type L3EventKind = L3Event["kind"]; diff --git a/apps/memos-local-plugin/core/pipeline/deps.ts b/apps/memos-local-plugin/core/pipeline/deps.ts index 79714b35e..f3048fbbe 100644 --- a/apps/memos-local-plugin/core/pipeline/deps.ts +++ b/apps/memos-local-plugin/core/pipeline/deps.ts @@ -217,7 +217,7 @@ export function buildPipelineSubscribers( const log = deps.log ?? rootLogger.child({ channel: "core.pipeline" }); const bgLlmSemaphore = createSemaphore(algorithm.session.bgLlmConcurrency); const bgLlm = rateLimitLlmClient(deps.llm, bgLlmSemaphore, resources); - const bgReflectLlm = rateLimitLlmClient(deps.reflectLlm, bgLlmSemaphore, resources); + const bgEvolverLlm = rateLimitLlmClient(deps.reflectLlm ?? deps.llm, bgLlmSemaphore, resources); const bgL3Llm = rateLimitLlmClient(deps.l3Llm ?? deps.llm, bgLlmSemaphore, resources); const bgEmbedder = resources ? prioritizeEmbedder(deps.embedder, resources, "background") @@ -233,9 +233,8 @@ export function buildPipelineSubscribers( // Issue #2148: capture batch reflection emits JSON, so it must use // the main model rather than the potentially thinking-enabled // skill-evolver model. Keep the background wrapper so capture also - // participates in the shared concurrency limit. `bgReflectLlm` - // remains read-only evaluator metadata below; the original - // `deps.reflectLlm` is also exposed to the Overview health card. + // participates in the shared concurrency limit. The dedicated + // evolver client is used by L2 induction and skill crystallization. reflectLlm: bgLlm, bus: buses.capture, cfg: algorithm.capture, @@ -276,8 +275,8 @@ export function buildPipelineSubscribers( bus: buses.reward, cfg: algorithm.reward, evaluator: { - reflectionProvider: bgReflectLlm?.provider, - reflectionModel: bgReflectLlm?.model, + reflectionProvider: bgLlm?.provider, + reflectionModel: bgLlm?.model, scorerProvider: bgLlm?.provider, scorerModel: bgLlm?.model, }, @@ -311,7 +310,7 @@ export function buildPipelineSubscribers( repos: deps.repos, rewardBus: buses.reward, l2Bus: buses.l2, - llm: bgLlm, + llm: bgEvolverLlm, log: log.child({ channel: "core.memory.l2" }), config: algorithm.l2Induction, thresholds: { @@ -336,7 +335,7 @@ export function buildPipelineSubscribers( const skillHandle = attachSkillSubscriber({ repos: deps.repos, embedder: bgEmbedder, - llm: bgLlm, + llm: bgEvolverLlm, bus: buses.skill, l2Bus: buses.l2, rewardBus: buses.reward, diff --git a/apps/memos-local-plugin/core/pipeline/memory-core.ts b/apps/memos-local-plugin/core/pipeline/memory-core.ts index 660ea737c..c7b2782e4 100644 --- a/apps/memos-local-plugin/core/pipeline/memory-core.ts +++ b/apps/memos-local-plugin/core/pipeline/memory-core.ts @@ -426,7 +426,7 @@ export async function bootstrapMemoryCoreFull( llm = null; } - // Build a dedicated LLM for the reflection phase from skillEvolver + // Build a dedicated LLM for L2 induction and skills from skillEvolver // config when the user has configured a stronger model there. Falls // back to the main `llm` when skillEvolver.model is blank. let reflectLlm: ReturnType | null = null; diff --git a/apps/memos-local-plugin/core/pipeline/types.ts b/apps/memos-local-plugin/core/pipeline/types.ts index 4e20bfada..05c8c2d7b 100644 --- a/apps/memos-local-plugin/core/pipeline/types.ts +++ b/apps/memos-local-plugin/core/pipeline/types.ts @@ -144,7 +144,7 @@ export interface PipelineDeps { repos: Repos; llm: LlmClient | null; /** - * Dedicated LLM for the topic-end reflection + α scoring pass. + * Dedicated LLM for L2 induction and skill crystallization. * Built from `config.skillEvolver.*` when the user configures a * stronger model for skill evolution; falls back to `llm` when * absent. Summarization and per-turn lite capture still use `llm`. @@ -181,7 +181,7 @@ export interface PipelineHandle { readonly repos: Repos; readonly llm: LlmClient | null; /** - * Dedicated client for skill-evolution reflection. When the operator + * Dedicated client for L2 induction and skill crystallization. When the operator * leaves `skillEvolver.*` blank, this is the same instance as `llm` * (so call sites can blindly read whichever is non-null). When they * configure their own model it carries its own `stats()` so the diff --git a/apps/memos-local-plugin/core/skill/ALGORITHMS.md b/apps/memos-local-plugin/core/skill/ALGORITHMS.md index e1560d926..f0f783553 100644 --- a/apps/memos-local-plugin/core/skill/ALGORITHMS.md +++ b/apps/memos-local-plugin/core/skill/ALGORITHMS.md @@ -285,7 +285,10 @@ still under trial. `cooldownMs` debounces repeat runs for the same policy triggered by rapid-fire upstream events (e.g. a burst of `reward.updated`). The -subscriber holds a simple in-memory `{policyId → lastRunAt}` table. +subscriber holds an in-memory `{policyId → lastRunAt}` table plus a pending +policy set. A reward event received during cooldown is coalesced and drained +when that policy becomes eligible; it is not silently discarded. Different +policies have independent cooldowns and queue entries. If `cooldownMs === 0` (as in unit tests), every event triggers a run. diff --git a/apps/memos-local-plugin/core/skill/README.md b/apps/memos-local-plugin/core/skill/README.md index 04a2de8bf..b061d5af9 100644 --- a/apps/memos-local-plugin/core/skill/README.md +++ b/apps/memos-local-plugin/core/skill/README.md @@ -53,6 +53,11 @@ event-driven and every triggered run is fully async. Listener errors are captured so a bad downstream consumer can never break the orchestrator. +For `reward.updated`, the subscriber resolves only policies linked to the +updated episode (falling back to `sourceEpisodeIds` for legacy rows). Each +resolved policy receives its own queue entry and cooldown, so unrelated +policies neither trigger a global scan nor block one another. + ## Key concepts ### Eligibility @@ -201,7 +206,7 @@ See `algorithm.skill` in | `minSupport` | `2` | Min distinct-episode support to crystallize. | | `minGain` | `0.02` | Min policy gain required (paired with the new shrinkage-anchored gain in `core/memory/l2/gain.ts`). | | `candidateTrials` | `3` | Trials required to transition out of `candidate`. NOTE: legacy docs called this `probationaryTrials`; the schema field is `candidateTrials`. | -| `cooldownMs` | `60000` | Debounce between runs triggered by the same policy. | +| `cooldownMs` | `21600000` | Per-policy debounce for reward-triggered runs (6 hours; `0` disables). | | `traceCharCap` | `600` | Char cap per evidence trace in the crystallize prompt.| | `evidenceLimit` | `4` | Max evidence traces per crystallize call. | | `useLlm` | `true` | Toggle the LLM off (tests / degraded mode). | diff --git a/apps/memos-local-plugin/core/skill/crystallize.ts b/apps/memos-local-plugin/core/skill/crystallize.ts index 04de81654..77505df40 100644 --- a/apps/memos-local-plugin/core/skill/crystallize.ts +++ b/apps/memos-local-plugin/core/skill/crystallize.ts @@ -146,12 +146,19 @@ export async function crystallizeDraft( // human-facing fields (display_title, summary, preconditions, steps, // examples) come out in the same language the user was using. The // `name` slug stays snake_case regardless — enforced by `sanitiseName`. - const evidenceLang = detectDominantLanguage([ - input.policy.title, - input.policy.trigger, - input.policy.procedure, - ...input.evidence.flatMap((t) => [t.userText, t.agentText, t.reflection]), + // Trace evidence is authoritative for rendering language. Including an + // old English policy here can force a Chinese episode back to English and + // recreate the verifier resonance failure during upgrades. + const evidenceSamples = input.evidence.flatMap((t) => [ + t.userText, + t.agentText, + t.reflection, ]); + const evidenceLang = detectDominantLanguage( + evidenceSamples.some((sample) => sample?.trim()) + ? evidenceSamples + : [input.policy.title, input.policy.trigger, input.policy.procedure], + ); try { const rsp = await llm.completeJson>( diff --git a/apps/memos-local-plugin/core/skill/skill.ts b/apps/memos-local-plugin/core/skill/skill.ts index d0cdc48ae..be9c0aea2 100644 --- a/apps/memos-local-plugin/core/skill/skill.ts +++ b/apps/memos-local-plugin/core/skill/skill.ts @@ -174,6 +174,7 @@ export async function runSkill( kind: "skill.verification.failed", at: nowMs(), skillId: "sk_placeholder" as SkillId, + policyId: decision.policy.id, reason: verdict.reason ?? "verify-failed", }); continue; diff --git a/apps/memos-local-plugin/core/skill/subscriber.ts b/apps/memos-local-plugin/core/skill/subscriber.ts index 6b30d9d5e..c827a8dfc 100644 --- a/apps/memos-local-plugin/core/skill/subscriber.ts +++ b/apps/memos-local-plugin/core/skill/subscriber.ts @@ -7,10 +7,10 @@ * - `l2.policy.induced` → `runSkill({ trigger, policyId })` * - `l2.policy.status_changed` → `runSkill({ trigger, policyId })` when * the new status is `active` - * - `reward.updated` → `runSkill({ trigger: "reward.updated" })` - * — evaluates every policy referenced by - * the updated episode. Also drives the η - * drift adjustment on existing skills. + * - `reward.updated` → one scoped run per policy linked to the + * updated episode. Each policy has its own + * cooldown and pending queue entry. Also + * drives η adjustment on existing skills. * * The handle returns `runOnce` for manual runs (used by the CLI / viewer * rebuild button) and `applyFeedback` for explicit skill feedback. @@ -33,7 +33,7 @@ import type { SkillFeedbackKind, SkillTrigger, } from "./types.js"; -import type { SkillId } from "../types.js"; +import type { PolicyId, SkillId } from "../types.js"; import { now as nowMs } from "../time.js"; import { IDLE_ARCHIVE_BATCH_LIMIT } from "../storage/repos/skills.js"; @@ -79,13 +79,18 @@ export function attachSkillSubscriber( }; let inflight: Promise | null = null; - let queued: { trigger: SkillTrigger; hint?: { policyId?: string; skillId?: SkillId } } | null = - null; + let disposed = false; + let rewardTimer: ReturnType | null = null; + const lastRewardRunAt = new Map(); + const pendingRewardPolicies = new Set(); + const queued: Array<{ + trigger: SkillTrigger; + hint?: { policyId?: PolicyId; skillId?: SkillId }; + }> = []; async function drain(): Promise { - while (queued) { - const next = queued; - queued = null; + while (queued.length > 0) { + const next = queued.shift()!; try { await runSkill( { trigger: next.trigger, policyId: next.hint?.policyId, skillId: next.hint?.skillId }, @@ -102,9 +107,9 @@ export function attachSkillSubscriber( function triggerRun( trigger: SkillTrigger, - hint?: { policyId?: string; skillId?: SkillId }, + hint?: { policyId?: PolicyId; skillId?: SkillId }, ): void { - queued = { trigger, hint }; + queued.push({ trigger, hint }); if (inflight) { log.debug("skill.run.queued", { trigger }); return; @@ -115,6 +120,41 @@ export function attachSkillSubscriber( inflight = promise; } + function scheduleRewardRuns(): void { + if (disposed) return; + if (rewardTimer) { + clearTimeout(rewardTimer); + rewardTimer = null; + } + const cooldownMs = Math.max(0, deps.config.cooldownMs); + const at = nowMs(); + let nextDelay: number | null = null; + for (const policyId of pendingRewardPolicies) { + const lastRunAt = lastRewardRunAt.get(policyId); + const remainingMs = lastRunAt === undefined ? 0 : cooldownMs - (at - lastRunAt); + if (remainingMs > 0) { + nextDelay = nextDelay === null ? remainingMs : Math.min(nextDelay, remainingMs); + log.debug("skill.run.cooldown", { + trigger: "reward.updated", + policyId, + remainingMs, + }); + continue; + } + pendingRewardPolicies.delete(policyId); + lastRewardRunAt.set(policyId, at); + triggerRun("reward.updated", { policyId }); + } + if (nextDelay !== null && pendingRewardPolicies.size > 0) { + rewardTimer = setTimeout(scheduleRewardRuns, nextDelay); + } + } + + function triggerRewardRuns(policyIds: readonly PolicyId[]): void { + for (const policyId of policyIds) pendingRewardPolicies.add(policyId); + scheduleRewardRuns(); + } + const offInduced = deps.l2Bus.on("l2.policy.induced", (evt: L2Event) => { if (evt.kind !== "l2.policy.induced") return; log.debug("trigger.l2.policy.induced", { policyId: evt.policyId }); @@ -134,10 +174,22 @@ export function attachSkillSubscriber( episodeId: evt.result.episodeId, }); resolveTrialsForReward(evt); - triggerRun("reward.updated"); + const linkedPolicyIds = deps.repos.tracePolicyLinks.getLinkedPolicyIds(evt.result.episodeId); + const sourceEpisodePolicyIds = deps.repos.policies + .list({ status: "active", limit: 200 }) + .filter((policy) => policy.sourceEpisodeIds.includes(evt.result.episodeId)) + .map((policy) => policy.id); + const relatedPolicyIds = Array.from( + new Set([...linkedPolicyIds, ...sourceEpisodePolicyIds]), + ); + triggerRewardRuns(relatedPolicyIds); }); function dispose(): void { + disposed = true; + if (rewardTimer) clearTimeout(rewardTimer); + rewardTimer = null; + pendingRewardPolicies.clear(); offInduced(); offStatus(); offReward(); diff --git a/apps/memos-local-plugin/core/skill/tool-names.ts b/apps/memos-local-plugin/core/skill/tool-names.ts index 4dd903a75..34532f63e 100644 --- a/apps/memos-local-plugin/core/skill/tool-names.ts +++ b/apps/memos-local-plugin/core/skill/tool-names.ts @@ -21,7 +21,20 @@ export function extractToolNames(traces: readonly TraceRow[]): Set { if (name && !IGNORED_NAMES.has(name)) out.add(name); if (typeof tc.input === "string") { - const first = tc.input.trim().split(/\s+/)[0]?.toLowerCase(); + const raw = tc.input.trim(); + // JSON tool arguments are payload, not shell commands. Taking the + // first whitespace token from them produces entries such as `{"code":` + // and poisons the EVIDENCE_TOOLS whitelist. + let parsed: unknown; + try { + parsed = JSON.parse(raw); + } catch { + parsed = undefined; + } + if (parsed !== undefined) { + continue; + } + const first = raw.split(/\s+/)[0]?.toLowerCase(); if (first && first.length >= 2) out.add(first); } } diff --git a/apps/memos-local-plugin/core/skill/types.ts b/apps/memos-local-plugin/core/skill/types.ts index b04e3bc35..5c1dd89ef 100644 --- a/apps/memos-local-plugin/core/skill/types.ts +++ b/apps/memos-local-plugin/core/skill/types.ts @@ -214,6 +214,7 @@ export interface SkillVerificationPassedEvent export interface SkillVerificationFailedEvent extends SkillEventBase<"skill.verification.failed"> { skillId: SkillId; + policyId: string; reason: string; } diff --git a/apps/memos-local-plugin/core/storage/index.ts b/apps/memos-local-plugin/core/storage/index.ts index 1cc7dc5e3..bb2f9baaf 100644 --- a/apps/memos-local-plugin/core/storage/index.ts +++ b/apps/memos-local-plugin/core/storage/index.ts @@ -20,6 +20,7 @@ export { type MigrationFile, type MigrationsResult, } from "./migrator.js"; +export { backfillLegacyPolicyMetadata } from "./policy-metadata-backfill.js"; export { withRetry, withSavepoint, diff --git a/apps/memos-local-plugin/core/storage/migrations/019-policy-metadata.sql b/apps/memos-local-plugin/core/storage/migrations/019-policy-metadata.sql new file mode 100644 index 000000000..b3e63e7e4 --- /dev/null +++ b/apps/memos-local-plugin/core/storage/migrations/019-policy-metadata.sql @@ -0,0 +1,4 @@ +-- Structured provenance for L2 policies. Nullable for backwards compatibility; +-- legacy rows continue to use the L3 text heuristic until re-induced. +ALTER TABLE policies ADD COLUMN metadata_json TEXT; +CREATE INDEX IF NOT EXISTS idx_policies_metadata ON policies(metadata_json); diff --git a/apps/memos-local-plugin/core/storage/migrator.ts b/apps/memos-local-plugin/core/storage/migrator.ts index b858f7aa3..85a5be3c8 100644 --- a/apps/memos-local-plugin/core/storage/migrator.ts +++ b/apps/memos-local-plugin/core/storage/migrator.ts @@ -18,6 +18,7 @@ import { fileURLToPath } from "node:url"; import { now } from "../time.js"; import { rootLogger } from "../logger/index.js"; import { markReady } from "./connection.js"; +import { backfillLegacyPolicyMetadata } from "./policy-metadata-backfill.js"; import type { StorageDb } from "./types.js"; const log = rootLogger.child({ channel: "storage.migration" }); @@ -34,6 +35,8 @@ export interface MigrationsResult { applied: Array<{ version: number; name: string; durationMs: number }>; skipped: number; total: number; + /** Number of legacy policy rows healed during this boot (bounded batch). */ + metadataBackfilled: number; } /** @@ -113,6 +116,19 @@ export function runMigrations(db: StorageDb, dir: string = defaultMigrationsDir( skipped++; continue; } + // Some recovery/compatibility tests (and a few hand-created legacy + // databases) contain only the migration bookkeeping tables. The policy + // metadata migration is additive and has nothing to do when the parent + // policies table is absent; record it as applied so boot can continue. + // Normal databases always have policies from 001-initial.sql, so this + // branch never bypasses the real ALTER TABLE upgrade. + if (file.version === 19 && !tableExists(db, "policies")) { + db.prepare( + `INSERT INTO schema_migrations (version, name, applied_at) VALUES (@version, @name, @applied_at)`, + ).run({ version: file.version, name: file.name, applied_at: now() }); + applied.push({ version: file.version, name: file.name, durationMs: 0 }); + continue; + } const t0 = now(); db.tx(() => { applyMigration(db, file); @@ -134,6 +150,10 @@ export function runMigrations(db: StorageDb, dir: string = defaultMigrationsDir( } ensureHubSharingSearchColumns(db); + const metadataBackfilled = backfillLegacyPolicyMetadata(db); + if (metadataBackfilled > 0) { + log.info("policy-metadata.backfilled", { count: metadataBackfilled }); + } markReady(db); log.info("migrations.summary", { @@ -142,7 +162,7 @@ export function runMigrations(db: StorageDb, dir: string = defaultMigrationsDir( skipped, }); - return { applied, skipped, total: allFiles.length }; + return { applied, skipped, total: allFiles.length, metadataBackfilled }; } /** @@ -210,6 +230,17 @@ function applyMigration(db: StorageDb, file: MigrationFile): void { } return; } + if (file.version === 19 && file.name === "policy-metadata") { + // 019 is intentionally idempotent beyond the schema_migrations marker: + // a process can crash after ALTER TABLE but before recording the marker. + // Re-running must detect the already-present column instead of failing + // with "duplicate column name". + if (tableExists(db, "policies")) { + ensureColumn(db, "policies", "metadata_json", "TEXT"); + db.exec(`CREATE INDEX IF NOT EXISTS idx_policies_metadata ON policies(metadata_json)`); + } + return; + } db.exec(fs.readFileSync(file.fullPath, "utf8")); } diff --git a/apps/memos-local-plugin/core/storage/policy-metadata-backfill.ts b/apps/memos-local-plugin/core/storage/policy-metadata-backfill.ts new file mode 100644 index 000000000..088c9212b --- /dev/null +++ b/apps/memos-local-plugin/core/storage/policy-metadata-backfill.ts @@ -0,0 +1,189 @@ +import type { PolicyMetadata } from "../types.js"; +import type { StorageDb } from "./types.js"; + +const DEFAULT_BATCH_SIZE = 100; + +interface LegacyPolicyRow { + id: string; + title: string; + trigger: string; + procedure: string; + verification: string; + boundary: string; + source_episodes_json: string | null; + source_trace_ids_json: string | null; +} + +interface LegacyTraceRow { + id: string; + episode_id: string; + user_text: string; + agent_text: string; + reflection: string | null; + tags_json: string | null; + tool_calls_json: string | null; +} + +/** + * Fill metadata for legacy policies without holding the migration transaction + * open. The query is deliberately bounded: a large 2.0.x database is healed + * over several normal boots instead of making the first upgrade block on a + * full-table rewrite. Rows are selected by `metadata_json IS NULL`, making the + * operation idempotent and safe to resume after an interrupted boot. + */ +export function backfillLegacyPolicyMetadata( + db: StorageDb, + options: { batchSize?: number } = {}, +): number { + if (!tableExists(db, "policies") || !hasColumn(db, "policies", "metadata_json")) return 0; + + const limit = Math.max(1, Math.min(500, Math.floor(options.batchSize ?? DEFAULT_BATCH_SIZE))); + const policies = db + .prepare<{ limit: number }, LegacyPolicyRow>( + `SELECT id, title, trigger, procedure, verification, boundary, + source_episodes_json, source_trace_ids_json + FROM policies + WHERE metadata_json IS NULL + ORDER BY updated_at ASC, id ASC + LIMIT @limit`, + ) + .all({ limit }); + if (policies.length === 0) return 0; + + const traceById = new Map(); + if (tableExists(db, "traces")) { + const traceIds = new Set(); + const episodeIds = new Set(); + for (const policy of policies) { + for (const id of parseStringArray(policy.source_trace_ids_json)) traceIds.add(id); + for (const id of parseStringArray(policy.source_episodes_json)) episodeIds.add(id); + } + if (traceIds.size > 0 || episodeIds.size > 0) { + const idPlaceholders = Array.from(traceIds, (_, i) => `@id${i}`); + const episodePlaceholders = Array.from(episodeIds, (_, i) => `@episode${i}`); + const predicates = []; + if (idPlaceholders.length > 0) predicates.push(`id IN (${idPlaceholders.join(",")})`); + if (episodePlaceholders.length > 0) predicates.push(`episode_id IN (${episodePlaceholders.join(",")})`); + const params = { + ...Object.fromEntries(Array.from(traceIds, (id, i) => [`id${i}`, id])), + ...Object.fromEntries(Array.from(episodeIds, (id, i) => [`episode${i}`, id])), + }; + const keyed = db + .prepare, LegacyTraceRow>( + `SELECT id, episode_id, user_text, agent_text, reflection, tags_json, tool_calls_json + FROM traces WHERE ${predicates.join(" OR ")}`, + ) + .all(params); + for (const row of keyed) traceById.set(row.id, row); + } + } + + const update = db.prepare<{ id: string; metadata_json: string }>( + `UPDATE policies SET metadata_json=@metadata_json WHERE id=@id AND metadata_json IS NULL`, + ); + return db.tx(() => { + let updated = 0; + for (const policy of policies) { + const traces = parseStringArray(policy.source_trace_ids_json) + .map((id) => traceById.get(id)) + .filter((trace): trace is LegacyTraceRow => trace !== undefined); + if (traces.length === 0) { + const episodes = new Set(parseStringArray(policy.source_episodes_json)); + for (const trace of traceById.values()) { + if (episodes.has(trace.episode_id)) traces.push(trace); + if (traces.length >= 8) break; + } + } + const metadata = deriveLegacyMetadata(policy, traces); + updated += update.run({ id: policy.id, metadata_json: JSON.stringify(metadata) }).changes; + } + return updated; + }); +} + +function deriveLegacyMetadata( + policy: LegacyPolicyRow, + traces: readonly LegacyTraceRow[], +): PolicyMetadata { + const domainTags = unique(traces.flatMap((trace) => parseStringArray(trace.tags_json))); + const toolNames = unique( + traces.flatMap((trace) => parseJsonArray(trace.tool_calls_json) + .map((call) => typeof call.name === "string" ? call.name : "") + .filter(Boolean)), + ); + const errorCodes = unique( + traces.flatMap((trace) => extractErrorCodes([ + trace.agent_text, + trace.reflection ?? "", + ...parseJsonArray(trace.tool_calls_json).map((call) => + typeof call.output === "string" ? call.output : "", + ), + ].join(" "))), + ); + const prose = traces.length > 0 + ? traces.flatMap((trace) => [trace.user_text, trace.agent_text, trace.reflection ?? ""]) + : [policy.title, policy.trigger, policy.procedure, policy.verification, policy.boundary]; + const language = detectLanguage(prose); + return { version: 1, language, domainTags, toolNames, errorCodes }; +} + +function detectLanguage(values: readonly string[]): PolicyMetadata["language"] { + let zh = 0; + let en = 0; + for (const value of values) { + for (const char of value) { + const code = char.charCodeAt(0); + if (code >= 0x4e00 && code <= 0x9fff) zh++; + else if ((code >= 0x41 && code <= 0x5a) || (code >= 0x61 && code <= 0x7a)) en++; + } + } + if (zh === 0 && en === 0) return "unknown"; + if (zh > 0 && en > 0) return zh / en >= 1.5 ? "zh" : en / zh >= 1.5 ? "en" : "mixed"; + return zh > 0 ? "zh" : "en"; +} + +function extractErrorCodes(text: string): string[] { + return Array.from(text.matchAll(/\b[A-Z][A-Z0-9]{2,}_[A-Z0-9_]+\b/g), (match) => match[0]!); +} + +function parseStringArray(raw: string | null): string[] { + if (!raw) return []; + try { + const value: unknown = JSON.parse(raw); + return Array.isArray(value) + ? value.filter((item): item is string => typeof item === "string" && item.trim().length > 0) + : []; + } catch { + return []; + } +} + +function parseJsonArray(raw: string | null): Array> { + if (!raw) return []; + try { + const value: unknown = JSON.parse(raw); + return Array.isArray(value) + ? value.filter((item): item is Record => + Boolean(item) && typeof item === "object" && !Array.isArray(item), + ) + : []; + } catch { + return []; + } +} + +function unique(values: readonly string[]): string[] { + return Array.from(new Set(values.map((value) => value.trim()).filter(Boolean))).slice(0, 32); +} + +function tableExists(db: StorageDb, table: string): boolean { + return Boolean(db.prepare<{ name: string }, { name: string }>( + `SELECT name FROM sqlite_master WHERE type='table' AND name=@name`, + ).get({ name: table })); +} + +function hasColumn(db: StorageDb, table: string, column: string): boolean { + return db.prepare(`PRAGMA table_info(${table})`) + .all() + .some((row) => row.name === column); +} diff --git a/apps/memos-local-plugin/core/storage/repos/policies.ts b/apps/memos-local-plugin/core/storage/repos/policies.ts index 29920f60a..cf769e2bc 100644 --- a/apps/memos-local-plugin/core/storage/repos/policies.ts +++ b/apps/memos-local-plugin/core/storage/repos/policies.ts @@ -1,4 +1,4 @@ -import type { EmbeddingVector, PolicyId, PolicyRow, ShareScope } from "../../types.js"; +import type { EmbeddingVector, PolicyId, PolicyMetadata, PolicyRow, ShareScope } from "../../types.js"; import type { PolicyListFilter, StorageDb } from "../types.js"; import { buildInsert, buildUpdate } from "../tx.js"; import { scanAndTopK, type VectorHit } from "../vector.js"; @@ -46,6 +46,7 @@ const COLUMNS = [ "share_target", "shared_at", "edited_at", + "metadata_json", ]; export interface PolicySearchMeta { @@ -420,6 +421,7 @@ interface RawPolicyRow { share_target: string | null; shared_at: number | null; edited_at: number | null; + metadata_json: string | null; } type RawPolicySearchRow = Pick< @@ -476,6 +478,7 @@ function rowToParams(row: PolicyRow): Record { share_target: row.share?.target ?? null, shared_at: row.share?.sharedAt ?? null, edited_at: row.editedAt ?? null, + metadata_json: row.metadata ? toJsonText(row.metadata) : null, }; } @@ -517,6 +520,30 @@ function mapRow(r: RawPolicyRow): PolicyRow { } : null, editedAt: r.edited_at, + metadata: parsePolicyMetadata(r.metadata_json), + }; +} + +function parsePolicyMetadata(raw: string | null | undefined): PolicyMetadata | undefined { + if (!raw) return undefined; + const value = fromJsonText | null>(raw, null); + if (!value || value.version !== 1) return undefined; + const language = value.language; + if (!['zh', 'en', 'other', 'mixed', 'unknown'].includes(String(language))) return undefined; + const list = (input: unknown): string[] => + Array.isArray(input) + ? Array.from(new Set(input.filter((v): v is string => typeof v === 'string' && !!v.trim()).map((v) => v.trim().slice(0, 64)))) + : []; + const sourceSignature = typeof value.sourceSignature === 'string' && value.sourceSignature.trim() + ? value.sourceSignature.trim().slice(0, 256) + : undefined; + return { + version: 1, + language: language as PolicyMetadata['language'], + domainTags: list(value.domainTags), + toolNames: list(value.toolNames), + errorCodes: list(value.errorCodes), + ...(sourceSignature ? { sourceSignature } : {}), }; } diff --git a/apps/memos-local-plugin/core/storage/repos/trace-policy-links.ts b/apps/memos-local-plugin/core/storage/repos/trace-policy-links.ts index a3c443a9c..6c20bdc86 100644 --- a/apps/memos-local-plugin/core/storage/repos/trace-policy-links.ts +++ b/apps/memos-local-plugin/core/storage/repos/trace-policy-links.ts @@ -24,6 +24,12 @@ export function makeTracePolicyLinksRepo(db: StorageDb) { WHERE policy_id=@policy_id ORDER BY episode_id`, ); + const selectPolicyIds = db.prepare<{ episode_id: EpisodeId }, { policy_id: PolicyId }>( + `SELECT DISTINCT policy_id + FROM trace_policy_links + WHERE episode_id=@episode_id + ORDER BY policy_id`, + ); return { link(args: { @@ -47,5 +53,9 @@ export function makeTracePolicyLinksRepo(db: StorageDb) { getLinkedEpisodeIds(policyId: PolicyId): EpisodeId[] { return selectEpisodeIds.all({ policy_id: policyId }).map((r) => r.episode_id); }, + + getLinkedPolicyIds(episodeId: EpisodeId): PolicyId[] { + return selectPolicyIds.all({ episode_id: episodeId }).map((r) => r.policy_id); + }, }; } diff --git a/apps/memos-local-plugin/core/types.ts b/apps/memos-local-plugin/core/types.ts index 0e4f87f52..1256326a8 100644 --- a/apps/memos-local-plugin/core/types.ts +++ b/apps/memos-local-plugin/core/types.ts @@ -203,6 +203,17 @@ export interface PolicyRow extends OwnedRow { } | null; /** Last user edit through the viewer's edit modal (migration 009). */ editedAt?: EpochMs | null; + /** Structured provenance used by L3 clustering; optional for legacy rows. */ + metadata?: PolicyMetadata; +} + +export interface PolicyMetadata { + version: 1; + language: "zh" | "en" | "other" | "mixed" | "unknown"; + domainTags: string[]; + toolNames: string[]; + errorCodes: string[]; + sourceSignature?: string; } /** diff --git a/apps/memos-local-plugin/docs/CONFIG-ADVANCED.md b/apps/memos-local-plugin/docs/CONFIG-ADVANCED.md index 82e03230d..e25743e34 100644 --- a/apps/memos-local-plugin/docs/CONFIG-ADVANCED.md +++ b/apps/memos-local-plugin/docs/CONFIG-ADVANCED.md @@ -138,6 +138,8 @@ algorithm: minPolicyGain: 0.02 # eligible L2 gain floor (paired with shrinkage-anchored gain) minPolicySupport: 1 # eligible L2 support floor clusterMinSimilarity: 0.6 # cosine cutoff for cluster-and-merge decisions + maxPoliciesPerCluster: 20 # policies per abstraction batch; overflow is retained + maxPromptChars: 32000 # hard total prompt cap; oversized legacy batches are quarantined policyCharCap: 800 # chars per policy in the prompt traceCharCap: 500 # chars per evidence trace in the prompt traceEvidencePerPolicy: 1 # evidence traces per policy in the prompt diff --git a/apps/memos-local-plugin/install.sh b/apps/memos-local-plugin/install.sh index 4e96ba2e9..eb30bff68 100755 --- a/apps/memos-local-plugin/install.sh +++ b/apps/memos-local-plugin/install.sh @@ -135,6 +135,12 @@ Usage: "all" keeps its existing meaning: installed OpenClaw + Hermes targets. Select DSH explicitly with --agent dsh. +Hermes desktop/custom installs (environment variables): + HERMES_INSTALL_DIR Actual Hermes backend source directory + HERMES_PYTHON Python executable used by that Hermes installation + HERMES_HOME Hermes data/config directory (default: ~/.hermes) +Paths containing spaces are supported; quote environment variable values. + Each agent runs its viewer on a fixed port: openclaw → http://127.0.0.1:${OPENCLAW_PORT} hermes → http://127.0.0.1:${HERMES_PORT} @@ -229,7 +235,10 @@ HAS_OPENCLAW="false" HAS_HERMES="false" HAS_DSH="false" [[ -d "${HOME}/.openclaw" ]] && HAS_OPENCLAW="true" -[[ -d "${HOME}/.hermes" ]] && HAS_HERMES="true" +if [[ -d "${HERMES_HOME:-${HOME}/.hermes}" || -n "${HERMES_INSTALL_DIR:-}" || -n "${HERMES_PYTHON:-}" ]] \ + || command -v hermes >/dev/null 2>&1; then + HAS_HERMES="true" +fi find_openclaw_cli() { command -v openclaw 2>/dev/null && return 0 @@ -301,8 +310,15 @@ SOURCE_KIND="" # "path" for a local file, "npm" otherwise SOURCE_SPEC="" GATEWAY_RECOVERY_BIN="" GATEWAY_RECOVERY_STATE="inactive" +HERMES_STAGED_PREFIX="" +HERMES_PREFIX_BACKUP_ROOT="" +HERMES_LIVE_PREFIX="" +HERMES_PREFIX_SWAPPED="false" cleanup_install_state() { + if declare -F rollback_hermes_install >/dev/null 2>&1; then + rollback_hermes_install + fi if [[ "${GATEWAY_RECOVERY_STATE:-inactive}" == "needs_recovery" \ && -n "${GATEWAY_RECOVERY_BIN:-}" ]]; then local recovery_out="" @@ -744,15 +760,263 @@ NODE return 1 } +# Probe only: do not stop processes or write host files until this succeeds. +# The bootstrap Python uses only the standard library; each candidate is +# validated in its own interpreter, independently of the bootstrap environment. +resolve_hermes_environment() { + local bootstrap="${HERMES_PYTHON:-$(command -v python3 || true)}" + [[ -x "${bootstrap}" ]] || { printf '%s\n' "Cannot run Python: ${bootstrap}. Set HERMES_PYTHON to the Hermes interpreter." >&2; return 1; } + "${bootstrap}" - <<'PY' +import os +import shlex +import shutil +import subprocess +import sys +from pathlib import Path + + +def absolute(value): + return Path(os.path.abspath(os.path.expanduser(value))) + + +explicit_root = os.environ.get("HERMES_INSTALL_DIR", "") +explicit_python = os.environ.get("HERMES_PYTHON", "") +host_home = absolute(os.environ.get("HERMES_HOME") or str(Path.home() / ".hermes")) +roots = [] +candidates = [] + + +def add_root(root): + root = absolute(str(root)) + if root not in roots: + roots.append(root) + + +def add_python(python, root=None): + # Do not resolve the Python symlink: that would escape its virtualenv. + pair = (str(absolute(str(python))), str(root) if root else "") + if pair not in candidates: + candidates.append(pair) + + +if explicit_root: + add_root(explicit_root) +else: + launchers = [shutil.which("hermes"), str(Path.home() / ".local/bin/hermes")] + for launcher in launchers: + if not launcher or not Path(launcher).is_file(): + continue + entry = Path(launcher).resolve() + for parent in entry.parents: + if (parent / "hermes_cli").is_dir(): + add_root(parent) + break + try: + lines = entry.read_text().splitlines() + except (OSError, UnicodeError): + continue + for line in lines: + try: + words = shlex.split(line[2:] if line.startswith("#!") else line) + except ValueError: + continue + # Recognize literal Python shebangs and official `exec "...python"` + # launchers. Never source/eval a wrapper or expand shell expressions. + value = words[1] if len(words) > 1 and words[0] == "exec" else ( + words[0] if words and line.startswith("#!") else "" + ) + if not value.startswith("/") or not Path(value).name.startswith("python"): + continue + python = absolute(value) + root = next((p for p in python.parents if (p / "hermes_cli").is_dir()), None) + if root: + add_root(root) + if not explicit_python: + add_python(python, root) + add_root(host_home / "hermes-agent") + +if explicit_python: + candidates = [] + for root in roots: + add_python(explicit_python, root) + if not explicit_root: + add_python(explicit_python) +else: + for root in roots: + for venv in ("venv", ".venv"): + for name in ("python", "python3"): + add_python(root / venv / "bin" / name, root) + # A PATH interpreter is acceptable only if it really imports the host API. + # Do not combine a system interpreter with a guessed source checkout. + if not explicit_root: + for name in ("python3", "python"): + value = shutil.which(name) + if value: + add_python(value) + +probe = ''' +import sys +from pathlib import Path +if sys.argv[1]: + sys.path.insert(0, sys.argv[1]) +import hermes_cli +import plugins.memory as memory +from plugins.memory import load_memory_provider +if not callable(load_memory_provider): + raise ImportError("Hermes load_memory_provider is not callable") +directory = Path(memory.__file__).resolve().parent +if sys.argv[1] and directory != Path(sys.argv[1]).resolve() / "plugins" / "memory": + raise ImportError("Memory API resolved outside the selected Hermes source directory") +print(directory) +print(directory.parent.parent) +''' +child_env = {k: v for k, v in os.environ.items() if k not in ("PYTHONPATH", "PYTHONHOME")} +for python, root in candidates: + if not os.access(python, os.X_OK): + print("Not executable: " + python, file=sys.stderr) + continue + try: + result = subprocess.run( + [python, "-c", probe, root], env=child_env, cwd="/", + capture_output=True, text=True, timeout=15, + ) + except (OSError, subprocess.TimeoutExpired) as exc: + print("Cannot probe {}: {}".format(python, exc), file=sys.stderr) + continue + paths = result.stdout.strip().splitlines() + if result.returncode == 0 and len(paths) == 2 and Path(paths[0]).is_dir(): + print(python) + print("\n".join(paths)) + sys.exit(0) + print("Hermes import failed using {} (root: {}):\n{}".format( + python, root or "interpreter default", result.stderr.strip() or result.stdout.strip() + ), file=sys.stderr) +print("Cannot locate a compatible Hermes environment. Set HERMES_INSTALL_DIR to the actual " + "backend source directory and HERMES_PYTHON to its Python executable. " + "HERMES_HOME is the data/config directory, not the application bundle. " + "If imports still fail, repair/update that Hermes environment first.", file=sys.stderr) +sys.exit(1) +PY +} + # ─── Hermes install ─────────────────────────────────────────────────────── +HERMES_PROVIDER_BACKUP_DIR="" +HERMES_PROVIDER_BACKUP_TARGETS=() +HERMES_PROVIDER_BACKUP_PATHS=() + +rollback_hermes_install() { + if ((${#HERMES_PROVIDER_BACKUP_TARGETS[@]} > 0)); then + restore_hermes_provider_targets + fi + if [[ "${HERMES_PREFIX_SWAPPED}" == "true" \ + && "${HERMES_LIVE_PREFIX:-}" == "${HOME}/.hermes/memos-plugin" ]]; then + if [[ -d "${HERMES_LIVE_PREFIX}" ]]; then + rm -rf "${HERMES_LIVE_PREFIX}" + fi + if [[ -d "${HERMES_PREFIX_BACKUP_ROOT:-}" && -d "${HERMES_PREFIX_BACKUP_ROOT}/live" ]]; then + mv "${HERMES_PREFIX_BACKUP_ROOT}/live" "${HERMES_LIVE_PREFIX}" \ + || warn "Failed to restore the previous Hermes plugin prefix." + fi + fi + if [[ -n "${HERMES_STAGED_PREFIX:-}" && -d "${HERMES_STAGED_PREFIX}" ]]; then + rm -rf "${HERMES_STAGED_PREFIX}" + fi + if [[ -n "${HERMES_PREFIX_BACKUP_ROOT:-}" && -d "${HERMES_PREFIX_BACKUP_ROOT}" ]]; then + rm -rf "${HERMES_PREFIX_BACKUP_ROOT}" + fi + HERMES_STAGED_PREFIX="" + HERMES_PREFIX_BACKUP_ROOT="" + HERMES_LIVE_PREFIX="" + HERMES_PREFIX_SWAPPED="false" +} + +restore_hermes_provider_targets() { + local index target backup + for ((index=${#HERMES_PROVIDER_BACKUP_TARGETS[@]}-1; index>=0; index--)); do + target="${HERMES_PROVIDER_BACKUP_TARGETS[index]}" + backup="${HERMES_PROVIDER_BACKUP_PATHS[index]}" + if [[ -L "${target}" ]]; then + rm -f "${target}" + elif [[ -e "${target}" ]]; then + warn "Cannot restore Hermes provider target because it changed after installation: ${target}" + continue + fi + if [[ -e "${backup}" || -L "${backup}" ]]; then + mv "${backup}" "${target}" || warn "Failed to restore Hermes provider target: ${target}" + fi + done + if [[ -n "${HERMES_PROVIDER_BACKUP_DIR}" && -d "${HERMES_PROVIDER_BACKUP_DIR}" ]]; then + rm -rf "${HERMES_PROVIDER_BACKUP_DIR}" + fi + HERMES_PROVIDER_BACKUP_DIR="" + HERMES_PROVIDER_BACKUP_TARGETS=() + HERMES_PROVIDER_BACKUP_PATHS=() +} + +prepare_hermes_provider_targets() { + local backup_root="$1" + shift + HERMES_PROVIDER_BACKUP_DIR="$(mktemp -d "${TMPDIR:-/tmp}/memos-hermes-provider.XXXXXX")" \ + || die "Unable to create a Hermes provider rollback directory." + local index=0 target backup + for target in "$@"; do + backup="" + if [[ -e "${target}" || -L "${target}" ]]; then + backup="${HERMES_PROVIDER_BACKUP_DIR}/target-${index}" + mv "${target}" "${backup}" \ + || { restore_hermes_provider_targets; die "Unable to back up Hermes provider target: ${target}"; } + fi + HERMES_PROVIDER_BACKUP_TARGETS+=("${target}") + HERMES_PROVIDER_BACKUP_PATHS+=("${backup}") + ln -sfn "${backup_root}" "${target}" \ + || { restore_hermes_provider_targets; die "Unable to link Hermes memtensor provider: ${target}"; } + success "Symlinked → ${target}" + index=$((index + 1)) + done +} + install_hermes() { STEP_CURRENT=0 header "Hermes Install" + step "Locating and validating Hermes Python environment" + local resolved python_bin plugin_dir hermes_source + resolved="$(resolve_hermes_environment)" || die "Hermes preflight failed; host files and processes were not changed." + python_bin="$(printf '%s\n' "${resolved}" | sed -n '1p')" + plugin_dir="$(printf '%s\n' "${resolved}" | sed -n '2p')" + hermes_source="$(printf '%s\n' "${resolved}" | sed -n '3p')" + success "Python: ${python_bin}" + success "plugins/memory: ${plugin_dir}" + # Keep all subsequent Python operations in the validated host context. + local PYTHONPATH="${hermes_source}" PYTHONHOME="" + export PYTHONPATH PYTHONHOME + local hermes_host_home="${HERMES_HOME:-${HOME}/.hermes}" local prefix="${HOME}/.hermes/memos-plugin" + local config_file="${hermes_host_home}/config.yaml" + HERMES_LIVE_PREFIX="${prefix}" + mkdir -p "${HOME}/.hermes" + HERMES_STAGED_PREFIX="$(mktemp -d "${HOME}/.hermes/.memos-plugin-stage.XXXXXX")" \ + || die "Unable to create a staged Hermes plugin directory." + local staged_prefix="${HERMES_STAGED_PREFIX}" local home="${prefix}" - local config_file="${HOME}/.hermes/config.yaml" local adapter_dir="${prefix}/adapters/hermes" - mkdir -p "${HOME}/.hermes" + local staged_adapter_dir="${staged_prefix}/adapters/hermes" + deploy_tarball_to_prefix "${staged_prefix}" + + step "Preparing staged Hermes runtime" + ensure_runtime_home "hermes" "${staged_prefix}" "${staged_prefix}" + local staged_bridge_entry="${staged_prefix}/dist/bridge.cjs" + [[ -f "${staged_bridge_entry}" ]] || staged_bridge_entry="${staged_prefix}/bridge.cts" + echo "${prefix}/dist/bridge.cjs" > "${staged_adapter_dir}/bridge_path.txt" + [[ -f "${staged_adapter_dir}/plugin.yaml" ]] || die "Staged Hermes adapter manifest is missing." + local staged_version_sync="${staged_prefix}/scripts/sync-hermes-version.cjs" + if [[ -f "${staged_version_sync}" ]]; then + node "${staged_version_sync}" "${staged_prefix}" >/dev/null \ + || die "Failed to synchronize staged Hermes plugin version metadata." + else + warn "Hermes version sync helper missing; using packaged plugin.yaml as-is." + fi + cp "${staged_adapter_dir}/plugin.yaml" "${staged_adapter_dir}/memos_provider/plugin.yaml" 2>/dev/null \ + || die "Failed to prepare the staged memtensor provider." step "Stopping existing bridge daemon" local bridge_pids="" @@ -792,49 +1056,41 @@ install_hermes() { fi fi - deploy_tarball_to_prefix "${prefix}" - - step "Configuring runtime environment" - ensure_runtime_home "hermes" "${home}" "${prefix}" + # Copy user data only after the host processes are stopped. The old prefix + # remains intact until the staged prefix has been handed off successfully. + step "Handing off staged Hermes runtime" + local preserved_item + local preserved_items=(data logs skills daemon .migrations config.yaml .auth.json) + mkdir -p "${prefix%/*}" + for preserved_item in "${preserved_items[@]}"; do + if [[ -e "${prefix}/${preserved_item}" ]]; then + rm -rf "${staged_prefix:?}/${preserved_item}" + cp -a "${prefix}/${preserved_item}" "${staged_prefix}/${preserved_item}" \ + || die "Failed to stage Hermes runtime data: ${preserved_item}" + fi + done + HERMES_PREFIX_BACKUP_ROOT="$(mktemp -d "${HOME}/.hermes/.memos-plugin-backup.XXXXXX")" \ + || die "Unable to create a Hermes rollback directory." + if [[ -d "${prefix}" ]]; then + mv "${prefix}" "${HERMES_PREFIX_BACKUP_ROOT}/live" \ + || die "Unable to move the existing Hermes plugin prefix into rollback storage." + fi + HERMES_PREFIX_SWAPPED="true" + mv "${staged_prefix}" "${prefix}" \ + || { rollback_hermes_install; die "Unable to activate the staged Hermes plugin prefix."; } + HERMES_STAGED_PREFIX="" + adapter_dir="${prefix}/adapters/hermes" + home="${prefix}" + staged_bridge_entry="" + success "Staged Hermes runtime activated" local bridge_entry="${prefix}/dist/bridge.cjs" [[ -f "${bridge_entry}" ]] || bridge_entry="${prefix}/bridge.cts" echo "${bridge_entry}" > "${adapter_dir}/bridge_path.txt" success "Bridge path recorded" - step "Locating Hermes Python environment" - local python_bin="" - if command -v hermes >/dev/null 2>&1; then - local shebang; shebang="$(head -1 "$(command -v hermes)" 2>/dev/null || true)" - [[ "${shebang}" == "#!"*python* ]] && python_bin="$(echo "${shebang}" | sed 's/^#!\s*//')" - fi - if [[ -z "${python_bin}" || ! -x "${python_bin}" ]] \ - && [[ -x "${HOME}/.hermes/hermes-agent/venv/bin/python3" ]]; then - python_bin="${HOME}/.hermes/hermes-agent/venv/bin/python3" - fi - [[ -z "${python_bin}" || ! -x "${python_bin}" ]] && python_bin="$(command -v python3 || true)" - [[ -n "${python_bin}" && -x "${python_bin}" ]] || die "Cannot locate Python for Hermes." - success "Python: ${python_bin}" - - local plugin_dir="" - plugin_dir="$("${python_bin}" -c " -from pathlib import Path -try: - import plugins.memory as pm - print(Path(pm.__file__).parent) -except Exception: - pass -" 2>/dev/null || true)" - if [[ -z "${plugin_dir}" || ! -d "${plugin_dir}" ]]; then - for d in "${HOME}/.hermes/hermes-agent/plugins/memory"; do - [[ -d "${d}" && -f "${d}/__init__.py" ]] && { plugin_dir="${d}"; break; } - done - fi - [[ -n "${plugin_dir}" && -d "${plugin_dir}" ]] || die "plugins/memory not found" - success "plugins/memory: ${plugin_dir}" - step "Linking memtensor provider" - local user_plugin_dir="${HOME}/.hermes/plugins/memory" + local user_plugin_dir="${hermes_host_home}/plugins/memory" mkdir -p "${user_plugin_dir}" local version_sync="${prefix}/scripts/sync-hermes-version.cjs" if [[ -f "${version_sync}" ]]; then @@ -843,20 +1099,13 @@ except Exception: else warn "Hermes version sync helper missing; using packaged plugin.yaml as-is." fi - # Ensure the provider directory is fully populated before symlinking so - # the second symlink (user-level) already points at a complete tree. - cp "${adapter_dir}/plugin.yaml" "${adapter_dir}/memos_provider/plugin.yaml" 2>/dev/null || true + # The provider directory was populated while the package was staged, before + # any host link was changed. local provider_targets=( "${plugin_dir}/memtensor" "${user_plugin_dir}/memtensor" ) - local target - for target in "${provider_targets[@]}"; do - # Use `ln -sfn` for atomic, idempotent replace; matches install.hermes.sh. - if [[ -e "${target}" && ! -L "${target}" ]]; then rm -rf "${target}"; fi - ln -sfn "${adapter_dir}/memos_provider" "${target}" - success "Symlinked → ${target}" - done + prepare_hermes_provider_targets "${adapter_dir}/memos_provider" "${provider_targets[@]}" step "Verifying provider & patching config" local verify @@ -865,8 +1114,13 @@ from plugins.memory import load_memory_provider p = load_memory_provider('memtensor') print('OK' if p and p.name == 'memtensor' else 'FAIL') " 2>/dev/null || true)" - [[ "${verify}" == "OK" ]] && success "Provider verification passed" \ - || warn "Provider verification didn't return OK" + if [[ "${verify}" == "OK" ]]; then + success "Provider verification passed" + else + restore_hermes_provider_targets + rollback_hermes_install + die "Hermes memtensor provider verification failed. Existing provider targets and plugin prefix were restored." + fi step "Installing Hermes profile defaults hook" "${python_bin}" - <<'PYEOF' || warn "Hermes profile defaults hook install failed" @@ -1015,7 +1269,7 @@ PYEOF if [[ -f "${config_file}" ]]; then local patched_configs - patched_configs="$("${python_bin}" - "${HOME}/.hermes" 2>/dev/null <<'PYEOF' + patched_configs="$("${python_bin}" - "${hermes_host_home}" 2>/dev/null <<'PYEOF' import sys from pathlib import Path @@ -1133,6 +1387,20 @@ CFGEOF else printf " ${DIM}Next:${NC} ${BOLD}hermes chat${NC}\n" fi + # Keep rollback material until the daemon smoke test and all post-swap + # configuration work have completed successfully. + if [[ -n "${HERMES_PROVIDER_BACKUP_DIR}" ]]; then + rm -rf "${HERMES_PROVIDER_BACKUP_DIR}" + fi + if [[ -n "${HERMES_PREFIX_BACKUP_ROOT}" && -d "${HERMES_PREFIX_BACKUP_ROOT}" ]]; then + rm -rf "${HERMES_PREFIX_BACKUP_ROOT}" + fi + HERMES_PROVIDER_BACKUP_DIR="" + HERMES_PROVIDER_BACKUP_TARGETS=() + HERMES_PROVIDER_BACKUP_PATHS=() + HERMES_PREFIX_BACKUP_ROOT="" + HERMES_LIVE_PREFIX="" + HERMES_PREFIX_SWAPPED="false" return 0 } diff --git a/apps/memos-local-plugin/tests/integration/skill-l3-full-chain.test.ts b/apps/memos-local-plugin/tests/integration/skill-l3-full-chain.test.ts new file mode 100644 index 000000000..f5c593bef --- /dev/null +++ b/apps/memos-local-plugin/tests/integration/skill-l3-full-chain.test.ts @@ -0,0 +1,279 @@ +/** + * Full local-memory verification for the repaired Skill → L3 path. + * + * This deliberately keeps the same Chinese, JSON-tool evidence that used to + * fail Skill verification, then feeds three compatible active L2 policies to + * L3. It proves the repaired boundaries in one deterministic SQLite run: + * language steering, tool-name extraction, Skill persistence, L3 creation, + * and idempotent L3 merge. + */ + +import { afterEach, describe, expect, it } from "vitest"; + +import { rootLogger } from "../../core/logger/index.js"; +import { + createL3EventBus, + runL3, + type L3Config, + type L3Event, +} from "../../core/memory/l3/index.js"; +import { L3_ABSTRACTION_PROMPT } from "../../core/llm/prompts/l3-abstraction.js"; +import { runL2, type L2Config } from "../../core/memory/l2/index.js"; +import { + createSkillEventBus, + runSkill, + type SkillEvent, +} from "../../core/skill/index.js"; +import type { EpisodeId } from "../../core/types.js"; +import { fakeLlm } from "../helpers/fake-llm.js"; +import { makeTmpDb, type TmpDbHandle } from "../helpers/tmp-db.js"; +import { + makeDraft, + makeSkillConfig, + seedSessionOnly, + seedTrace, +} from "../unit/skill/_helpers.js"; + +const L3_OP = `${L3_ABSTRACTION_PROMPT.id}.v${L3_ABSTRACTION_PROMPT.version}`; +const log = rootLogger.child({ channel: "integration.skill-l3" }); + +function l3Config(): L3Config { + return { + minPolicies: 3, + minPolicyGain: 0.1, + minPolicySupport: 2, + clusterMinSimilarity: 0.3, + maxPoliciesPerCluster: 20, + maxPromptChars: 32_000, + looseMinCohesion: 0.55, + policyCharCap: 800, + traceCharCap: 500, + traceEvidencePerPolicy: 1, + useLlm: true, + cooldownDays: 0, + confidenceDelta: 0.1, + minConfidenceForRetrieval: 0.2, + }; +} + +function l2Config(): L2Config { + return { + minSimilarity: 0.95, + candidateTtlDays: 30, + gamma: 0.9, + tauSoftmax: 0.4, + useLlm: true, + minTraceValue: 0.1, + minEpisodesForInduction: 2, + inductionTraceCharCap: 2_000, + gainEmaAlpha: 0.4, + }; +} + +describe("integration: Chinese Skill crystallization → L3 world model", () => { + let handle: TmpDbHandle | null = null; + + afterEach(() => { + handle?.cleanup(); + handle = null; + }); + + it("completes Skill verification and creates then merges one L3 model", async () => { + handle = makeTmpDb(); + const h = handle; + const groups = [ + { packageName: "cryptography", errorCode: "EXIT_1", vector: [1, 0, 0] }, + { packageName: "lxml", errorCode: "EXIT_2", vector: [0, 1, 0] }, + { packageName: "psycopg2", errorCode: "EXIT_3", vector: [0, 0, 1] }, + ] as const; + for (const [groupIndex, group] of groups.entries()) { + for (let repetition = 0; repetition < 2; repetition += 1) { + const episodeId = `ep_chain_${groupIndex + 1}_${repetition + 1}` as EpisodeId; + const sessionId = `s_chain_${groupIndex}_${repetition}`; + const toolCalls = [ + { + name: "pip", + input: JSON.stringify({ package: group.packageName }), + output: `Error: ${group.errorCode}`, + }, + { name: "execute_code", input: '{"code":"print(1)"}' }, + ]; + seedSessionOnly(h, sessionId); + seedTrace(h, { + episodeId, + sessionId, + userText: `在 Alpine 镜像中安装 ${group.packageName} 失败`, + agentText: `先执行 apk add ${group.packageName}-dev,再执行 pip install ${group.packageName}`, + reflection: "先安装系统库,再重试 pip 安装", + value: 0.9, + tags: ["alpine", "pip"], + vec: new Float32Array(group.vector), + toolCalls, + }); + } + } + + const skillEvents: SkillEvent[] = []; + const skillBus = createSkillEventBus(); + skillBus.onAny((event) => skillEvents.push(event)); + const skillPrompts: unknown[] = []; + let inductionCount = 0; + const llm = fakeLlm({ + completeJson: { + "l2.l2.induction.v2": () => { + const group = groups[inductionCount++] ?? groups[0]; + return { + title: `Alpine ${group.packageName} dependency repair`, + trigger: `pip install ${group.packageName} fails in Alpine`, + procedure: `apk add ${group.packageName}-dev then pip install ${group.packageName}`, + verification: `${group.packageName} installs successfully`, + boundary: "仅适用于 Alpine 镜像", + rationale: "多次观察到 Alpine 原生依赖缺失", + caveats: ["Alpine 使用 musl libc"], + confidence: 0.8, + }; + }, + "skill.crystallize": (input) => { + skillPrompts.push(input); + return makeDraft({ + name: "alpine_cryptography_fix", + displayTitle: "Alpine cryptography 安装修复", + summary: "在 Alpine 中先安装系统库,再重试 pip 安装 cryptography", + preconditions: ["当前环境是 Alpine 镜像"], + steps: [ + { title: "检查错误", body: "确认 pip install cryptography 安装失败" }, + { title: "安装系统库", body: "执行 apk add openssl-dev" }, + { title: "重新安装", body: "再次执行 pip install cryptography" }, + ], + tools: ["execute_code"], + tags: ["alpine", "pip"], + }); + }, + [L3_OP]: () => ({ + title: "Alpine Python dependency environment", + domain_tags: ["alpine", "python", "pip"], + environment: [ + { + label: "Alpine uses musl libc", + description: "Alpine images commonly lack glibc-linked build dependencies by default.", + }, + ], + inference: [ + { + label: "Native extensions need system headers", + description: "Python packages with native extensions fail to build when matching headers are absent.", + }, + ], + constraints: [], + body: "# Alpine Python dependency environment", + confidence: 0.75, + supersedes_world_ids: [], + }), + }, + }); + + const l2Runs = []; + const l2Deps = { + db: h.db, + repos: h.repos, + llm, + log, + config: l2Config(), + thresholds: { minSupport: 2, minGain: 0.1, archiveGain: -0.05 }, + }; + for (const group of groups) { + const traces = h.repos.traces + .list({ limit: 20 }) + .filter((trace) => trace.userText.includes(group.packageName)); + expect(traces).toHaveLength(2); + for (const trace of traces) { + l2Runs.push( + await runL2( + { + episodeId: trace.episodeId, + sessionId: trace.sessionId, + traces: [trace], + trigger: "reward.updated", + }, + l2Deps, + ), + ); + } + } + + const inducedPolicies = h.repos.policies.list({ status: "active" }); + expect(l2Runs.flatMap((run) => run.inductions).filter((induction) => induction.policyId)) + .toHaveLength(3); + expect(inducedPolicies).toHaveLength(3); + const p1 = inducedPolicies.find((policy) => policy.title.includes("cryptography")); + expect(p1).toBeDefined(); + expect(p1!.metadata?.domainTags).toEqual(["alpine", "pip"]); + expect(p1!.metadata?.toolNames).toEqual(["pip", "execute_code"]); + + const skillResult = await runSkill( + { trigger: "manual", policyId: p1!.id, episodeId: p1!.sourceEpisodeIds[0] }, + { + repos: h.repos, + embedder: null, + llm, + log, + bus: skillBus, + config: makeSkillConfig({ minSupport: 2, minGain: 0.1 }), + }, + ); + + expect(skillResult).toMatchObject({ evaluated: 1, crystallized: 1, rejected: 0 }); + expect(skillEvents).toContainEqual( + expect.objectContaining({ kind: "skill.crystallized", policyId: p1!.id }), + ); + expect(skillEvents.some((event) => event.kind === "skill.verification.failed")).toBe(false); + + const messages = skillPrompts[0] as Array<{ role: string; content: string }>; + expect(messages).toEqual( + expect.arrayContaining([ + expect.objectContaining({ role: "system", content: expect.stringContaining("简体中文") }), + ]), + ); + const payload = JSON.parse(messages.find((message) => message.role === "user")!.content) as { + evidence_tools: string[]; + }; + expect(payload.evidence_tools).toEqual(["pip", "execute_code"]); + expect(payload.evidence_tools).not.toContain('{"code":'); + + const storedSkill = h.repos.skills.list()[0]!; + expect(storedSkill.status).toBe("candidate"); + expect(storedSkill.sourcePolicyIds).toContain(p1!.id); + + const l3Events: L3Event[] = []; + const l3Bus = createL3EventBus(); + l3Bus.onAny((event) => l3Events.push(event)); + const l3Deps = { + repos: h.repos, + llm, + log, + bus: l3Bus, + config: l3Config(), + }; + + const created = await runL3( + { trigger: "l2.policy.induced", episodeId: p1!.sourceEpisodeIds[0] }, + l3Deps, + ); + expect(created.abstractions).toHaveLength(1); + expect(created.abstractions[0]).toMatchObject({ + clusterKey: "alpine|pip", + policyCount: 3, + createdNew: true, + skippedReason: null, + }); + expect(h.repos.worldModel.list()).toHaveLength(1); + expect(l3Events).toContainEqual(expect.objectContaining({ kind: "l3.world-model.created" })); + + const merged = await runL3({ trigger: "manual" }, l3Deps); + expect(merged.abstractions).toHaveLength(1); + expect(merged.abstractions[0]!.createdNew).toBe(false); + expect(merged.abstractions[0]!.mergedIntoWorldId).toBe(created.abstractions[0]!.worldModelId); + expect(h.repos.worldModel.list()).toHaveLength(1); + expect(l3Events).toContainEqual(expect.objectContaining({ kind: "l3.world-model.updated" })); + }); +}); diff --git a/apps/memos-local-plugin/tests/python/test_hermes_install_environment.py b/apps/memos-local-plugin/tests/python/test_hermes_install_environment.py new file mode 100644 index 000000000..978c5677f --- /dev/null +++ b/apps/memos-local-plugin/tests/python/test_hermes_install_environment.py @@ -0,0 +1,169 @@ +"""Execute Unix installer preflight against isolated Hermes layouts.""" + +import os +import subprocess +import sys + +from pathlib import Path + + +INSTALLER = Path(__file__).resolve().parents[2] / "install.sh" + + +def fixture_host(root: Path, env_name: str = "venv") -> Path: + for module in ("hermes_cli", "plugins", "plugins/memory"): + directory = root / module + directory.mkdir(parents=True, exist_ok=True) + (directory / "__init__.py").write_text( + "def load_memory_provider(name): return None\n" if module == "plugins/memory" else "" + ) + python = root / env_name / "bin/python" + python.parent.mkdir(parents=True) + python.symlink_to(sys.executable) + return python + + +def probe(tmp_path: Path, **overrides: str) -> subprocess.CompletedProcess: + source = INSTALLER.read_text() + start = source.find("resolve_hermes_environment() {") + assert start >= 0, "installer needs a validated Hermes environment resolver" + end = source.index("# ─── Hermes install", start) + env = {k: v for k, v in os.environ.items() if not k.startswith("HERMES_")} + env.update(HOME=str(tmp_path), PATH="/usr/bin:/bin") + env.update(overrides) + script = source[start:end] + "\nresolve_hermes_environment\n" + return subprocess.run( + ["/bin/bash", "-c", script], + cwd=tmp_path, + env=env, + capture_output=True, + text=True, + check=False, + timeout=30, + ) + + +def test_default_layout(tmp_path: Path) -> None: + root = tmp_path / ".hermes/hermes-agent" + python = fixture_host(root) + result = probe(tmp_path) + assert result.returncode == 0, result.stderr + assert result.stdout.splitlines() == [str(python), str(root / "plugins/memory"), str(root)] + + +def test_explicit_desktop_root_with_spaces(tmp_path: Path) -> None: + root = tmp_path / "Library/Application Support/Desktop/backend" + python = fixture_host(root, ".venv") + result = probe(tmp_path, HERMES_INSTALL_DIR=str(root)) + assert result.returncode == 0, result.stderr + assert result.stdout.splitlines()[0] == str(python) + + +def test_bash_launcher_with_spaces(tmp_path: Path) -> None: + root = tmp_path / "Desktop backend" + python = fixture_host(root) + cli = tmp_path / ".local/bin/hermes" + cli.parent.mkdir(parents=True) + cli.write_text(f'#!/usr/bin/env bash\nexec "{python}" "{root}/hermes" "$@"\n') + cli.chmod(0o755) + result = probe(tmp_path) + assert result.returncode == 0, result.stderr + assert result.stdout.splitlines()[0] == str(python) + + +def test_custom_data_home(tmp_path: Path) -> None: + home = tmp_path / "custom home" + root = home / "hermes-agent" + fixture_host(root) + result = probe(tmp_path, HERMES_HOME=str(home)) + assert result.returncode == 0, result.stderr + assert result.stdout.splitlines()[2] == str(root) + + +def test_explicit_python_and_root(tmp_path: Path) -> None: + root = tmp_path / "desktop source" + fixture_host(root) + result = probe(tmp_path, HERMES_INSTALL_DIR=str(root), HERMES_PYTHON=sys.executable) + assert result.returncode == 0, result.stderr + assert result.stdout.splitlines()[0] == sys.executable + + +def test_invalid_override_does_not_fall_back(tmp_path: Path) -> None: + fixture_host(tmp_path / ".hermes/hermes-agent") + result = probe(tmp_path, HERMES_PYTHON=str(tmp_path / "missing-python")) + assert result.returncode != 0 + assert "missing-python" in result.stderr + assert "HERMES_PYTHON" in result.stderr + + +def test_missing_host_api_keeps_import_error(tmp_path: Path) -> None: + root = tmp_path / "old-hermes" + fixture_host(root) + (root / "plugins/memory/__init__.py").write_text( + "raise ImportError('missing host dependency')\n" + ) + result = probe(tmp_path, HERMES_INSTALL_DIR=str(root)) + assert result.returncode != 0 + assert "missing host dependency" in result.stderr + assert "HERMES_INSTALL_DIR" in result.stderr + + +def test_environment_is_checked_before_install_side_effects() -> None: + source = INSTALLER.read_text() + body = source[source.index("install_hermes() {") :] + assert body.index("resolve_hermes_environment") < body.index("Stopping existing bridge daemon") + assert body.index("resolve_hermes_environment") < body.index("deploy_tarball_to_prefix") + + +def test_incompatible_host_api_is_rejected(tmp_path: Path) -> None: + root = tmp_path / "old-hermes" + fixture_host(root) + (root / "plugins/memory/__init__.py").write_text("") + result = probe(tmp_path, HERMES_INSTALL_DIR=str(root)) + assert result.returncode != 0 + assert "load_memory_provider" in result.stderr + + +def test_python_shebang_launcher(tmp_path: Path) -> None: + root = tmp_path / "custom-host" + python = fixture_host(root) + cli = tmp_path / ".local/bin/hermes" + cli.parent.mkdir(parents=True) + cli.write_text(f"#!{python}\n") + cli.chmod(0o755) + result = probe(tmp_path) + assert result.returncode == 0, result.stderr + assert result.stdout.splitlines()[0] == str(python) + + +def test_failed_preflight_does_not_stop_or_deploy(tmp_path: Path) -> None: + source = INSTALLER.read_text() + start = source.index("resolve_hermes_environment() {") + end = source.index("# ─── Main", start) + script = ( + """ +header() { :; } +step() { :; } +die() { printf '%s\\n' "$*" >&2; exit 1; } +pgrep() { echo unexpected-process-probe >&2; exit 99; } +deploy_tarball_to_prefix() { echo unexpected-deployment >&2; exit 99; } +""" + + source[start:end] + + "\ninstall_hermes\n" + ) + result = subprocess.run( + ["/bin/bash", "-c", script], + env={ + "HOME": str(tmp_path), + "PATH": "/usr/bin:/bin", + "HERMES_PYTHON": str(tmp_path / "missing"), + }, + capture_output=True, + text=True, + check=False, + timeout=30, + ) + assert result.returncode != 0 + assert "preflight failed" in result.stderr + assert "unexpected-" not in result.stderr + assert not (tmp_path / ".hermes").exists() diff --git a/apps/memos-local-plugin/tests/unit/install/hermes-provider-link.test.ts b/apps/memos-local-plugin/tests/unit/install/hermes-provider-link.test.ts index 6943c4807..ffa33f0a0 100644 --- a/apps/memos-local-plugin/tests/unit/install/hermes-provider-link.test.ts +++ b/apps/memos-local-plugin/tests/unit/install/hermes-provider-link.test.ts @@ -8,11 +8,33 @@ describe("Hermes provider install links", () => { it("main Unix installer links both checkout-local and user-level provider paths", () => { const source = readFileSync(path.join(repoRoot, "install.sh"), "utf8"); - expect(source).toContain('${HOME}/.hermes/plugins/memory'); + expect(source).toContain('local hermes_host_home="${HERMES_HOME:-${HOME}/.hermes}"'); + expect(source).toContain('${hermes_host_home}/plugins/memory'); expect(source).toContain('"${plugin_dir}/memtensor"'); expect(source).toContain('"${user_plugin_dir}/memtensor"'); }); + it("fails closed and restores provider targets when live verification fails", () => { + const source = readFileSync(path.join(repoRoot, "install.sh"), "utf8"); + + expect(source).toContain("restore_hermes_provider_targets"); + expect(source).toContain('die "Hermes memtensor provider verification failed.'); + expect(source).not.toContain('if [[ -e "${target}" && ! -L "${target}" ]]; then rm -rf "${target}"; fi'); + }); + + it("stages the Hermes package before stopping the host and can roll back the prefix", () => { + const source = readFileSync(path.join(repoRoot, "install.sh"), "utf8"); + const hermesStart = source.indexOf("install_hermes() {"); + const hermesEnd = source.indexOf("# ─── DSH install", hermesStart); + const hermesSource = source.slice(hermesStart, hermesEnd); + + expect(hermesSource).toContain('deploy_tarball_to_prefix "${staged_prefix}"'); + expect(hermesSource).toContain("rollback_hermes_install"); + expect(hermesSource.indexOf('deploy_tarball_to_prefix "${staged_prefix}"')).toBeLessThan( + hermesSource.indexOf('step "Stopping existing bridge daemon"'), + ); + }); + it("adapter Unix installer keeps a user-level provider link", () => { const source = readFileSync( path.join(repoRoot, "adapters/hermes/install.hermes.sh"), @@ -142,12 +164,14 @@ describe("Hermes provider install links", () => { it("main Unix installer uses atomic ln -sfn and prepares provider dir first", () => { const source = readFileSync(path.join(repoRoot, "install.sh"), "utf8"); - // cp runs BEFORE the loop now so the provider dir is populated before - // the second symlink is created. - const cpPos = source.indexOf('cp "${adapter_dir}/plugin.yaml"'); + // The staged provider is populated before any host symlink is created. + const cpPos = source.indexOf('cp "${staged_adapter_dir}/plugin.yaml"'); const loopPos = source.indexOf("provider_targets=("); expect(cpPos).toBeGreaterThan(0); expect(loopPos).toBeGreaterThan(cpPos); - expect(source).toContain('ln -sfn "${adapter_dir}/memos_provider" "${target}"'); + expect(source).toContain('ln -sfn "${backup_root}" "${target}"'); + expect(source).toContain( + 'prepare_hermes_provider_targets "${adapter_dir}/memos_provider" "${provider_targets[@]}"', + ); }); }); diff --git a/apps/memos-local-plugin/tests/unit/llm/client.test.ts b/apps/memos-local-plugin/tests/unit/llm/client.test.ts index cf891e376..51051ab00 100644 --- a/apps/memos-local-plugin/tests/unit/llm/client.test.ts +++ b/apps/memos-local-plugin/tests/unit/llm/client.test.ts @@ -58,11 +58,17 @@ class FakeProvider implements LlmProvider { class StreamingProvider implements LlmProvider { readonly name: LlmProviderName = "openai_compatible"; + public lastInput: ProviderCallInput | null = null; + async complete(): Promise { return { text: "full", durationMs: 1 }; } // eslint-disable-next-line require-yield - async *stream(): AsyncGenerator { + async *stream( + _messages: LlmMessage[], + opts: ProviderCallInput, + ): AsyncGenerator { + this.lastInput = opts; yield { delta: "he", done: false }; yield { delta: "llo", done: false }; yield { @@ -335,6 +341,60 @@ describe("llm/client", () => { ); }); + // ─── op propagation to providers (issue #2308) ──────────────────────── + // + // The facade must forward `opts.op` onto the `ProviderCallInput` handed to + // `provider.complete()` / `provider.stream()` so per-op provider behavior + // (e.g. request-body tweaks, routing overrides, reasoning kill-switches + // keyed on `capture.summarize`) can fire. Dropping it silently makes + // those switches unreachable. + describe("op propagation (issue #2308)", () => { + it("complete forwards opts.op onto the provider input", async () => { + const fake = new FakeProvider("openai_compatible", () => ({ text: "ok", durationMs: 1 })); + const client = createLlmClientWithProvider(cfg(), fake); + await client.complete("hi", { op: "capture.summarize" }); + expect(fake.lastInput?.op).toBe("capture.summarize"); + }); + + it("completeJson forwards opts.op onto the provider input", async () => { + const fake = new FakeProvider("openai_compatible", () => ({ + text: '{"a":1}', + durationMs: 1, + })); + const client = createLlmClientWithProvider(cfg(), fake); + await client.completeJson<{ a: number }>("score", { op: "retrieval.filter" }); + expect(fake.lastInput?.op).toBe("retrieval.filter"); + }); + + it("stream forwards opts.op onto the provider input (non-streaming provider)", async () => { + // FakeProvider has no stream(); the facade wraps complete() in a + // single-chunk iterable, which still exercises buildCallInput. + const fake = new FakeProvider("openai_compatible", () => ({ text: "one", durationMs: 1 })); + const client = createLlmClientWithProvider(cfg(), fake); + const parts: string[] = []; + for await (const c of client.stream("x", { op: "skill.evolve" })) { + if (!c.done) parts.push(c.delta); + } + expect(fake.lastInput?.op).toBe("skill.evolve"); + }); + + it("stream forwards opts.op onto a native streaming provider", async () => { + const provider = new StreamingProvider(); + const client = createLlmClientWithProvider(cfg(), provider); + for await (const _chunk of client.stream("x", { op: "capture.summarize" })) { + // Consume the stream so the provider receives the cooked input. + } + expect(provider.lastInput?.op).toBe("capture.summarize"); + }); + + it("leaves op undefined when the caller supplies no op", async () => { + const fake = new FakeProvider("openai_compatible", () => ({ text: "ok", durationMs: 1 })); + const client = createLlmClientWithProvider(cfg(), fake); + await client.complete("hi"); + expect(fake.lastInput?.op).toBeUndefined(); + }); + }); + // ─── Circuit breaker (issue #1897) ────────────────────────────────────── describe("circuit breaker", () => { function statusSink(): { rows: LlmStatusDetail[]; push: (d: LlmStatusDetail) => void } { diff --git a/apps/memos-local-plugin/tests/unit/llm/prompts.test.ts b/apps/memos-local-plugin/tests/unit/llm/prompts.test.ts index bee52730a..1e35b5e7f 100644 --- a/apps/memos-local-plugin/tests/unit/llm/prompts.test.ts +++ b/apps/memos-local-plugin/tests/unit/llm/prompts.test.ts @@ -43,6 +43,13 @@ describe("llm/prompts", () => { it("detectDominantLanguage only chooses Chinese when CJK dominates", () => { expect(detectDominantLanguage(["请修复这个问题,并解释原因"])).toBe("zh"); + expect( + detectDominantLanguage([ + "在 Alpine 镜像中安装 cryptography 失败", + "先执行 apk add openssl-dev,再执行 pip install cryptography", + "先安装系统库,再重试 pip 安装", + ]), + ).toBe("zh"); expect(detectDominantLanguage(["Excelファイルの欠落値を復元してください"])).toBe("en"); expect(detectDominantLanguage(["저는 GRPO를 사용하여 모델을 훈련시키고 있습니다"])).toBe("en"); expect(detectDominantLanguage(["GRPO / TRL / reward_fn.py"])).toBe("en"); diff --git a/apps/memos-local-plugin/tests/unit/memory/l2/induce.test.ts b/apps/memos-local-plugin/tests/unit/memory/l2/induce.test.ts index 06d61459f..52b3b3e73 100644 --- a/apps/memos-local-plugin/tests/unit/memory/l2/induce.test.ts +++ b/apps/memos-local-plugin/tests/unit/memory/l2/induce.test.ts @@ -182,6 +182,7 @@ describe("memory/l2/induce", () => { episodeIds: ["ep_1", "ep_2"] as EpisodeId[], evidenceTraces: [mkTrace("tr_a", "ep_1", vec([1, 0])), mkTrace("tr_b", "ep_2", vec([0, 1]))], inducedBy: "l2.l2.induction.v1", + sourceSignature: "docker|pip|pip.install|MODULE_NOT_FOUND", now: 42, }); expect(row.status).toBe("candidate"); @@ -191,5 +192,13 @@ describe("memory/l2/induce", () => { expect(row.vec).not.toBeNull(); expect(row.createdAt).toBe(42); expect(row.title).toBe("t"); + expect(row.metadata).toEqual(expect.objectContaining({ + version: 1, + language: "en", + domainTags: ["docker", "pip"], + toolNames: ["pip.install"], + errorCodes: ["module_not_found"], + sourceSignature: "docker|pip|pip.install|MODULE_NOT_FOUND", + })); }); }); diff --git a/apps/memos-local-plugin/tests/unit/memory/l2/l2.integration.test.ts b/apps/memos-local-plugin/tests/unit/memory/l2/l2.integration.test.ts index d9d062d1f..7cbd6b303 100644 --- a/apps/memos-local-plugin/tests/unit/memory/l2/l2.integration.test.ts +++ b/apps/memos-local-plugin/tests/unit/memory/l2/l2.integration.test.ts @@ -186,6 +186,12 @@ describe("memory/l2/integration", () => { const persisted = handle.repos.policies.getById(induced.policyId!)!; expect(persisted.status).toBe("candidate"); expect(persisted.sourceEpisodeIds.sort()).toEqual(["ep_A", "ep_B"]); + expect(persisted.metadata).toEqual(expect.objectContaining({ + version: 1, + domainTags: ["docker", "pip"], + toolNames: ["pip.install"], + sourceSignature: "docker|pip|pip.install|MODULE_NOT_FOUND", + })); // ── A third run with a trace that cosine-matches the new policy should // associate (not re-induce) and bump gain/support. diff --git a/apps/memos-local-plugin/tests/unit/memory/l3/abstract.test.ts b/apps/memos-local-plugin/tests/unit/memory/l3/abstract.test.ts index ff759b17f..2bbefa199 100644 --- a/apps/memos-local-plugin/tests/unit/memory/l3/abstract.test.ts +++ b/apps/memos-local-plugin/tests/unit/memory/l3/abstract.test.ts @@ -168,6 +168,20 @@ describe("memory/l3/abstract", () => { expect(res.reason).toBe("llm_disabled"); }); + it("skips oversized legacy batches before calling the LLM", async () => { + const llm = fakeLlm({ completeJson: { [OP]: {} } }); + const cluster = mkCluster(); + cluster.policies[0]!.procedure = "x".repeat(5_000); + const res = await abstractDraft( + { cluster, evidenceByPolicy: new Map() }, + { llm, log, config: cfg({ maxPromptChars: 4_000, policyCharCap: 4_000 }) }, + ); + expect(res.ok).toBe(false); + if (res.ok) return; + expect(res.reason).toBe("prompt_too_large"); + expect(llm.stats().requests).toBe(0); + }); + it("returns llm_failed when the LLM throws — never rethrows", async () => { const llm = throwingLlm(new Error("boom")); const res = await abstractDraft( @@ -227,4 +241,3 @@ describe("memory/l3/abstract", () => { expect(row.body).toContain("Environment"); }); }); - diff --git a/apps/memos-local-plugin/tests/unit/memory/l3/cluster.test.ts b/apps/memos-local-plugin/tests/unit/memory/l3/cluster.test.ts index 737b602d9..32e57684a 100644 --- a/apps/memos-local-plugin/tests/unit/memory/l3/cluster.test.ts +++ b/apps/memos-local-plugin/tests/unit/memory/l3/cluster.test.ts @@ -33,6 +33,7 @@ function mkPolicy(partial: Partial & { id: PolicyId }): PolicyRow { vec: partial.vec ?? vec([1, 0, 0]), createdAt: NOW, updatedAt: NOW, + metadata: partial.metadata, }; } @@ -61,6 +62,63 @@ describe("memory/l3/cluster", () => { expect(tags).toEqual([]); }); + it("prefers structured trace metadata over English-only policy prose", () => { + const p = mkPolicy({ + id: "po_zh" as PolicyId, + title: "处理依赖安装失败", + procedure: "执行包管理器并重试", + metadata: { + version: 1, + language: "zh", + domainTags: ["python", "alpine"], + toolNames: ["pip.install"], + errorCodes: ["module_not_found"], + sourceSignature: "python|alpine|pip.install|MODULE_NOT_FOUND", + }, + }); + expect(domainKeyOf(p)).toEqual({ + key: "python|pip.install", + tags: expect.arrayContaining(["python", "alpine", "pip.install", "module_not_found"]), + }); + }); + + it("keeps the legacy text fallback when backfilled metadata has no tags", () => { + const p = mkPolicy({ + id: "po_legacy" as PolicyId, + title: "修复 Docker 中的 pip 安装失败", + trigger: "pip install fails in Alpine container", + procedure: "apk add build tools before pip install", + metadata: { + version: 1, + language: "mixed", + domainTags: [], + toolNames: [], + errorCodes: [], + }, + }); + const { key, tags } = domainKeyOf(p); + expect(key).toContain("docker"); + expect(key).toContain("pip"); + expect(tags).toEqual(expect.arrayContaining(["docker", "alpine", "pip"])); + }); + + it("does not create a cluster from an untagged bucket", () => { + const policies = [ + [1, 0, 0], + [0, 1, 0], + [0, 0, 1], + ].map((v, n) => mkPolicy({ + id: `po_u${n}` as PolicyId, + title: `中文策略 ${n}`, + vec: vec(v), + })); + const clusters = clusterPolicies( + { policies }, + { config: { clusterMinSimilarity: 0.6, minPolicies: 1, maxPoliciesPerCluster: 20 } }, + ); + expect(clusters).toHaveLength(0); + }); + it("groups network-related text under 'network'", () => { const p = mkPolicy({ id: "po_n" as PolicyId, @@ -108,6 +166,20 @@ describe("memory/l3/cluster", () => { ); }); + it("preserves every policy so prompt batching can process the full cluster", () => { + const policies = [1, 2, 3].map((n) => mkPolicy({ + id: `po_n${n}` as PolicyId, + title: `network retry ${n}`, + trigger: "proxy DNS failure", + vec: vec([1, 0, 0]), + })); + const clusters = clusterPolicies( + { policies }, + { config: { clusterMinSimilarity: 0.99, minPolicies: 1, maxPoliciesPerCluster: 2 } }, + ); + expect(clusters[0]!.policies).toHaveLength(3); + }); + it("skips a bucket that doesn't meet minPolicies", () => { const policies = [ mkPolicy({ @@ -184,30 +256,25 @@ describe("memory/l3/cluster", () => { expect(keys.some((k) => k.includes("node") || k.includes("npm"))).toBe(true); }); - it("falls back to loose admission when strict subset is too small but bucket survives", () => { - // All three policies share the same domain key (`python|_`) but - // their vectors point in mutually-orthogonal directions, so the - // strict (cosine ≥ minSimilarity) subset would be empty. The - // bucket itself satisfies minPolicies, so `cluster.ts` should - // fall back to admitting the WHOLE bucket as a `loose` cluster. + it("does not cluster unrelated untagged policies", () => { const policies = [ mkPolicy({ id: "po_validate" as PolicyId, title: "validate python syntax", - trigger: "after writing python files", - procedure: "python -m py_compile ", + trigger: "after writing source files", + procedure: "run the syntax checker on the changed file", vec: vec([1, 0, 0]), }), mkPolicy({ id: "po_cli" as PolicyId, - title: "register python CLI subcommand", + title: "register a command line subcommand", trigger: "adding a new task verb", procedure: "register(subparsers) + handler() -> int", vec: vec([0, 1, 0]), }), mkPolicy({ id: "po_storage" as PolicyId, - title: "implement python storage backend", + title: "implement a storage backend", trigger: "new persistence format requested", procedure: "implement load/save with UTF-8", vec: vec([0, 0, 1]), @@ -217,16 +284,22 @@ describe("memory/l3/cluster", () => { { policies }, { config: { clusterMinSimilarity: 0.6, minPolicies: 2 } }, ); - expect(clusters.length).toBe(1); - const c = clusters[0]!; - expect(c.admission).toBe("loose"); - expect(c.policies.length).toBe(3); - // Centroid of three orthogonal unit vectors gives mean cosine - // 1/sqrt(3) ≈ 0.577 — strictly less than 0.6 (the strict floor), - // confirming we landed in the loose fallback for the right - // reason and not because of a bug elsewhere. - expect(c.cohesion).toBeLessThan(0.6); - expect(c.cohesion).toBeGreaterThan(0.49); + expect(clusters).toEqual([]); + }); + + it("clusters similar untagged policies without dropping prompt overflow", () => { + const policies = [1, 2, 3].map((n) => mkPolicy({ + id: `po_uv${n}` as PolicyId, + title: `中文策略 ${n}`, + vec: vec([1, n * 0.01, 0]), + })); + const clusters = clusterPolicies( + { policies }, + { config: { clusterMinSimilarity: 0.6, minPolicies: 2, maxPoliciesPerCluster: 2 } }, + ); + expect(clusters).toHaveLength(1); + expect(clusters[0]!.policies).toHaveLength(3); + expect(clusters[0]!.admission).toBe("strict"); }); it("filters outliers below clusterMinSimilarity", () => { diff --git a/apps/memos-local-plugin/tests/unit/memory/l3/l3.integration.test.ts b/apps/memos-local-plugin/tests/unit/memory/l3/l3.integration.test.ts index 51ba5d4b8..7dc1cd811 100644 --- a/apps/memos-local-plugin/tests/unit/memory/l3/l3.integration.test.ts +++ b/apps/memos-local-plugin/tests/unit/memory/l3/l3.integration.test.ts @@ -38,6 +38,16 @@ import { const OP = `${L3_ABSTRACTION_PROMPT.id}.v${L3_ABSTRACTION_PROMPT.version}`; const log = rootLogger.child({ channel: "core.memory.l3" }); +const validDraft = { + title: "Alpine python dependency model", + domain_tags: ["docker", "alpine", "pip"], + environment: [{ label: "musl libc", description: "no glibc" }], + inference: [{ label: "binary wheels fail", description: "compile from source" }], + constraints: [], + body: "# summary", + confidence: 0.75, + supersedes_world_ids: [], +}; function cfg(overrides: Partial = {}): L3Config { return { @@ -155,6 +165,63 @@ describe("memory/l3/integration", () => { ); }); + it("batches a large cluster without dropping policies or creating duplicate world models", async () => { + for (let index = 1; index <= 5; index++) { + const episodeId = `ep_batch_${index}`; + seedPolicy(handle, { + id: `po_batch_${index}` as PolicyId, + title: `Alpine pip dependency ${index}`, + trigger: "pip install fails in Alpine container", + procedure: `apk add dependency-${index} then pip install`, + sourceEpisodeIds: [episodeId as EpisodeId], + vec: vec([1, index * 0.01, 0]), + }); + seedTrace(handle, { + id: `tr_batch_${index}`, + episodeId, + tags: ["docker", "alpine", "pip"], + }); + } + let calls = 0; + const llm = fakeLlm({ + completeJson: { + [OP]: () => { + calls += 1; + return { + ...validDraft, + environment: [{ label: `batch ${calls}`, description: "covered" }], + }; + }, + }, + }); + + const result = await runL3( + { trigger: "manual" }, + { + repos: { + policies: handle.repos.policies, + traces: handle.repos.traces, + worldModel: handle.repos.worldModel, + kv: handle.repos.kv, + }, + llm, + log, + config: cfg({ minPolicies: 1, maxPoliciesPerCluster: 2 }), + }, + ); + + expect(calls).toBe(3); + expect(result.abstractions).toHaveLength(1); + expect(handle.repos.worldModel.list()).toHaveLength(1); + expect(handle.repos.worldModel.list()[0]!.policyIds.map(String).sort()).toEqual([ + "po_batch_1", + "po_batch_2", + "po_batch_3", + "po_batch_4", + "po_batch_5", + ]); + }); + it("merges into an existing WM that covers the same domain", async () => { seedTriplet(); // Seed a prior WM that shares domain tags + vector, so merge kicks in. @@ -237,6 +304,110 @@ describe("memory/l3/integration", () => { ); expect(res.abstractions.every((a) => a.skippedReason === "llm_disabled")).toBe(true); expect(handle.repos.worldModel.list().length).toBe(0); + expect(handle.repos.kv.all().filter((row) => row.key.startsWith("l3.retry."))).toEqual([]); + }); + + it("backs off a failed abstraction, retries after expiry, and clears retry state", async () => { + seedTriplet(); + let calls = 0; + const llm = fakeLlm({ + completeJson: { + [OP]: () => { + calls++; + if (calls === 1) throw new Error("temporary failure"); + return validDraft; + }, + }, + }); + const deps = { + repos: { + policies: handle.repos.policies, + traces: handle.repos.traces, + worldModel: handle.repos.worldModel, + kv: handle.repos.kv, + }, + llm, + log, + config: cfg(), + }; + + const failed = await runL3({ trigger: "manual", now: NOW }, deps); + expect(failed.abstractions[0]!.skippedReason).toBe("llm_failed"); + expect(calls).toBe(1); + expect(handle.repos.kv.all().some((row) => row.key.startsWith("l3.retry."))).toBe(true); + + const deferred = await runL3({ trigger: "manual", now: NOW + 299_999 }, deps); + expect(deferred.abstractions[0]!.skippedReason).toBe("retry_cooldown"); + expect(calls).toBe(1); + + const retried = await runL3({ trigger: "manual", now: NOW + 300_000 }, deps); + expect(retried.abstractions[0]!.skippedReason).toBeNull(); + expect(calls).toBe(2); + expect(handle.repos.kv.all().filter((row) => row.key.startsWith("l3.retry."))).toEqual([]); + }); + + it("quarantines a deterministically failing legacy cluster after bounded attempts", async () => { + seedTriplet(); + let calls = 0; + const llm = fakeLlm({ + completeJson: { + [OP]: () => { + calls++; + throw new Error("malformed legacy response"); + }, + }, + }); + const deps = { + repos: { + policies: handle.repos.policies, + traces: handle.repos.traces, + worldModel: handle.repos.worldModel, + kv: handle.repos.kv, + }, + llm, + log, + config: cfg(), + }; + + for (const at of [0, 300_000, 2_100_000, 9_300_000]) { + const result = await runL3({ trigger: "manual", now: NOW + at }, deps); + expect(result.abstractions[0]!.skippedReason).toBe("llm_failed"); + } + expect(calls).toBe(4); + + const quarantined = await runL3({ trigger: "manual", now: NOW + 100_000_000 }, deps); + expect(quarantined.abstractions[0]!.skippedReason).toBe("quarantined"); + expect(calls).toBe(4); + const state = handle.repos.kv.all().find((row) => row.key.startsWith("l3.retry.")); + expect(state?.value).toMatchObject({ version: 2, failures: 4, quarantined: true }); + }); + + it("records retry state instead of success cooldown when persistence fails", async () => { + seedTriplet(); + const worldModel = { + ...handle.repos.worldModel, + insert: () => { + throw new Error("disk full"); + }, + }; + const result = await runL3( + { trigger: "manual", now: NOW }, + { + repos: { + policies: handle.repos.policies, + traces: handle.repos.traces, + worldModel, + kv: handle.repos.kv, + }, + llm: fakeLlm({ completeJson: { [OP]: validDraft } }), + log, + config: cfg({ cooldownDays: 1 }), + }, + ); + + expect(result.warnings.some((warning) => warning.stage === "insert")).toBe(true); + expect(handle.repos.kv.all().some((row) => row.key.startsWith("l3.retry."))).toBe(true); + expect(handle.repos.kv.all().some((row) => row.key.startsWith("l3.lastRun."))).toBe(false); }); it("adjustConfidence clamps in [0,1] and emits an event", async () => { diff --git a/apps/memos-local-plugin/tests/unit/pipeline/capture-reflect-llm-wiring.test.ts b/apps/memos-local-plugin/tests/unit/pipeline/capture-reflect-llm-wiring.test.ts index 3eede9e1c..20013e980 100644 --- a/apps/memos-local-plugin/tests/unit/pipeline/capture-reflect-llm-wiring.test.ts +++ b/apps/memos-local-plugin/tests/unit/pipeline/capture-reflect-llm-wiring.test.ts @@ -29,6 +29,46 @@ const captureRunnerCalls: Array<{ llm: LlmClient | null; reflectLlm: LlmClient | null; }> = []; +const evolutionSubscriberCalls: Array<{ + l2Llm: LlmClient | null; + skillLlm: LlmClient | null; +}> = []; + +vi.mock("../../../core/memory/l2/index.js", async () => { + const actual = await vi.importActual< + typeof import("../../../core/memory/l2/index.js") + >("../../../core/memory/l2/index.js"); + return { + ...actual, + attachL2Subscriber: (deps: { llm: LlmClient | null; [k: string]: unknown }) => { + evolutionSubscriberCalls.push({ + l2Llm: deps.llm, + skillLlm: null, + }); + return actual.attachL2Subscriber( + deps as Parameters[0], + ); + }, + }; +}); + +vi.mock("../../../core/skill/index.js", async () => { + const actual = await vi.importActual< + typeof import("../../../core/skill/index.js") + >("../../../core/skill/index.js"); + return { + ...actual, + attachSkillSubscriber: (deps: { llm: LlmClient | null; [k: string]: unknown }) => { + evolutionSubscriberCalls.push({ + l2Llm: null, + skillLlm: deps.llm, + }); + return actual.attachSkillSubscriber( + deps as Parameters[0], + ); + }, + }; +}); vi.mock("../../../core/capture/index.js", async () => { const actual = await vi.importActual< @@ -141,6 +181,7 @@ function buildDepsWithDistinctLlms( beforeEach(() => { dbHandle = makeTmpDb(); captureRunnerCalls.length = 0; + evolutionSubscriberCalls.length = 0; }); afterEach(() => { @@ -189,4 +230,27 @@ describe("pipeline/deps captureRunner wiring (issue #2148)", () => { expect(call.reflectLlm).toBe(call.llm); expect(call.reflectLlm?.model).toBe("main-llm"); }); + + it("passes the dedicated skill-evolver model to L2 induction and skill crystallization", () => { + const buses = buildPipelineBuses(); + const deps = buildDepsWithDistinctLlms(dbHandle!, false); + const algorithm = extractAlgorithmConfig(deps); + const session = buildPipelineSession(deps, buses.session); + buildPipelineSubscribers(deps, buses, algorithm, session); + + expect(evolutionSubscriberCalls[0].l2Llm?.model).toBe("skill-evolver-llm"); + expect(evolutionSubscriberCalls[1].skillLlm?.model).toBe("skill-evolver-llm"); + expect(evolutionSubscriberCalls[0].l2Llm).toBe(evolutionSubscriberCalls[1].skillLlm); + }); + + it("inherits the main model when no dedicated evolver client is available", () => { + const buses = buildPipelineBuses(); + const deps = buildDepsWithDistinctLlms(dbHandle!, false); + deps.reflectLlm = null; + buildPipelineSubscribers(deps, buses, extractAlgorithmConfig(deps)); + + expect(evolutionSubscriberCalls[0].l2Llm?.model).toBe("main-llm"); + expect(evolutionSubscriberCalls[1].skillLlm?.model).toBe("main-llm"); + }); + }); diff --git a/apps/memos-local-plugin/tests/unit/skill/events.test.ts b/apps/memos-local-plugin/tests/unit/skill/events.test.ts index 2643cff16..8b1bf6425 100644 --- a/apps/memos-local-plugin/tests/unit/skill/events.test.ts +++ b/apps/memos-local-plugin/tests/unit/skill/events.test.ts @@ -48,4 +48,18 @@ describe("skill/events", () => { }); expect(called).toEqual(["second"]); }); + + it("preserves the policy that produced a verification failure", () => { + const bus = createSkillEventBus(); + const seen: SkillEvent[] = []; + bus.on("skill.verification.failed", (event) => seen.push(event)); + bus.emit({ + kind: "skill.verification.failed", + at: 1, + skillId: "sk_1" as SkillId, + policyId: "po_1", + reason: "resonance=0.00<0.5", + }); + expect(seen[0]).toMatchObject({ policyId: "po_1" }); + }); }); diff --git a/apps/memos-local-plugin/tests/unit/skill/skill.integration.test.ts b/apps/memos-local-plugin/tests/unit/skill/skill.integration.test.ts index 476134abc..6f1c22c80 100644 --- a/apps/memos-local-plugin/tests/unit/skill/skill.integration.test.ts +++ b/apps/memos-local-plugin/tests/unit/skill/skill.integration.test.ts @@ -14,6 +14,7 @@ import { createL2EventBus } from "../../../core/memory/l2/index.js"; import { fakeLlm } from "../../helpers/fake-llm.js"; import { makeTmpDb, type TmpDbHandle } from "../../helpers/tmp-db.js"; import type { EpisodeId, PolicyId, SkillId } from "../../../core/types.js"; +import type { ToolCallDTO } from "../../../agent-contract/dto.js"; import { makeDraft, makeSkillConfig, @@ -92,6 +93,98 @@ function seedFullCandidate(h: TmpDbHandle): { } describe("skill/runSkill (integration)", () => { + it("crystallizes a Chinese multi-episode tool workflow without resonance or tool coverage loss", async () => { + const h = open(); + const episodeIds = ["ep_cn_1", "ep_cn_2"] as EpisodeId[]; + const toolCall: ToolCallDTO = { + name: "execute_code", + input: '{"code":"print(1)"}', + }; + + for (const [index, episodeId] of episodeIds.entries()) { + const sessionId = `s_cn_${index}`; + seedSessionOnly(h, sessionId); + seedTrace(h, { + episodeId, + sessionId, + userText: "在 Alpine 镜像中安装 cryptography 失败", + agentText: "先执行 apk add openssl-dev,再执行 pip install cryptography", + reflection: "先安装系统库,再重试 pip 安装", + value: 0.9, + tags: ["alpine", "pip"], + toolCalls: [toolCall], + }); + } + + const policy = seedPolicy(h, { + id: "po_cn_workflow" as PolicyId, + title: "Alpine 中先安装系统库再重试 pip", + trigger: "Alpine 镜像中的 cryptography 安装失败", + procedure: "先执行 apk add openssl-dev,再执行 pip install cryptography", + verification: "cryptography 安装成功", + boundary: "仅适用于 Alpine 镜像", + support: 2, + gain: 0.3, + sourceEpisodeIds: episodeIds, + }); + + const prompts: unknown[] = []; + const { deps, events } = makeDeps(h, { + llm: fakeLlm({ + completeJson: { + "skill.crystallize": (input) => { + prompts.push(input); + return makeDraft({ + name: "alpine_cryptography_fix", + displayTitle: "Alpine cryptography 安装修复", + summary: "在 Alpine 中先安装系统库,再重试 pip 安装 cryptography", + preconditions: ["当前环境是 Alpine 镜像"], + steps: [ + { title: "检查错误", body: "确认 pip install cryptography 安装失败" }, + { title: "安装系统库", body: "执行 apk add openssl-dev" }, + { title: "重新安装", body: "再次执行 pip install cryptography" }, + ], + tools: ["execute_code"], + tags: ["alpine", "pip"], + }); + }, + }, + }), + config: makeSkillConfig({ minSupport: 2, minGain: 0.1 }), + }); + + const result = await runSkill({ trigger: "manual", policyId: policy.id }, deps); + + expect(result).toMatchObject({ evaluated: 1, crystallized: 1, rejected: 0 }); + expect(events).toContainEqual( + expect.objectContaining({ + kind: "skill.crystallized", + policyId: policy.id, + }), + ); + expect(events.some((event) => event.kind === "skill.verification.failed")).toBe(false); + + const messages = prompts[0] as Array<{ role: string; content: string }>; + expect(messages).toEqual( + expect.arrayContaining([ + expect.objectContaining({ + role: "system", + content: expect.stringContaining("简体中文"), + }), + ]), + ); + const userMessage = messages.find((message) => message.role === "user"); + expect(userMessage).toBeDefined(); + const payload = JSON.parse(userMessage!.content) as { evidence_tools: string[] }; + expect(payload.evidence_tools).toEqual(["execute_code"]); + expect(payload.evidence_tools).not.toContain('{"code":'); + + const stored = h.repos.skills.list()[0]!; + expect(stored.status).toBe("candidate"); + expect(stored.sourcePolicyIds).toContain(policy.id); + expect(stored.procedureJson?.steps).toHaveLength(3); + }); + it("crystallizes a fresh skill for an eligible policy", async () => { const h = open(); const { policyId } = seedFullCandidate(h); diff --git a/apps/memos-local-plugin/tests/unit/skill/subscriber.test.ts b/apps/memos-local-plugin/tests/unit/skill/subscriber.test.ts index a8c35468a..ba4463298 100644 --- a/apps/memos-local-plugin/tests/unit/skill/subscriber.test.ts +++ b/apps/memos-local-plugin/tests/unit/skill/subscriber.test.ts @@ -2,6 +2,7 @@ import { describe, it, expect, afterEach, vi } from "vitest"; import { createL2EventBus } from "../../../core/memory/l2/events.js"; import { createRewardEventBus } from "../../../core/reward/events.js"; +import type { RewardResult } from "../../../core/reward/types.js"; import { attachSkillSubscriber, createSkillEventBus, @@ -49,6 +50,33 @@ function seedTracesForPolicy(h: TmpDbHandle, id: PolicyId) { return { episodeId }; } +function minimalRewardResult(episodeId: EpisodeId): RewardResult { + return { + episodeId, + sessionId: "s_reward", + rHuman: 0.8, + humanScore: { + rHuman: 0.8, + axes: { goalAchievement: 1, processQuality: 1, userSatisfaction: 1 }, + reason: "test", + source: "heuristic", + model: null, + }, + feedbackCount: 0, + backprop: { + updates: [], + meanAbsValue: 0, + maxPriority: 0, + echoParams: { gamma: 0.9, decayHalfLifeDays: 30, now: Date.now() }, + }, + traceIds: [], + timings: { summary: 0, score: 0, backprop: 0, persist: 0, total: 0 }, + warnings: [], + startedAt: Date.now(), + completedAt: Date.now(), + }; +} + describe("skill/subscriber", () => { it("triggers runSkill on l2.policy.induced", async () => { handle = makeTmpDb(); @@ -156,6 +184,128 @@ describe("skill/subscriber", () => { sub.dispose(); }); + it("runs reward updates only for policies linked to that episode", async () => { + handle = makeTmpDb(); + const h = handle; + const episodeA = seedTracesForPolicy(h, "po_reward_a" as PolicyId).episodeId; + const episodeB = seedTracesForPolicy(h, "po_reward_b" as PolicyId).episodeId; + const policyA = seedPolicy(h, { id: "po_reward_a" as PolicyId, sourceEpisodeIds: [episodeA] }); + seedPolicy(h, { id: "po_reward_b" as PolicyId, sourceEpisodeIds: [episodeB] }); + const traceA = h.repos.traces.listAllForEpisode(episodeA)[0]!; + h.repos.tracePolicyLinks.link({ traceId: traceA.id, policyId: policyA.id, episodeId: episodeA }); + + const rewardBus = createRewardEventBus(); + let draftNumber = 0; + const sub = attachSkillSubscriber({ + l2Bus: createL2EventBus(), + rewardBus, + bus: createSkillEventBus(), + repos: h.repos, + embedder: null, + llm: fakeLlm({ + completeJson: { + "skill.crystallize": () => { + draftNumber += 1; + return makeDraft({ name: `reward_policy_${draftNumber}` }); + }, + }, + }), + log: rootLogger.child({ channel: "core.skill.subscriber" }), + config: makeSkillConfig({ cooldownMs: 6 * 60 * 60 * 1000 }), + }); + + rewardBus.emit({ + kind: "reward.updated", + result: minimalRewardResult(episodeA), + }); + await sub.flush(); + + expect(h.repos.skills.list().map((skill) => skill.sourcePolicyIds.flat()).flat()).toEqual([policyA.id]); + sub.dispose(); + }); + + it("does not let one policy cooldown block another policy reward update", async () => { + handle = makeTmpDb(); + const h = handle; + const episodeA = seedTracesForPolicy(h, "po_reward_c" as PolicyId).episodeId; + const episodeB = seedTracesForPolicy(h, "po_reward_d" as PolicyId).episodeId; + const policyA = seedPolicy(h, { id: "po_reward_c" as PolicyId, sourceEpisodeIds: [episodeA] }); + const policyB = seedPolicy(h, { id: "po_reward_d" as PolicyId, sourceEpisodeIds: [episodeB] }); + for (const [policy, episode] of [[policyA, episodeA], [policyB, episodeB]] as const) { + const trace = h.repos.traces.listAllForEpisode(episode)[0]!; + h.repos.tracePolicyLinks.link({ traceId: trace.id, policyId: policy.id, episodeId: episode }); + } + + const rewardBus = createRewardEventBus(); + let draftNumber = 0; + const sub = attachSkillSubscriber({ + l2Bus: createL2EventBus(), + rewardBus, + bus: createSkillEventBus(), + repos: h.repos, + embedder: null, + llm: fakeLlm({ + completeJson: { + "skill.crystallize": () => { + draftNumber += 1; + return makeDraft({ name: `reward_cooldown_${draftNumber}` }); + }, + }, + }), + log: rootLogger.child({ channel: "core.skill.subscriber" }), + config: makeSkillConfig({ cooldownMs: 6 * 60 * 60 * 1000 }), + }); + + rewardBus.emit({ kind: "reward.updated", result: minimalRewardResult(episodeA) }); + await sub.flush(); + rewardBus.emit({ kind: "reward.updated", result: minimalRewardResult(episodeB) }); + await sub.flush(); + + expect(h.repos.skills.list().flatMap((skill) => skill.sourcePolicyIds.map(String)).sort()).toEqual([ + String(policyA.id), + String(policyB.id), + ]); + sub.dispose(); + }); + + it("retains a same-policy reward event received during cooldown", async () => { + handle = makeTmpDb(); + const h = handle; + const episode = seedTracesForPolicy(h, "po_reward_pending" as PolicyId).episodeId; + const policy = seedPolicy(h, { + id: "po_reward_pending" as PolicyId, + sourceEpisodeIds: [episode], + }); + const trace = h.repos.traces.listAllForEpisode(episode)[0]!; + h.repos.tracePolicyLinks.link({ traceId: trace.id, policyId: policy.id, episodeId: episode }); + + const rewardBus = createRewardEventBus(); + const bus = createSkillEventBus(); + let eligibilityRuns = 0; + bus.on("skill.eligibility.checked", () => { + eligibilityRuns += 1; + }); + const sub = attachSkillSubscriber({ + l2Bus: createL2EventBus(), + rewardBus, + bus, + repos: h.repos, + embedder: null, + llm: fakeLlm({ completeJson: { "skill.crystallize": makeDraft() } }), + log: rootLogger.child({ channel: "core.skill.subscriber" }), + config: makeSkillConfig({ cooldownMs: 20 }), + }); + + rewardBus.emit({ kind: "reward.updated", result: minimalRewardResult(episode) }); + await sub.flush(); + rewardBus.emit({ kind: "reward.updated", result: minimalRewardResult(episode) }); + await new Promise((resolve) => setTimeout(resolve, 30)); + await sub.flush(); + + expect(eligibilityRuns).toBe(2); + sub.dispose(); + }); + it("archives each stale low-η active skill once without regressing candidate promotion", async () => { handle = makeTmpDb(); const h = handle; diff --git a/apps/memos-local-plugin/tests/unit/skill/verifier.test.ts b/apps/memos-local-plugin/tests/unit/skill/verifier.test.ts index d74ac9d09..930c95aff 100644 --- a/apps/memos-local-plugin/tests/unit/skill/verifier.test.ts +++ b/apps/memos-local-plugin/tests/unit/skill/verifier.test.ts @@ -119,6 +119,22 @@ describe("skill/verifier", () => { expect(r.coverage).toBe(1); }); + it("does not treat JSON string arguments as command names", () => { + const draft = makeDraft({ + summary: "Execute code safely", + tools: ["execute_code"], + steps: [{ title: "execute", body: "execute the supplied code" }], + }); + const evidence = [ + trace("tr_json", "run code", "execution completed", [ + { name: "execute_code", input: '{"code": "print(1)"}' }, + ]), + ]; + const r = verifyDraft({ draft, evidence }, { log }); + expect(r.coverage).toBe(1); + expect(r.unmappedTokens).toEqual([]); + }); + it("partial coverage below threshold fails", () => { const draft = makeDraft({ summary: "Use several tools", diff --git a/apps/memos-local-plugin/tests/unit/storage/migrator.test.ts b/apps/memos-local-plugin/tests/unit/storage/migrator.test.ts index 933f59977..88e1af788 100644 --- a/apps/memos-local-plugin/tests/unit/storage/migrator.test.ts +++ b/apps/memos-local-plugin/tests/unit/storage/migrator.test.ts @@ -59,6 +59,31 @@ describe("storage/migrator", () => { } }); + it("replays the policy metadata migration after an interrupted marker write", () => { + const { dbPath, cleanup } = tmpDb(); + cleanups.push(cleanup); + const db = openDb({ filepath: dbPath, agent: "openclaw" }); + try { + runMigrations(db); + const metadataColumn = db + .prepare(`PRAGMA table_info(policies)`) + .all() + .find((row) => row.name === "metadata_json"); + expect(metadataColumn).toBeDefined(); + db.prepare<{ version: number }>(`DELETE FROM schema_migrations WHERE version=@version`) + .run({ version: 19 }); + + const replay = runMigrations(db); + expect(replay.applied.map((m) => m.version)).toContain(19); + expect(db + .prepare(`PRAGMA table_info(policies)`) + .all() + .filter((row) => row.name === "metadata_json")).toHaveLength(1); + } finally { + db.close(); + } + }); + it("rejects duplicate migration versions in a custom dir", () => { const dir = fs.mkdtempSync(path.join(os.tmpdir(), "memos-mig-dup-")); cleanups.push(() => fs.rmSync(dir, { recursive: true, force: true })); diff --git a/apps/memos-local-plugin/tests/unit/storage/policy-metadata-backfill.test.ts b/apps/memos-local-plugin/tests/unit/storage/policy-metadata-backfill.test.ts new file mode 100644 index 000000000..36dc03cb5 --- /dev/null +++ b/apps/memos-local-plugin/tests/unit/storage/policy-metadata-backfill.test.ts @@ -0,0 +1,51 @@ +import fs from "node:fs"; +import os from "node:os"; +import path from "node:path"; + +import { afterEach, describe, expect, it } from "vitest"; + +import { backfillLegacyPolicyMetadata } from "../../../core/storage/policy-metadata-backfill.js"; +import { openDb, runMigrations } from "../../../core/storage/index.js"; + +describe("policy metadata backfill", () => { + const cleanups: Array<() => void> = []; + + afterEach(() => { + while (cleanups.length) cleanups.pop()!(); + }); + + it("backfills old rows in bounded, idempotent batches", () => { + const dir = fs.mkdtempSync(path.join(os.tmpdir(), "memos-meta-")); + const dbPath = path.join(dir, "m.db"); + cleanups.push(() => fs.rmSync(dir, { recursive: true, force: true })); + const db = openDb({ filepath: dbPath, agent: "openclaw" }); + try { + runMigrations(db); + db.exec(` + INSERT INTO sessions (id, agent, started_at, last_seen_at) VALUES ('s1', 'openclaw', 1, 1); + INSERT INTO episodes (id, session_id, started_at) VALUES ('e1', 's1', 1); + INSERT INTO traces (id, episode_id, session_id, ts, user_text, agent_text, reflection, tool_calls_json, tags_json, turn_id) + VALUES ('t1', 'e1', 's1', 1, '请修复部署', '部署失败 DEPLOY_TIMEOUT', '重试服务', + '[{"name":"kubectl","input":"{\\"code\\":1}"}]', '["k8s","部署"]', 1); + INSERT INTO policies (id, title, trigger, procedure, verification, boundary, source_trace_ids_json, created_at, updated_at) + VALUES ('p1', 'legacy policy', 'trigger', 'procedure', 'verify', 'boundary', '["t1"]', 1, 1); + `); + + const replay = runMigrations(db); + expect(replay.metadataBackfilled).toBe(1); + const row = db.prepare<{ id: string }, { metadata_json: string | null }>( + `SELECT metadata_json FROM policies WHERE id=@id`, + ).get({ id: "p1" }); + expect(JSON.parse(row!.metadata_json!)).toMatchObject({ + version: 1, + language: "mixed", + domainTags: ["k8s", "部署"], + toolNames: ["kubectl"], + errorCodes: ["DEPLOY_TIMEOUT"], + }); + expect(backfillLegacyPolicyMetadata(db, { batchSize: 1 })).toBe(0); + } finally { + db.close(); + } + }); +});