Skip to content
Open
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
2 changes: 1 addition & 1 deletion .github/workflows/build-ingest-binary.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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 }}-
Expand Down
11 changes: 8 additions & 3 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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/**'
Expand Down Expand Up @@ -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, ' ') }}
Expand Down
3 changes: 3 additions & 0 deletions .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -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
5 changes: 3 additions & 2 deletions apps/cli/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -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": {
Expand Down
6 changes: 6 additions & 0 deletions apps/cli/src/server/assets.d.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
}
55 changes: 55 additions & 0 deletions apps/cli/src/server/otlp/ai-stamp.ts
Original file line number Diff line number Diff line change
@@ -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<string, Record<string, () => 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<Uint8Array, string> {
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))
}
36 changes: 35 additions & 1 deletion apps/cli/src/server/otlp/proto.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<string, unknown>)
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<string, unknown>)
Expand Down
17 changes: 14 additions & 3 deletions apps/cli/src/server/serve.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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":
Expand All @@ -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":
Expand Down
118 changes: 118 additions & 0 deletions apps/cli/test/ai-stamp.test.ts
Original file line number Diff line number Diff line change
@@ -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<string, string>) =>
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)
})
})
19 changes: 19 additions & 0 deletions apps/cli/test/native-checkpoint-smoke.sh
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down
5 changes: 5 additions & 0 deletions apps/ingest/.cargo/config.toml
Original file line number Diff line number Diff line change
@@ -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"']
18 changes: 18 additions & 0 deletions apps/ingest/Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

Loading
Loading