- shared TS collection manager: self-heal creates missing collection (1024/cosine) and missing keyword payload indexes; never mutates incompatible contracts (semantic_index_incompatible); async index visibility polled with bounded deadline - session admission (qdrantEnsure) uses the manager in self-heal mode; operator path keeps require_existing semantics - runtime lease exposes semanticQdrantUrl to the operator - operator commands vector-inspect/vector-rebuild with exact confirmation guards - thothctl workspace vector inspect|rebuild (Go) with --collection/--confirm/--destroy - p4 acceptance runner: real Qdrant (v1.18.2) lifecycle checks, 11/11 PASS - docs: CLI contract, manual walkthrough P4 (PENDING), PROJECT_STATE
73 lines
2.3 KiB
TypeScript
73 lines
2.3 KiB
TypeScript
import type {
|
|
OllamaEnsureResult,
|
|
SemanticReadinessCode,
|
|
ThtRunner,
|
|
} from "../tht/tht-runner.js";
|
|
import type { PrincipalContext } from "../auth/principal.js";
|
|
import type { WorkspaceDescriptor } from "../workspaces/schema.js";
|
|
|
|
export type ReadinessResult = OllamaEnsureResult & { code?: SemanticReadinessCode };
|
|
|
|
interface ReadyEntry {
|
|
expiresAt: number;
|
|
result: ReadinessResult;
|
|
}
|
|
|
|
/**
|
|
* Deduplicates embedding readiness checks and keeps only short-lived successes.
|
|
* Failures are deliberately not cached so a submit can retry after a transient outage.
|
|
*/
|
|
export class ReadinessManager {
|
|
private inFlight = new Map<string, Promise<ReadinessResult>>();
|
|
private ready = new Map<string, ReadyEntry>();
|
|
|
|
constructor(
|
|
private tht: ThtRunner,
|
|
private timeoutSec: number,
|
|
private ttlMs = 60_000,
|
|
private now: () => number = Date.now,
|
|
) {}
|
|
|
|
ensure(
|
|
workspace = "",
|
|
principal?: PrincipalContext,
|
|
descriptor?: WorkspaceDescriptor,
|
|
): Promise<ReadinessResult> {
|
|
const key = `${principal?.issuer ?? ""}\0${principal?.subject ?? ""}\0${workspace}`;
|
|
const cached = this.ready.get(key);
|
|
if (cached && cached.expiresAt > this.now()) return Promise.resolve(cached.result);
|
|
if (cached) this.ready.delete(key);
|
|
|
|
const current = this.inFlight.get(key);
|
|
if (current) return current;
|
|
|
|
const runner = principal && typeof (this.tht as any).withPrincipal === "function"
|
|
? this.tht.withPrincipal(principal) : this.tht;
|
|
const pending = (async (): Promise<ReadinessResult> => {
|
|
try {
|
|
if (descriptor) {
|
|
const qdrant = await runner.qdrantEnsure(descriptor, this.timeoutSec, "self_heal");
|
|
if (!qdrant.ok) return qdrant;
|
|
}
|
|
const ollama = await runner.ollamaEnsure(workspace, this.timeoutSec);
|
|
return ollama.ok
|
|
? ollama
|
|
: { ...ollama, code: "workspace_not_activatable" };
|
|
} catch {
|
|
return { ok: false, code: "workspace_not_activatable" };
|
|
}
|
|
})()
|
|
.then((result) => {
|
|
if (result.ok) {
|
|
this.ready.set(key, { result, expiresAt: this.now() + this.ttlMs });
|
|
}
|
|
return result;
|
|
})
|
|
.finally(() => {
|
|
if (this.inFlight.get(key) === pending) this.inFlight.delete(key);
|
|
});
|
|
this.inFlight.set(key, pending);
|
|
return pending;
|
|
}
|
|
}
|