From a6b2c0a04f90790a3ff41d9d8a4a7dbd557192c5 Mon Sep 17 00:00:00 2001 From: juemimgcd Date: Sun, 30 Aug 2026 21:09:21 +0800 Subject: [PATCH] feat(memoria): add consented ad personalization --- .env-example | 2 + README.md | 4 + app/mneme/domains/tasks/outbox.py | 17 ++- app/mneme/memoria/api/memory.py | 7 +- app/mneme/memoria/clients/memory_agent.py | 18 +++ app/mneme/memoria/memory_gateway.py | 25 +++- app/mneme/memoria/schemas/memory_agent.py | 68 +++++++++- .../20260830_01_add_ad_personalization.py | 72 +++++++++++ app/mneme/memoria/server/api/memories.py | 6 +- .../memoria/server/api/recommendations.py | 25 ++++ app/mneme/memoria/server/app.py | 2 + .../memoria/server/contracts/memories.py | 1 + .../server/contracts/recommendations.py | 47 +++++++ app/mneme/memoria/server/memory/extraction.py | 6 +- .../memoria/server/memory/reconciliation.py | 18 ++- app/mneme/memoria/server/memory/schemas.py | 18 ++- .../memoria/server/memory/sensitivity.py | 47 ++++++- .../memoria/server/models/canonical_memory.py | 6 + .../memoria/server/models/memory_settings.py | 7 + .../memoria/server/security/service_tokens.py | 1 + .../server/services/ad_recommendations.py | 121 ++++++++++++++++++ .../server/services/memory_commands.py | 4 +- .../memoria/server/services/memory_events.py | 28 +++- docker-compose.yml | 10 +- docs/architecture.md | 9 ++ docs/runtime-contracts.md | 19 +++ requirements/ai.txt | 3 +- 27 files changed, 557 insertions(+), 34 deletions(-) create mode 100644 app/mneme/memoria/server/alembic/versions/20260830_01_add_ad_personalization.py create mode 100644 app/mneme/memoria/server/api/recommendations.py create mode 100644 app/mneme/memoria/server/contracts/recommendations.py create mode 100644 app/mneme/memoria/server/services/ad_recommendations.py diff --git a/.env-example b/.env-example index 6ac829f..579a99a 100644 --- a/.env-example +++ b/.env-example @@ -146,3 +146,5 @@ EXTERNAL_RETRY_MAX_DELAY_SECONDS=4.0 CIRCUIT_BREAKER_FAILURE_THRESHOLD=3 CIRCUIT_BREAKER_RECOVERY_TIMEOUT_SECONDS=30 + +MEMORY_AGENT_SERVICE_JWT_SECRET=<另一组随机密钥> diff --git a/README.md b/README.md index f9bfc9a..370feaf 100644 --- a/README.md +++ b/README.md @@ -122,6 +122,8 @@ cd Mneme cp .env-example .env ``` +Windows PowerShell: + ```powershell # Windows PowerShell Copy-Item .env-example .env @@ -405,6 +407,8 @@ CI 会分别执行前端、后端和集成检查;只有三个阶段全部通 | [Operations Runbook](docs/operations-runbook.md) | 监控、告警、备份、恢复与故障处理 | | [Deployment](deploy/DEPLOY.md) | 生产部署、发布、回滚及 Memoria 运维 | +规范路径:[docs/architecture.md](docs/architecture.md)、[docs/runtime-contracts.md](docs/runtime-contracts.md)、[docs/current-state.md](docs/current-state.md)。 + ## 参与贡献 欢迎通过 [Issues](https://github.com/juemimgcd/Reminder/issues) 报告问题或提出建议。提交 Pull Request 前,请运行与改动范围对应的质量检查,并保持以下边界: diff --git a/app/mneme/domains/tasks/outbox.py b/app/mneme/domains/tasks/outbox.py index 7d0fb4b..3e15642 100644 --- a/app/mneme/domains/tasks/outbox.py +++ b/app/mneme/domains/tasks/outbox.py @@ -471,21 +471,32 @@ async def enqueue_user_memory_settings_changed( db: AsyncSession, *, owner_id: int, - automatic_conversation_memory: bool, + automatic_conversation_memory: bool | None = None, + ad_personalization_enabled: bool | None = None, occurred_at: datetime, ) -> OutboxEvent: + payload = { + key: value + for key, value in { + "automatic_conversation_memory": automatic_conversation_memory, + "ad_personalization_enabled": ad_personalization_enabled, + }.items() + if value is not None + } + if not payload: + raise ValueError("at least one memory setting is required") event = MemoryAgentEvent( event_id=_memory_event_id( "memory-settings-changed", str(owner_id), occurred_at.isoformat(), - str(automatic_conversation_memory), + repr(sorted(payload.items())), ), event_type="user.memory_settings.changed", occurred_at=occurred_at, owner_id=owner_id, knowledge_base_id=None, - payload={"automatic_conversation_memory": automatic_conversation_memory}, + payload=payload, ) return await _enqueue_memory_agent_event( db, diff --git a/app/mneme/memoria/api/memory.py b/app/mneme/memoria/api/memory.py index 3f95eff..fc2bb39 100644 --- a/app/mneme/memoria/api/memory.py +++ b/app/mneme/memoria/api/memory.py @@ -314,5 +314,10 @@ async def patch_settings( current_user: User = Depends(get_current_user), db: AsyncSession = Depends(get_write_database), ): - data = await service.update_settings(db, owner_id=current_user.id, enabled=payload.automatic_conversation_memory) + data = await service.update_settings( + db, + owner_id=current_user.id, + automatic_conversation_memory=payload.automatic_conversation_memory, + ad_personalization_enabled=payload.ad_personalization_enabled, + ) return success_response(data=data, message="memory settings accepted") diff --git a/app/mneme/memoria/clients/memory_agent.py b/app/mneme/memoria/clients/memory_agent.py index 6e1d122..a63c600 100644 --- a/app/mneme/memoria/clients/memory_agent.py +++ b/app/mneme/memoria/clients/memory_agent.py @@ -13,6 +13,8 @@ from app.mneme.conf.config import settings from app.mneme.memoria.schemas.memory_agent import ( + AdRecommendationRequest, + AdRecommendationResponse, CanonicalMemoryData, ConversationMemorySettingsData, EventReceipt, @@ -116,6 +118,22 @@ async def create_answer(self, request: MemoryAgentAnswerRequest) -> MemoryAgentA except ValueError as exc: raise MemoryAgentPermanentFailure("memory agent returned an invalid answer response") from exc + async def recommend_ads(self, request: AdRecommendationRequest) -> AdRecommendationResponse: + """Rerank caller-filtered ads without exposing governed memory content.""" + response = await self._post_json( + path="/v1/ad-recommendations", + payload=request.model_dump(mode="json"), + request_id=request.request_id, + scope="ads:recommend", + owner_id=request.owner_id, + knowledge_base_id=request.knowledge_base_id, + retry_transient=True, + ) + try: + return AdRecommendationResponse.model_validate(response.json()) + except ValueError as exc: + raise MemoryAgentPermanentFailure("memory agent returned an invalid ad recommendation") from exc + async def stream_answer(self, request: MemoryAgentAnswerRequest) -> AsyncIterator[MemoryAgentStreamEvent]: """Yield validated events from the Memoria answer stream. diff --git a/app/mneme/memoria/memory_gateway.py b/app/mneme/memoria/memory_gateway.py index 5ce01a3..9a1cc02 100644 --- a/app/mneme/memoria/memory_gateway.py +++ b/app/mneme/memoria/memory_gateway.py @@ -17,7 +17,6 @@ from app.mneme.memoria.clients.memory_agent import MemoryAgentClient from app.mneme.memoria.schemas.memory_agent import ( CanonicalMemoryData, - ConversationMemorySettingsData, GovernedMemoryPage, MemoryCandidateData, MemoryCandidatePage, @@ -293,11 +292,29 @@ async def list_candidates( return MemoryCandidatePage(items=items, next_cursor=next_cursor, total=total, pending_count=total) -async def update_settings(db: AsyncSession, *, owner_id: int, enabled: bool) -> ConversationMemorySettingsData: +async def update_settings( + db: AsyncSession, + *, + owner_id: int, + automatic_conversation_memory: bool | None, + ad_personalization_enabled: bool | None, +) -> dict[str, bool]: await enqueue_user_memory_settings_changed( - db, owner_id=owner_id, automatic_conversation_memory=enabled, occurred_at=datetime.now(UTC) + db, + owner_id=owner_id, + automatic_conversation_memory=automatic_conversation_memory, + ad_personalization_enabled=ad_personalization_enabled, + occurred_at=datetime.now(UTC), ) - return ConversationMemorySettingsData(automatic_conversation_memory=enabled, applied=False) + return { + key: value + for key, value in { + "automatic_conversation_memory": automatic_conversation_memory, + "ad_personalization_enabled": ad_personalization_enabled, + "applied": False, + }.items() + if value is not None + } def _secret() -> str: diff --git a/app/mneme/memoria/schemas/memory_agent.py b/app/mneme/memoria/schemas/memory_agent.py index 203d2cc..0332c9e 100644 --- a/app/mneme/memoria/schemas/memory_agent.py +++ b/app/mneme/memoria/schemas/memory_agent.py @@ -4,9 +4,17 @@ """ from datetime import datetime -from typing import Any, Literal - -from pydantic import BaseModel, ConfigDict, Field, SecretStr, field_validator, model_validator +from typing import Annotated, Any, Literal + +from pydantic import ( + BaseModel, + ConfigDict, + Field, + SecretStr, + StringConstraints, + field_validator, + model_validator, +) AnswerMode = Literal[ "kb_qa", @@ -154,6 +162,49 @@ class MemoryAgentStreamEvent(BaseModel): response: MemoryAgentAnswerResponse | None = None +AdTag = Annotated[str, StringConstraints(strip_whitespace=True, min_length=1, max_length=64)] + + +class AdCandidateData(BaseModel): + model_config = ConfigDict(extra="forbid") + + ad_id: str = Field(min_length=1, max_length=128) + title: str = Field(min_length=1, max_length=240) + description: str = Field(default="", max_length=2000) + tags: list[AdTag] = Field(default_factory=list, max_length=20) + business_score: float = Field(default=0.0, ge=0, le=1) + + +class AdRecommendationRequest(BaseModel): + model_config = ConfigDict(extra="forbid") + + request_id: str = Field(min_length=1, max_length=128) + owner_id: int = Field(gt=0) + knowledge_base_id: str | None = Field(default=None, max_length=128) + placement: str = Field(min_length=1, max_length=64) + candidates: list[AdCandidateData] = Field(min_length=1, max_length=100) + limit: int = Field(default=1, ge=1, le=10) + + @model_validator(mode="after") + def candidate_ids_must_be_unique(self) -> "AdRecommendationRequest": + ids = [candidate.ad_id for candidate in self.candidates] + if len(ids) != len(set(ids)): + raise ValueError("candidate ad IDs must be unique") + return self + + +class AdRecommendationItem(BaseModel): + ad_id: str + score: float = Field(ge=0, le=1) + matched_topics: list[str] = Field(default_factory=list) + + +class AdRecommendationResponse(BaseModel): + request_id: str + personalized: bool + items: list[AdRecommendationItem] + + class CanonicalMemoryData(BaseModel): memory_id: str knowledge_base_id: str | None @@ -162,6 +213,7 @@ class CanonicalMemoryData(BaseModel): predicate: str value: str confidence: float + sensitivity: str = "unknown" status: str active_revision_id: str created_at: datetime @@ -288,9 +340,17 @@ def exactly_one_selector(self) -> "MemoryPurgeRequest": class ConversationMemorySettingsUpdate(BaseModel): - automatic_conversation_memory: bool + automatic_conversation_memory: bool | None = None + ad_personalization_enabled: bool | None = None + + @model_validator(mode="after") + def require_one_setting(self) -> "ConversationMemorySettingsUpdate": + if self.automatic_conversation_memory is None and self.ad_personalization_enabled is None: + raise ValueError("at least one memory setting is required") + return self class ConversationMemorySettingsData(BaseModel): automatic_conversation_memory: bool + ad_personalization_enabled: bool = False applied: bool diff --git a/app/mneme/memoria/server/alembic/versions/20260830_01_add_ad_personalization.py b/app/mneme/memoria/server/alembic/versions/20260830_01_add_ad_personalization.py new file mode 100644 index 0000000..a018cc0 --- /dev/null +++ b/app/mneme/memoria/server/alembic/versions/20260830_01_add_ad_personalization.py @@ -0,0 +1,72 @@ +"""persist memory sensitivity and ad personalization consent + +Revision ID: 20260830_01 +Revises: 20260718_03 +Create Date: 2026-08-30 +""" + +from collections.abc import Sequence + +import sqlalchemy as sa + +from alembic import op + +revision: str = "20260830_01" +down_revision: str | Sequence[str] | None = "20260718_03" +branch_labels: str | Sequence[str] | None = None +depends_on: str | Sequence[str] | None = None + + +def upgrade() -> None: + op.add_column( + "canonical_memories", + sa.Column("sensitivity", sa.String(length=16), server_default="unknown", nullable=False), + ) + op.create_check_constraint( + "ck_canonical_memories_sensitivity", + "canonical_memories", + "sensitivity IN ('unknown', 'low', 'sensitive')", + ) + op.execute( + """ + UPDATE canonical_memories AS memory + SET sensitivity = 'sensitive' + FROM ( + SELECT + owner_id, + knowledge_base_id, + fingerprint + FROM memory_candidates + WHERE status = 'promoted' + AND sensitivity = 'sensitive' + GROUP BY owner_id, knowledge_base_id, fingerprint + ) AS candidates + WHERE memory.owner_id = candidates.owner_id + AND memory.knowledge_base_id IS NOT DISTINCT FROM candidates.knowledge_base_id + AND memory.fingerprint = candidates.fingerprint + """ + ) + op.add_column( + "memory_settings", + sa.Column("ad_personalization_enabled", sa.Boolean(), server_default="false", nullable=False), + ) + op.add_column( + "memory_settings", + sa.Column("ad_personalization_last_event_occurred_at", sa.DateTime(timezone=True), nullable=True), + ) + op.add_column( + "memory_settings", + sa.Column("ad_personalization_last_event_id", sa.String(length=128), nullable=True), + ) + + +def downgrade() -> None: + op.drop_column("memory_settings", "ad_personalization_last_event_id") + op.drop_column("memory_settings", "ad_personalization_last_event_occurred_at") + op.drop_column("memory_settings", "ad_personalization_enabled") + op.drop_constraint( + "ck_canonical_memories_sensitivity", + "canonical_memories", + type_="check", + ) + op.drop_column("canonical_memories", "sensitivity") diff --git a/app/mneme/memoria/server/api/memories.py b/app/mneme/memoria/server/api/memories.py index 16f3da0..d2725c3 100644 --- a/app/mneme/memoria/server/api/memories.py +++ b/app/mneme/memoria/server/api/memories.py @@ -41,6 +41,7 @@ class CanonicalMemoryDTO(BaseModel): predicate: str value: str confidence: float + sensitivity: str status: str active_revision_id: str created_at: datetime @@ -309,6 +310,7 @@ async def get_memory_settings( row = await db.get(MemorySettings, owner_id) return { "automatic_conversation_memory": bool(row and row.automatic_conversation_memory), + "ad_personalization_enabled": bool(row and row.ad_personalization_enabled), "applied": row is not None, } @@ -371,7 +373,8 @@ async def command_memory( ) else: assert command.subject is not None and command.predicate is not None and command.value is not None - if classify_sensitivity(command.subject, command.predicate, command.value) == "secret": + sensitivity = classify_sensitivity(command.subject, command.predicate, command.value) + if sensitivity == "secret": raise ValueError("secret values cannot be persisted") row = await memory_commands.revise( db, @@ -382,6 +385,7 @@ async def command_memory( predicate=command.predicate, value=command.value, confidence=command.confidence, + sensitivity=sensitivity, actor_id=command.actor_id, reason=command.reason, ) diff --git a/app/mneme/memoria/server/api/recommendations.py b/app/mneme/memoria/server/api/recommendations.py new file mode 100644 index 0000000..3c775ff --- /dev/null +++ b/app/mneme/memoria/server/api/recommendations.py @@ -0,0 +1,25 @@ +"""Expose scoped memory-personalized ad reranking.""" + +from typing import Annotated, Any + +from fastapi import APIRouter, Depends + +from app.mneme.memoria.server.api.dependencies import require_claimed_scope, require_service_scope +from app.mneme.memoria.server.contracts.recommendations import AdRecommendationRequest, AdRecommendationResponse +from app.mneme.memoria.server.security.service_tokens import AD_RECOMMENDATIONS_SCOPE +from app.mneme.memoria.server.services.ad_recommendations import recommend_ads + +router = APIRouter() + + +@router.post("/ad-recommendations", response_model=AdRecommendationResponse) +async def create_ad_recommendation( + request: AdRecommendationRequest, + claims: Annotated[dict[str, Any], Depends(require_service_scope(AD_RECOMMENDATIONS_SCOPE))], +) -> AdRecommendationResponse: + require_claimed_scope( + claims, + owner_id=request.owner_id, + knowledge_base_id=request.knowledge_base_id, + ) + return await recommend_ads(request) diff --git a/app/mneme/memoria/server/app.py b/app/mneme/memoria/server/app.py index bf04c88..b80e1f1 100644 --- a/app/mneme/memoria/server/app.py +++ b/app/mneme/memoria/server/app.py @@ -14,6 +14,7 @@ from app.mneme.memoria.server.api.events import router as event_router from app.mneme.memoria.server.api.health import router as health_router from app.mneme.memoria.server.api.memories import router as memories_router +from app.mneme.memoria.server.api.recommendations import router as recommendations_router from app.mneme.memoria.server.api.runs import router as runs_router from app.mneme.memoria.server.config import settings from app.mneme.memoria.server.observability.context import safe_log @@ -46,4 +47,5 @@ def create_memory_agent_app() -> FastAPI: app.include_router(answers_router, prefix="/v1") app.include_router(runs_router, prefix="/v1") app.include_router(memories_router, prefix="/v1") + app.include_router(recommendations_router, prefix="/v1") return app diff --git a/app/mneme/memoria/server/contracts/memories.py b/app/mneme/memoria/server/contracts/memories.py index e8ceb86..fe0ebce 100644 --- a/app/mneme/memoria/server/contracts/memories.py +++ b/app/mneme/memoria/server/contracts/memories.py @@ -28,6 +28,7 @@ class MemoryData(BaseModel): predicate: str value: str confidence: float + sensitivity: str status: MemoryStatus created_at: datetime updated_at: datetime diff --git a/app/mneme/memoria/server/contracts/recommendations.py b/app/mneme/memoria/server/contracts/recommendations.py new file mode 100644 index 0000000..946844e --- /dev/null +++ b/app/mneme/memoria/server/contracts/recommendations.py @@ -0,0 +1,47 @@ +"""Define the bounded contract for memory-personalized ad reranking.""" + +from typing import Annotated + +from pydantic import BaseModel, ConfigDict, Field, StringConstraints, model_validator + +AdTag = Annotated[str, StringConstraints(strip_whitespace=True, min_length=1, max_length=64)] + + +class AdCandidate(BaseModel): + model_config = ConfigDict(extra="forbid") + + ad_id: str = Field(min_length=1, max_length=128) + title: str = Field(min_length=1, max_length=240) + description: str = Field(default="", max_length=2000) + tags: list[AdTag] = Field(default_factory=list, max_length=20) + business_score: float = Field(default=0.0, ge=0, le=1) + + +class AdRecommendationRequest(BaseModel): + model_config = ConfigDict(extra="forbid") + + request_id: str = Field(min_length=1, max_length=128) + owner_id: int = Field(gt=0) + knowledge_base_id: str | None = Field(default=None, max_length=128) + placement: str = Field(min_length=1, max_length=64) + candidates: list[AdCandidate] = Field(min_length=1, max_length=100) + limit: int = Field(default=1, ge=1, le=10) + + @model_validator(mode="after") + def candidate_ids_must_be_unique(self) -> "AdRecommendationRequest": + ids = [candidate.ad_id for candidate in self.candidates] + if len(ids) != len(set(ids)): + raise ValueError("candidate ad IDs must be unique") + return self + + +class AdRecommendationItem(BaseModel): + ad_id: str + score: float = Field(ge=0, le=1) + matched_topics: list[str] = Field(default_factory=list) + + +class AdRecommendationResponse(BaseModel): + request_id: str + personalized: bool + items: list[AdRecommendationItem] diff --git a/app/mneme/memoria/server/memory/extraction.py b/app/mneme/memoria/server/memory/extraction.py index 696a446..a2adca6 100644 --- a/app/mneme/memoria/server/memory/extraction.py +++ b/app/mneme/memoria/server/memory/extraction.py @@ -81,7 +81,11 @@ async def extract_candidates(evidence: EvidenceInput) -> list[ExtractedCandidate "content": ( "Extract only durable user memories explicitly supported by the excerpt. " "Use only the fixed schema types. Copy an exact supporting quote and its " - "zero-based start/end offsets. Do not decide persistence or status. Return JSON." + "zero-based start/end offsets. Mark sensitivity_signals for identity, health, " + "finance, authentication, political or religious beliefs, sexual orientation, " + "race or ethnicity, trade-union membership, minors, precise location, biometric, " + "genetic, credential, secret, or other sensitive data. Do not decide persistence " + "or status. Return JSON." ), }, {"role": "user", "content": evidence.excerpt}, diff --git a/app/mneme/memoria/server/memory/reconciliation.py b/app/mneme/memoria/server/memory/reconciliation.py index 69660ae..0f4c26d 100644 --- a/app/mneme/memoria/server/memory/reconciliation.py +++ b/app/mneme/memoria/server/memory/reconciliation.py @@ -18,7 +18,7 @@ normalize_memory_text, ) from app.mneme.memoria.server.memory.policy import PolicyDecision, classify_candidate -from app.mneme.memoria.server.models.canonical_memory import CanonicalMemory +from app.mneme.memoria.server.models.canonical_memory import CanonicalMemory, CanonicalSensitivity from app.mneme.memoria.server.models.evidence import Evidence from app.mneme.memoria.server.models.memory_candidate import ( MEMORY_TYPES, @@ -125,6 +125,11 @@ def _combined_confidence(current: float, additional: float) -> float: return min(1.0, 1.0 - ((1.0 - current) * (1.0 - additional))) +def _combined_sensitivity(current: str, additional: Sensitivity) -> CanonicalSensitivity: + rank = {"low": 0, "unknown": 1, "sensitive": 2} + return max((current, additional), key=rank.__getitem__) # type: ignore[return-value] + + async def _create_memory( db: AsyncSession, *, @@ -136,6 +141,7 @@ async def _create_memory( value: str, fingerprint: str, confidence: float, + sensitivity: Sensitivity, reason: str, actor: str, evidence_ids: list[str], @@ -152,6 +158,7 @@ async def _create_memory( value=value, fingerprint=fingerprint, confidence=confidence, + sensitivity=sensitivity, status="active", active_revision_id=revision_id, ) @@ -298,6 +305,7 @@ async def reconcile_candidate( ) if attached: compatible.confidence = _combined_confidence(compatible.confidence, confidence) + compatible.sensitivity = _combined_sensitivity(compatible.sensitivity, sensitivity) await db.flush() return ReconciliationResult( decision="promote", @@ -366,6 +374,7 @@ async def reconcile_candidate( value=normalized_value, fingerprint=fingerprint, confidence=confidence, + sensitivity=sensitivity, reason="explicit_request" if explicit_request else "automatic_promotion", actor=actor, evidence_ids=evidence_ids, @@ -384,9 +393,9 @@ async def revise_memory( value: str, reason: str, actor: str, + sensitivity: Sensitivity, evidence_ids: list[str] | None = None, confidence: float | None = None, - sensitivity: Sensitivity = "low", ) -> CanonicalMemory: """Replace the active value of a scoped memory by creating a new revision. @@ -477,6 +486,7 @@ async def revise_memory( memory.predicate = normalized_predicate memory.value = normalized_value memory.fingerprint = fingerprint + memory.sensitivity = sensitivity memory.active_revision_id = revision.revision_id memory.status = "active" if confidence is not None: @@ -550,6 +560,9 @@ async def confirm_candidate( compatible.confidence = _combined_confidence( compatible.confidence, candidate.confidence ) + compatible.sensitivity = _combined_sensitivity( + compatible.sensitivity, candidate.sensitivity + ) memory = compatible else: current_conflict = await find_conflicting_memory( @@ -603,6 +616,7 @@ async def confirm_candidate( value=candidate.value, fingerprint=candidate.fingerprint, confidence=candidate.confidence, + sensitivity=candidate.sensitivity, reason="user_confirmation", actor=actor, evidence_ids=evidence_ids, diff --git a/app/mneme/memoria/server/memory/schemas.py b/app/mneme/memoria/server/memory/schemas.py index 65ad003..3a06e57 100644 --- a/app/mneme/memoria/server/memory/schemas.py +++ b/app/mneme/memoria/server/memory/schemas.py @@ -26,6 +26,15 @@ "health", "finance", "authentication", + "political_belief", + "religious_belief", + "sexual_orientation", + "race_ethnicity", + "trade_union", + "minor", + "precise_location", + "biometric", + "genetic", "credential", "secret", "password", @@ -134,4 +143,11 @@ class DocumentMemoryObservedPayload(BaseModel): class MemorySettingsChangedPayload(BaseModel): model_config = ConfigDict(extra="forbid") - automatic_conversation_memory: bool + automatic_conversation_memory: bool | None = None + ad_personalization_enabled: bool | None = None + + @model_validator(mode="after") + def require_one_setting(self) -> "MemorySettingsChangedPayload": + if self.automatic_conversation_memory is None and self.ad_personalization_enabled is None: + raise ValueError("at least one memory setting is required") + return self diff --git a/app/mneme/memoria/server/memory/sensitivity.py b/app/mneme/memoria/server/memory/sensitivity.py index 45cd036..ae1a138 100644 --- a/app/mneme/memoria/server/memory/sensitivity.py +++ b/app/mneme/memoria/server/memory/sensitivity.py @@ -40,10 +40,49 @@ ) _SENSITIVE_PATTERNS = ( - re.compile(r"\b(?:social security|ssn|passport|national id|身份证|护照)\b", re.IGNORECASE), - re.compile(r"\b(?:diagnos(?:is|ed)|medical|medication|disease|病历|诊断|药物)\b", re.IGNORECASE), - re.compile(r"\b(?:bank account|credit card|routing number|iban|银行卡|信用卡)\b", re.IGNORECASE), - re.compile(r"\b(?:login|authentication|two-factor|2fa|登录|认证)\b", re.IGNORECASE), + re.compile(r"\b(?:social security|ssn|passport|national id)\b", re.IGNORECASE), + re.compile(r"(?:身份证|护照)"), + re.compile( + r"\b(?:diagnos(?:is|ed)|medical|medication|disease|pregnan\w*|disability|mental health)\b", + re.IGNORECASE, + ), + re.compile(r"(?:病历|诊断|药物|怀孕|残疾|心理健康)"), + re.compile( + r"\b(?:bank account|credit card|routing number|iban|income|salary|debt|credit score)\b", + re.IGNORECASE, + ), + re.compile(r"(?:银行卡|信用卡|收入|工资|负债|信用评分)"), + re.compile(r"\b(?:login|authentication|two-factor|2fa)\b", re.IGNORECASE), + re.compile(r"(?:登录凭据|身份认证|双重认证)"), + re.compile( + r"\b(?:politic(?:al|s)(?: belief| view| affiliation)?|party membership|voting preference)\b", + re.IGNORECASE, + ), + re.compile(r"(?:政治倾向|政治观点|政治立场|政党成员|党员|投票偏好)"), + re.compile( + r"\b(?:religion|religious belief|faith|christian(?:ity)?|muslim|islam|jewish|judaism|" + r"hindu(?:ism)?|buddhis[mt]|atheis[mt])\b", + re.IGNORECASE, + ), + re.compile(r"(?:宗教|信仰|基督徒|基督教|穆斯林|伊斯兰教|犹太教|佛教徒|佛教|印度教|无神论)"), + re.compile( + r"\b(?:sexual orientation|gay|lesbian|bisexual|transgender|lgbtq?\+?)\b", + re.IGNORECASE, + ), + re.compile(r"(?:性取向|同性恋|双性恋|跨性别)"), + re.compile(r"\b(?:race|racial identity|ethnicity|ethnic origin|national origin)\b", re.IGNORECASE), + re.compile(r"(?:种族|民族|族裔)"), + re.compile(r"\b(?:trade union|labor union|union membership)\b", re.IGNORECASE), + re.compile(r"(?:工会成员|工会会籍)"), + re.compile(r"\b(?:underage|minor)\b", re.IGNORECASE), + re.compile(r"(?:未成年)"), + re.compile( + r"\b(?:home address|residential address|precise location|gps coordinates?|latitude|longitude)\b", + re.IGNORECASE, + ), + re.compile(r"(?:家庭住址|住宅地址|精确位置|实时位置|经纬度)"), + re.compile(r"\b(?:biometric|facial recognition|genetic|dna profile)\b", re.IGNORECASE), + re.compile(r"(?:生物特征|人脸识别|基因信息|DNA信息)", re.IGNORECASE), ) diff --git a/app/mneme/memoria/server/models/canonical_memory.py b/app/mneme/memoria/server/models/canonical_memory.py index 490c930..b8990ad 100644 --- a/app/mneme/memoria/server/models/canonical_memory.py +++ b/app/mneme/memoria/server/models/canonical_memory.py @@ -22,6 +22,7 @@ from app.mneme.memoria.server.models.base import Base MemoryStatus = Literal["active", "superseded", "invalidated"] +CanonicalSensitivity = Literal["unknown", "low", "sensitive"] class CanonicalMemory(Base): @@ -35,6 +36,10 @@ class CanonicalMemory(Base): "status IN ('active', 'superseded', 'invalidated')", name="ck_canonical_memories_status", ), + CheckConstraint( + "sensitivity IN ('unknown', 'low', 'sensitive')", + name="ck_canonical_memories_sensitivity", + ), CheckConstraint("confidence >= 0 AND confidence <= 1", name="ck_canonical_memories_confidence"), ForeignKeyConstraint( ["active_revision_id", "memory_id"], @@ -70,6 +75,7 @@ class CanonicalMemory(Base): value: Mapped[str] = mapped_column(Text, nullable=False) fingerprint: Mapped[str] = mapped_column(String(64), nullable=False) confidence: Mapped[float] = mapped_column(Float, nullable=False) + sensitivity: Mapped[str] = mapped_column(String(16), nullable=False, default="unknown", server_default="unknown") retrieval_weight: Mapped[float] = mapped_column(Float, nullable=False, server_default="1") status: Mapped[str] = mapped_column(String(16), nullable=False, server_default="active") active_revision_id: Mapped[str] = mapped_column(String(64), nullable=False) diff --git a/app/mneme/memoria/server/models/memory_settings.py b/app/mneme/memoria/server/models/memory_settings.py index 9be0fd2..8b0f7fd 100644 --- a/app/mneme/memoria/server/models/memory_settings.py +++ b/app/mneme/memoria/server/models/memory_settings.py @@ -18,8 +18,15 @@ class MemorySettings(Base): automatic_conversation_memory: Mapped[bool] = mapped_column( Boolean, nullable=False, default=False, server_default="false" ) + ad_personalization_enabled: Mapped[bool] = mapped_column( + Boolean, nullable=False, default=False, server_default="false" + ) last_event_occurred_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True)) last_event_id: Mapped[str | None] = mapped_column(String(128)) + ad_personalization_last_event_occurred_at: Mapped[datetime | None] = mapped_column( + DateTime(timezone=True) + ) + ad_personalization_last_event_id: Mapped[str | None] = mapped_column(String(128)) created_at: Mapped[datetime] = mapped_column( DateTime(timezone=True), nullable=False, server_default=func.now() ) diff --git a/app/mneme/memoria/server/security/service_tokens.py b/app/mneme/memoria/server/security/service_tokens.py index b1022fc..06f7c9e 100644 --- a/app/mneme/memoria/server/security/service_tokens.py +++ b/app/mneme/memoria/server/security/service_tokens.py @@ -15,6 +15,7 @@ RUNS_READ_SCOPE = "runs:read" MEMORIES_READ_SCOPE = "memories:read" MEMORIES_WRITE_SCOPE = "memories:write" +AD_RECOMMENDATIONS_SCOPE = "ads:recommend" SERVICE_TOKEN_ALGORITHM = "HS256" diff --git a/app/mneme/memoria/server/services/ad_recommendations.py b/app/mneme/memoria/server/services/ad_recommendations.py new file mode 100644 index 0000000..9b090ca --- /dev/null +++ b/app/mneme/memoria/server/services/ad_recommendations.py @@ -0,0 +1,121 @@ +"""Rerank caller-supplied eligible ads using consented low-sensitivity preferences.""" + +import logging + +from sqlalchemy import select + +from app.mneme.memoria.server.contracts.recommendations import ( + AdCandidate, + AdRecommendationItem, + AdRecommendationRequest, + AdRecommendationResponse, +) +from app.mneme.memoria.server.database import open_read_session +from app.mneme.memoria.server.models.canonical_memory import CanonicalMemory +from app.mneme.memoria.server.models.memory_settings import MemorySettings +from app.mneme.memoria.server.observability.context import safe_log +from app.mneme.memoria.server.services.embeddings import embed_texts + +PROFILE_MEMORY_LIMIT = 8 +SEMANTIC_WEIGHT = 0.75 +logger = logging.getLogger(__name__) + + +async def recommend_ads(request: AdRecommendationRequest) -> AdRecommendationResponse: + preferences = await _load_preferences(request) + if not preferences: + return _fallback(request) + + profile_text = " ".join( + f"{memory.subject} {memory.predicate} {memory.value}" + for memory in preferences + ) + candidate_texts = [ + " ".join((candidate.title, candidate.description, *candidate.tags)).strip() + for candidate in request.candidates + ] + try: + profile_vector, *candidate_vectors = await embed_texts([profile_text, *candidate_texts]) + except Exception: + safe_log( + logger, + logging.WARNING, + "ad_recommendation", + status="degraded", + error_code="AD_EMBEDDING_FAILED", + ) + return _fallback(request) + + profile_folded = profile_text.casefold() + scored = [ + ( + _score(profile_vector, vector, candidate.business_score), + index, + candidate, + [tag for tag in candidate.tags if tag.casefold() in profile_folded], + ) + for index, (candidate, vector) in enumerate(zip(request.candidates, candidate_vectors, strict=True)) + ] + scored.sort(key=lambda item: (-item[0], item[1])) + return AdRecommendationResponse( + request_id=request.request_id, + personalized=True, + items=[ + AdRecommendationItem( + ad_id=candidate.ad_id, + score=round(score, 6), + matched_topics=matched_topics, + ) + for score, _, candidate, matched_topics in scored[: request.limit] + ], + ) + + +async def _load_preferences(request: AdRecommendationRequest) -> list[CanonicalMemory]: + scope = ( + CanonicalMemory.knowledge_base_id.is_(None) + if request.knowledge_base_id is None + else CanonicalMemory.knowledge_base_id == request.knowledge_base_id + ) + async with open_read_session() as db: + settings = await db.get(MemorySettings, request.owner_id) + if settings is None or not settings.ad_personalization_enabled: + return [] + return list( + await db.scalars( + select(CanonicalMemory) + .where( + CanonicalMemory.owner_id == request.owner_id, + scope, + CanonicalMemory.status == "active", + CanonicalMemory.sensitivity == "low", + CanonicalMemory.memory_type == "preference", + ) + .order_by( + CanonicalMemory.confidence.desc(), + CanonicalMemory.retrieval_weight.desc(), + CanonicalMemory.updated_at.desc(), + ) + .limit(PROFILE_MEMORY_LIMIT) + ) + ) + + +def _score(profile: list[float], candidate: list[float], business_score: float) -> float: + semantic = (sum(left * right for left, right in zip(profile, candidate, strict=True)) + 1.0) / 2.0 + return max(0.0, min(1.0, SEMANTIC_WEIGHT * semantic + (1.0 - SEMANTIC_WEIGHT) * business_score)) + + +def _fallback(request: AdRecommendationRequest) -> AdRecommendationResponse: + ranked: list[tuple[int, AdCandidate]] = sorted( + enumerate(request.candidates), + key=lambda item: (-item[1].business_score, item[0]), + ) + return AdRecommendationResponse( + request_id=request.request_id, + personalized=False, + items=[ + AdRecommendationItem(ad_id=candidate.ad_id, score=candidate.business_score) + for _, candidate in ranked[: request.limit] + ], + ) diff --git a/app/mneme/memoria/server/services/memory_commands.py b/app/mneme/memoria/server/services/memory_commands.py index 8c71cd3..de99687 100644 --- a/app/mneme/memoria/server/services/memory_commands.py +++ b/app/mneme/memoria/server/services/memory_commands.py @@ -15,7 +15,7 @@ from app.mneme.memoria.server.models.canonical_memory import CanonicalMemory from app.mneme.memoria.server.models.evidence import Evidence, revision_evidence from app.mneme.memoria.server.models.memory_audit import MemoryActionAudit -from app.mneme.memoria.server.models.memory_candidate import MemoryCandidate +from app.mneme.memoria.server.models.memory_candidate import MemoryCandidate, Sensitivity from app.mneme.memoria.server.models.memory_revision import MemoryRevision from app.mneme.memoria.server.repositories.memories import ( hard_delete_memory, @@ -180,6 +180,7 @@ async def revise( predicate: str, value: str, confidence: float | None, + sensitivity: Sensitivity, actor_id: str, reason: str, ) -> CanonicalMemory: @@ -224,6 +225,7 @@ async def revise( predicate=predicate, value=value, confidence=confidence, + sensitivity=sensitivity, actor=actor_id, reason=reason, ) diff --git a/app/mneme/memoria/server/services/memory_events.py b/app/mneme/memoria/server/services/memory_events.py index d42eea7..2cc3751 100644 --- a/app/mneme/memoria/server/services/memory_events.py +++ b/app/mneme/memoria/server/services/memory_events.py @@ -80,16 +80,32 @@ async def _apply_settings(db: AsyncSession, event: AgentEventEnvelope) -> None: row = MemorySettings(owner_id=event.owner_id) db.add(row) await db.flush() - current_order = ( + automatic_order = ( (row.last_event_occurred_at, row.last_event_id) if row.last_event_occurred_at is not None and row.last_event_id is not None else None ) - if current_order is not None and _event_order(event) <= current_order: - return - row.automatic_conversation_memory = payload.automatic_conversation_memory - row.last_event_occurred_at = event.occurred_at - row.last_event_id = event.event_id + if payload.automatic_conversation_memory is not None and ( + automatic_order is None or _event_order(event) > automatic_order + ): + row.automatic_conversation_memory = payload.automatic_conversation_memory + row.last_event_occurred_at = event.occurred_at + row.last_event_id = event.event_id + ad_order = ( + ( + row.ad_personalization_last_event_occurred_at, + row.ad_personalization_last_event_id, + ) + if row.ad_personalization_last_event_occurred_at is not None + and row.ad_personalization_last_event_id is not None + else None + ) + if payload.ad_personalization_enabled is not None and ( + ad_order is None or _event_order(event) > ad_order + ): + row.ad_personalization_enabled = payload.ad_personalization_enabled + row.ad_personalization_last_event_occurred_at = event.occurred_at + row.ad_personalization_last_event_id = event.event_id await db.flush() diff --git a/docker-compose.yml b/docker-compose.yml index a77f505..15e01ff 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -67,11 +67,11 @@ x-memory-agent-base: &memory-agent-base MEMORY_AGENT_CELERY_BROKER_URL: redis://redis:6379/2 MEMORY_AGENT_CELERY_RESULT_BACKEND: redis://redis:6379/3 MEMORY_AGENT_CELERY_QUEUE: ${MEMORY_AGENT_CELERY_QUEUE:-memory_agent} - MEMORY_AGENT_EMBEDDING_MODEL_NAME: ${MEMORY_AGENT_EMBEDDING_MODEL_NAME:-BAAI/bge-m3} - MEMORY_AGENT_EMBEDDING_MODEL_PATH: ${MEMORY_AGENT_EMBEDDING_MODEL_PATH:-} - MEMORY_AGENT_EMBEDDING_CACHE_DIR: ${MEMORY_AGENT_EMBEDDING_CACHE_DIR:-/app/storage/model_cache/sentence_transformers} - MEMORY_AGENT_EMBEDDING_LOCAL_FILES_ONLY: "${MEMORY_AGENT_EMBEDDING_LOCAL_FILES_ONLY:-false}" - MEMORY_AGENT_EMBEDDING_PRELOAD_ON_STARTUP: "${MEMORY_AGENT_EMBEDDING_PRELOAD_ON_STARTUP:-true}" + MEMORY_AGENT_EMBEDDING_MODEL_NAME: ${MEMORY_AGENT_EMBEDDING_MODEL_NAME:-${EMBEDDING_MODEL_NAME:-BAAI/bge-m3}} + MEMORY_AGENT_EMBEDDING_MODEL_PATH: ${MEMORY_AGENT_EMBEDDING_MODEL_PATH:-${EMBEDDING_MODEL_PATH:-}} + MEMORY_AGENT_EMBEDDING_CACHE_DIR: ${MEMORY_AGENT_EMBEDDING_CACHE_DIR:-${EMBEDDING_CACHE_DIR:-/app/storage/model_cache/sentence_transformers}} + MEMORY_AGENT_EMBEDDING_LOCAL_FILES_ONLY: "${MEMORY_AGENT_EMBEDDING_LOCAL_FILES_ONLY:-${EMBEDDING_LOCAL_FILES_ONLY:-false}}" + MEMORY_AGENT_EMBEDDING_PRELOAD_ON_STARTUP: "${MEMORY_AGENT_EMBEDDING_PRELOAD_ON_STARTUP:-${EMBEDDING_PRELOAD_ON_STARTUP:-false}}" MEMORY_AGENT_UVICORN_WORKERS: ${MEMORY_AGENT_UVICORN_WORKERS:-1} MEMORY_AGENT_WORKER_CONCURRENCY: ${MEMORY_AGENT_WORKER_CONCURRENCY:-1} CELERY_BROKER_URL: redis://redis:6379/2 diff --git a/docs/architecture.md b/docs/architecture.md index ae69688..29ee0be 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -59,6 +59,7 @@ retrieval, generation and citation phases. Responsibilities: - retrieve owner-scoped document, memory, profile and relation evidence; +- rerank caller-supplied eligible ads from explicitly enabled, low-sensitivity preferences; - execute bounded single-agent or multi-agent reasoning; - expose only approval proposals for write-class tools; - call the configured primary model and optional fallback model; @@ -114,6 +115,14 @@ loads a snapshot, dispatches by target backend, records retry state, and moves e dead-letter state. Projection consumers are idempotent and must remain rebuildable from the PostgreSQL source records. +### Ad recommendation + +Mneme owns the call boundary; it receives campaign-eligible candidates from its ad source and +supplies them to `POST /v1/ad-recommendations`. Memoria checks the owner-scoped setting and reranks +only from active, low-sensitivity preference memories. Disabled consent, missing preferences, or an +embedding failure returns the candidates in business-score order with `personalized=false`. +Memoria does not own campaigns, impressions, clicks, budgets, or landing-page content. + ### Run control `POST /kb/chat/runs/{run_id}/control` supports: diff --git a/docs/runtime-contracts.md b/docs/runtime-contracts.md index 0b13666..98f8d39 100644 --- a/docs/runtime-contracts.md +++ b/docs/runtime-contracts.md @@ -161,6 +161,25 @@ runtime uses the original bounded conversation context rather than sending an em - Generated answers record selected provider/model, attempts, fallback use, token usage, stop reason and degraded multi-agent state where applicable. +## Ad recommendation + +- `POST /v1/ad-recommendations` requires the `ads:recommend` service scope and an exact matching + owner and optional knowledge-base claim. +- The caller supplies already eligible candidates; Memoria does not decide campaign, budget, + geographic, frequency, or inventory eligibility. +- Personalization is disabled by default and reads only active `preference` memories whose + canonical sensitivity is `low`. Political or religious beliefs, sexual orientation, race or + ethnicity, trade-union membership, minor status, precise location, biometric and genetic data + are sensitive and ineligible. Historical or unclassified memories remain `unknown` and are + ineligible until reclassified. +- The response contains ranked ad IDs, scores, and candidate-owned matched tags. It never returns + memory text or document evidence. +- Disabled consent, no eligible preference, or embedding failure produces a deterministic + business-score fallback with `personalized=false`. +- `user.memory_settings.changed` accepts optional partial updates so older queued events remain + valid; omitted settings retain their current value. Each setting keeps an independent event + order so a delayed consent revocation cannot be discarded by a newer unrelated setting change. + ## Tools and approvals Tool execution is bounded by `app/mneme/memoria/server/runtime/tools.py`. diff --git a/requirements/ai.txt b/requirements/ai.txt index 2079e48..258f53d 100644 --- a/requirements/ai.txt +++ b/requirements/ai.txt @@ -18,6 +18,7 @@ sympy==1.14.0 threadpoolctl==3.6.0 tiktoken==0.12.0 tokenizers==0.22.2 -torch @ https://download.pytorch.org/whl/cpu/torch-2.11.0%2Bcpu-cp312-cp312-manylinux_2_28_x86_64.whl +torch @ https://download.pytorch.org/whl/cpu/torch-2.11.0%2Bcpu-cp312-cp312-manylinux_2_28_x86_64.whl ; platform_machine == "x86_64" +torch @ https://download.pytorch.org/whl/cpu/torch-2.11.0%2Bcpu-cp312-cp312-manylinux_2_28_aarch64.whl ; platform_machine == "aarch64" tqdm==4.67.3 transformers==5.5.0