mirror of
https://github.com/volcengine/OpenViking.git
synced 2026-10-01 01:38:07 +08:00
* feat(ingest): replay local agent-harness logs into OV sessions / 本地 agent harness 日志重放入库
Add openviking/ingest/: parse Claude Code / Codex / OpenCode / Hermes / OpenClaw conversation logs into normalized messages and replay them through OpenViking's existing session pipeline (create_session -> batch_add_messages -> commit -> async memory extraction), instead of a bespoke ETL.
Supports one-shot backfill ("存量") and cursor-driven incremental polling ("新增", WatchScheduler-style, no fs-event dependency), per-harness enable/mode/paths config, and meaningful peer_id on every turn (assistant = {harness}/{model}; user = git identity for single-user harnesses, original username for group-chat harnesses). Cursor IDE is a registered but deferred stub.
Read-position cursors persist under ~/.openviking/ingest/state.db for crash-safe, idempotent resume. New openviking-ingest CLI (backfill/watch/run/status/list-sources) and an "ingest" section on OpenVikingConfig. Verified end-to-end against a local server: 3-message fixture -> session commit -> 10 memories extracted -> idempotent re-run.
Inspired by / supersedes volcengine/OpenViking#2674.
Co-authored-by: baobaodae <2014596548@qq.com>
* docs(ingest): bilingual guide + ov.conf.example for openviking-ingest / 本地日志入库双语文档与配置示例
Add docs/{zh,en}/agent-integrations/09-log-ingestion.md (auto-registered in the VitePress sidebar) and an `ingest` section in examples/ov.conf.example (off by default).
* fix(ingest): address review — gating, crash-safe batch replay, commit recovery, single-instance lock / 修复评审问题
Fixes the merge-blockers from the adversarial review:
- master switch ingest.enabled now actually gates enabled_harnesses();
- idempotent per-batch append with a durable pending-intent reconciled against the server message count on restart (no duplicate imports after a mid-append crash);
- bounded reads (<=100 msgs/call) so huge sessions don't materialize at once;
- needs_commit flag + commit_if_needed so appended-but-uncommitted sessions still get extracted (commit even when no new source rows);
- poller keeps dirty sessions until a commit actually succeeds;
- OpenCode advances its SQLite cursor only past complete rows (late part text no longer skipped);
- single-instance file lock guards concurrent ingest processes;
- positive-value config validation (no poll busy-loop); malformed ov.conf surfaces instead of silently defaulting.
Adds 6 tests (config gating/validation, crash reconcile both ways, commit recovery).
* refactor(ingest): expose as 'openviking-server ingest' subcommand; English-only code/docs
- Route the ingest CLI through 'openviking-server ingest ...' (same dispatch as 'init'/'doctor') and drop the separate 'openviking-ingest' console_script.
- Remove mixed-in Chinese terms (存量/新增) from source docstrings, CLI help, and the English doc; the Chinese doc keeps them.
* style(ingest): ruff format + import sort (isort I)
Run ruff 0.15.16 (from the uv cache) with the repo config: fixes 5 I001 import-order errors in tests and reformats 9 files. 'ruff check' and 'ruff format --check' now pass on all added/edited files.
---------
Co-authored-by: baobaodae <2014596548@qq.com>
149 lines
6.0 KiB
Python
149 lines
6.0 KiB
Python
# Copyright (c) 2026 Beijing Volcano Engine Technology Co., Ltd.
|
|
# SPDX-License-Identifier: AGPL-3.0
|
|
"""Incremental ingest (watch mode): a WatchScheduler-style asyncio poll loop.
|
|
|
|
Mirrors ``openviking/resource/watch_scheduler.py`` (interval polling, graceful start/stop)
|
|
rather than depending on filesystem events: a durable read-cursor makes polling correct
|
|
and self-healing (a missed tick / sleep / restart just reads cursor->EOF next time).
|
|
|
|
Each tick, per enabled harness: rescan sessions -> incremental ``read_messages`` from the
|
|
stored cursor -> append + advance cursor -> commit on idle or token threshold.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import asyncio
|
|
import time
|
|
from typing import Dict, List, Tuple
|
|
|
|
from openviking.ingest.models import Cursor
|
|
from openviking.ingest.replay import SessionReplayer
|
|
from openviking.ingest.sources.base import DEFAULT_READ_LIMIT, LogSource, NotSupportedError
|
|
from openviking_cli.utils import get_logger
|
|
from openviking_cli.utils.config.ingest_config import IngestHarnessConfig
|
|
|
|
logger = get_logger(__name__)
|
|
|
|
|
|
class _Dirty:
|
|
__slots__ = ("name", "cfg", "source", "ref", "last_activity")
|
|
|
|
def __init__(self, name, cfg, source, ref, last_activity):
|
|
self.name = name
|
|
self.cfg = cfg
|
|
self.source = source
|
|
self.ref = ref
|
|
self.last_activity = last_activity
|
|
|
|
|
|
class IngestPoller:
|
|
def __init__(
|
|
self,
|
|
sources: List[Tuple[str, IngestHarnessConfig, LogSource]],
|
|
replayer: SessionReplayer,
|
|
):
|
|
self._sources = sources
|
|
self.replayer = replayer
|
|
self.store = replayer.store
|
|
self._running = False
|
|
self._stop = asyncio.Event()
|
|
self._dirty: Dict[Tuple[str, str], _Dirty] = {}
|
|
self.poll_interval = min((cfg.poll_interval_seconds for _, cfg, _ in sources), default=5.0)
|
|
|
|
def stop(self) -> None:
|
|
self._running = False
|
|
self._stop.set()
|
|
|
|
async def run(self) -> None:
|
|
if not self._sources:
|
|
logger.info("[ingest] no harnesses enabled for watch; nothing to do")
|
|
return
|
|
self._running = True
|
|
logger.info(
|
|
"[ingest] watch loop started (poll=%.1fs, harnesses=%s)",
|
|
self.poll_interval,
|
|
[n for n, _, _ in self._sources],
|
|
)
|
|
try:
|
|
while self._running:
|
|
await self._tick()
|
|
await self._commit_idle()
|
|
try:
|
|
await asyncio.wait_for(self._stop.wait(), timeout=self.poll_interval)
|
|
except asyncio.TimeoutError:
|
|
pass
|
|
finally:
|
|
await self._flush_all()
|
|
logger.info("[ingest] watch loop stopped")
|
|
|
|
async def _tick(self) -> None:
|
|
for name, cfg, source in self._sources:
|
|
try:
|
|
refs = list(source.discover_sessions())
|
|
except NotSupportedError as exc:
|
|
logger.warning("[ingest:%s] %s", name, exc)
|
|
continue
|
|
except Exception: # noqa: BLE001
|
|
logger.exception("[ingest:%s] discovery failed", name)
|
|
continue
|
|
for ref in refs:
|
|
await self._poll_session(name, cfg, source, ref)
|
|
|
|
async def _poll_session(self, name, cfg, source, ref) -> None:
|
|
await self.replayer.reconcile(name, ref)
|
|
appended = False
|
|
while True:
|
|
cursor = self.store.get_cursor(name, ref.native_session_id, source.cursor_kind)
|
|
try:
|
|
messages, new_cursor = source.read_messages(ref, cursor, DEFAULT_READ_LIMIT)
|
|
except Exception: # noqa: BLE001
|
|
logger.exception("[ingest:%s] read failed for %s", name, ref.native_session_id)
|
|
return
|
|
if not messages:
|
|
from_value = (cursor or Cursor.zero(source.cursor_kind)).value
|
|
if new_cursor.value != from_value:
|
|
sid = self.replayer.ov_session_id(name, ref.native_session_id)
|
|
self.store.advance_cursor(
|
|
name, ref.native_session_id, sid, new_cursor, locator=ref.locator
|
|
)
|
|
break
|
|
from_cursor = cursor or Cursor.zero(source.cursor_kind)
|
|
added = await self.replayer.append_batch(name, ref, messages, from_cursor, new_cursor)
|
|
if added > 0:
|
|
appended = True
|
|
|
|
rec = self.store.get(name, ref.native_session_id)
|
|
if appended or (rec and rec.needs_commit):
|
|
self._dirty[(name, ref.native_session_id)] = _Dirty(
|
|
name, cfg, source, ref, time.monotonic()
|
|
)
|
|
committed = await self.replayer.maybe_commit_on_threshold(
|
|
name, ref, cfg.commit.commit_token_threshold, cfg.commit.keep_recent_count
|
|
)
|
|
if committed:
|
|
self._dirty.pop((name, ref.native_session_id), None)
|
|
|
|
async def _commit_idle(self) -> None:
|
|
now = time.monotonic()
|
|
for key, d in list(self._dirty.items()):
|
|
if now - d.last_activity >= d.cfg.commit.commit_idle_seconds:
|
|
try:
|
|
await self.replayer.commit_if_needed(
|
|
d.name, d.ref, d.cfg.commit.keep_recent_count
|
|
)
|
|
self._dirty.pop(key, None) # only drop on success
|
|
except Exception: # noqa: BLE001 - keep dirty and retry after a backoff
|
|
logger.exception("[ingest:%s] idle commit failed; will retry", d.name)
|
|
d.last_activity = time.monotonic()
|
|
|
|
async def _flush_all(self) -> None:
|
|
"""Final commit for every dirty session on shutdown (needs_commit persists if it fails)."""
|
|
for key, d in list(self._dirty.items()):
|
|
try:
|
|
await self.replayer.commit_if_needed(d.name, d.ref, d.cfg.commit.keep_recent_count)
|
|
except Exception: # noqa: BLE001
|
|
logger.exception(
|
|
"[ingest:%s] shutdown commit failed (needs_commit persisted)", d.name
|
|
)
|
|
self._dirty.pop(key, None)
|