diff --git a/.github/workflows/build-ingest-binary.yml b/.github/workflows/build-ingest-binary.yml index 278d5b814c..74c559d954 100644 --- a/.github/workflows/build-ingest-binary.yml +++ b/.github/workflows/build-ingest-binary.yml @@ -49,7 +49,7 @@ jobs: path: | ~/.cargo-ingest ${{ runner.temp }}/ingest-target - key: ingest-cargo-${{ runner.os }}-${{ runner.arch }}-${{ hashFiles('apps/ingest/Cargo.lock') }}-${{ hashFiles('apps/ingest/src/**') }} + key: ingest-cargo-${{ runner.os }}-${{ runner.arch }}-${{ hashFiles('apps/ingest/Cargo.lock') }}-${{ hashFiles('apps/ingest/src/**', 'apps/ingest/crates/**') }} restore-keys: | ingest-cargo-${{ runner.os }}-${{ runner.arch }}-${{ hashFiles('apps/ingest/Cargo.lock') }}- ingest-cargo-${{ runner.os }}-${{ runner.arch }}- diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index c92104aa74..e5bf341057 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -78,6 +78,10 @@ jobs: - *base - 'apps/cli/**' - 'apps/local-ui/**' + # The binary embeds the gateway's AI stamping as wasm. + - 'apps/ingest/crates/**' + - 'apps/ingest/Cargo.*' + - 'apps/ingest/.cargo/**' - 'packages/effect-sdk/**' - 'packages/domain/**' - 'packages/query-engine/**' @@ -361,9 +365,10 @@ jobs: - uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v6 - uses: jdx/mise-action@c2a87611a18de5b3828c5652fe268e992400cb5c # v4.3.0 with: - # Rust and Python are for the ingest and Tinybird workflows; skipping - # them keeps this job off their ~1 GB of toolchain cache. - install_args: bun node + # Python is for the Tinybird workflows and Rust for ingest, except that + # @maple/cli's tests build the gateway's AI stamping to wasm. Skipping + # them elsewhere keeps lanes off their ~1 GB of toolchain cache. + install_args: ${{ contains(matrix.filters, '@maple/cli') && 'bun node rust' || 'bun node' }} - uses: ./.github/actions/bun-install with: filters: ${{ join(matrix.filters, ' ') }} diff --git a/.gitignore b/.gitignore index f03cea3d1f..34034c83d7 100644 --- a/.gitignore +++ b/.gitignore @@ -75,3 +75,6 @@ apps/ios/build/ # Vitest Browser Mode failure screenshots and traces .vitest/ + +# AI stamping wasm, built from apps/ingest by `bun run --cwd apps/cli build:ai-stamp` +apps/cli/src/server/otlp/ai-stamp.wasm diff --git a/apps/cli/package.json b/apps/cli/package.json index 17d25e0505..8a5a9d84b2 100644 --- a/apps/cli/package.json +++ b/apps/cli/package.json @@ -6,8 +6,9 @@ }, "type": "module", "scripts": { - "start": "bun run src/bin.ts", - "test": "bun test --max-concurrency=1 test/*.test.ts src/server/otlp/encode.test.ts src/core/*.test.ts", + "build:ai-stamp": "cd ../ingest && cargo build --locked --profile wasm --target wasm32-unknown-unknown -p maple-ai-session-wasm && cp target/wasm32-unknown-unknown/wasm/maple_ai_session_wasm.wasm ../cli/src/server/otlp/ai-stamp.wasm", + "start": "bun run build:ai-stamp && bun run src/bin.ts", + "test": "bun run build:ai-stamp && bun test --max-concurrency=1 test/*.test.ts src/server/otlp/encode.test.ts src/core/*.test.ts", "typecheck": "tsc --noEmit" }, "dependencies": { diff --git a/apps/cli/src/server/assets.d.ts b/apps/cli/src/server/assets.d.ts index 01490304b9..2fbcb5d24e 100644 --- a/apps/cli/src/server/assets.d.ts +++ b/apps/cli/src/server/assets.d.ts @@ -9,3 +9,9 @@ declare module "*.proto" { const content: string export default content } +// `with { type: "file" }` imports are embedded the same way and resolve to a +// path readable with `node:fs`. +declare module "*.wasm" { + const path: string + export default path +} diff --git a/apps/cli/src/server/otlp/ai-stamp.ts b/apps/cli/src/server/otlp/ai-stamp.ts new file mode 100644 index 0000000000..ed0927fd8a --- /dev/null +++ b/apps/cli/src/server/otlp/ai-stamp.ts @@ -0,0 +1,55 @@ +import { Result } from "effect" +import { readFileSync } from "node:fs" +import wasmPath from "./ai-stamp.wasm" with { type: "file" } + +/** + * The ingest gateway's `maple_ai.*` stamping (`apps/ingest/crates/ai-session`), + * compiled to WebAssembly by `bun run build:ai-stamp`. Local ingest runs the + * same code as the gateway, so Agent Sessions reads the same stamps locally. + * See `apps/ingest/crates/ai-session-wasm` for the calling protocol. + */ +interface AiStampExports { + readonly memory: WebAssembly.Memory + readonly input: (length: number) => number + readonly stamp: () => number + readonly output: () => number + readonly output_len: () => number +} + +const module = new WebAssembly.Module(readFileSync(wasmPath)) +// `opentelemetry-proto` depends on `opentelemetry`, which links wasm-bindgen's +// `js-sys` on wasm32. The stamping never reaches it, so its imports are no-ops. +const imports: Record void>> = {} +for (const { module: namespace, name } of WebAssembly.Module.imports(module)) + (imports[namespace] ??= {})[name] = () => {} +const instantiate = () => + // SAFETY: these are the exports of `apps/ingest/crates/ai-session-wasm`, built alongside this file. + new WebAssembly.Instance(module, imports).exports as unknown as AiStampExports +let wasm = instantiate() + +/** Wasm memory never shrinks; past this, the next request gets a fresh instance. */ +const MAX_RETAINED_MEMORY_BYTES = 64 * 1024 * 1024 + +/** Stamp a protobuf `ExportTraceServiceRequest`; fails with the decode error + * when it is not one, which the gateway rejects too. */ +export function stampTraceRequest(request: Uint8Array): Result.Result { + const stamped = Result.try({ + try: () => { + // `input` may grow memory, detaching any earlier view of it. + const at = wasm.input(request.length) + new Uint8Array(wasm.memory.buffer, at, request.length).set(request) + const ok = wasm.stamp() === 1 + const output = wasm.output() + const length = wasm.output_len() + return { ok, bytes: new Uint8Array(wasm.memory.buffer, output, length).slice() } + }, + // A panic is a stamping bug, not bad input, and a trap can leave the + // instance unusable: keep the batch unstamped and start over. + catch: () => { + wasm = instantiate() + return { ok: true, bytes: request } + }, + }).pipe(Result.merge) + if (wasm.memory.buffer.byteLength > MAX_RETAINED_MEMORY_BYTES) wasm = instantiate() + return stamped.ok ? Result.succeed(stamped.bytes) : Result.fail(new TextDecoder().decode(stamped.bytes)) +} diff --git a/apps/cli/src/server/otlp/proto.ts b/apps/cli/src/server/otlp/proto.ts index fa29e3364e..ef5f16fd2f 100644 --- a/apps/cli/src/server/otlp/proto.ts +++ b/apps/cli/src/server/otlp/proto.ts @@ -419,12 +419,46 @@ export function decodeMetricsRequest(bytes: Uint8Array): unknown { return ExportMetricsServiceRequest.toObject(message, toObjectOptions) } -/** Test helper: encode a trace request object to protobuf bytes. */ +/** Encode a trace request object to protobuf bytes. */ export function encodeTraceRequest(obj: unknown): Uint8Array { const message = ExportTraceServiceRequest.fromObject(obj as Record) return ExportTraceServiceRequest.encode(message).finish() } +interface IdsJson { + [field: string]: unknown + links?: IdsJson[] +} +interface TraceRequestJson { + resourceSpans?: { scopeSpans?: { spans?: IdsJson[] }[] }[] +} + +const ID_BYTES = { traceId: 16, spanId: 8, parentSpanId: 8 } +const HEX = /^[0-9a-fA-F]+$/ + +/** OTLP/JSON spells ids in hex where protobufjs reads base64. Only hex at an + * id's exact length converts; anything else is read as base64, and a wrong + * length is rejected by `idHex` after the round trip. */ +function hexIdsToBytes(holder: IdsJson): void { + for (const [field, bytes] of Object.entries(ID_BYTES)) { + const value = holder[field] + if (typeof value === "string" && value.length === bytes * 2 && HEX.test(value)) + holder[field] = Buffer.from(value, "hex") + } +} + +/** Encode a parsed OTLP/JSON trace request to protobuf (for the AI stamping), + * rewriting its ids in place. */ +export function encodeTraceJsonRequest(request: unknown): Uint8Array { + for (const resourceSpans of (request as TraceRequestJson | null)?.resourceSpans ?? []) + for (const scopeSpans of resourceSpans.scopeSpans ?? []) + for (const span of scopeSpans.spans ?? []) { + hexIdsToBytes(span) + for (const link of span.links ?? []) hexIdsToBytes(link) + } + return encodeTraceRequest(request) +} + /** Test helper: encode a metrics request object to protobuf bytes. */ export function encodeMetricsRequest(obj: unknown): Uint8Array { const message = ExportMetricsServiceRequest.fromObject(obj as Record) diff --git a/apps/cli/src/server/serve.ts b/apps/cli/src/server/serve.ts index 886036f93b..e107c90e3b 100644 --- a/apps/cli/src/server/serve.ts +++ b/apps/cli/src/server/serve.ts @@ -31,12 +31,14 @@ import { } from "./eventing/control-store" import { ensureEventConsumerToken, eventConsumerTokenMatches } from "./eventing/consumer-auth" import { LocalEventingRuntime, type LocalEventingRuntimeApi } from "./eventing/runtime" +import { stampTraceRequest } from "./otlp/ai-stamp" import { encodeLogs, encodeMetrics, encodeTraces, type EncodedBatch, OtlpFieldError } from "./otlp/encode" import { decodeLogsRequest, decodeMetricsRequest, decodeTraceRequest, encodeExportResponse, + encodeTraceJsonRequest, } from "./otlp/proto" import { CURRENT_LOCAL_SCHEMA, LOCAL_SCHEMA_SQL, SCHEMA_FINGERPRINT } from "./schema-identity" import { assertCurrentPhysicalSchema } from "./schema-physical" @@ -297,13 +299,15 @@ function decodeOtlp( const decodedBytes = bytes return Result.try({ try: (): unknown => { - if (contentType.includes("json")) - return Schema.decodeUnknownSync(Schema.fromJsonString(Schema.Unknown))( + if (contentType.includes("json")) { + const request = Schema.decodeUnknownSync(Schema.fromJsonString(Schema.Unknown))( new TextDecoder().decode(decodedBytes), ) + return signal === "traces" ? stampedTraces(encodeTraceJsonRequest(request)) : request + } switch (signal) { case "traces": - return decodeTraceRequest(decodedBytes) + return stampedTraces(decodedBytes) case "logs": return decodeLogsRequest(decodedBytes) case "metrics": @@ -314,6 +318,13 @@ function decodeOtlp( }) } +/** Traces carry the ingest gateway's `maple_ai.*` stamps before anything + * reads them. Throws into `decodeOtlp`'s `Result.try`. */ +const stampedTraces = (bytes: Uint8Array): unknown => + decodeTraceRequest( + Result.getOrThrowWith(stampTraceRequest(bytes), (message) => new OtlpDecodeFailed({ message })), + ) + function encodeFor(signal: Signal, req: unknown): EncodedBatch[] { switch (signal) { case "traces": diff --git a/apps/cli/test/ai-stamp.test.ts b/apps/cli/test/ai-stamp.test.ts new file mode 100644 index 0000000000..2d64d85a43 --- /dev/null +++ b/apps/cli/test/ai-stamp.test.ts @@ -0,0 +1,118 @@ +import { deepStrictEqual, ok, strictEqual } from "node:assert" +import { Result } from "effect" +import { describe, it } from "vitest" +import { spanIdHex, traceIdHex } from "../src/server/otlp/encode" +import { encodeTraceRequest } from "../src/server/otlp/proto" +import { __testables } from "../src/server/serve" + +const TRACE_ID = "5b8efff798038103d269b633813fc60c" +const SPAN_ID = "eee19b7ec3c1b174" +const PARENT_ID = "0102030405060708" +const LINKED_TRACE_ID = "0af7651916cd43dd8448eb211c80319c" +// Several MiB, so the wasm memory has to grow mid-request. +const LARGE = "x".repeat(8 * 1024 * 1024) + +const attrs = (pairs: Record) => + Object.entries(pairs).map(([key, stringValue]) => ({ key, value: { stringValue } })) + +// The gateway's own fixture, `vercel_v7_agent_call_and_tool` in +// apps/ingest/crates/ai-session/src/facts.rs, plus a forged stamp. `id` spells +// the ids for the wire format. +const request = (id: (hex: string) => string | Uint8Array) => ({ + resourceSpans: [ + { + scopeSpans: [ + { + scope: { name: "gen_ai" }, + spans: [ + { + traceId: id(TRACE_ID), + spanId: id(SPAN_ID), + parentSpanId: id(PARENT_ID), + links: [{ traceId: id(LINKED_TRACE_ID), spanId: id(PARENT_ID) }], + name: "chat openai/gpt-4o-mini", + startTimeUnixNano: "1700000000000000000", + endTimeUnixNano: "1700000001000000000", + attributes: attrs({ + "gen_ai.operation.name": "chat", + "gen_ai.request.model": "openai/gpt-4o-mini", + "gen_ai.response.model": "openai/gpt-4o-mini-2024-07-18", + "gen_ai.response.id": "gen-1", + "maple_ai.tool_call": "1", + "app.note": LARGE, + }), + }, + ], + }, + ], + }, + ], +}) + +interface DecodedSpan { + traceId: string + spanId: string + parentSpanId: string + links: { traceId: string; spanId: string }[] + attributes: { key: string; value: { stringValue?: string } }[] +} + +const stampsOf = (decoded: unknown) => { + const [resourceSpans] = (decoded as { resourceSpans: { scopeSpans: { spans: DecodedSpan[] }[] }[] }) + .resourceSpans + const span = resourceSpans!.scopeSpans[0]!.spans[0]! + const stamps = span.attributes.filter(({ key }) => key.startsWith("maple_ai.")) + const [link] = span.links + return { + ids: [ + traceIdHex(span.traceId, "traceId"), + spanIdHex(span.spanId, "spanId"), + spanIdHex(span.parentSpanId, "parentSpanId"), + traceIdHex(link!.traceId, "link.traceId"), + spanIdHex(link!.spanId, "link.spanId"), + ], + note: span.attributes.find(({ key }) => key === "app.note")?.value.stringValue?.length, + stamps: Object.fromEntries(stamps.map(({ key, value }) => [key, value.stringValue])), + } +} + +const expected = { + ids: [TRACE_ID, SPAN_ID, PARENT_ID, LINKED_TRACE_ID, PARENT_ID], + note: LARGE.length, + stamps: { + "maple_ai.vendor.id": "vercel_ai_sdk", + "maple_ai.vendor.version": "0", + "maple_ai.llm_call": "1", + "maple_ai.model": "openai/gpt-4o-mini-2024-07-18", + "maple_ai.response.id": "gen-1", + }, +} + +describe("local ingest AI stamping", () => { + it("stamps an OTLP/protobuf request like the gateway", () => { + const bytes = encodeTraceRequest(request((hex) => Buffer.from(hex, "hex"))) + const decoded = __testables.decodeOtlp("traces", bytes, "application/x-protobuf", null) + ok(Result.isSuccess(decoded)) + deepStrictEqual(stampsOf(decoded.success), expected) + }) + + it("stamps an OTLP/JSON request and keeps its hex ids", () => { + const body = new TextEncoder().encode(JSON.stringify(request((hex) => hex))) + const decoded = __testables.decodeOtlp("traces", body, "application/json", null) + ok(Result.isSuccess(decoded)) + deepStrictEqual(stampsOf(decoded.success), expected) + }) + + it("rejects a body the gateway cannot decode", () => { + const decoded = __testables.decodeOtlp( + "traces", + new Uint8Array([0x0a, 0xff]), + "application/x-protobuf", + null, + ) + ok(Result.isFailure(decoded)) + strictEqual(decoded.failure._tag, "@maple/cli/OtlpDecodeFailed") + // prost's error, so the rejection came from the stamping, as in the gateway. + ok(decoded.failure.message.includes("failed to decode Protobuf message"), decoded.failure.message) + }) +}) diff --git a/apps/cli/test/native-checkpoint-smoke.sh b/apps/cli/test/native-checkpoint-smoke.sh index 6cf7950269..e1bd46cf7b 100755 --- a/apps/cli/test/native-checkpoint-smoke.sh +++ b/apps/cli/test/native-checkpoint-smoke.sh @@ -193,6 +193,25 @@ wait "$SERVER_PID" 2>/dev/null || true SERVER_PID="" start_server + +# The binary embeds the ingest gateway's AI stamping: an AI SDK span reaches +# ai_trace_index, which projects stamped spans only. +now_ns="$(date +%s)000000000" +curl --fail-with-body -sS --max-time 30 "http://127.0.0.1:$PORT/v1/traces" \ + -H 'content-type: application/json' \ + --data "$(jq -nc --arg now "$now_ns" '{resourceSpans: [{ + resource: {attributes: [{key: "service.name", value: {stringValue: "ai-stamp-smoke"}}]}, + scopeSpans: [{scope: {name: "gen_ai"}, spans: [{ + traceId: "5b8efff798038103d269b633813fc60c", spanId: "eee19b7ec3c1b174", + name: "chat gpt-4o-mini", startTimeUnixNano: $now, endTimeUnixNano: $now, + attributes: [ + {key: "gen_ai.operation.name", value: {stringValue: "chat"}}, + {key: "gen_ai.request.model", value: {stringValue: "gpt-4o-mini"}} + ]}]}]}]}')" >/dev/null +ai_row="$(query "SELECT VendorId, IsLlmCall FROM ai_trace_index WHERE ServiceName = 'ai-stamp-smoke'" | + jq -r '.[0] | "\(.VendorId) \(.IsLlmCall)"')" +[[ "$ai_row" == "vercel_ai_sdk 1" ]] || fail "AI span not stamped at ingest: ai_trace_index row '$ai_row'" + insert_marker A C1="$(checkpoint)" [[ "$C1" =~ ^[0-9a-f-]{36}$ ]] || fail "invalid C1 ID: $C1" diff --git a/apps/ingest/.cargo/config.toml b/apps/ingest/.cargo/config.toml new file mode 100644 index 0000000000..cffb0d2fb0 --- /dev/null +++ b/apps/ingest/.cargo/config.toml @@ -0,0 +1,5 @@ +# opentelemetry-proto's `trace` feature pulls opentelemetry_sdk -> rand -> +# getrandom, which has no default wasm32-unknown-unknown backend. The AI +# stamping module never draws randomness. +[target.wasm32-unknown-unknown] +rustflags = ['--cfg', 'getrandom_backend="unsupported"'] diff --git a/apps/ingest/Cargo.lock b/apps/ingest/Cargo.lock index 8c01ff6ff3..d2435732c0 100644 --- a/apps/ingest/Cargo.lock +++ b/apps/ingest/Cargo.lock @@ -1563,6 +1563,23 @@ version = "0.1.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "112b39cec0b298b6c1999fee3e31427f74f676e4cb9879ed1a121b43661a4154" +[[package]] +name = "maple-ai-session" +version = "0.1.0" +dependencies = [ + "opentelemetry-proto", + "serde_json", +] + +[[package]] +name = "maple-ai-session-wasm" +version = "0.1.0" +dependencies = [ + "maple-ai-session", + "opentelemetry-proto", + "prost", +] + [[package]] name = "maple-ingest" version = "0.1.0" @@ -1581,6 +1598,7 @@ dependencies = [ "flate2", "hmac 0.12.1", "hotpath", + "maple-ai-session", "moka", "num_cpus", "opentelemetry", diff --git a/apps/ingest/Cargo.toml b/apps/ingest/Cargo.toml index 7e75b968ce..3a85635583 100644 --- a/apps/ingest/Cargo.toml +++ b/apps/ingest/Cargo.toml @@ -4,7 +4,15 @@ version = "0.1.0" edition = "2021" default-run = "maple-ingest" +# crates/ai-session is the AI stamping the gateway runs at decode time; +# crates/ai-session-wasm compiles it for the local CLI (apps/cli) so local +# ingest stamps with the same code (`bun run --cwd apps/cli build:ai-stamp`). +[workspace] +members = ["crates/*"] +default-members = [".", "crates/ai-session"] + [dependencies] +maple-ai-session = { path = "crates/ai-session" } axum = "0.8.8" aes-gcm = "0.10" async-trait = "0.1" @@ -75,13 +83,16 @@ default = [] hotpath = ["hotpath/hotpath"] hotpath-alloc = ["hotpath/hotpath-alloc"] -[lints.rust] +[lints] +workspace = true + +[workspace.lints.rust] unsafe_code = "deny" missing_debug_implementations = "warn" unreachable_pub = "warn" unused_qualifications = "warn" -[lints.clippy] +[workspace.lints.clippy] # groups first, negative priority so specific lints below can override pedantic = { level = "warn", priority = -1 } @@ -132,3 +143,13 @@ harness = false [[bench]] name = "ai_session_bench" harness = false + +# `cargo build --profile wasm --target wasm32-unknown-unknown -p maple-ai-session-wasm`: +# size over speed, and no unwinding in a module the CLI instantiates once. +[profile.wasm] +inherits = "release" +opt-level = "s" +lto = true +codegen-units = 1 +panic = "abort" +strip = true diff --git a/apps/ingest/crates/ai-session-wasm/Cargo.toml b/apps/ingest/crates/ai-session-wasm/Cargo.toml new file mode 100644 index 0000000000..1142d3b971 --- /dev/null +++ b/apps/ingest/crates/ai-session-wasm/Cargo.toml @@ -0,0 +1,15 @@ +[package] +name = "maple-ai-session-wasm" +version = "0.1.0" +edition = "2021" + +[lib] +crate-type = ["cdylib"] + +[dependencies] +maple-ai-session = { path = "../ai-session" } +opentelemetry-proto = { version = "0.32", default-features = false, features = ["gen-tonic-messages", "trace"] } +prost = "0.14.3" + +[lints] +workspace = true diff --git a/apps/ingest/crates/ai-session-wasm/src/lib.rs b/apps/ingest/crates/ai-session-wasm/src/lib.rs new file mode 100644 index 0000000000..74ce1deb98 --- /dev/null +++ b/apps/ingest/crates/ai-session-wasm/src/lib.rs @@ -0,0 +1,56 @@ +//! The gateway's AI stamping as a WebAssembly module, so the local CLI's +//! ingest (`apps/cli/src/server/otlp/ai-stamp.ts`) runs the same code: a +//! protobuf `ExportTraceServiceRequest` in, the stamped request out. +//! +//! The host writes the request into the buffer `input(len)` returns, calls +//! `stamp()`, then reads `output_len()` bytes at `output()`. The pointers are +//! only valid until the next call, which may grow memory. + +use std::cell::RefCell; + +use maple_ai_session::stamp_trace_request; +use opentelemetry_proto::tonic::collector::trace::v1::ExportTraceServiceRequest; +use prost::Message; + +thread_local! { + static INPUT: RefCell> = const { RefCell::new(Vec::new()) }; + static OUTPUT: RefCell> = const { RefCell::new(Vec::new()) }; +} + +#[expect(unsafe_code, reason = "wasm export")] +#[no_mangle] +pub extern "C" fn input(len: usize) -> *mut u8 { + INPUT.with_borrow_mut(|input| { + input.clear(); + input.resize(len, 0); + input.as_mut_ptr() + }) +} + +/// False when the input is not a protobuf `ExportTraceServiceRequest`; the +/// output then holds prost's decode error. +#[expect(unsafe_code, reason = "wasm export")] +#[no_mangle] +pub extern "C" fn stamp() -> bool { + let (stamped, output) = match ExportTraceServiceRequest::decode(INPUT.take().as_slice()) { + Ok(mut request) => { + stamp_trace_request(&mut request); + (true, request.encode_to_vec()) + } + Err(error) => (false, error.to_string().into_bytes()), + }; + OUTPUT.set(output); + stamped +} + +#[expect(unsafe_code, reason = "wasm export")] +#[no_mangle] +pub extern "C" fn output() -> *const u8 { + OUTPUT.with_borrow(Vec::as_ptr) +} + +#[expect(unsafe_code, reason = "wasm export")] +#[no_mangle] +pub extern "C" fn output_len() -> usize { + OUTPUT.with_borrow(Vec::len) +} diff --git a/apps/ingest/crates/ai-session/Cargo.toml b/apps/ingest/crates/ai-session/Cargo.toml new file mode 100644 index 0000000000..7a98e4796b --- /dev/null +++ b/apps/ingest/crates/ai-session/Cargo.toml @@ -0,0 +1,11 @@ +[package] +name = "maple-ai-session" +version = "0.1.0" +edition = "2021" + +[dependencies] +opentelemetry-proto = { version = "0.32", default-features = false, features = ["gen-tonic-messages", "trace"] } +serde_json = "1.0.145" + +[lints] +workspace = true diff --git a/apps/ingest/src/ai_session/claude_code.rs b/apps/ingest/crates/ai-session/src/claude_code.rs similarity index 100% rename from apps/ingest/src/ai_session/claude_code.rs rename to apps/ingest/crates/ai-session/src/claude_code.rs diff --git a/apps/ingest/src/ai_session/facts.rs b/apps/ingest/crates/ai-session/src/facts.rs similarity index 99% rename from apps/ingest/src/ai_session/facts.rs rename to apps/ingest/crates/ai-session/src/facts.rs index 5e907fa7cb..129f6831c9 100644 --- a/apps/ingest/src/ai_session/facts.rs +++ b/apps/ingest/crates/ai-session/src/facts.rs @@ -36,7 +36,7 @@ use opentelemetry_proto::tonic::trace::v1::status::StatusCode; use opentelemetry_proto::tonic::trace::v1::Span; use super::{owned_string_attribute, usage}; -use crate::telemetry::any_value_string; +use crate::value::any_value_string; const TOOL_CALL_ATTR: &str = "maple_ai.tool_call"; const ERROR_ATTR: &str = "maple_ai.error"; @@ -427,7 +427,7 @@ pub(super) fn mark_tool_failed(span: &mut Span) { #[cfg(test)] mod tests { use super::*; - use crate::ai_session::{stamp_trace_request, value_str}; + use crate::{stamp_trace_request, value_str}; use opentelemetry_proto::tonic::collector::trace::v1::ExportTraceServiceRequest; use opentelemetry_proto::tonic::common::v1::InstrumentationScope; use opentelemetry_proto::tonic::resource::v1::Resource; diff --git a/apps/ingest/src/ai_session.rs b/apps/ingest/crates/ai-session/src/lib.rs similarity index 99% rename from apps/ingest/src/ai_session.rs rename to apps/ingest/crates/ai-session/src/lib.rs index 642d364795..68e7fd58df 100644 --- a/apps/ingest/src/ai_session.rs +++ b/apps/ingest/crates/ai-session/src/lib.rs @@ -3,8 +3,10 @@ //! Runs at decode time (`enrich_trace_request`), before the write path forks, //! so the stamps ride inside the OTLP payload itself and reach both the native //! row encoder and the forward-to-collector path — same as the `maple_org_id` -//! resource enrichment. A span that matches a vendor (or, failing that, a -//! generic AI dialect) gets these span attributes appended: +//! resource enrichment. The local CLI's ingest runs the same code, compiled to +//! WebAssembly by `crates/ai-session-wasm`, so this crate stays pure: no I/O, +//! clocks, threads or randomness. A span that matches a vendor (or, failing +//! that, a generic AI dialect) gets these span attributes appended: //! //! - `maple_ai.vendor.id` — vendor slug //! - `maple_ai.vendor.version` — identified vendor version, currently always `"0"` @@ -24,14 +26,14 @@ //! //! One vendor's dialect is restated as well as stamped: Claude Code's native //! keys become the `gen_ai.*` keys every reader keys on, and the phases of its -//! tool calls are left unstamped — see `ai_session/claude_code.rs`. +//! tool calls are left unstamped — see `claude_code.rs`. //! //! Every stamped span also carries the facts Agent Sessions aggregates and //! filters on, decided here once: whether it is a model call //! (`maple_ai.llm_call`) or a tool call, whether it failed, its model, agent //! and tool, and a model call's token usage as five disjoint //! `maple_ai.usage.*` buckets, whatever convention its emitter reported under -//! — see `ai_session/facts.rs` and `ai_session/usage.rs`. +//! — see `facts.rs` and `usage.rs`. //! //! One vendor's evaluations are left unstamped entirely: a Mastra scorer run //! grades a finished agent run and is not a conversation (see @@ -74,6 +76,7 @@ use opentelemetry_proto::tonic::trace::v1::span::Event; mod claude_code; mod facts; mod usage; +pub mod value; pub const ATTR_NAMESPACE: &str = "maple_ai."; pub const VENDOR_ID_ATTR: &str = "maple_ai.vendor.id"; diff --git a/apps/ingest/src/ai_session/usage.rs b/apps/ingest/crates/ai-session/src/usage.rs similarity index 99% rename from apps/ingest/src/ai_session/usage.rs rename to apps/ingest/crates/ai-session/src/usage.rs index a0b8498010..aa79b2fe8e 100644 --- a/apps/ingest/src/ai_session/usage.rs +++ b/apps/ingest/crates/ai-session/src/usage.rs @@ -320,7 +320,7 @@ fn tokens(value: f64) -> u64 { #[cfg(test)] mod tests { use super::*; - use crate::ai_session::{stamp_trace_request, value_str, VENDOR_ID_ATTR}; + use crate::{stamp_trace_request, value_str, VENDOR_ID_ATTR}; use opentelemetry_proto::tonic::collector::trace::v1::ExportTraceServiceRequest; use opentelemetry_proto::tonic::common::v1::{any_value, AnyValue, InstrumentationScope}; use opentelemetry_proto::tonic::resource::v1::Resource; diff --git a/apps/ingest/crates/ai-session/src/value.rs b/apps/ingest/crates/ai-session/src/value.rs new file mode 100644 index 0000000000..428b4ea15d --- /dev/null +++ b/apps/ingest/crates/ai-session/src/value.rs @@ -0,0 +1,65 @@ +//! AnyValue rendering shared by the gateway's row encoder and the AI stamps, +//! so a stamped value reads exactly like the attribute it came from. + +use opentelemetry_proto::tonic::common::v1::{any_value, AnyValue}; +use serde_json::Value; + +pub fn any_value_string(value: &AnyValue) -> String { + match value.value.as_ref() { + Some(any_value::Value::StringValue(value)) => value.clone(), + Some(any_value::Value::BoolValue(value)) => value.to_string(), + Some(any_value::Value::IntValue(value)) => value.to_string(), + Some(any_value::Value::DoubleValue(value)) => value.to_string(), + // Text sent as bytes (LangSmith's prompt and completion) stays + // readable. Binary is hex, including binary that happens to be valid + // UTF-8: a control character other than whitespace marks it. + Some(any_value::Value::BytesValue(value)) => match std::str::from_utf8(value) { + Ok(text) + if !text + .chars() + .any(|c| c.is_control() && !matches!(c, '\t' | '\n' | '\r')) => + { + text.to_owned() + } + _ => bytes_hex(value), + }, + Some(any_value::Value::ArrayValue(_) | any_value::Value::KvlistValue(_)) => { + serde_json::to_string(&any_value_json(value)).unwrap_or_default() + } + // String-table references (OTLP 1.9 experimental encoding) cannot be resolved + // without the sender's dictionary, which the gateway does not accept yet. + Some(any_value::Value::StringValueStrindex(_)) | None => String::new(), + } +} + +pub fn any_value_json(value: &AnyValue) -> Value { + match value.value.as_ref() { + Some(any_value::Value::ArrayValue(array)) => { + Value::Array(array.values.iter().map(any_value_json).collect()) + } + Some(any_value::Value::KvlistValue(kvlist)) => Value::Object( + kvlist + .values + .iter() + .map(|kv| { + let value = kv.value.as_ref().map_or(Value::from(""), any_value_json); + (kv.key.clone(), value) + }) + .collect(), + ), + _ => Value::String(any_value_string(value)), + } +} + +pub fn bytes_hex(bytes: &[u8]) -> String { + if bytes.is_empty() || bytes.iter().all(|byte| *byte == 0) { + return String::new(); + } + const HEX: &[u8; 16] = b"0123456789abcdef"; + let mut out = String::with_capacity(bytes.len() * 2); + for byte in bytes { + out.push(HEX[(byte >> 4) as usize] as char); + out.push(HEX[(byte & 0x0f) as usize] as char); + } + out +} diff --git a/apps/ingest/src/lib.rs b/apps/ingest/src/lib.rs index ceba8a7d6a..ce3e0d2e92 100644 --- a/apps/ingest/src/lib.rs +++ b/apps/ingest/src/lib.rs @@ -1,4 +1,4 @@ -pub mod ai_session; +pub use maple_ai_session as ai_session; pub mod aws; pub mod clickhouse_insert_mappings; pub mod metrics; diff --git a/apps/ingest/src/telemetry.rs b/apps/ingest/src/telemetry.rs index 2ca2762ffd..e8f1847316 100644 --- a/apps/ingest/src/telemetry.rs +++ b/apps/ingest/src/telemetry.rs @@ -18,11 +18,12 @@ use crc32fast::Hasher as Crc32; use dashmap::DashMap; use flate2::write::GzEncoder; use flate2::Compression; +use maple_ai_session::value::{any_value_json, any_value_string, bytes_hex}; use opentelemetry::trace::{SpanContext, TraceContextExt}; use opentelemetry_proto::tonic::collector::logs::v1::ExportLogsServiceRequest; use opentelemetry_proto::tonic::collector::metrics::v1::ExportMetricsServiceRequest; use opentelemetry_proto::tonic::collector::trace::v1::ExportTraceServiceRequest; -use opentelemetry_proto::tonic::common::v1::{any_value, AnyValue, KeyValue}; +use opentelemetry_proto::tonic::common::v1::KeyValue; use opentelemetry_proto::tonic::logs::v1::LogRecord; use opentelemetry_proto::tonic::metrics::v1::{ metric, number_data_point, Exemplar, NumberDataPoint, @@ -4010,56 +4011,9 @@ fn attr_map(attributes: &[KeyValue]) -> Map { out } -pub(crate) fn any_value_string(value: &AnyValue) -> String { - match value.value.as_ref() { - Some(any_value::Value::StringValue(value)) => value.clone(), - Some(any_value::Value::BoolValue(value)) => value.to_string(), - Some(any_value::Value::IntValue(value)) => value.to_string(), - Some(any_value::Value::DoubleValue(value)) => value.to_string(), - // Text sent as bytes (LangSmith's prompt and completion) stays - // readable. Binary is hex, including binary that happens to be valid - // UTF-8: a control character other than whitespace marks it. - Some(any_value::Value::BytesValue(value)) => match std::str::from_utf8(value) { - Ok(text) - if !text - .chars() - .any(|c| c.is_control() && !matches!(c, '\t' | '\n' | '\r')) => - { - text.to_owned() - } - _ => bytes_hex(value), - }, - Some(any_value::Value::ArrayValue(_) | any_value::Value::KvlistValue(_)) => { - serde_json::to_string(&any_value_json(value)).unwrap_or_default() - } - // String-table references (OTLP 1.9 experimental encoding) cannot be resolved - // without the sender's dictionary, which the gateway does not accept yet. - Some(any_value::Value::StringValueStrindex(_)) | None => String::new(), - } -} - /// An array or map as JSON. Scalars keep their string form, as flat arrays and /// maps always have, but a nested array or map stays JSON instead of becoming /// a string of escaped JSON (structured `gen_ai.input.messages`). -fn any_value_json(value: &AnyValue) -> Value { - match value.value.as_ref() { - Some(any_value::Value::ArrayValue(array)) => { - Value::Array(array.values.iter().map(any_value_json).collect()) - } - Some(any_value::Value::KvlistValue(kvlist)) => Value::Object( - kvlist - .values - .iter() - .map(|kv| { - let value = kv.value.as_ref().map_or(Value::from(""), any_value_json); - (kv.key.clone(), value) - }) - .collect(), - ), - _ => Value::String(any_value_string(value)), - } -} - fn span_kind(kind: i32) -> &'static str { match kind { x if x == span::SpanKind::Internal as i32 => "Internal", @@ -4091,19 +4045,6 @@ fn severity_number_to_text(n: i32) -> &'static str { } } -fn bytes_hex(bytes: &[u8]) -> String { - if bytes.is_empty() || bytes.iter().all(|byte| *byte == 0) { - return String::new(); - } - const HEX: &[u8; 16] = b"0123456789abcdef"; - let mut out = String::with_capacity(bytes.len() * 2); - for byte in bytes { - out.push(HEX[(byte >> 4) as usize] as char); - out.push(HEX[(byte & 0x0f) as usize] as char); - } - out -} - fn format_timestamp_nano(unix_nano: u64) -> String { if unix_nano == 0 { return "1970-01-01 00:00:00.000000000".to_owned(); @@ -4154,6 +4095,7 @@ mod tests { use axum::routing::post; use axum::Router; use flate2::read::GzDecoder; + use opentelemetry_proto::tonic::common::v1::{any_value, AnyValue}; use opentelemetry_proto::tonic::common::v1::{ArrayValue, InstrumentationScope, KeyValueList}; use opentelemetry_proto::tonic::metrics::v1::{ exponential_histogram_data_point, metric, AggregationTemporality, ExponentialHistogram, diff --git a/apps/web/src/lib/agent-sessions/vendor-label.ts b/apps/web/src/lib/agent-sessions/vendor-label.ts index 5ee6188c68..aab64b0acb 100644 --- a/apps/web/src/lib/agent-sessions/vendor-label.ts +++ b/apps/web/src/lib/agent-sessions/vendor-label.ts @@ -1,6 +1,6 @@ /** * Vendor ids the ingest gateway stamps (`AI_VENDORS` in - * apps/ingest/src/ai_session.rs) → brand names. Listed here are the ids whose + * apps/ingest/crates/ai-session/src/lib.rs) → brand names. Listed here are the ids whose * brand casing the title-case fallback below can't derive — acronyms (SDK, * ADK), camel brands (LiteLLM, DSPy), and deliberately lowercase ones (eve, * smolagents). diff --git a/docs/local-mode.md b/docs/local-mode.md index 1cbacc3acb..56153dd528 100644 --- a/docs/local-mode.md +++ b/docs/local-mode.md @@ -156,13 +156,14 @@ There is a single binary, `maple`, compiled from **`apps/cli`** (package the server. It talks to the embedded ClickHouse engine **directly via `bun:ffi`**, with no subprocess and no second language in front: -| Concern | Where | How | -| -------------------- | ------------------------------ | -------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | -| CLI commands | `apps/cli/src/commands` | `maple services`, `traces`, `errors`, … run against **either** the local server **or** a remote workspace. `apps/cli/src/core/operations.ts` picks the path per [mode](#local-vs-remote-mode). | -| `maple start` server | `apps/cli/src/server/serve.ts` | A `Bun.serve` hosting OTLP/HTTP ingest (`POST /v1/{traces,logs,metrics}`), the query API (`POST /local/query`), and the bundled SPA, all on one port. | -| Embedded ClickHouse | `apps/cli/src/server/chdb.ts` | `dlopen`s `libchdb` via `bun:ffi` (the `chdb_*` accessor C API) and holds a single connection for the process. | -| OTLP → rows | `apps/cli/src/server/otlp/` | Decodes OTLP protobuf/JSON (protobufjs) and encodes each signal to per-table NDJSON, matching the generated `local-inserts.json` schema exactly. Ported from the production Rust encoders so row shapes can't diverge. | -| UI (SPA) | `apps/local-ui` (Vite + React) | Hooks compile queries with `CH.compile(...)` and POST to `/local/query`. The same build is deployed to `local.maple.dev` (the default) **and** inlined into the binary as the `--offline` fallback (see [release bundle](#release-bundle)); it picks its query base URL from `window.location` at runtime (see [Where the UI comes from](#where-the-ui-comes-from)). | +| Concern | Where | How | +| -------------------- | -------------------------------------- | -------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | +| CLI commands | `apps/cli/src/commands` | `maple services`, `traces`, `errors`, … run against **either** the local server **or** a remote workspace. `apps/cli/src/core/operations.ts` picks the path per [mode](#local-vs-remote-mode). | +| `maple start` server | `apps/cli/src/server/serve.ts` | A `Bun.serve` hosting OTLP/HTTP ingest (`POST /v1/{traces,logs,metrics}`), the query API (`POST /local/query`), and the bundled SPA, all on one port. | +| Embedded ClickHouse | `apps/cli/src/server/chdb.ts` | `dlopen`s `libchdb` via `bun:ffi` (the `chdb_*` accessor C API) and holds a single connection for the process. | +| OTLP → rows | `apps/cli/src/server/otlp/` | Decodes OTLP protobuf/JSON (protobufjs) and encodes each signal to per-table NDJSON, matching the generated `local-inserts.json` schema exactly. Ported from the production Rust encoders so row shapes can't diverge. | +| AI stamping | `apps/cli/src/server/otlp/ai-stamp.ts` | Runs the ingest gateway's `maple_ai.*` stamping (`apps/ingest/crates/ai-session`, compiled to wasm and embedded) on every trace request before it is encoded, so Agent Sessions reads the same stamps as in the cloud. | +| UI (SPA) | `apps/local-ui` (Vite + React) | Hooks compile queries with `CH.compile(...)` and POST to `/local/query`. The same build is deployed to `local.maple.dev` (the default) **and** inlined into the binary as the `--offline` fallback (see [release bundle](#release-bundle)); it picks its query base URL from `window.location` at runtime (see [Where the UI comes from](#where-the-ui-comes-from)). | chDB allows exactly one connection per process and isn't safe to call concurrently. The long-lived `maple start` process owns the connection, and @@ -422,9 +423,13 @@ their matching loopback address instead of the unusable `0.0.0.0` or `::`. ## Dev workflow -No Rust toolchain needed. Run the server and the SPA dev server in two terminals: +Run the server and the SPA dev server in two terminals. The server embeds the +gateway's AI stamping as wasm, built with the mise Rust toolchain; rebuild it +after changing `apps/ingest/crates/ai-session`: ```bash +bun run --cwd apps/cli build:ai-stamp + # Terminal 1: the server (OTLP ingest + query API + chDB) on :4318. # Needs libchdb: set MAPLE_LIBCHDB, or keep libchdb.so in ~/.maple/bin. bun run apps/cli/src/bin.ts start diff --git a/docs/otel-spec/resource-and-config.md b/docs/otel-spec/resource-and-config.md index 852c1d5419..982c57d597 100644 --- a/docs/otel-spec/resource-and-config.md +++ b/docs/otel-spec/resource-and-config.md @@ -6,7 +6,7 @@ A compliance reference for the OpenTelemetry **Resource** model (what a resource - Maple is mainly a **consumer/backend** of OTel data (Rust ingest gateway `apps/ingest` → Tinybird/ClickHouse, or a collector in forward mode → web dashboard). We also **self-instrument** our own services (`apps/ingest`, the API, etc.). So we must parse resources correctly from arbitrary SDKs _and_ emit spec-correct resources ourselves. - **Dashboards key off `service.name` + `deployment.environment.name`.** The legacy `deployment.environment` attribute is deprecated (replaced by `deployment.environment.name`). The ingest gateway's own self-telemetry resource dual-emits both keys today (`apps/ingest/src/otel.rs`). Every materialized view coalesces the two (`DEPLOYMENT_ENV_SQL` in `packages/domain/src/tinybird/semconv-renames.ts`, ClickHouse migration `packages/domain/src/clickhouse/migrations/0020_semconv_key_renames.ts`), so telemetry carrying either key resolves to the same environment. The deprecation citation below is the spec basis for dropping the legacy emission. That now waits only on pre-0020 rows aging out and on BYO-ClickHouse schemas catching up. -- The ingest gateway stamps a Maple-internal `maple_org_id` resource attribute (and `maple_ai.*` span attributes on AI agent spans, see `apps/ingest/src/ai_session.rs`). These are intentionally outside semconv and must never collide with a real OTel attribute name. The dotted `maple.*` namespace belongs to Maple's SDKs, not the gateway. +- The ingest gateway stamps a Maple-internal `maple_org_id` resource attribute (and `maple_ai.*` span attributes on AI agent spans, see `apps/ingest/crates/ai-session/src/lib.rs`). These are intentionally outside semconv and must never collide with a real OTel attribute name. The dotted `maple.*` namespace belongs to Maple's SDKs, not the gateway. - We do not consume declarative/file configuration or the Entities data model yet. Both are tracked because they affect how upstream SDKs will emit resources (Entities) and how our self-instrumentation could be configured (file config). --- diff --git a/mise.toml b/mise.toml index ef58e6e517..7e0a899c4e 100644 --- a/mise.toml +++ b/mise.toml @@ -16,7 +16,9 @@ [tools] bun = "1.4.0" # keep in sync with `packageManager` in package.json node = "24.18.0" # some scripts shell out to `node` -rust = "1.94.1" # apps/ingest (Rust gateway); was a floating `stable` in CI +# apps/ingest (Rust gateway); wasm32 builds its AI stamping for the local CLI +# (`bun run --cwd apps/cli build:ai-stamp`). Was a floating `stable` in CI. +rust = { version = "1.94.1", targets = ["wasm32-unknown-unknown"] } python = "3.12.13" # Tinybird tooling + .github/scripts/format-ingest-benchmark-comment.py # go = "1.x.y" # otel-collector exporter pins via go.mod; add here only if we want it unified diff --git a/packages/backend/src/services/warehouse/clickhouse-e2e-support.ts b/packages/backend/src/services/warehouse/clickhouse-e2e-support.ts index 09387fdfe7..2e01613d59 100644 --- a/packages/backend/src/services/warehouse/clickhouse-e2e-support.ts +++ b/packages/backend/src/services/warehouse/clickhouse-e2e-support.ts @@ -214,7 +214,7 @@ export const looksLikeIdentityColumn = (name: string): boolean => /** * What the ingest gateway stamps on a span it classifies - * (`apps/ingest/src/ai_session/facts.rs`, `usage.rs`), spelled out per seed as + * (`apps/ingest/crates/ai-session/src/facts.rs`, `usage.rs`), spelled out per seed as * the gateway would have written it: `ai_trace_index_mv` reads these and * nothing else, so a seed without them materializes as neither a call nor a * tool. Seeds keep their dialect attributes so detail reads see whole spans. diff --git a/packages/domain/src/gen-ai.test.ts b/packages/domain/src/gen-ai.test.ts index 4f9b5d1b41..ba280cba05 100644 --- a/packages/domain/src/gen-ai.test.ts +++ b/packages/domain/src/gen-ai.test.ts @@ -8,7 +8,7 @@ import { MAPLE_AI_STAMP_ATTRS } from "./gen-ai" // literals are pinned against the Rust sources that write them. const gatewaySource = ["facts.rs", "usage.rs"] .map((file) => - readFileSync(new URL(`../../../apps/ingest/src/ai_session/${file}`, import.meta.url), "utf8"), + readFileSync(new URL(`../../../apps/ingest/crates/ai-session/src/${file}`, import.meta.url), "utf8"), ) .join("\n") diff --git a/packages/domain/src/gen-ai.ts b/packages/domain/src/gen-ai.ts index ce131d4295..8fedc49443 100644 --- a/packages/domain/src/gen-ai.ts +++ b/packages/domain/src/gen-ai.ts @@ -70,7 +70,7 @@ export const traceSessionTraceId = (sessionId: string): string | undefined => { // The gateway strips `maple_ai.*` on the way in and re-stamps its own verdict, // with the exceptions it does not own and lets through: the turn id, the // dropped-messages counter and the model-duration clock (`PRESERVED_ATTRS` in -// `apps/ingest/src/ai_session.rs`). +// `apps/ingest/crates/ai-session/src/lib.rs`). /** * Groups every span of one conversation/investigation into one agent session. @@ -105,7 +105,7 @@ export const AI_AGENT_OPERATIONS = [ ] as const /** The convention's memory-store operations: the agent's own bookkeeping, never * a model turn or a tool call. The ingest gateway decides this for the spans it - * stamps (`KNOWN_OPS` in `apps/ingest/src/ai_session/usage.rs`); the span + * stamps (`KNOWN_OPS` in `apps/ingest/crates/ai-session/src/usage.rs`); the span * classifier reads it for the spans ingested before. */ export const AI_MEMORY_OPERATIONS = [ "search_memory", @@ -134,7 +134,7 @@ export const MAPLE_GENAI_MODEL_DURATION_MS_ATTR = "maple_ai.model_duration_ms" /** * Every fact Agent Sessions aggregates or filters on, as the ingest gateway - * decided it for one span (`apps/ingest/src/ai_session/facts.rs`, `usage.rs`). + * decided it for one span (`apps/ingest/crates/ai-session/src/facts.rs`, `usage.rs`). * `ai_trace_index_mv` projects these and holds no vendor rule, and the detail * page reads the same ones, so the list and the page cannot disagree. * diff --git a/packages/domain/src/tinybird/gen-ai-columns.ts b/packages/domain/src/tinybird/gen-ai-columns.ts index 8a47670f88..c9efad47ce 100644 --- a/packages/domain/src/tinybird/gen-ai-columns.ts +++ b/packages/domain/src/tinybird/gen-ai-columns.ts @@ -2,7 +2,7 @@ // // Every fact the index aggregates or filters on is decided by the ingest // gateway, once per span, and written onto the span as a `maple_ai.*` stamp -// (`MAPLE_AI_STAMP_ATTRS`, decided in `apps/ingest/src/ai_session/facts.rs` +// (`MAPLE_AI_STAMP_ATTRS`, decided in `apps/ingest/crates/ai-session/src/facts.rs` // and `usage.rs`): whether the span is a model call or a tool call, whether it // failed, its model, agent and tool, and a model call's usage as five disjoint // buckets. So the view is a projection of those stamps plus generic OTel diff --git a/packages/query-engine-integrations/src/ai/ai-integrations.ts b/packages/query-engine-integrations/src/ai/ai-integrations.ts index 9e676901f0..8036a31bc1 100644 --- a/packages/query-engine-integrations/src/ai/ai-integrations.ts +++ b/packages/query-engine-integrations/src/ai/ai-integrations.ts @@ -9,7 +9,7 @@ // a different value. Vendor entries are keyed on the `maple_ai.vendor.id` the // ingest gateway stamped at decode time, from evidence the read path no longer // has (instrumentation scope, resource SDK name, span events — see -// `apps/ingest/src/ai_session.rs`). Attributes arrive as `Map(String, String)`, +// `apps/ingest/crates/ai-session/src/lib.rs`). Attributes arrive as `Map(String, String)`, // so a missing key reads back as `''` and an undecodable value yields no field. import { Effect, Option } from "effect" diff --git a/packages/query-engine-integrations/src/ai/ai-sessions.ts b/packages/query-engine-integrations/src/ai/ai-sessions.ts index de405e0a9b..f666f1f904 100644 --- a/packages/query-engine-integrations/src/ai/ai-sessions.ts +++ b/packages/query-engine-integrations/src/ai/ai-sessions.ts @@ -1,7 +1,7 @@ // AI agent sessions — read side // // The ingest gateway stamps three attributes on AI-agent spans at decode time -// (`apps/ingest/src/ai_session.rs`): `maple_ai.vendor.id`, +// (`apps/ingest/crates/ai-session/src/lib.rs`): `maple_ai.vendor.id`, // `maple_ai.vendor.version` and `maple_ai.session.id`. Only the last one is // sparse — a vendor exposes a session key on the spans that own the turn // (`ai.eve.turn`, `invoke_agent`), never on the sibling `chat`, `execute_tool`, diff --git a/scripts/build-local-binary.sh b/scripts/build-local-binary.sh index c71650d101..4a4dcfc4bc 100755 --- a/scripts/build-local-binary.sh +++ b/scripts/build-local-binary.sh @@ -1,12 +1,14 @@ #!/usr/bin/env bash # Build the distributable `maple` local binary — a single Bun-compiled -# executable plus libchdb. No Rust/cargo involved. +# executable plus libchdb. Cargo only builds the embedded AI stamping wasm. # # Pipeline: # 1. Build the workspace packages consumed through gitignored `dist/` paths. # 2. Build the lightweight SPA (`apps/local-ui` → its `dist/`). # 3. Inline that dist into apps/cli/src/server/ui-embed.gen.ts so # `bun build --compile` bakes the SPA into the binary. +# 3b. Build the ingest gateway's AI stamping to wasm (apps/ingest/crates/ +# ai-session-wasm) for the binary to embed; needs the mise rust toolchain. # 4. Compile apps/cli (the CLI + the OTLP-ingest/query server) into a single # executable with `bun build --compile`. The schema artifacts and SPA are # embedded; the OTLP encoders run in-process; chDB is reached via bun:ffi. @@ -48,6 +50,9 @@ restore_stub() { git -C "$REPO_ROOT" checkout -- "$UI_EMBED" 2>/dev/null || true trap restore_stub EXIT bun run "$REPO_ROOT/scripts/gen-ui-embed.ts" +echo "==> Building the AI stamping wasm" +bun run --cwd "$REPO_ROOT/apps/cli" build:ai-stamp + echo "==> Compiling maple binary (bun build --compile) — version $MAPLE_BUILD_VERSION" ( cd "$REPO_ROOT" && bun build apps/cli/src/bin.ts --compile \ --define "__MAPLE_VERSION__=\"$MAPLE_BUILD_VERSION\"" \ diff --git a/turbo.json b/turbo.json index 80cd565e2d..1f98719da8 100644 --- a/turbo.json +++ b/turbo.json @@ -71,7 +71,18 @@ "@maple/domain#test": { "dependsOn": ["^build"], "outputs": [], - "inputs": ["$TURBO_DEFAULT$", "$TURBO_ROOT$/apps/ingest/src/ai_session/*.rs"] + "inputs": ["$TURBO_DEFAULT$", "$TURBO_ROOT$/apps/ingest/crates/ai-session/src/*.rs"] + }, + // The test script builds the ingest gateway's AI stamping to wasm first. + "@maple/cli#test": { + "dependsOn": ["^build"], + "outputs": [], + "inputs": [ + "$TURBO_DEFAULT$", + "$TURBO_ROOT$/apps/ingest/crates/**", + "$TURBO_ROOT$/apps/ingest/Cargo.*", + "$TURBO_ROOT$/apps/ingest/.cargo/**" + ] }, "eval": { "dependsOn": ["^build"],