Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
8.4 kB · 115 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116import { createHash } from "node:crypto";import fs from "node:fs/promises";import path from "node:path";import type { JsonObject } from "../core/json.js";import type { EventCandidate, PrivacyClass, ThoughtEvent } from "../events/types.js";import type { JazzThoughtStore } from "../jazz/store.js";import { readWorkspaceArtifact, storeBlob } from "./blob.js";import { rebuildArtifactCatalog } from "./catalog.js";import type { ArtifactKind, ArtifactMediaType, ArtifactPayload, ArtifactProvenance, ArtifactRelation, ArtifactRequestPayload, ArtifactRequestStatus, ArtifactVisibility } from "./types.js";import { ARTIFACT_EVENT_TYPE, ARTIFACT_REQUEST_EVENT_TYPE, ARTIFACT_REQUEST_SCHEMA_VERSION, ARTIFACT_SCHEMA_VERSION } from "./types.js";
export interface ArtifactRequestInput { requestId: string; workspaceLeaseId: string; workspaceRootLabel: string; relativeFilePath: string; artifactId: string; artifactVersion: number; kind: ArtifactKind; title: string; summary: string; mediaType: ArtifactMediaType; provenance: ArtifactProvenance; visibility: ArtifactVisibility; requestingAgentId?: string; requestingRunId?: string; sourceEventId?: string; supersedesArtifactEventId?: string; status?: ArtifactRequestStatus; reasonCode?: string;}export interface AppendArtifactRequestOptions extends ArtifactRequestInput { store: JazzThoughtStore; occurredAt?: string }export interface MaterializeArtifactOptions { store: JazzThoughtStore; requestEventId: string; workspaceRoot: string; workspaceLeaseId: string; occurredAt?: string }export interface AppendArtifactResult { requestEvent?: ThoughtEvent; event: ThoughtEvent; inserted: boolean; blobInserted: boolean; bodySha256: string; byteCount: number }export interface RequestAndMaterializeOptions extends ArtifactRequestInput { store: JazzThoughtStore; workspaceRoot: string; occurredAt?: string }
export async function appendArtifactRequest(options: AppendArtifactRequestOptions): Promise<{ event: ThoughtEvent; inserted: boolean }> { const status = options.status ?? "pending"; const payload: ArtifactRequestPayload = { requestId: options.requestId, workspaceLeaseId: options.workspaceLeaseId, workspaceRootLabel: options.workspaceRootLabel, relativeFilePath: options.relativeFilePath, artifactId: options.artifactId, artifactVersion: options.artifactVersion, kind: options.kind, title: options.title, summary: options.summary, mediaType: options.mediaType, provenance: options.provenance, visibility: options.visibility, publicationEligible: false, publicationAuthority: false, status, ...(options.reasonCode ? { reasonCode: options.reasonCode } : {}), ...(options.requestingAgentId ? { requestingAgentId: options.requestingAgentId } : {}), ...(options.requestingRunId ? { requestingRunId: options.requestingRunId } : {}), ...(options.sourceEventId ? { sourceEventId: options.sourceEventId } : {}), ...(options.supersedesArtifactEventId ? { supersedesArtifactEventId: options.supersedesArtifactEventId } : {}), }; const privacy = privacyFor(options.visibility); const result = await options.store.appendEvent({ type: ARTIFACT_REQUEST_EVENT_TYPE, schemaVersion: ARTIFACT_REQUEST_SCHEMA_VERSION, source: "artifact-request:local-broker", sourceKind: "system", externalId: `artifact-request:${options.requestId}:${status}`, idempotencyKey: `artifact-request:${options.requestId}:${status}`, occurredAt: options.occurredAt ?? new Date().toISOString(), actor: options.requestingAgentId ?? "trusted-artifact-broker", correlationId: `artifact-request:${options.requestId}`, privacy, payload: payload as unknown as JsonObject, }); return result;}
export async function materializeArtifactRequest(options: MaterializeArtifactOptions): Promise<AppendArtifactResult> { const request = await options.store.getEvent(options.requestEventId); if (!request || request.type !== ARTIFACT_REQUEST_EVENT_TYPE) throw new Error("Artifact request event not found"); const payload = request.payload as unknown as ArtifactRequestPayload; if (payload.status !== "pending") throw new Error("Only pending artifact requests can be materialized"); if (payload.workspaceLeaseId !== options.workspaceLeaseId) throw new Error("Workspace lease identity does not match artifact request"); const bytes = await readWorkspaceArtifact(options.workspaceRoot, payload.relativeFilePath, payload.mediaType); const hash = hashBytes(bytes); const blobPath = `sha256/${hash.slice(0, 2)}/${hash}`; let blobInserted = false; try { await fs.access(path.join(options.store.getArtifactRoot(), ...blobPath.split("/"))); } catch { blobInserted = true; } const blob = await storeBlob(options.store.getArtifactRoot(), bytes); const relations: ArtifactRelation[] = [{ type: "materialized-from-request", eventId: request.id }]; if (payload.sourceEventId) relations.push({ type: "source-event", eventId: payload.sourceEventId }); if (payload.requestingRunId) relations.push({ type: "requesting-run", eventId: payload.requestingRunId }); const artifactPayload: ArtifactPayload = { artifactId: payload.artifactId, artifactVersion: payload.artifactVersion, kind: payload.kind, title: payload.title, summary: payload.summary, mediaType: payload.mediaType, blob, provenance: payload.provenance, relations, visibility: payload.visibility, publicationEligible: false, ...(payload.supersedesArtifactEventId ? { supersedesArtifactEventId: payload.supersedesArtifactEventId } : {}), }; const identity = `artifact:${payload.artifactId}:v${payload.artifactVersion}`; const candidate: EventCandidate = { type: ARTIFACT_EVENT_TYPE, schemaVersion: ARTIFACT_SCHEMA_VERSION, source: "artifact-materializer:local-broker", sourceKind: "system", externalId: identity, idempotencyKey: identity, occurredAt: options.occurredAt ?? request.occurredAt, actor: "trusted-artifact-materializer", correlationId: request.id, privacy: privacyFor(payload.visibility), payload: artifactPayload as unknown as JsonObject, }; const result = await options.store.appendEvent(candidate); await rebuildArtifactCatalog(options.store); return { requestEvent: request, event: result.event, inserted: result.inserted, blobInserted, bodySha256: blob.sha256, byteCount: blob.byteCount };}
export async function requestArtifactStorage(options: RequestAndMaterializeOptions): Promise<AppendArtifactResult> { const request = await appendArtifactRequest(options); const result = await materializeArtifactRequest({ store: options.store, requestEventId: request.event.id, workspaceRoot: options.workspaceRoot, workspaceLeaseId: options.workspaceLeaseId, ...(options.occurredAt ? { occurredAt: options.occurredAt } : {}) }); return { ...result, requestEvent: request.event };}
/** Compatibility wrapper retained only for callers migrating to the explicit request seam. */export async function appendArtifactFromFile(options: { store: JazzThoughtStore; rootPath: string; filePath: string; artifactId: string; artifactVersion: number; kind: ArtifactKind; title: string; summary: string; mediaType: ArtifactMediaType; provenance: ArtifactProvenance; visibility: ArtifactVisibility; source: string; actor: string; relations?: ArtifactRelation[]; supersedesArtifactEventId?: string; occurredAt?: string;}): Promise<AppendArtifactResult> { const relativeFilePath = path.relative(options.rootPath, options.filePath).split(path.sep).join("/"); return requestArtifactStorage({ store: options.store, workspaceRoot: options.rootPath, workspaceLeaseId: "legacy-explicit-root", workspaceRootLabel: "legacy explicit root", relativeFilePath, requestId: `${options.artifactId}-v${options.artifactVersion}`, artifactId: options.artifactId, artifactVersion: options.artifactVersion, kind: options.kind, title: options.title, summary: options.summary, mediaType: options.mediaType, provenance: options.provenance, visibility: options.visibility, ...(options.supersedesArtifactEventId ? { supersedesArtifactEventId: options.supersedesArtifactEventId } : {}), ...(options.occurredAt ? { occurredAt: options.occurredAt } : {}) });}
function privacyFor(visibility: ArtifactVisibility): PrivacyClass { return visibility === "sensitive" ? "sensitive" : "private" }function hashBytes(bytes: Buffer): string { return createHash("sha256").update(bytes).digest("hex") }