Skip to content
Open
8 changes: 8 additions & 0 deletions .env.example
Original file line number Diff line number Diff line change
Expand Up @@ -109,6 +109,14 @@ POWERCONTEXT_SERVER_RUNTIME_DREAM_MAX_PENDING_PER_SCOPE=32
# Memory write gate: opt-in evidence-sufficiency check before a write commits. Disabled by
# default; when unset no gate runs and no extra model call is made.
# POWERCONTEXT_SERVER_RUNTIME_MEMORY_WRITE_GATE_ENABLED=true
# `disabled` does not construct a gate. `shadow` records an internal observation but always lets
# the Memory write proceed. `advisory` turns an insufficient-evidence result into FLAG, while
# `enforcing` preserves the visible ACCEPT/FLAG/HOLD outcomes.
# POWERCONTEXT_SERVER_RUNTIME_MEMORY_WRITE_GATE_MODE=shadow
# The default makes no decision-model call. Set `local_only` only for an in-process or loopback
# decision backend controlled by this deployment. `hosted_redacted` is rejected until a shared
# PowerContext content sanitizer is separately reviewed.
# POWERCONTEXT_SERVER_RUNTIME_MEMORY_WRITE_GATE_PRIVACY_BOUNDARY=local_only

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[P2] The example value contradicts the actual default for the privacy boundary

RuntimeConfig.memory_write_gate_privacy_boundary defaults to no_external_call (runtime/config.py, and test_memory_write_gate_defaults_to_no_external_call pins it), and the comment above this line correctly says "The default makes no decision-model call". But by this file's own convention the commented line shows the default (..._MODE=shadow two lines up matches its actual default), and here it shows local_only - a value that DOES make a decision-model call on a loopback backend.

An operator who uncomments the line to "keep the documented default" actually enables model calls. Suggest changing the example to # POWERCONTEXT_SERVER_RUNTIME_MEMORY_WRITE_GATE_PRIVACY_BOUNDARY=no_external_call (and optionally a follow-up line showing the local_only opt-in), or rewording so the example is clearly an opt-in and not the default.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed in 5098f5d. .env.example now shows no_external_call, matching RuntimeConfig and the surrounding documentation. local_only remains described as the explicit loopback opt-in.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed in 5098f5d. .env.example now shows no_external_call, matching RuntimeConfig and the surrounding documentation. local_only remains described as the explicit loopback opt-in.

# Direction only (which verdict means the cited evidence is insufficient): "yes" or "no".
# POWERCONTEXT_SERVER_RUNTIME_MEMORY_WRITE_GATE_HOLD_ON=yes
# Optional confidence floor below which a hold becomes a written-but-annotated change.
Expand Down
2 changes: 1 addition & 1 deletion docs/en/docs/operate/troubleshoot.md
Original file line number Diff line number Diff line change
Expand Up @@ -267,7 +267,7 @@ so the previous database remains available for recovery:

```bash
obloader <connection-options> -D <new-database> --csv \
--table 'pc_scopes,pc_source_journal_heads,pc_sources,pc_artifacts,pc_source_cursors,pc_memory_source_windows,pc_artifact_processing_leases,pc_artifact_processing_binding_states,pc_artifact_processing_pending,pc_artifact_processing_auto_wave_targets,pc_artifact_processing_sequences,pc_artifact_processing_intents,pc_topic_memory_processing_targets,pc_artifact_processing_schema,pc_artifact_processing_migration_receipts,pc_topic_memory_work_budgets,pc_topic_memory_retrieval_shape,pc_connector_checkpoints,pc_source_definition_manifests,pc_external_skill_registrations,pc_skill_packages,pc_agent_skill_targets,pc_skill_publications,pc_model_usage_daily,pc_recall_token_daily,pc_receipt_migration_review' \
--table 'pc_scopes,pc_source_journal_heads,pc_sources,pc_artifacts,pc_source_cursors,pc_memory_source_windows,pc_artifact_processing_leases,pc_artifact_processing_binding_states,pc_artifact_processing_pending,pc_artifact_processing_auto_wave_targets,pc_artifact_processing_sequences,pc_artifact_processing_intents,pc_topic_memory_processing_targets,pc_artifact_processing_schema,pc_artifact_processing_migration_receipts,pc_topic_memory_work_budgets,pc_topic_memory_retrieval_shape,pc_connector_checkpoints,pc_source_definition_manifests,pc_external_skill_registrations,pc_skill_packages,pc_agent_skill_targets,pc_skill_publications,pc_model_usage_daily,pc_recall_token_daily,pc_receipt_migration_review,pc_decision_observations' \
-f <export-directory>
```

Expand Down
2 changes: 1 addition & 1 deletion docs/zh/docs/operate/troubleshoot.md
Original file line number Diff line number Diff line change
Expand Up @@ -259,7 +259,7 @@ collation,但不会包含数据库 URL 或凭据。

```bash
obloader <connection-options> -D <new-database> --csv \
--table 'pc_scopes,pc_source_journal_heads,pc_sources,pc_artifacts,pc_source_cursors,pc_memory_source_windows,pc_artifact_processing_leases,pc_artifact_processing_binding_states,pc_artifact_processing_pending,pc_artifact_processing_auto_wave_targets,pc_artifact_processing_sequences,pc_artifact_processing_intents,pc_topic_memory_processing_targets,pc_artifact_processing_schema,pc_artifact_processing_migration_receipts,pc_topic_memory_work_budgets,pc_topic_memory_retrieval_shape,pc_connector_checkpoints,pc_source_definition_manifests,pc_external_skill_registrations,pc_skill_packages,pc_agent_skill_targets,pc_skill_publications,pc_model_usage_daily,pc_recall_token_daily,pc_receipt_migration_review' \
--table 'pc_scopes,pc_source_journal_heads,pc_sources,pc_artifacts,pc_source_cursors,pc_memory_source_windows,pc_artifact_processing_leases,pc_artifact_processing_binding_states,pc_artifact_processing_pending,pc_artifact_processing_auto_wave_targets,pc_artifact_processing_sequences,pc_artifact_processing_intents,pc_topic_memory_processing_targets,pc_artifact_processing_schema,pc_artifact_processing_migration_receipts,pc_topic_memory_work_budgets,pc_topic_memory_retrieval_shape,pc_connector_checkpoints,pc_source_definition_manifests,pc_external_skill_registrations,pc_skill_packages,pc_agent_skill_targets,pc_skill_publications,pc_model_usage_daily,pc_recall_token_daily,pc_receipt_migration_review,pc_decision_observations' \
-f <export-directory>
```

Expand Down
31 changes: 30 additions & 1 deletion src/powercontext/builtin/artifacts/memory/protocols.py
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,7 @@
from contextlib import AbstractAsyncContextManager
from dataclasses import dataclass
from enum import StrEnum
from typing import Protocol
from typing import TYPE_CHECKING, Protocol, runtime_checkable

from pydantic import BaseModel, ConfigDict

Expand All @@ -39,6 +39,18 @@
from powercontext.builtin.tags import TagFilter
from powercontext.sources import Source

if TYPE_CHECKING:
from powercontext.builtin.decision_observations import DecisionObservation


class MemoryWriteObservationSink(Protocol):
"""Accept a bounded decision sidecar without owning the Memory write."""

async def record(self, observation: DecisionObservation, /) -> None:
"""Store one decision observation after the gate has made its judgement."""

...


class MemoryCandidateRequest(BaseModel):
"""Canonical evidence and bounded current entries offered to a pipeline."""
Expand Down Expand Up @@ -110,6 +122,11 @@ class MemoryWriteGateRequest:
candidates: tuple[str, ...]
evidence: tuple[str, ...]
expected_revision: int | None = None
scope_id: str = "unscoped"
operation_id: str | None = None
subject_refs: tuple[str, ...] = ()
evidence_refs: tuple[str, ...] = ()
observation_sink: MemoryWriteObservationSink | None = None


class MemoryWriteGate(Protocol):
Expand All @@ -127,6 +144,18 @@ async def assess(self, request: MemoryWriteGateRequest, /) -> MemoryWriteAssessm
...


@runtime_checkable
class MemoryWriteGatePreflight(Protocol):
"""Optionally map a deterministic service-side preflight rejection through a gate policy."""

async def assess_preflight(
self, request: MemoryWriteGateRequest, rejection: MemoryWriteAssessment, /
) -> MemoryWriteAssessment:
"""Return the mode-aware decision for a bounded preflight rejection."""

...


class MemoryWritePlan(BaseModel):
"""A side-effect-free result that can be committed in an outer transaction."""

Expand Down
45 changes: 38 additions & 7 deletions src/powercontext/builtin/artifacts/memory/service.py
Original file line number Diff line number Diff line change
Expand Up @@ -92,7 +92,9 @@
MemorySearchRequest,
MemoryWriteAssessment,
MemoryWriteGate,
MemoryWriteGatePreflight,
MemoryWriteGateRequest,
MemoryWriteObservationSink,
MemoryWritePlan,
MemoryWriteRejectionCode,
MemoryWriteVerdict,
Expand Down Expand Up @@ -259,7 +261,9 @@ def __init__(
artifact_resolver: _ArtifactResolver | None = None,
id_factory: IdFactory | None = None,
prompt_context: ScopedPrompts | None = None,
scope_id: str = "unscoped",
write_gate: MemoryWriteGate | None = None,
write_gate_observation_sink: MemoryWriteObservationSink | None = None,
capacity_budget: MemoryCapacityBudget | None = None,
compaction: MemoryCompactionPolicy | None = None,
max_history_revisions: int = 100,
Expand All @@ -268,6 +272,8 @@ def __init__(
self._prompt_context = prompt_context
self._candidate_pipeline = candidate_pipeline
self._write_gate = write_gate
self._write_gate_observation_sink = write_gate_observation_sink
self._scope_id = scope_id
self._embedding_model = embedding_model
if rerank_candidate_limit < 1:
raise _InvalidMemoryOperationError("search-limit")
Expand Down Expand Up @@ -1375,17 +1381,23 @@ async def _assess_write(
if self._write_gate is None:
return None
projection = await self._gate_evidence(base, candidates, evidence, current_entries)
request = MemoryWriteGateRequest(
candidates=tuple(candidate.text for candidate in candidates),
evidence=projection.entries,
expected_revision=None if base is None else base.revision,
scope_id=self._scope_id,
operation_id=_gate_operation_id(base),
subject_refs=_gate_subject_refs(candidates),
evidence_refs=tuple(_gate_evidence_ref(entry) for entry in projection.entries),
observation_sink=self._write_gate_observation_sink,
Comment thread
Copilot marked this conversation as resolved.
Outdated
)
if projection.rejection is not None:
if isinstance(self._write_gate, MemoryWriteGatePreflight):
return await self._write_gate.assess_preflight(request, projection.rejection)

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[P2] assess_preflight is not covered by the fail-open fallback that protects assess

_assess_write wraps only the assess() call in try/except Exception -> ACCEPT (a few lines below), but the new assess_preflight delegation returns directly with no exception handling. A gate that implements the optional MemoryWriteGatePreflight contract and raises during preflight (a bug, or a transient error inside a custom implementation) propagates out of _assess_write and fails the whole plan_remember, while the same gate's assess path fails open by design (failure_policy=FAIL_OPEN, and the PR's stated behavior is that gate problems never block writes).

Before this PR the projection-rejection path never called gate code, so this is a new failure mode introduced by the preflight hook. Suggest wrapping the preflight call in the same fail-open fallback (return an ACCEPT assessment with used_fallback=True and log the error type), so both entry points of a preflight-capable gate share one failure policy. A focused test: a preflight gate whose assess_preflight raises should still yield a committable plan.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed in 5098f5d. assess_preflight now shares the fail-open behavior of assess(): any gate exception returns an ACCEPT assessment with used_fallback=True, leaving the plan committable. Added a focused raising-preflight regression test.

_log_gate_assessment(projection.rejection)
return projection.rejection
try:
return await self._write_gate.assess(
MemoryWriteGateRequest(
candidates=tuple(candidate.text for candidate in candidates),
evidence=projection.entries,
expected_revision=None if base is None else base.revision,
)
)
return await self._write_gate.assess(request)
except Exception:
return MemoryWriteAssessment(
verdict=MemoryWriteVerdict.ACCEPT,
Expand Down Expand Up @@ -1839,6 +1851,25 @@ def _candidate_gate_identity(candidate_index: int, identity: str) -> str:
return f"candidate:{candidate_index} {identity}"


def _gate_operation_id(base: Memory | None) -> str:
if base is None:
return "memory-write:new"
return f"memory-write:{base.artifact_id}@{base.revision + 1}"


def _gate_subject_refs(candidates: tuple[MemoryEntryInput, ...]) -> tuple[str, ...]:
return tuple(
f"candidate:{index}"
if candidate.entry is None
else f"entry:{candidate.entry.entry_id}@{candidate.entry.entry_version_id}"
for index, candidate in enumerate(candidates, start=1)

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed in 5098f5d. Applied writes replace candidate placeholders with exact committed entry/version refs and the committed Memory revision operation ID. Held, unchanged, and failed candidates cannot be reconstructed without retaining raw input, so their sidecars explicitly clear subject_refs and set incomplete_subject_count rather than claiming replayability.

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[P2] Retain references for the complete assessed candidate batch

Applied single-candidate writes now have exact committed refs, but batch coverage is still incomplete on cd2bc83. Through MemoryService.remember() with real SQLite, storing First fact. and then writing [First fact., Second fact.] evaluates two candidates but persists only the newly created Second fact. entry in subject_refs; metadata says candidate_count=2 and has no incomplete_subject_count.

apply() uses only plan.commit.entry_versions, which excludes deduplicated/unchanged candidates and loses their ordinal mapping to evidence. Preserve exact refs for those candidates too, or explicitly record incomplete subjects and their positions, so re-evaluation can reconstruct the input actually judged.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed as part of the Atomic Memory migration in abeed49c. The collection Memory service no longer exists on the current base. The Atomic window gate assesses the complete prepared change set, and publish derives one final, encoded primary Artifact reference from every corresponding mutation result before persisting the observation. It does not derive subjects from only newly created collection entries.

)


def _gate_evidence_ref(value: str) -> str:
return value.partition("\n")[0]


def _source_gate_content(source: Source, resolver: _SourceResolver | None) -> str | None:
content = getattr(source, "content", None)
if isinstance(content, str):
Expand Down
159 changes: 159 additions & 0 deletions src/powercontext/builtin/decision_observations.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,159 @@
# Copyright (c) 2026 OceanBase.
#
# Licensed under the Apache License, Version 2.0 (the "License");
# you may not use this file except in compliance with the License.
# You may obtain a copy of the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS,
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.

"""Runtime-independent contracts for bounded decision-policy sidecars."""

from __future__ import annotations

from collections.abc import Mapping
from datetime import UTC, datetime
from enum import StrEnum
from typing import Annotated
from uuid import uuid4

from pydantic import BaseModel, ConfigDict, Field, JsonValue, model_validator

from powercontext.builtin.inference import InferenceUsage


class DecisionPolicyMode(StrEnum):
"""Runtime mode for one versioned decision policy."""

DISABLED = "disabled"
SHADOW = "shadow"
ADVISORY = "advisory"
ENFORCING = "enforcing"


class DecisionFailurePolicy(StrEnum):
"""How a consumer treats backend failure when it owns a domain action."""

FAIL_OPEN = "fail_open"
FAIL_CLOSED = "fail_closed"


class DecisionPrivacyBoundary(StrEnum):
"""Content boundary declared by a policy before backend calls."""

LOCAL_ONLY = "local_only"
HOSTED_REDACTED = "hosted_redacted"
REFERENCES_ONLY = "references_only"
NO_EXTERNAL_CALL = "no_external_call"


class DecisionCoverage(StrEnum):
"""Whether a policy evaluation actually judged the supplied content."""

ADJUDICATED = "adjudicated"
UNADJUDICATED = "unadjudicated"


class DecisionVerdict(StrEnum):
"""Domain-neutral verdict before an owning service maps it to an action."""

ALLOW = "allow"
DENY = "deny"
REVIEW = "review"
UNKNOWN = "unknown"


class DecisionAssessmentSource(StrEnum):
"""The source that produced the assessment's substantive judgement."""

LOCAL_RULE = "local_rule"
DECISION_MODEL = "decision_model"
NONE = "none"


class _StrictModel(BaseModel):
model_config = ConfigDict(extra="forbid", frozen=True)


class DecisionPolicy(_StrictModel):
"""Versioned, reviewable policy manifest for one bounded runtime question."""

policy_id: Annotated[str, Field(min_length=1, max_length=256)]
decision_kind: Annotated[str, Field(min_length=1, max_length=128)]
version: Annotated[str, Field(min_length=1, max_length=64)]
consumer: Annotated[str, Field(min_length=1, max_length=128)]
mode: DecisionPolicyMode
failure_policy: DecisionFailurePolicy
privacy_boundary: DecisionPrivacyBoundary
local_rules: tuple[Annotated[str, Field(min_length=1, max_length=128)], ...] = ()
question: Annotated[str, Field(min_length=1, max_length=8192)]
subject_selector: Annotated[str, Field(min_length=1, max_length=512)]
evidence_selector: Annotated[str, Field(min_length=1, max_length=512)]
outcome_mapping: Mapping[str, Annotated[str, Field(min_length=1, max_length=128)]] = Field(default_factory=dict)
promotion_criteria: tuple[Annotated[str, Field(min_length=1, max_length=512)], ...] = ()


class DecisionAssessment(_StrictModel):
"""Policy-level assessment before domain-specific action mapping."""

policy_id: Annotated[str, Field(min_length=1, max_length=256)]
policy_version: Annotated[str, Field(min_length=1, max_length=64)]
mode: DecisionPolicyMode
coverage: DecisionCoverage
verdict: DecisionVerdict
source: DecisionAssessmentSource
reason: Annotated[str | None, Field(min_length=1, max_length=4096)] = None
confidence: Annotated[float | None, Field(ge=0.0, le=1.0)] = None
used_fallback: bool = False
usage: InferenceUsage = Field(default_factory=lambda: InferenceUsage(requests=0))
latency_ms: Annotated[float | None, Field(ge=0.0, allow_inf_nan=False)] = None

@model_validator(mode="after")
def validate_fallback_coverage(self) -> DecisionAssessment:
if self.used_fallback and self.coverage is not DecisionCoverage.UNADJUDICATED:
raise ValueError("fallback assessments must be unadjudicated") # noqa: TRY003
return self


class DecisionObservation(_StrictModel):
"""Audit/replay sidecar for one policy evaluation attempt."""

observation_id: Annotated[str, Field(min_length=1, max_length=128)] = Field(
default_factory=lambda: f"decision-observation-{uuid4().hex}"
)
created_at: datetime = Field(default_factory=lambda: datetime.now(UTC))
operation_id: Annotated[str, Field(min_length=1, max_length=256)]
scope_id: Annotated[str, Field(min_length=1, max_length=256)]
consumer: Annotated[str, Field(min_length=1, max_length=128)]
policy_id: Annotated[str, Field(min_length=1, max_length=256)]
policy_version: Annotated[str, Field(min_length=1, max_length=64)]
mode: DecisionPolicyMode
subject_refs: tuple[Annotated[str, Field(min_length=1, max_length=512)], ...] = ()
evidence_refs: tuple[Annotated[str, Field(min_length=1, max_length=512)], ...] = ()
privacy_boundary: DecisionPrivacyBoundary
privacy_outcome: Annotated[str | None, Field(min_length=1, max_length=128)] = None
provider_id: Annotated[str | None, Field(min_length=1, max_length=128)] = None
backend_model_id: Annotated[str | None, Field(min_length=1, max_length=256)] = None
model_policy_id: Annotated[str | None, Field(min_length=1, max_length=256)] = None
assessment: DecisionAssessment
final_action: Annotated[str, Field(min_length=1, max_length=128)]
fallback_reason: Annotated[str | None, Field(min_length=1, max_length=256)] = None
metadata: Mapping[str, JsonValue] = Field(default_factory=dict)


__all__ = [
"DecisionAssessment",
"DecisionAssessmentSource",
"DecisionCoverage",
"DecisionFailurePolicy",
"DecisionObservation",
"DecisionPolicy",
"DecisionPolicyMode",
"DecisionPrivacyBoundary",
"DecisionVerdict",
]
2 changes: 2 additions & 0 deletions src/powercontext/builtin/persistence/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@
from powercontext.builtin.persistence.candidates import CandidateRepository
from powercontext.builtin.persistence.connectors import ConnectorCheckpointRepository
from powercontext.builtin.persistence.database import AsyncDatabase
from powercontext.builtin.persistence.decision_observations import DecisionObservationRepository
from powercontext.builtin.persistence.errors import (
ArtifactProcessingLeadershipLostError,
ArtifactProcessingWaveIncompleteError,
Expand Down Expand Up @@ -78,6 +79,7 @@
"CompositeTopicMemoryIndex",
"ConnectorCheckpointRepository",
"DatabaseClosedError",
"DecisionObservationRepository",
"ExternalSkillRepository",
"GenerationConflictError",
"IdentityMismatchError",
Expand Down
Loading