Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
4.8 kB · 186 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187import { schema as s } from "jazz-tools";
export const inferenceAccountingEnabled = process.env.THOUGHTSTREAM_INFERENCE_ACCOUNTING !== "disabled";
const coreSchema = { events: s.table({ key: s.string(), sourceSequence: s.int().optional(), type: s.string(), schemaVersion: s.int(), source: s.string(), sourceKind: s.string(), externalId: s.string(), idempotencyKey: s.string(), occurredAt: s.string(), observedAt: s.string(), actor: s.string(), rootEventKey: s.string(), parentEventKey: s.string(), correlationId: s.string(), privacy: s.string(), payloadJson: s.string(), payloadHash: s.string(), traceKey: s.string(), createdByRuntime: s.string(), }), documentVersions: s.table({ key: s.string(), source: s.string(), documentKey: s.string(), path: s.string(), contentType: s.string(), sha256: s.string(), content: s.string(), sizeBytes: s.int(), mtimeMs: s.float(), createdAt: s.string(), }), documents: s.table({ key: s.string(), source: s.string(), documentKey: s.string(), path: s.string(), versionKey: s.string(), sha256: s.string(), contentType: s.string(), sizeBytes: s.int(), mtimeMs: s.float(), deleted: s.boolean(), updatedAt: s.string(), }), sources: s.table({ key: s.string(), kind: s.string(), enabled: s.boolean(), configJson: s.string(), lastSequence: s.int().optional(), createdAt: s.string(), updatedAt: s.string(), }), sourceCursors: s.table({ key: s.string(), source: s.string(), cursorJson: s.string(), lastSuccessAt: s.string(), lastFailureAt: s.string(), lastError: s.string(), updatedAt: s.string(), }), agents: s.table({ key: s.string(), version: s.int(), enabled: s.boolean(), specJson: s.string(), specHash: s.string(), updatedAt: s.string(), }), runs: s.table({ key: s.string(), executionKey: s.string().optional(), triggerEventKey: s.string().optional(), agentKey: s.string(), agentVersion: s.int(), status: s.string(), inputEventIdsJson: s.string(), outputEventIdsJson: s.string(), attempt: s.int(), provider: s.string(), model: s.string(), adapterRevision: s.string(), promptHash: s.string(), contextManifestJson: s.string(), resultJson: s.string(), errorText: s.string(), createdAt: s.string(), startedAt: s.string(), completedAt: s.string(), updatedAt: s.string(), }), traceChunks: s.table({ key: s.string(), runKey: s.string(), sequence: s.int(), type: s.string(), payloadJson: s.string(), createdAt: s.string(), }), consumerProgress: s.table({ key: s.string(), consumerKey: s.string(), consumerVersion: s.int(), source: s.string(), lastSequence: s.int(), lastEventKey: s.string(), updatedAt: s.string(), }), projections: s.table({ key: s.string(), payloadJson: s.string(), lastEventKey: s.string(), projectionVersion: s.int(), updatedAt: s.string(), }),};
const accountingSchema = { ...coreSchema, runs: s.table({ key: s.string(), executionKey: s.string().optional(), triggerEventKey: s.string().optional(), agentKey: s.string(), agentVersion: s.int(), status: s.string(), inputEventIdsJson: s.string(), outputEventIdsJson: s.string(), attempt: s.int(), provider: s.string(), model: s.string(), adapterRevision: s.string(), promptHash: s.string(), contextManifestJson: s.string(), accountingReservationKey: s.string(), resultJson: s.string(), errorText: s.string(), createdAt: s.string(), startedAt: s.string(), completedAt: s.string(), updatedAt: s.string(), }), inferenceBudgetAccounts: s.table({ key: s.string(), scopeType: s.string(), scopeKey: s.string(), windowsJson: s.string(), activeLeasesJson: s.string(), updatedAt: s.string(), }), inferenceAccounting: s.table({ key: s.string(), scopeType: s.string(), scopeKey: s.string(), runKey: s.string(), agentKey: s.string(), agentVersion: s.int(), provider: s.string(), model: s.string(), status: s.string(), usageStatus: s.string(), estimateJson: s.string(), chargedJson: s.string(), actualUsageJson: s.string(), windowKeysJson: s.string(), reservedAt: s.string(), leaseExpiresAt: s.string(), settledAt: s.string(), denialReason: s.string(), limitingWindowKeysJson: s.string(), updatedAt: s.string(), }),};
const accountingApp = s.defineApp(accountingSchema);export const thoughtstreamApp = ( inferenceAccountingEnabled ? accountingApp : s.defineApp(coreSchema)) as typeof accountingApp;