Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -0,0 +1,55 @@
import { createHash } from "node:crypto";
import { broadcast } from "../wire";
import { JobRegistry } from "../jobs/registry";
import { JobService } from "../jobs/service";
import { researchAnswerDefinition, researchEvidenceDefinition } from "../jobs/research-workspace";
import { ResearchWorkspaceService } from "../research/service";
import { REPORT } from "../research/test-support";
import * as Plan from "../plan/service";
import type { headingMemory } from "../chat/job-heading-memory.test-fixtures";

async function evidence() {
return { findings: [], sources: [] };
}
async function privateEvidence() {
return { findings: [] };
}
async function report() {
return REPORT;
}
async function answer() {
return { text: "Offline answer", sourceUrls: [] };
}

export function researchStorage(h: Awaited<ReturnType<typeof headingMemory>>) {
let jobs = new JobService({
storage: h.opened.storage,
registry: new JobRegistry([
researchEvidenceDefinition({ config: { agent: true, model: "offline" }, engine: evidence }),
researchAnswerDefinition({
config: { agent: true, model: "offline" },
engines: { private: privateEvidence, synthesize: report, answer },
}),
]),
lease: () => h.opened.lease,
});
let research = new ResearchWorkspaceService({
storage: h.opened.storage,
jobs,
lease: () => h.opened.lease,
current: async () => ({
channelId: h.plan.id,
revision: h.plan.revision,
source: Plan.source(h.plan),
sourceHash: `sha256:${createHash("sha256").update(Plan.source(h.plan)).digest("hex")}`,
}),
publish: (channelId, workspaceId, revision) =>
broadcast(h.opened.server, channelId, {
kind: "research:changed",
ts: 0,
workspaceId,
revision,
}),
});
return { research, jobs };
}
136 changes: 115 additions & 21 deletions apps/server/src/main.ts
Original file line number Diff line number Diff line change
Expand Up @@ -37,7 +37,8 @@ import * as Inject from "./questions/inject";
import * as Marks from "./comments/inject";
import * as Questions from "./questions/service";
import { registerResearchWorkspaceRoutes } from "./research/routes";
import { ResearchWorkspaceService } from "./research/service";
import { ResearchWorkspaceError, ResearchWorkspaceService } from "./research/service";
import { placeResearchReference as placeResearch } from "./research/placement";
import * as Rooms from "./rooms";
import { admit } from "./socket/admission";
import { StorageError } from "./storage/errors";
Expand Down Expand Up @@ -70,6 +71,7 @@ const LEASE_RENEW_MS = 10_000;
const LEASE_SAFETY_MS = 5_000;
const SESSION_CLEANUP_MS = 5 * 60_000;
const ACCESS_RECHECK_MS = 60_000;
const RESEARCH_RECOVERY_RETRY_MS = 10_000;

let server: Server<SocketData>;
let heldLease: Lease | undefined;
Expand All @@ -81,6 +83,8 @@ let cleaningSessions: Promise<void> | undefined;
let ownerBindings: ActiveOwnerBindings | undefined;
let jobRunner: JobRunner | undefined;
let researchService: ResearchWorkspaceService | undefined;
let researchRecoveryTimer: ReturnType<typeof setTimeout> | undefined;
let recoveringResearch: Promise<void> | undefined;
let referenceService: ReferenceService | undefined;
let summaryCoordinator: DocumentSummaryCoordinator | undefined;
let descriptionProjector: DocumentDescriptionProjector | undefined;
Expand Down Expand Up @@ -119,17 +123,8 @@ function presence(server: Server<SocketData>, room: Rooms.Room): void {
});
}

/**
* Attach the document to a room, once.
*
* Two clients opening at the same moment must not build two documents, so the
* first stores its promise and the second waits on it.
*/
async function plan(room: Rooms.Room, server: Server<SocketData>): Promise<Service.Plan> {
if (room.closing) await room.closing;
if (deletingChannels.has(room.id)) throw new Error("document is unavailable");
if (room.plan) return room.plan;
let backend: Service.Backend = {
function documentBackend(): Service.Backend {
return {
storage,
lease: () => {
if (!heldLease) throw new Error("storage writer lease is unavailable");
Expand All @@ -141,6 +136,56 @@ async function plan(room: Rooms.Room, server: Server<SocketData>): Promise<Servi
},
onDocumentPersisted: target => summaryCoordinator?.schedule(target),
};
}

function placeResearchReference(
channelId: string,
workspaceId: string,
): Promise<"placed" | "deferred"> {
return placeResearch(channelId, workspaceId, {
get: Rooms.get,
exclusive: withDocumentLock,
detached: async (id, action) => {
let detached = await Service.open(id, documentBackend(), server);
try {
return await action(detached);
} finally {
await Service.close(detached);
}
},
});
}

function scheduleResearchRecovery(deferred: number): void {
if (deferred === 0 || researchRecoveryTimer || draining) return;
researchRecoveryTimer = setTimeout(() => {
researchRecoveryTimer = undefined;
if (draining) return;
recoveringResearch = researchService!.recoverPendingPlannerInline(placeResearchReference).then(
result => {
scheduleResearchRecovery(result.deferred);
},
err => {
console.error("chopin: research reference recovery failed -", err);
signal();
},
).finally(() => {
recoveringResearch = undefined;
});
}, RESEARCH_RECOVERY_RETRY_MS);
}

/**
* Attach the document to a room, once.
*
* Two clients opening at the same moment must not build two documents, so the
* first stores its promise and the second waits on it.
*/
async function plan(room: Rooms.Room, server: Server<SocketData>): Promise<Service.Plan> {
if (room.closing) await room.closing;
if (deletingChannels.has(room.id)) throw new Error("document is unavailable");
if (room.plan) return room.plan;
let backend = documentBackend();
let opening = room.opening ??= withDocumentLock(room.id, async () => {
if (deletingChannels.has(room.id)) throw new Error("document is unavailable");
if (room.plan) return room.plan;
Expand Down Expand Up @@ -186,14 +231,23 @@ function chat(room: Rooms.Room, ws: Socket): Chat.Room {
? async request => {
let service = researchService;
if (!service) throw new Error("research workspaces are unavailable");
let created = await service.startPlanner({
channelId: room.id,
question: request.question,
originMessageId: request.entryId,
requestedBy: request.userId,
requestedByHandle: request.handle,
beforeStart: () => jobRunner?.ownerAvailable(room.id) ?? Promise.resolve(),
});
let created;
try {
created = await service.startPlannerInline({
channelId: room.id,
question: request.question,
originMessageId: request.entryId,
requestedBy: request.userId,
requestedByHandle: request.handle,
beforeStart: () => jobRunner?.ownerAvailable(room.id) ?? Promise.resolve(),
placeReference: id => placeResearchReference(room.id, id),
});
} catch (err) {
if (err instanceof ResearchWorkspaceError && err.code === "not-ready") {
scheduleResearchRecovery(1);
}
throw err;
}
return {
workspaceId: created.request.id,
state: created.request.state,
Expand Down Expand Up @@ -644,6 +698,8 @@ function drain(): Promise<void> {
}
};
await attempt(() => server.stop(true));
if (researchRecoveryTimer) clearTimeout(researchRecoveryTimer);
if (recoveringResearch) await attempt(() => recoveringResearch!);
if (sessionCleanup) clearInterval(sessionCleanup);
for (let result of await Promise.allSettled([cleaningSessions])) {
if (result.status === "rejected") record(result.reason);
Expand Down Expand Up @@ -808,6 +864,11 @@ async function restoreChannelLocked(channelId: string, now: Date) {
channelId,
() => storage.channels.restore({ id: channelId, now }),
);
let recovery = await researchService?.recoverPendingPlannerInline(
placeResearchReference,
channelId,
);
scheduleResearchRecovery(recovery?.deferred ?? 0);
summaryCoordinator?.resume(channelId);
announceChannel(result.channel);
if (summaryCoordinator) void summaryCoordinator.ensure(channelId).catch(() => {});
Expand Down Expand Up @@ -907,6 +968,32 @@ function announceResearchChanged(channelId: string, workspaceId: string, revisio
});
}

function announceResearchTerminal(channelId: string, id: string, text: string): Promise<void> {
return withDocumentLock(channelId, async () => {
let active = Rooms.get(channelId)?.plan;
if (active) {
let existing = active.chat.entries.find(entry => entry.id === id);
await Chat.noticeOnce(
{ chat: active.chat, plan: active, server, room: channelId },
id,
existing?.text ?? text,
);
return;
}
let detached = await Service.open(channelId, documentBackend(), server);
try {
let existing = detached.chat.entries.find(entry => entry.id === id);
await Chat.noticeOnce(
{ chat: detached.chat, plan: detached, server, room: channelId },
id,
existing?.text ?? text,
);
} finally {
await Service.close(detached);
}
});
}

async function currentDocumentTarget(
channelId: string,
): Promise<Service.DocumentTarget | undefined> {
Expand Down Expand Up @@ -1047,6 +1134,7 @@ researchService = new ResearchWorkspaceService({
},
current: currentDocumentTarget,
publish: announceResearchChanged,
terminalNotice: announceResearchTerminal,
});
referenceService = new ReferenceService({
storage,
Expand Down Expand Up @@ -1185,10 +1273,16 @@ leaseRenewal = setInterval(renewLease, LEASE_RENEW_MS);
cleanSessions();
sessionCleanup = setInterval(cleanSessions, SESSION_CLEANUP_MS);

let listening: Server<SocketData> | undefined;
try {
server = listen();
server = listening = listen();
let recovery = await researchService.recoverPendingPlannerInline(placeResearchReference);
scheduleResearchRecovery(recovery.deferred);
await researchService.recoverTerminalPlannerInline();
jobRunner.start();
} catch (err) {
if (listening) await listening.stop(true).catch(() => {});
if (researchRecoveryTimer) clearTimeout(researchRecoveryTimer);
if (sessionCleanup) clearInterval(sessionCleanup);
if (leaseRenewal) clearInterval(leaseRenewal);
if (leaseWatchdog) clearTimeout(leaseWatchdog);
Expand Down
Loading
Loading