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 .github/workflows/plugin-api-v3.yml
Original file line number Diff line number Diff line change
Expand Up @@ -36,7 +36,7 @@ jobs:
- uses: actions/checkout@v4
with:
repository: kachofugetsu09/akashic-agent
ref: 9da3a988a2bf62b0f550bd4f6bb98c4eeb1f56f5
ref: 0607e546de923e9377b04fb841d01384b1a666c1
path: .akashic-core
- uses: actions/setup-python@v5
with:
Expand Down
7 changes: 3 additions & 4 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -36,10 +36,9 @@ history reader 始终使用 SQLite `mode=ro`。数据库不存在表示合法空
若尚无三个 preview 列,history 会把这三个稳定字段投影为 `null` 后参与 canonical hash;
旧表经既有 `open_db` 逐列 `ALTER` 后形成的 SQLite schema 也属于同一已知 lineage。

非引用评分使用 Core 正式运行时的共享 HTTP resources。嵌入配置从
`AKASHIC_CONFIG` 指向的 Core 配置加载,不从插件 checkout 的当前目录猜测配置。
embedding 继续使用既有 Core provider 数据流;API key 只作为运行时认证,不进入
inbox、projection 或 typed event,完整正文也不进入这些持久/发布边界。
非引用评分只注入公共 `EMBEDDINGS` service,不读取 Core 配置、provider 或凭据。
后台评分每次都在当前 generation 的 `runtime_scope` 内执行 `bind → embed`,因此热更新
前后的任务不会混用模型快照。向量和完整正文不进入 inbox、projection 或 history。

旧 v2 `scripts/backfill_proactive_feedback.py` 已移除:它直接操作
`workspace/proactive_feedback/proactive_feedback.db`,而 `--clear` 会删除旧 DB、WAL
Expand Down
44 changes: 17 additions & 27 deletions plugin.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,14 +3,14 @@
import asyncio
import json
import logging
import os
from collections.abc import Callable, Iterable, Iterator
from pathlib import Path
from typing import Any, cast

from agent.config_models import Config as CoreConfig
from agent.plugin_composition import (
Context,
EMBEDDINGS,
Embeddings,
MobileUiDefinition,
MobileUiNavigation,
MobileUiRpcInvalidRequest,
Expand All @@ -20,8 +20,6 @@
)
from agent.turn_events.after_turn import AFTER_TURN_COMMITTED
from bus.events_lifecycle import TurnCommitted
from core.net.http import get_default_http_requester
from memory2.embedder import Embedder

from .dashboard import ProactiveFeedbackDashboardReader
from .db import (
Expand All @@ -37,6 +35,7 @@
pending_feedback_inputs,
)
from .scorer import (
EmbedBatch,
MessageRow,
message_rows_from_snapshot,
parse_quote_parts,
Expand All @@ -61,7 +60,7 @@
version = "3.0.0"
desc = "记录主动消息被继续的反馈,并提供桌面与移动只读投影。"
author = "Akashic"
inject = (SESSION_READ, UI_SLOTS)
inject = (SESSION_READ, UI_SLOTS, EMBEDDINGS)
skill_roots: tuple[str, ...] = ()
drift_skill_roots: tuple[str, ...] = ()
workspace_roots: tuple[str, ...] = ()
Expand All @@ -75,10 +74,11 @@ async def apply(ctx: Context, config: object) -> None:
_ = config
session_read = ctx.require(SESSION_READ)
ui_slots = ctx.require(UI_SLOTS)
embeddings = ctx.require(EMBEDDINGS)
db_path = ctx.data_root / _FEEDBACK_DB_NAME
runtime = ProactiveFeedbackRuntime(
session_read=session_read,
workspace=ctx.runtime.workspace,
embed_batch=_bind_embeddings(embeddings, ctx),
db_path=db_path,
)
_ = await ctx.provide(
Expand Down Expand Up @@ -111,18 +111,17 @@ def __init__(
self,
*,
session_read: SessionReadService,
workspace: Path,
embed_batch: EmbedBatch,
db_path: Path,
session_keys: Callable[[], Iterable[str]] | None = None,
) -> None:
self._session_read = session_read
self._workspace = workspace
self._embed_batch = embed_batch
self._db_path = db_path
self._session_keys = session_keys or (
lambda: _formal_session_keys_from_read_service(session_read)
)
self._queue: asyncio.Queue[int] = asyncio.Queue(maxsize=_QUEUE_MAX)
self._embedder: Embedder | None = None
self._discovery_done = False

def observe_committed(self, event: TurnCommitted) -> None:
Expand Down Expand Up @@ -270,7 +269,7 @@ async def _process(
# 3. Persist one deduplicated projection, including bounded display text.
try:
scored = await score_followup(
embed_batch=self._get_embedder().embed_batch if allow_pua else _no_embed,
embed_batch=self._embed_batch if allow_pua else _no_embed,
user=user,
assistant=assistant,
candidates=candidates,
Expand Down Expand Up @@ -449,11 +448,6 @@ async def _process_input_record(self, record: FeedbackInputRecord) -> None:
input_row_id=record.row_id,
)

def _get_embedder(self) -> Embedder:
if self._embedder is None:
self._embedder = _build_embedder(self._workspace)
return self._embedder

def _discover_committed_inputs(self) -> None:
"""Discover bounded eligible Turns committed before the callback fanout."""

Expand Down Expand Up @@ -687,18 +681,14 @@ def _bounded_session_keys(
return (*unique[: limit - 1], unique[-1])


def _build_embedder(workspace: Path) -> Embedder:
config_path = os.environ.get("AKASHIC_CONFIG", "").strip()
if not config_path:
raise RuntimeError("proactive_feedback 需要 Core 的 AKASHIC_CONFIG")
embedding = CoreConfig.load(path=config_path, workspace=workspace).memory.embedding
return Embedder(
base_url=embedding.base_url,
api_key=embedding.api_key,
model=embedding.model,
output_dimensionality=embedding.output_dimensionality,
requester=get_default_http_requester("external_default"),
)
def _bind_embeddings(embeddings: Embeddings, ctx: Context) -> EmbedBatch:
async def embed_batch(texts: list[str]) -> list[list[float]]:
async with ctx.runtime_scope():
async with embeddings.bind() as bound:
result = await bound.embed(texts)
return [list(vector) for vector in result.vectors]

return embed_batch


def _decode_outbox_payload(payload_json: str) -> dict[str, object]:
Expand Down
Loading
Loading