mirror of
https://github.com/volcengine/OpenViking.git
synced 2026-09-30 09:17:51 +08:00
Co-authored-by: ZaynJarvis <31875147+ZaynJarvis@users.noreply.github.com>
37 lines
1.1 KiB
TypeScript
37 lines
1.1 KiB
TypeScript
import type { OVClient } from "./client.js";
|
|
import type { OVConfig } from "./config.js";
|
|
import type { SyncManager } from "./sync.js";
|
|
import { TakeoverCore } from "./lib/takeover-core.mjs";
|
|
|
|
export function createTakeoverManager(opts: {
|
|
pi: any;
|
|
client: OVClient;
|
|
sync: SyncManager;
|
|
config: OVConfig;
|
|
log?: (message: string) => void;
|
|
}): TakeoverCore {
|
|
const { pi, client, sync, config } = opts;
|
|
return new TakeoverCore({
|
|
config,
|
|
io: {
|
|
flush: () => sync.flushForTakeover(),
|
|
commit: (commitOpts?: { queueOnFailure?: boolean; keepRecentCount?: number }) => sync.commit(commitOpts),
|
|
fetchOverview: async (tokenBudget?: number) => {
|
|
if (!sync.sessionId) return "";
|
|
const ctx = await client.getSessionContext(
|
|
sync.sessionId,
|
|
tokenBudget ?? config.takeoverOverviewBudget * 4,
|
|
);
|
|
return ctx?.latest_archive_overview ?? "";
|
|
},
|
|
persistEntry: (customType: string, data: any) => {
|
|
if (typeof pi?.appendEntry === "function") {
|
|
pi.appendEntry(customType, data);
|
|
}
|
|
},
|
|
getWatermark: () => sync.syncedCount,
|
|
log: opts.log,
|
|
},
|
|
});
|
|
}
|