From a0337e9434b041b72e43f4607ef06a4d8c01000c Mon Sep 17 00:00:00 2001 From: knqiufan Date: Mon, 28 Sep 2026 16:49:28 +0800 Subject: [PATCH] fix(memory): persist smaller windows after extraction timeouts --- docs/en/docs/operate/configuration.md | 2 +- docs/en/docs/operate/troubleshoot.md | 18 +- .../1515_artifact_processing_supervisor.md | 7 + docs/zh/docs/operate/configuration.md | 2 +- docs/zh/docs/operate/troubleshoot.md | 14 +- .../1515_artifact_processing_supervisor.md | 5 + .../builtin/persistence/memory_windows.py | 46 +++++ .../builtin/persistence/tables.py | 10 + .../builtin/runtime/relational.py | 68 ++++++- .../builtin/runtime/test_family_processing.py | 57 ++++++ .../runtime/test_memory_window_recovery.py | 179 ++++++++++++++++++ tests/e2e/real_experience_skill/harness.py | 1 + 12 files changed, 400 insertions(+), 9 deletions(-) create mode 100644 src/powercontext/builtin/persistence/memory_windows.py create mode 100644 tests/builtin/runtime/test_memory_window_recovery.py diff --git a/docs/en/docs/operate/configuration.md b/docs/en/docs/operate/configuration.md index ae0dcba1e6..cb2c5b8c0b 100644 --- a/docs/en/docs/operate/configuration.md +++ b/docs/en/docs/operate/configuration.md @@ -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 | diff --git a/docs/en/docs/operate/troubleshoot.md b/docs/en/docs/operate/troubleshoot.md index 74fe0144f8..8cbfdcc4c8 100644 --- a/docs/en/docs/operate/troubleshoot.md +++ b/docs/en/docs/operate/troubleshoot.md @@ -267,7 +267,7 @@ so the previous database remains available for recovery: ```bash obloader -D --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 ``` @@ -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. diff --git a/docs/en/rfcs/1515_artifact_processing_supervisor.md b/docs/en/rfcs/1515_artifact_processing_supervisor.md index 363dc16ecd..e16e73a00d 100644 --- a/docs/en/rfcs/1515_artifact_processing_supervisor.md +++ b/docs/en/rfcs/1515_artifact_processing_supervisor.md @@ -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 diff --git a/docs/zh/docs/operate/configuration.md b/docs/zh/docs/operate/configuration.md index b5db43b09d..e94ba2ff01 100644 --- a/docs/zh/docs/operate/configuration.md +++ b/docs/zh/docs/operate/configuration.md @@ -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 | diff --git a/docs/zh/docs/operate/troubleshoot.md b/docs/zh/docs/operate/troubleshoot.md index b76bbd95a8..f7ddacc160 100644 --- a/docs/zh/docs/operate/troubleshoot.md +++ b/docs/zh/docs/operate/troubleshoot.md @@ -259,7 +259,7 @@ collation,但不会包含数据库 URL 或凭据。 ```bash obloader -D --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 ``` @@ -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 被终止前记录缩窗结果。 diff --git a/docs/zh/rfcs/1515_artifact_processing_supervisor.md b/docs/zh/rfcs/1515_artifact_processing_supervisor.md index 3fe235f17e..81b80aadc0 100644 --- a/docs/zh/rfcs/1515_artifact_processing_supervisor.md +++ b/docs/zh/rfcs/1515_artifact_processing_supervisor.md @@ -407,6 +407,11 @@ Supervisor 不替处理器作业务授权。显式触发在领域 API 接受意 退避只放内存。重启或换主后可以立即重试已接受的失败调用,随后重新建立退避序列;这不会跳过其业务失败位置。 永久错误可能持续消耗调用成本,需要通过可观测性发现,首版不自动弃置。 +Memory 在生成超时后按 Scope 持久化缩窗上限。失败的提取可以将下一次窗口减半,但不推进 Source 游标, +也不确认调用;缩窗上限跨 Worker 更换保留,在已观察到的积压消费完成时清除。该更新受游标 CAS 和当前 +任期 fencing 保护。它是输入大小提示,不代表处理成功的 checkpoint,也不是 Supervisor 的失败任务状态。 +单个日志位置失败时仍按原策略重试,Supervisor 的截止时间与退避规则保持不变。 + ## 8. 任期、后端和进程角色 一条 Lease 对应一个逻辑 Supervisor,而不是一个 binding、Scope、Worker 或 Worker 预算分区。`global` 模式只有 diff --git a/src/powercontext/builtin/persistence/memory_windows.py b/src/powercontext/builtin/persistence/memory_windows.py new file mode 100644 index 0000000000..1228ba5b1e --- /dev/null +++ b/src/powercontext/builtin/persistence/memory_windows.py @@ -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)) diff --git a/src/powercontext/builtin/persistence/tables.py b/src/powercontext/builtin/persistence/tables.py index 530d3f13de..e38d0d099f 100644 --- a/src/powercontext/builtin/persistence/tables.py +++ b/src/powercontext/builtin/persistence/tables.py @@ -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, @@ -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, diff --git a/src/powercontext/builtin/runtime/relational.py b/src/powercontext/builtin/runtime/relational.py index e402860222..ff080768a8 100644 --- a/src/powercontext/builtin/runtime/relational.py +++ b/src/powercontext/builtin/runtime/relational.py @@ -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, @@ -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, @@ -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 @@ -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 @@ -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: @@ -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( @@ -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) @@ -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, diff --git a/tests/builtin/runtime/test_family_processing.py b/tests/builtin/runtime/test_family_processing.py index 0922561908..823fcba957 100644 --- a/tests/builtin/runtime/test_family_processing.py +++ b/tests/builtin/runtime/test_family_processing.py @@ -29,6 +29,7 @@ import powercontext.builtin.runtime.family_processing as family_processing from powercontext.builtin.artifacts.experience import ExperienceCandidateInput, ExperienceContent from powercontext.builtin.artifacts.memory import MemoryEntryInput +from powercontext.builtin.inference import InferenceTimeoutError from powercontext.builtin.inference.models import GenerationResult, InferenceUsage from powercontext.builtin.inference.usage import UsageReportingStructuredGenerator from powercontext.builtin.persistence.cursors import SourceCursorRepository @@ -160,6 +161,62 @@ def security_spec(*, allowed=True): ).model_dump(mode="json") +def test_memory_timeout_retries_shrink_without_acknowledging_or_skipping_input(tmp_path, monkeypatch): + windows = [] + original = MemoryPipeline.extract + + async def bounded_extract(self, request): + windows.append(tuple(source.name for source in request.sources)) + if len(request.sources) > 1: + raise InferenceTimeoutError("generate", 60) + return await original(self, request) + + monkeypatch.setattr(MemoryPipeline, "extract", bounded_extract) + + async def scenario(): + config = BuiltinConfig(database=SQLiteConfig(url=f"sqlite+aiosqlite:///{tmp_path / 'timeout.db'}")) + async with _open_sqlite(config, tables=BUILTIN_TABLES) as profile: + contexts, assignment = await prepare(profile, "memory") + for index in range(3): + await contexts.records.create_source(assignment.scope_id, "content", f"Additional fact {index}") + for _ in range(2): + with pytest.raises(InferenceTimeoutError): + await process_family_invocation(contexts, assignment, config=config) + async with profile.database.transaction() as connection: + cursor = await SourceCursorRepository().load( + connection, assignment.scope_id, assignment.binding_name + ) + intent = await ArtifactProcessingIntentRepository().load( + connection, assignment.scope_id, assignment.binding_name + ) + assert cursor is not None and cursor.cursor.sequence == 0 + assert intent is not None and intent.handled_generation == 0 + for position in range(1, 5): + assert (await process_family_invocation(contexts, assignment, config=config)).outcome == "succeeded" + # Replaying an acknowledged invocation cannot process its successor. + await process_family_invocation(contexts, assignment, config=config) + async with profile.database.transaction() as connection: + cursor = await SourceCursorRepository().load( + connection, assignment.scope_id, assignment.binding_name + ) + assert cursor is not None and cursor.cursor.sequence == position + intent = await ArtifactProcessingIntentRepository().load( + connection, assignment.scope_id, assignment.binding_name + ) + assert intent is not None and intent.handled_generation == assignment.claimed_request_generation + if position < 4: + assert intent.dirty_generation > intent.clean_generation + request = await ArtifactProcessingIntentRepository().request( + connection, assignment.scope_id, assignment.binding_name + ) + assignment = replace(assignment, claimed_request_generation=request.requested_generation) + else: + assert intent.dirty_generation == intent.clean_generation + assert tuple(item for window in windows[2:] for item in window) == windows[0] + + asyncio.run(scenario()) + + @pytest.mark.parametrize("family", ["memory", "experience", "profile"]) def test_owner_failure_rolls_back_domain_cursor_and_ack_then_retry_owns_result(tmp_path, family): async def scenario(): diff --git a/tests/builtin/runtime/test_memory_window_recovery.py b/tests/builtin/runtime/test_memory_window_recovery.py new file mode 100644 index 0000000000..22bed295d0 --- /dev/null +++ b/tests/builtin/runtime/test_memory_window_recovery.py @@ -0,0 +1,179 @@ +# 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 timeouts shrink durable windows without consuming failed evidence.""" + +import asyncio +from uuid import uuid4 + +import pytest + +from powercontext.builtin.artifacts.memory import MemoryEntryInput +from powercontext.builtin.inference import InferenceTimeoutError, InferenceUnavailableError +from powercontext.builtin.persistence.errors import GenerationConflictError +from powercontext.builtin.persistence.sqlite import SQLiteConfig, SQLiteProfile +from powercontext.builtin.persistence.tables import BUILTIN_TABLES +from powercontext.builtin.runtime.relational import RelationalContexts +from powercontext.builtin.scope import ScopeDraft + + +class Pipeline: + def __init__(self, maximum=1, error=None): + self.maximum = maximum + self.error = error + self.windows = [] + + async def extract(self, request): + self.windows.append(tuple(source.name for source in request.sources)) + if self.error is not None: + raise self.error + if len(request.sources) > self.maximum: + raise InferenceTimeoutError("generate", 60) + return tuple( + MemoryEntryInput(kind="fact", text=source.content, sources=(source,)) for source in request.sources + ) + + +async def create_scope(contexts, count): + scope = ( + await contexts.scopes.create(ScopeDraft(title="Recovery", summary="Recovery", idempotency_key=str(uuid4()))) + ).scope_id + for index in range(count): + await contexts.records.create_source(scope, "content", f"Fact {index}") + return scope + + +def test_timeout_reduction_survives_reopen_and_resets_after_backlog(tmp_path): + async def scenario(): + config = SQLiteConfig(url=f"sqlite+aiosqlite:///{tmp_path / 'recovery.db'}") + pipeline = Pipeline() + async with SQLiteProfile.open(config, tables=BUILTIN_TABLES) as profile: + contexts = RelationalContexts(database=profile.database, candidate_pipeline=pipeline) + scope = await create_scope(contexts, 4) + context = await contexts.get(scope) + with pytest.raises(InferenceTimeoutError): + await context.triggers.flush(limit=100) + assert (await context.triggers.cursor()).sequence == 0 + + # Reopen the database and rebuild the processor, as a new Worker does. + async with SQLiteProfile.open(config, tables=BUILTIN_TABLES) as profile: + contexts = RelationalContexts(database=profile.database, candidate_pipeline=pipeline) + context = await contexts.get(scope) + with pytest.raises(InferenceTimeoutError): + await context.triggers.flush(limit=100) + assert (await context.triggers.cursor()).sequence == 0 + assert len(pipeline.windows[-1]) < len(pipeline.windows[0]) + for position in range(1, 5): + result = await context.triggers.flush(limit=100) + assert result.current_cursor == position + assert result.memory_ref is not None + assert tuple(item for window in pipeline.windows[2:] for item in window) == pipeline.windows[0] + + pipeline.maximum = 4 + for index in range(4): + await contexts.records.create_source(scope, "content", f"New fact {index}") + result = await context.triggers.flush(limit=100) + assert result.source_count == 4 + assert result.current_cursor == 8 + + asyncio.run(scenario()) + + +def test_single_source_timeout_preserves_cursor_and_can_recover(): + async def scenario(): + pipeline = Pipeline(error=InferenceTimeoutError("generate", 60)) + async with SQLiteProfile.open(SQLiteConfig(), tables=BUILTIN_TABLES) as profile: + contexts = RelationalContexts(database=profile.database, candidate_pipeline=pipeline) + scope = await create_scope(contexts, 1) + context = await contexts.get(scope) + for _ in range(2): + with pytest.raises(InferenceTimeoutError): + await context.triggers.flush(limit=100) + assert (await context.triggers.cursor()).sequence == 0 + pipeline.error = None + assert (await context.triggers.flush(limit=100)).current_cursor == 1 + assert pipeline.windows[0] == pipeline.windows[-1] + + asyncio.run(scenario()) + + +@pytest.mark.parametrize("error", [InferenceTimeoutError("embed", 60), InferenceUnavailableError("generate")]) +def test_non_extraction_failures_do_not_reduce_source_window(error): + async def scenario(): + pipeline = Pipeline(maximum=4, error=error) + async with SQLiteProfile.open(SQLiteConfig(), tables=BUILTIN_TABLES) as profile: + contexts = RelationalContexts(database=profile.database, candidate_pipeline=pipeline) + scope = await create_scope(contexts, 4) + context = await contexts.get(scope) + with pytest.raises(type(error)): + await context.triggers.flush(limit=100) + pipeline.error = None + result = await context.triggers.flush(limit=100) + assert result.current_cursor == 4 + assert pipeline.windows[0] == pipeline.windows[-1] + + asyncio.run(scenario()) + + +def test_stale_timeout_cannot_rewind_concurrently_committed_cursor(): + async def scenario(): + entered, release = asyncio.Event(), asyncio.Event() + + class DelayedTimeout: + async def extract(self, request): + entered.set() + await release.wait() + raise InferenceTimeoutError("generate", 60) + + async with SQLiteProfile.open(SQLiteConfig(), tables=BUILTIN_TABLES) as profile: + first = RelationalContexts(database=profile.database, candidate_pipeline=DelayedTimeout()) + second = RelationalContexts(database=profile.database, candidate_pipeline=Pipeline(maximum=4)) + scope = await create_scope(first, 4) + first_context, second_context = await first.get(scope), await second.get(scope) + task = asyncio.create_task(first_context.triggers.flush(limit=4)) + await asyncio.wait_for(entered.wait(), timeout=5) + try: + assert (await second_context.triggers.flush(limit=2)).current_cursor == 2 + finally: + release.set() + with pytest.raises(GenerationConflictError): + await task + assert (await second_context.triggers.cursor()).sequence == 2 + assert (await second_context.triggers.flush(limit=4)).current_cursor == 4 + + asyncio.run(scenario()) + + +def test_failed_commit_retains_reduction_until_consumption_succeeds(): + async def fail_commit(*_args): + raise OSError("ownership commit failed") # noqa: TRY003 - injected transaction failure + + async def scenario(): + pipeline = Pipeline() + async with SQLiteProfile.open(SQLiteConfig(), tables=BUILTIN_TABLES) as profile: + contexts = RelationalContexts(database=profile.database, candidate_pipeline=pipeline) + scope = await create_scope(contexts, 2) + context = await contexts.get(scope) + with pytest.raises(InferenceTimeoutError): + await context.triggers.flush(limit=100) + assert (await context.triggers.flush(limit=1)).current_cursor == 1 + with pytest.raises(OSError, match="ownership commit failed"): + await contexts.process_memory(scope, 100, on_commit=fail_commit) + assert (await context.triggers.cursor()).sequence == 1 + await contexts.records.create_source(scope, "content", "New input while retrying") + # The failed final commit must roll back clearing the reduction too. + assert (await context.triggers.flush(limit=100)).current_cursor == 2 + assert (await context.triggers.flush(limit=100)).current_cursor == 3 + + asyncio.run(scenario()) diff --git a/tests/e2e/real_experience_skill/harness.py b/tests/e2e/real_experience_skill/harness.py index c03de72cc8..f26a8e9050 100644 --- a/tests/e2e/real_experience_skill/harness.py +++ b/tests/e2e/real_experience_skill/harness.py @@ -116,6 +116,7 @@ "pc_artifacts", "pc_skill_packages", "pc_source_cursors", + "pc_memory_source_windows", "pc_external_skill_registrations", "pc_sources", "pc_source_journal_heads",