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/en/docs/operate/configuration.md
Original file line number Diff line number Diff line change
Expand Up @@ -71,7 +71,7 @@ Server settings use the `POWERCONTEXT_SERVER_` prefix.
| `POWERCONTEXT_SERVER_DATABASE_PATH` | user data `seekdb` directory | Embedded seekdb path; used only when `DATABASE_KIND=seekdb` |
| `POWERCONTEXT_SERVER_DATABASE_BUSY_TIMEOUT_MS` | `5000` | Milliseconds a business connection waits for SQLite's single write lock before raising; it does not govern usage accounting, which has its own bounded budget |
| `POWERCONTEXT_SERVER_RUNTIME_SCOPE_CACHE_SIZE` | `128` | Inactive scope compositions retained by the Runtime; in-flight scopes are never evicted |
| `POWERCONTEXT_SERVER_RUNTIME_SOURCE_WINDOW_LIMIT` | `100` | Maximum Sources processed in one activation |
| `POWERCONTEXT_SERVER_RUNTIME_SOURCE_WINDOW_LIMIT` | `100` | Maximum Source journal positions per activation; Memory reduces its window after generation timeouts |
| `POWERCONTEXT_SERVER_RUNTIME_CONTEXT_ASSEMBLY_MAX_ENTRIES` | `8` | Maximum sum of explicit `assembly.sections[].limit`; positive integer. Per-family limits still apply. |
| `POWERCONTEXT_SERVER_RUNTIME_MEMORY_EXTRACTION_PROFILE` | `coding` | Memory selection policy: `coding` or `conversation` |
| `POWERCONTEXT_SERVER_RUNTIME_MEMORY_RERANK_ENABLED` | `false` | Apply listwise reranking after coarse Memory retrieval |
Expand Down
18 changes: 17 additions & 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_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' \
-f <export-directory>
```

Expand Down Expand Up @@ -450,3 +450,19 @@ connections, session charset changes, or other applications sharing the environm
claim that the older driver is free of other vulnerabilities. Use an isolated environment and retain
the required charset. Remove this restriction only after the async driver supports the updated
binary encoding and Source write/read/replay tests pass against both backends.

## Memory extraction repeatedly times out

A generation timeout leaves the Memory Source cursor unchanged. The Memory processor halves the
failed journal window for the next attempt, down to one position. The reduced limit is stored per
Scope and survives Worker replacement and Server restart; it never exceeds `SOURCE_WINDOW_LIMIT`
or an explicit flush limit. Each invocation makes one extraction attempt, so Supervisor backoff and
Worker deadlines continue to bound retries. The reduction is retained while consuming the backlog
visible at the timeout and cleared atomically when that backlog is consumed.

The processor never skips failed Sources or acknowledges a failed invocation. A single Source can
still time out, and an unavailable provider, embedding failure, or hard Worker termination does not
trigger this generation-timeout recovery. Other Scopes continue independently. Inspect the
`memory.window_reduced` log event and the provider error; if one Source still fails, adjust the model
or generation timeout and retry. Keep the Worker timeout above the model request timeout plus
startup and commit overhead so the processor can record a reduction before the Worker is stopped.
7 changes: 7 additions & 0 deletions docs/en/rfcs/1515_artifact_processing_supervisor.md
Original file line number Diff line number Diff line change
Expand Up @@ -496,6 +496,13 @@ Backoff remains in memory. Restart or leadership change may retry an accepted fa
rebuild the backoff sequence; it never skips the failed business position. Permanent errors may continue consuming
invocation resources and must be found through observability; the first version does not discard them automatically.

Memory owns a persisted per-Scope window reduction for generation timeouts. A failed extraction can
halve the next window without advancing its Source cursor or acknowledging the invocation; the
reduction survives Worker replacement and is cleared when the observed backlog is consumed.
The cursor CAS and current fence guard this update. It is an input-size hint, not a successful
processing checkpoint or a Supervisor failed-job state. Single-position failures retain the same
retry policy. Supervisor deadlines and backoff remain unchanged.

## 8. Leadership terms, backends, and process roles

One Lease corresponds to one logical Supervisor, not to a binding, Scope, Worker, or Worker-budget partition. Because
Expand Down
2 changes: 1 addition & 1 deletion docs/zh/docs/operate/configuration.md
Original file line number Diff line number Diff line change
Expand Up @@ -67,7 +67,7 @@ Server 配置使用 `POWERCONTEXT_SERVER_` 前缀。
| `POWERCONTEXT_SERVER_DATABASE_PATH` | 用户数据目录下的 `seekdb` 目录 | 嵌入式 seekdb 路径;仅在 `DATABASE_KIND=seekdb` 时使用 |
| `POWERCONTEXT_SERVER_DATABASE_BUSY_TIMEOUT_MS` | `5000` | 业务连接等待 SQLite 单一写锁的毫秒数;它不约束用量记账,后者有自己的有界预算 |
| `POWERCONTEXT_SERVER_RUNTIME_SCOPE_CACHE_SIZE` | `128` | Runtime 保留的非活动 scope composition 数量;进行中的 scope 不会被驱逐 |
| `POWERCONTEXT_SERVER_RUNTIME_SOURCE_WINDOW_LIMIT` | `100` | 单次 activation 最多处理的 Source 数量 |
| `POWERCONTEXT_SERVER_RUNTIME_SOURCE_WINDOW_LIMIT` | `100` | 每次处理的 Source 日志位置上限;Memory 在生成超时后会缩小窗口 |
| `POWERCONTEXT_SERVER_RUNTIME_CONTEXT_ASSEMBLY_MAX_ENTRIES` | `8` | 显式 `assembly.sections[].limit` 之和的上限;正整数,各类别单独上限仍适用 |
| `POWERCONTEXT_SERVER_RUNTIME_MEMORY_EXTRACTION_PROFILE` | `coding` | Memory 选择策略:`coding` 或 `conversation` |
| `POWERCONTEXT_SERVER_RUNTIME_MEMORY_RERANK_ENABLED` | `false` | 在 Memory 粗召回后应用 listwise rerank |
Expand Down
14 changes: 13 additions & 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_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' \
-f <export-directory>
```

Expand Down Expand Up @@ -424,3 +424,15 @@ PyMySQL 1.2.1、1.2.2 会导致 aiomysql 导入失败;1.2.3 虽然能够导入
此评估不覆盖自定义连接、会话字符集变更或共用环境中的其他应用,也不代表旧驱动不存在其他漏洞。
请使用隔离环境并保留要求的字符集。解除版本限制前,需要异步驱动适配新的二进制编码,并在两种后端上通过
Source 写入、读取和幂等重放测试。

## Memory 提取持续超时

生成超时不会推进 Memory 的 Source 游标。Memory 处理器会将失败的日志窗口减半,供下一次尝试使用,
最小为一个日志位置。缩小后的上限按 Scope 持久化,跨 Worker 更换和 Server 重启保留,且不会超过
`SOURCE_WINDOW_LIMIT` 或显式 flush 的 limit。每次调用只尝试一次提取,仍受 Supervisor 退避和 Worker
截止时间约束。缩窗上限在处理超时时已可见的积压期间保留,并在积压消费完成时随业务提交原子清除。

处理器不会跳过失败的 Source,也不会确认失败调用。单条 Source 仍可能超时;服务不可用、Embedding 失败
或 Worker 被强制终止不会触发这一生成超时恢复策略。其他 Scope 可以独立继续。请检查
`memory.window_reduced` 日志事件及模型错误;如果单条仍失败,可调整模型或生成超时后重试。
Worker 超时应大于模型请求超时与启动、提交开销之和,让处理器能在 Worker 被终止前记录缩窗结果。
5 changes: 5 additions & 0 deletions docs/zh/rfcs/1515_artifact_processing_supervisor.md
Original file line number Diff line number Diff line change
Expand Up @@ -407,6 +407,11 @@ Supervisor 不替处理器作业务授权。显式触发在领域 API 接受意
退避只放内存。重启或换主后可以立即重试已接受的失败调用,随后重新建立退避序列;这不会跳过其业务失败位置。
永久错误可能持续消耗调用成本,需要通过可观测性发现,首版不自动弃置。

Memory 在生成超时后按 Scope 持久化缩窗上限。失败的提取可以将下一次窗口减半,但不推进 Source 游标,
也不确认调用;缩窗上限跨 Worker 更换保留,在已观察到的积压消费完成时清除。该更新受游标 CAS 和当前
任期 fencing 保护。它是输入大小提示,不代表处理成功的 checkpoint,也不是 Supervisor 的失败任务状态。
单个日志位置失败时仍按原策略重试,Supervisor 的截止时间与退避规则保持不变。

## 8. 任期、后端和进程角色

一条 Lease 对应一个逻辑 Supervisor,而不是一个 binding、Scope、Worker 或 Worker 预算分区。`global` 模式只有
Expand Down
46 changes: 46 additions & 0 deletions src/powercontext/builtin/persistence/memory_windows.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,46 @@
# 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.

"""Memory extraction window reductions retained until the backlog is consumed."""

from sqlalchemy import delete, insert, select
from sqlalchemy.ext.asyncio import AsyncConnection

from powercontext.builtin.persistence.tables import MEMORY_SOURCE_WINDOWS_TABLE


class MemorySourceWindowRepository:
"""Keep one bounded recovery hint per Scope, serialized by the Memory cursor CAS."""

async def limit(self, connection: AsyncConnection, scope_id: str, after: int, configured_limit: int) -> int:
table = MEMORY_SOURCE_WINDOWS_TABLE
limit = await connection.scalar(
select(table.c.window_limit).where(table.c.scope_id == scope_id, table.c.source_through > after)
)
return configured_limit if limit is None else min(configured_limit, int(limit))

async def reduce(
self, connection: AsyncConnection, scope_id: str, *, source_through: int, window_limit: int
) -> None:
"""Replace a hint only after acquiring the cursor CAS in this transaction."""

table = MEMORY_SOURCE_WINDOWS_TABLE
await connection.execute(delete(table).where(table.c.scope_id == scope_id))
await connection.execute(
insert(table).values(scope_id=scope_id, source_through=source_through, window_limit=window_limit)
)

async def clear_consumed(self, connection: AsyncConnection, scope_id: str, through: int) -> None:
table = MEMORY_SOURCE_WINDOWS_TABLE
await connection.execute(delete(table).where(table.c.scope_id == scope_id, table.c.source_through <= through))
10 changes: 10 additions & 0 deletions src/powercontext/builtin/persistence/tables.py
Original file line number Diff line number Diff line change
Expand Up @@ -431,6 +431,15 @@ def _entry_text_type():
CheckConstraint("version > 0", name="ck_pc_profile_policies_version_positive"),
)

MEMORY_SOURCE_WINDOWS_TABLE = Table(
"pc_memory_source_windows",
SHARED_METADATA,
Column("scope_id", identity_string(MAX_SCOPE_ID_LENGTH), primary_key=True),
Column("source_through", BigInteger, nullable=False),
Column("window_limit", BigInteger, nullable=False),
CheckConstraint("source_through > 0 AND window_limit > 0", name="ck_pc_memory_source_window_positive"),
)

SOURCE_CURSORS_TABLE = Table(
"pc_source_cursors",
SHARED_METADATA,
Expand Down Expand Up @@ -907,6 +916,7 @@ def _entry_text_type():
ARTIFACT_CANDIDATE_HEADS_TABLE,
PROFILE_POLICIES_TABLE,
SOURCE_CURSORS_TABLE,
MEMORY_SOURCE_WINDOWS_TABLE,
ARTIFACT_PROCESSING_LEASES_TABLE,
ARTIFACT_PROCESSING_BINDING_STATES_TABLE,
TOPIC_MEMORY_WORK_BUDGETS_TABLE,
Expand Down
68 changes: 63 additions & 5 deletions src/powercontext/builtin/runtime/relational.py
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,7 @@
from sqlalchemy import select
from sqlalchemy.ext.asyncio import AsyncConnection

from powercontext._logging import log_safely
from powercontext.artifacts import Artifact, ArtifactRef, MemoryCitation
from powercontext.builtin.artifacts.experience import (
EXPERIENCE_INCUBATION_CURSOR_NAME,
Expand Down Expand Up @@ -114,7 +115,13 @@
from powercontext.builtin.dream.models import DreamBudget, DreamOperation
from powercontext.builtin.dream.service import CandidateAttester, DreamAuthorizer, DreamService
from powercontext.builtin.evidence.resolver import AuthorizationContext, EvidenceAuthorizer, EvidenceResolver
from powercontext.builtin.inference import EmbeddingModel, InferenceUsage, InvalidInferenceOutputError, TokenEstimator
from powercontext.builtin.inference import (
EmbeddingModel,
InferenceTimeoutError,
InferenceUsage,
InvalidInferenceOutputError,
TokenEstimator,
)
from powercontext.builtin.persistence.agent_skill_targets import RemoteAgentSkillTargetRepository
from powercontext.builtin.persistence.artifact_governance import (
ArtifactGovernance,
Expand All @@ -125,7 +132,7 @@
from powercontext.builtin.persistence.artifacts import ArtifactRepository
from powercontext.builtin.persistence.candidates import CandidateRepository
from powercontext.builtin.persistence.connectors import ConnectorCheckpointRepository
from powercontext.builtin.persistence.cursors import SourceCursorRepository
from powercontext.builtin.persistence.cursors import SourceCursorRepository, StoredSourceCursor
from powercontext.builtin.persistence.database import AsyncDatabase
from powercontext.builtin.persistence.errors import RepositoryNotFoundError, StoredPayloadConflictError
from powercontext.builtin.persistence.experience_index import ExperienceIndex, NoExperienceIndex
Expand All @@ -145,6 +152,7 @@
)
from powercontext.builtin.persistence.memory import RelationalMemoryBackend
from powercontext.builtin.persistence.memory_index import MemoryIndex, NoMemoryIndex
from powercontext.builtin.persistence.memory_windows import MemorySourceWindowRepository
from powercontext.builtin.persistence.processing import ArtifactProcessingPendingRepository
from powercontext.builtin.persistence.processing_intents import ArtifactProcessingIntentRepository
from powercontext.builtin.persistence.records import RelationalRecordService
Expand Down Expand Up @@ -1707,7 +1715,12 @@ async def flush(
connection,
self._services.scope_id,
)
# Validate the caller's limit before applying a persisted reduction.
signal = SourceHighWatermark(sequence=high_watermark, limit=limit)
windows = MemorySourceWindowRepository()
signal = signal.model_copy(
update={"limit": await windows.limit(connection, self._services.scope_id, state.sequence, limit)}
)
transition = self._trigger.activate(signal, state)
sources = () if not transition.actions else await self._sources(connection, transition.actions[0])
if not transition.actions and processing is not None:
Expand All @@ -1722,9 +1735,14 @@ async def flush(
)

action = transition.actions[0]
prepared = (
None if not sources else await self._prepare_memory(sources, authorize_snapshot=authorize_snapshot)
)
try:
prepared = (
None if not sources else await self._prepare_memory(sources, authorize_snapshot=authorize_snapshot)
)
except InferenceTimeoutError as error:
if error.operation == "generate" and action.through - action.after > 1:
await self._reduce_memory_window(action, high_watermark, state_row, processing)
raise
held = _is_held_write(prepared)
commit = None if prepared is None else prepared.commit
with self._stage(
Expand All @@ -1748,6 +1766,7 @@ async def flush(
transition.state,
expected_generation=None if state_row is None else state_row.generation,
)
await windows.clear_consumed(connection, self._services.scope_id, action.through)
if on_commit is not None and prepared is not None:
before = prepared.result if prepared.commit is None else prepared.commit.base
await on_commit(connection, before, updated)
Expand All @@ -1763,6 +1782,45 @@ async def flush(
hold_codes=_hold_codes(prepared),
)

async def _reduce_memory_window(
self,
action: ProcessSourceWindow,
high_watermark: int,
state_row: StoredSourceCursor | None,
processing: ScopeInvocation | None,
) -> None:
reduced_limit = max(1, (action.through - action.after) // 2)
async with self._services.database.transaction() as connection:
if processing is not None:
await processing.guard(connection)
# Preserve the position, but invalidate stale publishers before
# changing the window used by a new Worker or SDK flush.
await self._services.repositories.cursors.save(
connection,
self._services.scope_id,
SOURCE_WINDOW_TRIGGER_NAME,
SourceCursor(sequence=action.after),
expected_generation=None if state_row is None else state_row.generation,
)
await MemorySourceWindowRepository().reduce(
connection,
self._services.scope_id,
source_through=high_watermark,
window_limit=reduced_limit,
)
log_safely(
logger,
logging.WARNING,
"Memory extraction timed out; reduced the next Source window",
extra={
"event": "memory.window_reduced",
"scope_id": self._services.scope_id,
"source_after": action.after,
"source_through": action.through,
"window_limit": reduced_limit,
},
)

async def _sources(
self,
connection: AsyncConnection,
Expand Down
Loading
Loading