Files
OpenViking/examples/pi-coding-agent-extension/sync.ts
T

125 lines
3.9 KiB
TypeScript

import type { OVClient } from "./client.js";
import type { OVConfig } from "./config.js";
import { deriveHarnessSessionId } from "./shared/session-model.mjs";
import { enqueue, listPending, replayPending } from "./shared/pending-queue.mjs";
import { extractBranchCapturePayloads } from "./lib/capture-adapter.mjs";
import { countUndeliveredForSession, estimatePayloadTokens } from "./lib/takeover-core.mjs";
// --- SyncManager ---
export interface AddPayloadResult {
accepted: boolean;
delivered: boolean;
}
export interface SyncBranchResult {
added: number;
tokens: number;
allDelivered: boolean;
}
export class SyncManager {
private client: OVClient;
private config: OVConfig;
private ovSessionId: string | null = null;
private syncedEntryCount = 0;
constructor(client: OVClient, config: OVConfig) {
this.client = client;
this.config = config;
}
get sessionId(): string | null { return this.ovSessionId; }
get syncedCount(): number { return this.syncedEntryCount; }
restoreWatermark(n: number): void {
const next = Math.max(0, Math.floor(Number(n) || 0));
this.syncedEntryCount = next;
}
async ensureSession(piSessionId: string): Promise<boolean> {
if (this.ovSessionId) return true;
const id = deriveHarnessSessionId("pi-", piSessionId);
this.ovSessionId = id;
return true;
}
async replayPending(): Promise<void> {
if (!this.client.connected) return;
await replayPending(
(path: string, init?: any) => this.client.fetchJSON(path, init, 10000),
() => {},
);
}
async flushForTakeover(): Promise<boolean> {
if (!this.ovSessionId) return false;
await this.replayPending();
const pending = await listPending();
return countUndeliveredForSession(pending, this.ovSessionId) === 0;
}
async syncBranch(branch: any[]): Promise<SyncBranchResult> {
if (!this.ovSessionId) return { added: 0, tokens: 0, allDelivered: true };
const extracted = extractBranchCapturePayloads(branch, this.syncedEntryCount, this.config);
if (extracted.resetWatermark) this.syncedEntryCount = 0;
let added = 0;
let tokens = 0;
let allDelivered = true;
for (const payload of extracted.payloads) {
const result = await this.addPayload(payload);
if (!result.accepted) break;
added++;
tokens += estimatePayloadTokens(payload);
allDelivered = allDelivered && result.delivered;
}
if (added === extracted.payloads.length) {
this.syncedEntryCount = extracted.nextEntryCount;
}
if (added > 0 && !this.config.takeoverEnabled) {
await this.commitIfNeeded();
}
return { added, tokens, allDelivered };
}
async addPayload(payload: any): Promise<AddPayloadResult> {
if (!this.ovSessionId) return { accepted: false, delivered: false };
const ok = await this.client.addMessagePayload(this.ovSessionId, payload);
if (ok) return { accepted: true, delivered: true };
await enqueue("addMessage", this.ovSessionId, payload);
return { accepted: true, delivered: false };
}
async commitIfNeeded(): Promise<void> {
if (!this.ovSessionId) return;
const meta = await this.client.getSession(this.ovSessionId);
const pending = Number(meta?.pending_tokens || 0);
if (pending >= this.config.commitTokenThreshold) {
await this.commit();
}
}
async commit(opts: { queueOnFailure?: boolean; keepRecentCount?: number } = {}): Promise<any | null> {
if (!this.ovSessionId) return null;
const result = await this.client.commitSession(
this.ovSessionId,
opts.keepRecentCount,
);
if (!result) {
if (opts.queueOnFailure !== false) {
await enqueue("commitSession", this.ovSessionId, {
keep_recent_count: opts.keepRecentCount ?? this.config.commitKeepRecentCount,
});
}
return null;
}
return result;
}
async shutdown(): Promise<void> {
return;
}
}