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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion docs/generated/release-truth.json
Original file line number Diff line number Diff line change
Expand Up @@ -135,7 +135,7 @@
"console_entrypoints": 8,
"mcp_tools": 51,
"ops_cli_commands": 5,
"pytest_test_functions": 4119
"pytest_test_functions": 4123
},
"feature_profile_matrix": {
"capture_hook": [
Expand Down
2 changes: 1 addition & 1 deletion docs/generated/release-truth.md
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,7 @@ Do not edit this file by hand. Run `python scripts/generate_release_truth.py`.
- Main CLI commands: **119**
- Operations CLI commands: **5**
- Console entrypoints: **8**
- Pytest source test functions: **4119**
- Pytest source test functions: **4123**

## MCP tools

Expand Down
72 changes: 60 additions & 12 deletions memorymaster/profile/engine.py
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,13 @@
from pathlib import Path
from typing import Any, Callable, Protocol

from memorymaster.profile.models import ProfileCandidate, ProfileDecision, ProfileFact, ProfileMessage
from memorymaster.profile.models import (
ProfileCandidate,
ProfileDecision,
ProfileFact,
ProfileMessage,
ProfileValidationError,
)
from memorymaster.profile.renderer import render_profile
from memorymaster.profile.repository import ProfileRepository

Expand Down Expand Up @@ -56,8 +62,17 @@ class ProfileConfig:
max_input_chars: int = DEFAULT_MAX_INPUT_CHARS
min_independent_sessions: int = 2
preference_ttl_days: int = 90
token_budget: int = 800
max_facts: int = 40
# 2026-08-30: con el run 3 destrabado el corpus paso a 52 hechos y el techo
# viejo (800 tok / 40 hechos) cortaba 20 en silencio — la proyeccion quedaba
# clavada en 800/800 justo cuando los hechos NUEVOS eran los buenos
# (memorymaster+UNMS sup=10, WISP, Task Scheduler) y los viejos los stale.
# Truncar por techo no avisa, asi que el techo sube con el corpus.
token_budget: int = 1400
max_facts: int = 60
# 68 candidatos particionaron bien (run 2); 234 no lo lograron ni una vez en
# diez dias (run 3). 40 deja margen bajo el limite observado sin volver el
# reduce innecesariamente charlatan.
reduce_batch_size: int = 40

@classmethod
def from_env(cls) -> "ProfileConfig":
Expand All @@ -72,6 +87,7 @@ def from_env(cls) -> "ProfileConfig":
preference_ttl_days=_env_int("MEMORYMASTER_PROFILE_PREFERENCE_TTL_DAYS", 90),
token_budget=_env_int("MEMORYMASTER_PROFILE_TOKEN_BUDGET", 800),
max_facts=_env_int("MEMORYMASTER_PROFILE_MAX_FACTS", 40),
reduce_batch_size=_env_int("MEMORYMASTER_PROFILE_REDUCE_BATCH", 40),
)


Expand Down Expand Up @@ -179,15 +195,46 @@ def _advance_mapping(
return None

def _reduce_and_complete(self, run_id: int, now: datetime) -> dict[str, Any]:
candidates = self.repo.candidates(run_id)
facts = self.repo.active_facts()
decisions = self.reducer.reduce(candidates, facts) if candidates else ()
stats = self.repo.apply_decisions(
run_id,
decisions,
now=now,
min_sessions=self.config.min_independent_sessions,
)
"""Reduce por LOTES acotados, no de una.

El validador exige que el modelo particione el lote perfectamente: cada
candidate_id exactamente una vez. Eso se sostiene con decenas de
candidatos y no con cientos — el run 2 completo con 68, el run 3 acumulo
234 y quedo clavado diez dias, fallando entre "candidates must appear
exactly once" y JSON malformado.

La semantica se conserva releyendo `active_facts()` entre lotes: el lote
N+1 ve los hechos que creo el N y puede fusionar contra ellos, que es el
mismo mecanismo incremental que ya opera entre runs. Lo que cambia es que
las fusiones se deciden con visibilidad parcial, asi que dos candidatos
de lotes distintos pueden quedar como dos hechos donde una particion
unica los unia. Ese es el precio, y es preferible a no reducir nunca.
"""
stats = {"applied": 0, "rejected": 0, "consumed": 0}
batches = 0
while True:
pending = self.repo.candidates(run_id, pending_only=True)
if not pending:
break
batch = pending[: self.config.reduce_batch_size]
facts = self.repo.active_facts() # releido: el lote previo creo hechos
decisions = self.reducer.reduce(batch, facts)
applied = self.repo.apply_decisions(
run_id,
decisions,
now=now,
min_sessions=self.config.min_independent_sessions,
)
for key, value in applied.items():
stats[key] = stats.get(key, 0) + value
batches += 1
if applied.get("consumed", 0) == 0:
# Ningun candidato quedo marcado: reintentar seria un bucle
# infinito sobre el mismo lote. Se corta y el run queda en
# `reducing` para el proximo ciclo, que es el estado honesto.
raise ProfileValidationError(
f"lote de {len(batch)} candidatos no consumio ninguno"
)
expired = self.repo.expire_preferences(
now=now, ttl_days=self.config.preference_ttl_days
)
Expand Down Expand Up @@ -241,6 +288,7 @@ def run_compiled_profile(
) -> dict[str, Any]:
from memorymaster.profile.providers import ProfileMapper, ProfileReducer


# MEMORYMASTER_PROFILE_OUTPUT_DIR existe para que los tests puedan sacar esta
# escritura del HOME real, igual que MEMORYMASTER_SNAPSHOT_DIR y
# MEMORYMASTER_SPOOL_DIR. Sin ella no habia forma: `scheduled_task._run_dream`
Expand Down
30 changes: 26 additions & 4 deletions memorymaster/profile/repository.py
Original file line number Diff line number Diff line change
Expand Up @@ -215,11 +215,22 @@ def mark_reducing(self, run_id: int, *, now: datetime) -> None:
)
conn.commit()

def candidates(self, run_id: int) -> tuple[ProfileCandidate, ...]:
def candidates(
self, run_id: int, *, pending_only: bool = False
) -> tuple[ProfileCandidate, ...]:
"""Candidatos del run; con ``pending_only`` solo los aun no consumidos.

El reduce va por lotes y marca cada candidato consumido al aplicar su
decision, asi que un run reanudado debe ver SOLO lo que falta. Sin ese
filtro, reanudar volveria a aplicar lo ya aplicado.
"""
where = "WHERE run_id=?" + (
" AND consumed_at IS NULL" if pending_only else ""
)
with closing(self.connect()) as conn:
rows = conn.execute(
"""SELECT * FROM compiled_profile_candidates
WHERE run_id=? ORDER BY candidate_id""",
f"""SELECT * FROM compiled_profile_candidates
{where} ORDER BY candidate_id""",
(run_id,),
).fetchall()
return tuple(
Expand Down Expand Up @@ -294,14 +305,25 @@ def apply_decisions(
min_sessions: int,
) -> dict[str, int]:
candidates = {item.candidate_id: item for item in self.candidates(run_id)}
stats = {"applied": 0, "rejected": 0}
stats = {"applied": 0, "rejected": 0, "consumed": 0}
stamp = now.isoformat()
with closing(self.connect()) as conn:
for decision in decisions:
support_ids = self._decision_supports(decision, candidates)
applied = self._apply_decision(
conn, decision, support_ids, now=now, min_sessions=min_sessions
)
stats["applied" if applied else "rejected"] += 1
# Se marca DENTRO de la misma transaccion que aplico la decision.
# Separarlo reabre la doble aplicacion: un commit del hecho sin el
# marcado deja el candidato listo para volver a aplicarse.
for candidate_id in decision.candidate_ids:
conn.execute(
"""UPDATE compiled_profile_candidates SET consumed_at=?
WHERE run_id=? AND candidate_id=? AND consumed_at IS NULL""",
(stamp, run_id, candidate_id),
)
stats["consumed"] += 1
conn.commit()
return stats

Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,56 @@
"""Marcar que candidato de perfil ya fue consumido por una decision.

El reduce pasaba TODOS los candidatos de un run en una sola llamada, y el
validador exige que el modelo devuelva una particion perfecta: cada
candidate_id exactamente una vez, sin duplicar ni omitir. Eso escala hasta
donde el modelo puede sostener la particion y no mas: el run 2 completo con 68
candidatos, el run 3 acumulo 234 y quedo clavado en `reducing` desde el
2026-08-20 — diez dias de intentos programados, mas cuatro reintentos medidos a
mano, todos fallando entre `profile candidates must appear exactly once` y
`profile provider returned malformed JSON`.

Lotear el reduce es la salida, pero sin esta columna es peor que el problema:
`apply_decisions` relee `candidates(run_id)` entero en cada llamada, asi que un
lote aplicado y un crash antes del siguiente dejaba los mismos candidatos listos
para aplicarse de nuevo — hechos duplicados en el perfil que se inyecta en cada
sesion. `consumed_at` es lo que hace que un run a medias sea reanudable en vez
de destructivo.

Se marca dentro de la MISMA transaccion que aplica la decision. Si eso se
separa, vuelve la doble aplicacion por otra puerta.
"""

from __future__ import annotations

from typing import Any


VERSION = 23
DESCRIPTION = "Track which profile candidates a reduce batch already consumed"


def _columns(conn: Any) -> set[str]:
return {
str(row[1])
for row in conn.execute("PRAGMA table_info(compiled_profile_candidates)")
}


def apply_sqlite(conn: Any) -> None:
if "consumed_at" not in _columns(conn):
conn.execute(
"ALTER TABLE compiled_profile_candidates ADD COLUMN consumed_at TEXT"
)
# Las filas previas quedan NULL a proposito. Un run ya completado no vuelve a
# reducirse (su estado terminal lo frena antes), y un run en vuelo DEBE ver
# sus candidatos como pendientes: son exactamente los que faltan aplicar.
conn.execute(
"""CREATE INDEX IF NOT EXISTS idx_compiled_profile_candidates_pending
ON compiled_profile_candidates(run_id, consumed_at)"""
)
conn.commit()


def apply_postgres(conn: Any) -> None:
"""Fail closed: el perfil compilado es SQLite-only, igual que la 21."""
raise RuntimeError("migration 23 is SQLite-only; the compiled profile is SQLite-only")
Loading
Loading