feat: add schema foundation for document pipeline, agent, and deployment

Add missing topics (librarian, knowledge, collection-management, flow),
pipeline message types (TextDocument, Chunk, Triples, EntityContexts),
service message types (Librarian, Knowledge, Collection, Flow CRUD),
and update AgentResponse for streaming chunk format.

Add RequestResponseSpec enabling flow-scoped request/response calls
(needed by knowledge extraction and agent services). Add requestor
registry to Flow class with proper lifecycle management.

Add end_of_dialog to gateway's isComplete() check for agent streaming.

Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
This commit is contained in:
elpresidank 2026-04-06 00:11:29 -05:00
parent 28747e1a92
commit 5ed3f0e2d8
6 changed files with 236 additions and 1 deletions

View file

@ -0,0 +1,36 @@
/**
* Request/response specification declares a request/response client for a flow.
*
* Enables FlowProcessor handlers to make request/response calls to other services
* (e.g., calling the prompt service or LLM from within a knowledge extraction handler).
*
* Python reference: trustgraph-base/trustgraph/base/prompt_client_spec.py
*/
import type { Spec } from "./types.js";
import type { PubSubBackend } from "../backend/types.js";
import type { Flow, FlowDefinition } from "../processor/flow.js";
import { RequestResponse } from "../messaging/request-response.js";
export class RequestResponseSpec<TReq, TRes> implements Spec {
constructor(
public readonly name: string,
private readonly requestTopicName: string,
private readonly responseTopicName: string,
) {}
async add(flow: Flow, pubsub: PubSubBackend, definition: FlowDefinition): Promise<void> {
const requestTopic = definition.topics?.[this.requestTopicName] ?? this.requestTopicName;
const responseTopic = definition.topics?.[this.responseTopicName] ?? this.responseTopicName;
const rr = new RequestResponse<TReq, TRes>({
pubsub,
requestTopic,
responseTopic,
subscription: `${flow.processorId}-${flow.name}-${this.name}`,
});
await rr.start();
flow.registerRequestor(this.name, rr as RequestResponse<unknown, unknown>);
}
}