Files
OpenViking/openviking/ingest/poller.py
T
t0sakiandbaobaodae 72d04cd488 feat(ingest): replay local agent-harness logs into OpenViking sessions (#2892)
* 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>
2026-06-30 12:14:31 +08:00

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)