mirror of
https://github.com/volcengine/OpenViking.git
synced 2026-09-29 04:02:57 +08:00
* fix(pi): honor conversation capture filters in faithful mode Apply the shared turn filter before capture decisions so configured rules cannot be bypassed by takeover or attached tool parts, while preserving tool values and behavior without applicable rules. (cherry picked from commite8fb6c9476) Signed-off-by: Hao Zhe <haozhe4547@gmail.com> * fix(opencode): wire captureFilters through config and the capture path captureFilters (sed-style redaction rules, e.g. s/sk-.../[REDACTED-SK]/g) was defined in the shared capture-utils library but never reached the OpenCode capture path: lib/config.mjs had no mapping for the key and memory-session's buildCapturePayload applied no filtering, so secrets typed into OpenCode sessions were captured verbatim. The Codex plugin applies the identical library correctly (issue #4984). - map captureFilters in config (file key, OPENVIKING_CAPTURE_FILTERS env override as JSON, non-string entries dropped) - apply filterCaptureParts() to extracted parts in buildCapturePayload, mirroring the Codex reference flow (role-level drops skip the turn; an empty filtered result gates capture) - cover config mapping, env override and validation in tests (cherry picked from commitb5bd21930d) Signed-off-by: Hao Zhe <haozhe4547@gmail.com> * refactor(capture): share sanitized filter path across adapters * test(capture): preserve mixed tool messages and cover OpenCode v2 * fix(capture): cap mixed text and preserve dated logs * fix(capture): use shared shaping in Claude Code hooks * fix(capture): sanitize fallback text once --------- Signed-off-by: Hao Zhe <haozhe4547@gmail.com> Co-authored-by: hemingzhe <hehesmilett@163.com>
537 lines
18 KiB
JavaScript
537 lines
18 KiB
JavaScript
import fs from "fs"
|
|
import path from "path"
|
|
import {
|
|
extractTextFromPayload,
|
|
isCaptureEnabled,
|
|
shapeCapturePayload,
|
|
} from "./shared/capture-utils.mjs"
|
|
import {
|
|
deriveHarnessSessionId,
|
|
} from "./shared/session-model.mjs"
|
|
import {
|
|
enqueue,
|
|
replayPending,
|
|
} from "./shared/pending-queue.mjs"
|
|
import {
|
|
sendSessionMessages,
|
|
} from "./shared/batch-send.mjs"
|
|
import {
|
|
isRetryableFailure,
|
|
} from "./shared/retryable.mjs"
|
|
import {
|
|
log,
|
|
effectivePeerId,
|
|
fetchJSON,
|
|
safeStringify,
|
|
} from "./utils.mjs"
|
|
|
|
export function createMemorySessionManager({ config, pluginRoot }) {
|
|
const sessions = new Map()
|
|
const statePath = path.join(pluginRoot, "openviking-session-state.json")
|
|
const oldSessionMapPath = path.join(pluginRoot, "openviking-session-map.json")
|
|
let saveTimer = null
|
|
let initBackground = Promise.resolve()
|
|
// Serialize saves: concurrent saveState() calls (a debounced save racing
|
|
// with flushAll / flushSession / session deletion) all share the same
|
|
// `${statePath}.tmp` temp file, so one rename can fail with ENOENT after
|
|
// another already moved it. Chaining through a promise queue keeps at most
|
|
// one write+rename in flight. Each queued save re-serializes the in-memory
|
|
// `sessions` map at execution time, so the last save always persists the
|
|
// latest state.
|
|
let saveChain = Promise.resolve()
|
|
|
|
function enqueueSave() {
|
|
const run = saveChain.then(() => saveState())
|
|
saveChain = run.catch(() => {})
|
|
return run
|
|
}
|
|
|
|
async function init({ deferNetwork = false } = {}) {
|
|
await loadState()
|
|
initBackground = Promise.resolve().then(async () => {
|
|
if (isCaptureEnabled(config)) await migrateLegacySessionMap()
|
|
const health = await fetchJSON(config, "/health", {}, { timeoutMs: 5000 })
|
|
if (!health.ok) return
|
|
await replayPending(
|
|
(endpoint, init = {}, options = {}) => fetchJSON(config, endpoint, init, options),
|
|
(stage, data) => log("DEBUG", "pending", stage, data),
|
|
)
|
|
}).catch((error) => {
|
|
log("WARN", "pending", "Pending replay failed during initialization", { error: error?.message })
|
|
})
|
|
if (!deferNetwork) await initBackground
|
|
}
|
|
|
|
async function waitForBackground() {
|
|
await initBackground
|
|
}
|
|
|
|
async function loadState() {
|
|
try {
|
|
if (!fs.existsSync(statePath)) {
|
|
log("INFO", "persistence", "No session state file found, starting fresh")
|
|
return
|
|
}
|
|
const data = JSON.parse(await fs.promises.readFile(statePath, "utf8"))
|
|
if (data.version !== 2) {
|
|
log("ERROR", "persistence", "Unsupported session map version", { version: data.version })
|
|
return
|
|
}
|
|
for (const [opencodeSessionId, persisted] of Object.entries(data.sessions ?? {})) {
|
|
sessions.set(opencodeSessionId, deserializeSessionState(persisted))
|
|
}
|
|
log("INFO", "persistence", "Session state loaded", { count: sessions.size })
|
|
} catch (error) {
|
|
log("ERROR", "persistence", "Failed to load session state", { error: error?.message })
|
|
if (fs.existsSync(statePath)) {
|
|
await fs.promises.rename(statePath, `${statePath}.corrupted.${Date.now()}`)
|
|
}
|
|
}
|
|
}
|
|
|
|
async function saveState() {
|
|
try {
|
|
const persisted = {}
|
|
for (const [opencodeSessionId, state] of sessions.entries()) {
|
|
persisted[opencodeSessionId] = serializeSessionState(state)
|
|
}
|
|
const tempPath = `${statePath}.tmp`
|
|
await fs.promises.writeFile(tempPath, JSON.stringify({ version: 2, sessions: persisted, lastSaved: Date.now() }, null, 2), "utf8")
|
|
await fs.promises.rename(tempPath, statePath)
|
|
log("DEBUG", "persistence", "Session state saved", { count: sessions.size })
|
|
} catch (error) {
|
|
log("ERROR", "persistence", "Failed to save session state", { error: error?.message })
|
|
}
|
|
}
|
|
|
|
function debouncedSaveState() {
|
|
if (saveTimer) clearTimeout(saveTimer)
|
|
saveTimer = setTimeout(() => {
|
|
enqueueSave().catch((error) => {
|
|
log("ERROR", "persistence", "Debounced save failed", { error: error?.message })
|
|
})
|
|
}, 300)
|
|
}
|
|
|
|
function serializeSessionState(state) {
|
|
return {
|
|
ovSessionId: state.ovSessionId,
|
|
createdAt: state.createdAt,
|
|
lastActivityAt: state.lastActivityAt,
|
|
lastCommitTime: state.lastCommitTime,
|
|
compactedAt: state.compactedAt,
|
|
messages: Array.from(state.messages.entries()).map(([messageId, message]) => ([
|
|
messageId,
|
|
{
|
|
role: message.role,
|
|
captured: message.captured,
|
|
// Captured messages are never read by flushPendingMessages again; retain
|
|
// only their metadata so completed payloads cannot grow the state file.
|
|
parts: message.captured ? [] : Array.from(message.parts.entries()),
|
|
},
|
|
])),
|
|
}
|
|
}
|
|
|
|
function deserializeSessionState(persisted) {
|
|
return {
|
|
ovSessionId: persisted.ovSessionId,
|
|
createdAt: persisted.createdAt,
|
|
lastActivityAt: persisted.lastActivityAt,
|
|
lastCommitTime: persisted.lastCommitTime,
|
|
compactedAt: persisted.compactedAt,
|
|
messages: new Map((persisted.messages ?? []).map(([messageId, message]) => ([
|
|
messageId,
|
|
{
|
|
role: message.role,
|
|
captured: Boolean(message.captured),
|
|
parts: new Map(message.parts ?? []),
|
|
},
|
|
]))),
|
|
}
|
|
}
|
|
|
|
function getMappedSessionId(opencodeSessionId) {
|
|
return getOrCreateSession(opencodeSessionId).ovSessionId
|
|
}
|
|
|
|
async function handleEvent(event) {
|
|
if (!event?.type || event.type === "session.diff") return
|
|
|
|
if (event.type === "session.created") {
|
|
await handleSessionCreated(event)
|
|
} else if (event.type === "session.deleted") {
|
|
await handleSessionDeleted(event)
|
|
} else if (event.type === "session.error") {
|
|
await handleSessionError(event)
|
|
} else if (event.type === "session.compacted") {
|
|
await handleSessionCompacted(event)
|
|
} else if (event.type === "session.idle") {
|
|
await handleSessionIdle(event)
|
|
} else if (event.type === "message.updated" && isCaptureEnabled(config)) {
|
|
await handleMessageUpdated(event)
|
|
} else if (event.type === "message.part.updated" && isCaptureEnabled(config)) {
|
|
await handleMessagePartUpdated(event)
|
|
}
|
|
}
|
|
|
|
async function handleSessionCreated(event) {
|
|
const sessionId = resolveEventSessionId(event)
|
|
if (!sessionId) {
|
|
log("ERROR", "event", "session.created event missing sessionId", { event: safeStringify(event) })
|
|
return
|
|
}
|
|
const state = getOrCreateSession(sessionId, event)
|
|
debouncedSaveState()
|
|
const health = await fetchJSON(config, "/health", {}, { timeoutMs: 5000 })
|
|
if (health.ok) {
|
|
await replayPending(
|
|
(endpoint, init = {}, options = {}) => fetchJSON(config, endpoint, init, options),
|
|
(stage, data) => log("DEBUG", "pending", stage, data),
|
|
)
|
|
}
|
|
log("INFO", "event", "OpenViking session derived", {
|
|
opencode_session: sessionId,
|
|
openviking_session: state.ovSessionId,
|
|
})
|
|
}
|
|
|
|
async function handleSessionDeleted(event) {
|
|
const sessionId = resolveEventSessionId(event)
|
|
if (!sessionId) return
|
|
await flushSession(sessionId, { commit: true, reason: event.type })
|
|
sessions.delete(sessionId)
|
|
await enqueueSave()
|
|
}
|
|
|
|
async function handleSessionError(event) {
|
|
const sessionId = resolveEventSessionId(event)
|
|
if (!sessionId) return
|
|
log("ERROR", "event", "OpenCode session error", { session_id: sessionId, error: safeStringify(event.error) })
|
|
await handleSessionDeleted(event)
|
|
}
|
|
|
|
async function handleSessionCompacted(event) {
|
|
await commitSessionBoundary(event, "session.compacted")
|
|
}
|
|
|
|
async function handleSessionIdle(event) {
|
|
const sessionId = resolveEventSessionId(event)
|
|
if (!sessionId) return
|
|
await flushSession(sessionId, { commit: false, reason: "session.idle" })
|
|
}
|
|
|
|
async function commitSessionBoundary(event, reason) {
|
|
const sessionId = resolveEventSessionId(event)
|
|
if (!sessionId) return
|
|
const state = getOrCreateSession(sessionId, event)
|
|
state.compactedAt = Date.now()
|
|
await flushSession(sessionId, { commit: true, reason })
|
|
}
|
|
|
|
async function handleMessageUpdated(event) {
|
|
const message = event.properties?.info
|
|
if (!message) return
|
|
|
|
const sessionId = message.sessionID
|
|
const messageId = message.id
|
|
const role = message.role
|
|
const finish = message.finish
|
|
if (!sessionId || !messageId) return
|
|
|
|
const state = getOrCreateSession(sessionId, event)
|
|
const captured = state.messages.get(messageId)
|
|
const next = captured ?? createMessageState()
|
|
if (role === "user") {
|
|
next.role = role
|
|
} else if (role === "assistant") {
|
|
next.role = role
|
|
}
|
|
state.messages.set(messageId, next)
|
|
state.lastActivityAt = Date.now()
|
|
debouncedSaveState()
|
|
}
|
|
|
|
async function handleMessagePartUpdated(event) {
|
|
const part = event.properties?.part
|
|
if (!part) return
|
|
|
|
const sessionId = part.sessionID
|
|
const messageId = part.messageID
|
|
if (!sessionId || !messageId) return
|
|
|
|
const state = getOrCreateSession(sessionId, event)
|
|
const message = state.messages.get(messageId) ?? createMessageState()
|
|
if (message.captured) return
|
|
const partId = part.id ?? `${messageId}:${message.parts.size}`
|
|
message.parts.set(partId, part)
|
|
state.messages.set(messageId, message)
|
|
state.lastActivityAt = Date.now()
|
|
debouncedSaveState()
|
|
}
|
|
|
|
async function flushAll({ commit = false } = {}) {
|
|
if (saveTimer) {
|
|
clearTimeout(saveTimer)
|
|
saveTimer = null
|
|
}
|
|
// Do not abort the whole dispose flush when one session fails (e.g. server
|
|
// already GC'd a stale session). Continue so later sessions still commit.
|
|
// See https://github.com/volcengine/OpenViking/issues/4490
|
|
for (const sessionId of sessions.keys()) {
|
|
try {
|
|
await flushSession(sessionId, { commit, reason: "flushAll" })
|
|
} catch (err) {
|
|
console.warn(
|
|
`[opencode-plugin] flushAll skipped session ${sessionId}:`,
|
|
err?.message || err,
|
|
)
|
|
}
|
|
}
|
|
await enqueueSave()
|
|
}
|
|
|
|
async function flushSession(opencodeSessionId, { commit = false, reason = "manual" } = {}) {
|
|
if (!opencodeSessionId) return false
|
|
const state = sessions.get(opencodeSessionId)
|
|
if (!state) return false
|
|
|
|
const added = await flushPendingMessages(opencodeSessionId, state)
|
|
if (commit && isCaptureEnabled(config)) {
|
|
await commitOvSession(state.ovSessionId, { force: true, reason })
|
|
} else if (added > 0) {
|
|
await maybeCommitByThreshold(state)
|
|
}
|
|
await enqueueSave()
|
|
return true
|
|
}
|
|
|
|
async function commitSession(sessionId, opencodeSessionId, abortSignal) {
|
|
if (opencodeSessionId) {
|
|
const state = sessions.get(opencodeSessionId)
|
|
if (state) await flushPendingMessages(opencodeSessionId, state)
|
|
}
|
|
return commitOvSession(sessionId, { force: true, abortSignal, reason: "tool" })
|
|
}
|
|
|
|
return {
|
|
init,
|
|
waitForBackground,
|
|
handleEvent,
|
|
getMappedSessionId,
|
|
commitSession,
|
|
flushAll,
|
|
flushSession,
|
|
}
|
|
|
|
function createSessionState(opencodeSessionId, event = {}) {
|
|
const parentId = event?.properties?.info?.parentID ?? event?.properties?.parentID ?? event?.parentID ?? ""
|
|
const ovSessionId = parentId
|
|
? deriveHarnessSessionId("oc-", parentId, `subagent-${opencodeSessionId}`)
|
|
: deriveHarnessSessionId("oc-", opencodeSessionId)
|
|
return {
|
|
ovSessionId,
|
|
createdAt: Date.now(),
|
|
lastActivityAt: Date.now(),
|
|
lastCommitTime: undefined,
|
|
compactedAt: undefined,
|
|
messages: new Map(),
|
|
}
|
|
}
|
|
|
|
function createMessageState() {
|
|
return {
|
|
role: "",
|
|
parts: new Map(),
|
|
captured: false,
|
|
}
|
|
}
|
|
|
|
function getOrCreateSession(opencodeSessionId, event = {}) {
|
|
let state = sessions.get(opencodeSessionId)
|
|
if (!state) {
|
|
state = createSessionState(opencodeSessionId, event)
|
|
sessions.set(opencodeSessionId, state)
|
|
}
|
|
return state
|
|
}
|
|
|
|
function resolveEventSessionId(event) {
|
|
return event?.properties?.info?.id ??
|
|
event?.properties?.info?.sessionID ??
|
|
event?.properties?.info?.sessionId ??
|
|
event?.properties?.sessionID ??
|
|
event?.properties?.sessionId ??
|
|
event?.sessionID ??
|
|
event?.sessionId ??
|
|
event?.id
|
|
}
|
|
|
|
function resolvePartRole(part, fallbackRole) {
|
|
if (fallbackRole) return fallbackRole
|
|
const type = String(part?.type || part?.kind || "").toLowerCase()
|
|
if (type.includes("tool") && type.includes("call")) return "assistant"
|
|
if (type.includes("tool")) return "user"
|
|
return ""
|
|
}
|
|
|
|
function buildCapturePayload(message) {
|
|
const partsRaw = Array.from(message.parts.values())
|
|
if (partsRaw.length === 0) return null
|
|
const role = resolvePartRole(partsRaw[0], message.role)
|
|
if (!role) return null
|
|
if (role === "assistant" && !config.captureAssistantTurns) return null
|
|
|
|
const shaped = shapeCapturePayload({ role, content: partsRaw }, role, config)
|
|
const toolParts = shaped.parts.filter((part) => part.type !== "text")
|
|
if (shaped.dropped || (!shaped.text && toolParts.length === 0)) return null
|
|
const body = toolParts.length > 0
|
|
? { role, parts: [
|
|
...(shaped.text && shaped.parts.some((part) => part.type === "text")
|
|
? [{ type: "text", text: shaped.text }]
|
|
: []),
|
|
...toolParts,
|
|
] }
|
|
: { role, content: shaped.text }
|
|
const peerId = effectivePeerId(config)
|
|
if (peerId) body.peer_id = peerId
|
|
return body
|
|
}
|
|
|
|
async function flushPendingMessages(opencodeSessionId, state) {
|
|
if (!isCaptureEnabled(config)) return 0
|
|
const toSend = []
|
|
for (const [messageId, message] of state.messages.entries()) {
|
|
if (message.captured) continue
|
|
const body = buildCapturePayload(message)
|
|
if (!body) {
|
|
message.captured = true
|
|
continue
|
|
}
|
|
toSend.push({ messageId, message, body })
|
|
}
|
|
if (toSend.length === 0) return 0
|
|
|
|
let added = 0
|
|
const health = await fetchJSON(config, "/health", {}, { timeoutMs: 5000 })
|
|
if (!health.ok) {
|
|
for (const item of toSend) {
|
|
const queued = await enqueue("addMessage", state.ovSessionId, item.body)
|
|
if (!queued.ok) break
|
|
item.message.captured = true
|
|
added += 1
|
|
}
|
|
} else {
|
|
const res = await sendSessionMessages(
|
|
(endpoint, init = {}, options = {}) => fetchJSON(config, endpoint, init, { timeoutMs: 10000, ...options }),
|
|
state.ovSessionId,
|
|
toSend.map((item) => item.body),
|
|
{ enqueueOnRetryable: true },
|
|
)
|
|
added = res.sent + res.queued
|
|
for (const item of toSend.slice(0, added)) {
|
|
item.message.captured = true
|
|
}
|
|
if (res.failed > 0 || res.enqueueFailed > 0) {
|
|
log("ERROR", "message", "Failed to add message to OpenViking session", {
|
|
openviking_session: state.ovSessionId,
|
|
status: res.lastError?.status,
|
|
error: res.lastError,
|
|
failed: res.failed,
|
|
enqueueFailed: res.enqueueFailed,
|
|
})
|
|
}
|
|
}
|
|
if (added > 0) {
|
|
state.lastActivityAt = Date.now()
|
|
debouncedSaveState()
|
|
}
|
|
return added
|
|
}
|
|
|
|
async function maybeCommitByThreshold(state) {
|
|
if (config.commitTokenThreshold <= 0) return { committed: false }
|
|
const meta = await fetchJSON(config, `/api/v1/sessions/${encodeURIComponent(state.ovSessionId)}`, {}, {
|
|
timeoutMs: 5000,
|
|
})
|
|
const pendingTokens = Number(meta.result?.pending_tokens || 0)
|
|
log("DEBUG", "session", "Pending token check", {
|
|
openviking_session: state.ovSessionId,
|
|
pendingTokens,
|
|
threshold: config.commitTokenThreshold,
|
|
})
|
|
if (!meta.ok || pendingTokens < config.commitTokenThreshold) return { committed: false, pendingTokens }
|
|
return commitOvSession(state.ovSessionId, { force: true, reason: "threshold" })
|
|
}
|
|
|
|
async function commitOvSession(ovSessionId, { force = false, reason = "manual", abortSignal } = {}) {
|
|
if (!force && config.commitTokenThreshold <= 0) return { status: "skipped" }
|
|
const body = { keep_recent_count: config.commitKeepRecentCount }
|
|
const res = await fetchJSON(config, `/api/v1/sessions/${encodeURIComponent(ovSessionId)}/commit`, {
|
|
method: "POST",
|
|
body: JSON.stringify(body),
|
|
signal: abortSignal,
|
|
}, { timeoutMs: 30000 })
|
|
if (res.ok) {
|
|
for (const state of sessions.values()) {
|
|
if (state.ovSessionId === ovSessionId) state.lastCommitTime = Date.now()
|
|
}
|
|
const traceId = res.traceId || res.result?.trace_id
|
|
log("INFO", "session", "Committed OpenViking session", {
|
|
openviking_session: ovSessionId,
|
|
reason,
|
|
trace_id: traceId,
|
|
})
|
|
return { status: "accepted", result: res.result, traceId }
|
|
}
|
|
if (isRetryableFailure(res)) {
|
|
await enqueue("commitSession", ovSessionId, body)
|
|
log("WARN", "session", "Queued OpenViking session commit", {
|
|
openviking_session: ovSessionId,
|
|
reason,
|
|
trace_id: res.traceId,
|
|
status: res.status,
|
|
})
|
|
return { status: "queued" }
|
|
}
|
|
log("ERROR", "session", "Failed to commit OpenViking session", {
|
|
openviking_session: ovSessionId,
|
|
reason,
|
|
trace_id: res.traceId,
|
|
status: res.status,
|
|
error: res.error?.message || res.error?.code,
|
|
})
|
|
throw new Error(
|
|
`Failed to commit OpenViking session ${ovSessionId}: ${res.error?.message || res.status}` +
|
|
(res.traceId ? ` (trace_id=${res.traceId})` : ""),
|
|
)
|
|
}
|
|
|
|
async function migrateLegacySessionMap() {
|
|
if (!fs.existsSync(oldSessionMapPath)) return
|
|
if (fs.existsSync(`${oldSessionMapPath}.migrated`)) return
|
|
try {
|
|
const data = JSON.parse(await fs.promises.readFile(oldSessionMapPath, "utf8"))
|
|
const ovSessionIds = new Set()
|
|
for (const persisted of Object.values(data.sessions ?? {})) {
|
|
if (persisted?.ovSessionId) ovSessionIds.add(persisted.ovSessionId)
|
|
}
|
|
for (const ovSessionId of ovSessionIds) {
|
|
try {
|
|
await commitOvSession(ovSessionId, { force: true, reason: "legacy-migration" })
|
|
} catch (error) {
|
|
log("WARN", "migration", "Legacy orphan session commit failed", {
|
|
openviking_session: ovSessionId,
|
|
error: error?.message,
|
|
})
|
|
}
|
|
}
|
|
await fs.promises.rename(oldSessionMapPath, `${oldSessionMapPath}.migrated`)
|
|
log("INFO", "migration", "Migrated legacy session map", { count: ovSessionIds.size })
|
|
} catch (error) {
|
|
log("ERROR", "migration", "Failed to migrate legacy session map", { error: error?.message })
|
|
}
|
|
}
|
|
|
|
}
|