Files
OpenViking/openviking/ingest/models.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

96 lines
3.3 KiB
Python

# Copyright (c) 2026 Beijing Volcano Engine Technology Co., Ltd.
# SPDX-License-Identifier: AGPL-3.0
"""Core data structures shared across the ingest subsystem.
Note on "cursor": throughout this package, ``Cursor`` / ``cursor_store`` / ``cursor_kind``
refer to the read-position POINTER (how far we have ingested a given log), NOT the Cursor
IDE harness (whose adapter is ``sources/cursor.py``).
"""
from __future__ import annotations
import json
from dataclasses import dataclass, field
from datetime import datetime, timezone
from typing import Any, Dict, List, Optional
def iso_from_epoch_ms(value: Any) -> Optional[str]:
"""Best-effort ISO-8601 from an epoch-millisecond integer (used by SQLite/JSONL adapters)."""
try:
ms = int(value)
except (TypeError, ValueError):
return None
if ms <= 0:
return None
# Heuristic: 10-digit values are seconds, 13-digit are milliseconds.
seconds = ms / 1000.0 if ms > 10_000_000_000 else float(ms)
try:
return datetime.fromtimestamp(seconds, tz=timezone.utc).isoformat()
except (OverflowError, OSError, ValueError):
return None
# Cursor kinds.
BYTE_OFFSET = "byte_offset" # append-only JSONL: byte offset of last consumed line
ROWID_TIME = "rowid_time" # SQLite: (time_created, id) of last consumed row
@dataclass
class NormalizedMessage:
"""A single conversation turn, harness-agnostic, ready for OV replay.
``peer_id`` is resolved by the adapter (see ``peer.py``): assistant turns ->
``{harness}/{model}``; user turns -> a human identifier (git identity for
single-user dev harnesses, original username for group-chat harnesses).
"""
role: str # "user" | "assistant"
text: str = ""
parts: List[Dict[str, Any]] = field(default_factory=list) # extra tool/context parts
created_at: Optional[str] = None # ISO-8601
peer_id: Optional[str] = None
meta: Dict[str, Any] = field(default_factory=dict) # model, provider, cwd, …
@dataclass
class SessionRef:
"""A discoverable conversation in a harness's storage."""
harness: str # registry name, e.g. "claude_code"
native_session_id: str # the harness's own session/conversation id
locator: str # file path (JSONL) or db session id (SQLite) used to read
title: Optional[str] = None
started_at: Optional[str] = None
meta: Dict[str, Any] = field(default_factory=dict) # session-level model, cwd, platform
@dataclass
class Cursor:
"""A durable read-position pointer for one (harness, session)."""
kind: str # BYTE_OFFSET | ROWID_TIME
value: Dict[str, Any] # {"offset": int} | {"time": int|float, "id": str}
@classmethod
def zero(cls, kind: str) -> "Cursor":
if kind == BYTE_OFFSET:
return cls(kind=kind, value={"offset": 0})
return cls(kind=kind, value={"time": 0, "id": ""})
@property
def offset(self) -> int:
return int(self.value.get("offset", 0))
def to_json(self) -> str:
return json.dumps(self.value, separators=(",", ":"))
@classmethod
def from_json(cls, kind: str, raw: Optional[str]) -> "Cursor":
if not raw:
return cls.zero(kind)
try:
return cls(kind=kind, value=json.loads(raw))
except (ValueError, TypeError):
return cls.zero(kind)