Skip to content

Commit 9c3a38e

Browse files
Add Langfuse ingestion and sync benchmark (#6)
1 parent 09684a2 commit 9c3a38e

13 files changed

Lines changed: 843 additions & 3 deletions

File tree

.changeset/brave-traces-sync.md

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
1+
---
2+
"agentpond": patch
3+
---
4+
5+
Add a manual Langfuse ingestion and DuckDB sync performance benchmark.

package.json

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -21,7 +21,8 @@
2121
"version-packages": "changeset version",
2222
"release": "changeset publish",
2323
"ingest": "tsx apps/ingest/src/index.ts",
24-
"cli": "tsx apps/cli/src/index.ts"
24+
"cli": "tsx apps/cli/src/index.ts",
25+
"perf:sync": "pnpm --filter @agentpond/perf sync"
2526
},
2627
"dependencies": {
2728
"@aws-sdk/client-s3": "^3.936.0",

packages/duckdb/src/cache.ts

Lines changed: 66 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,23 @@ export type SyncResult = {
1515
eventsProcessed: number;
1616
};
1717

18+
export type SyncProgress = SyncResult & {
19+
manifestsTotal: number;
20+
manifestsSeen: number;
21+
manifestsSkipped: number;
22+
objectsSkipped: number;
23+
phase:
24+
| "listed"
25+
| "manifest-skipped"
26+
| "manifest-processed"
27+
| "object-skipped"
28+
| "object-processed"
29+
| "events-processed"
30+
| "complete";
31+
currentManifestKey?: string;
32+
currentObjectKey?: string;
33+
};
34+
1835
export class AgentPondDuckDb {
1936
private instance?: DuckDBInstance;
2037
private connection?: DuckDBConnection;
@@ -115,22 +132,55 @@ export class AgentPondDuckDb {
115132
store: ObjectStore;
116133
projectId: string;
117134
prefix: string;
135+
onProgress?: (progress: SyncProgress) => void;
118136
}): Promise<SyncResult> {
119137
await this.init();
120138
const result: SyncResult = {
121139
manifestsProcessed: 0,
122140
objectsProcessed: 0,
123141
eventsProcessed: 0,
124142
};
143+
let manifestsSeen = 0;
144+
let manifestsSkipped = 0;
145+
let objectsSkipped = 0;
125146
const manifestKeys = await params.store.listKeys(
126147
manifestPrefix(params.prefix, params.projectId),
127148
);
149+
const emitProgress = (
150+
phase: SyncProgress["phase"],
151+
current?: Pick<SyncProgress, "currentManifestKey" | "currentObjectKey">,
152+
) => {
153+
params.onProgress?.({
154+
...result,
155+
manifestsTotal: manifestKeys.length,
156+
manifestsSeen,
157+
manifestsSkipped,
158+
objectsSkipped,
159+
phase,
160+
...current,
161+
});
162+
};
163+
emitProgress("listed");
128164

129165
for (const manifestKey of manifestKeys) {
130-
if (await this.exists("processed_manifests", manifestKey)) continue;
166+
manifestsSeen += 1;
167+
if (await this.exists("processed_manifests", manifestKey)) {
168+
manifestsSkipped += 1;
169+
emitProgress("manifest-skipped", {
170+
currentManifestKey: manifestKey,
171+
});
172+
continue;
173+
}
131174
const manifest = await params.store.getJson<BatchManifest>(manifestKey);
132175
for (const object of manifest.objects) {
133-
if (await this.exists("processed_objects", object.key)) continue;
176+
if (await this.exists("processed_objects", object.key)) {
177+
objectsSkipped += 1;
178+
emitProgress("object-skipped", {
179+
currentManifestKey: manifestKey,
180+
currentObjectKey: object.key,
181+
});
182+
continue;
183+
}
134184
const events =
135185
object.entityType === "otel"
136186
? otelResourceSpansToEvents(
@@ -152,13 +202,27 @@ export class AgentPondDuckDb {
152202
});
153203
await this.projectEvent(manifest.projectId, event);
154204
result.eventsProcessed += 1;
205+
if (result.eventsProcessed % 1000 === 0) {
206+
emitProgress("events-processed", {
207+
currentManifestKey: manifestKey,
208+
currentObjectKey: object.key,
209+
});
210+
}
155211
}
156212
await this.insertKey("processed_objects", object.key, manifestKey);
157213
result.objectsProcessed += 1;
214+
emitProgress("object-processed", {
215+
currentManifestKey: manifestKey,
216+
currentObjectKey: object.key,
217+
});
158218
}
159219
await this.insertKey("processed_manifests", manifestKey);
160220
result.manifestsProcessed += 1;
221+
emitProgress("manifest-processed", {
222+
currentManifestKey: manifestKey,
223+
});
161224
}
225+
emitProgress("complete");
162226

163227
return result;
164228
}

packages/perf/package.json

Lines changed: 21 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,21 @@
1+
{
2+
"name": "@agentpond/perf",
3+
"version": "0.1.0",
4+
"type": "module",
5+
"private": true,
6+
"scripts": {
7+
"sync": "tsx src/index.ts"
8+
},
9+
"dependencies": {
10+
"@agentpond/core": "workspace:*",
11+
"@agentpond/duckdb": "workspace:*",
12+
"@langfuse/otel": "^5.4.1",
13+
"@langfuse/tracing": "^5.4.1",
14+
"@opentelemetry/sdk-node": "^0.208.0"
15+
},
16+
"devDependencies": {
17+
"@types/node": "^24.10.1",
18+
"tsx": "^4.21.0",
19+
"typescript": "^5.9.3"
20+
}
21+
}

packages/perf/src/args.ts

Lines changed: 136 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,136 @@
1+
import { randomUUID } from "node:crypto";
2+
import { mkdtempSync } from "node:fs";
3+
import { tmpdir } from "node:os";
4+
import { join } from "node:path";
5+
import type { AgentPondConfig, S3ObjectStore } from "@agentpond/core";
6+
7+
export type PerfArgs = {
8+
traces: number;
9+
endpoint: string;
10+
bucket: string;
11+
accessKeyId: string;
12+
secretAccessKey: string;
13+
region: string;
14+
projectId: string;
15+
publicKey: string;
16+
secretKey: string;
17+
prefix: string;
18+
dbPath: string;
19+
};
20+
21+
const DEFAULT_TRACES = 100_000;
22+
23+
export function parseArgs(argv: string[]): PerfArgs {
24+
const values = parseFlags(argv.filter((arg) => arg !== "--"));
25+
const runId = values["run-id"] ?? randomUUID();
26+
const dbPath =
27+
values.db ??
28+
join(mkdtempSync(join(tmpdir(), "agentpond-perf-")), "cache.duckdb");
29+
30+
const traces = integerFlag(values, "traces", DEFAULT_TRACES);
31+
if (traces < 1) throw new Error("--traces must be at least 1");
32+
33+
return {
34+
traces,
35+
endpoint: values.endpoint ?? "http://localhost:9000",
36+
bucket: values.bucket ?? "agentpond",
37+
accessKeyId: values["access-key-id"] ?? "minio",
38+
secretAccessKey: values["secret-access-key"] ?? "minio123",
39+
region: values.region ?? "us-east-1",
40+
projectId: values["project-id"] ?? "default-project",
41+
publicKey: values["public-key"] ?? "pk-agentpond",
42+
secretKey: values["secret-key"] ?? "sk-agentpond",
43+
prefix: normalizePrefix(values.prefix ?? `perf/${runId}`),
44+
dbPath,
45+
};
46+
}
47+
48+
export function buildConfig(args: PerfArgs): AgentPondConfig {
49+
return {
50+
projectId: args.projectId,
51+
dbPath: args.dbPath,
52+
s3: {
53+
bucket: args.bucket,
54+
prefix: args.prefix,
55+
endpoint: args.endpoint,
56+
region: args.region,
57+
accessKeyId: args.accessKeyId,
58+
secretAccessKey: args.secretAccessKey,
59+
forcePathStyle: true,
60+
},
61+
auth: {
62+
projectId: args.projectId,
63+
publicKey: args.publicKey,
64+
secretKey: args.secretKey,
65+
},
66+
};
67+
}
68+
69+
export function configureLangfuseEnv(address: string, args: PerfArgs): void {
70+
process.env.LANGFUSE_BASE_URL = address;
71+
process.env.LANGFUSE_PUBLIC_KEY = args.publicKey;
72+
process.env.LANGFUSE_SECRET_KEY = args.secretKey;
73+
process.env.LANGFUSE_RELEASE = "agentpond-perf";
74+
process.env.LANGFUSE_ENVIRONMENT = "performance";
75+
}
76+
77+
export async function assertEmptyPrefix(
78+
store: S3ObjectStore,
79+
prefix: string,
80+
): Promise<void> {
81+
try {
82+
const keys = await store.listKeys(prefix);
83+
if (keys.length > 0) {
84+
throw new Error(
85+
`S3 prefix ${prefix} is not empty (${keys.length} existing objects)`,
86+
);
87+
}
88+
} catch (error) {
89+
if (error instanceof Error && error.message.includes("is not empty")) {
90+
throw error;
91+
}
92+
const message = error instanceof Error ? error.message : String(error);
93+
throw new Error(
94+
`Could not access local S3 storage. Start MinIO and create the bucket first: docker compose up -d minio create-bucket. Cause: ${message}`,
95+
);
96+
}
97+
}
98+
99+
export function runIdFromPrefix(prefix: string): string {
100+
const trimmed = prefix.replace(/\/$/, "");
101+
return trimmed.slice(trimmed.lastIndexOf("/") + 1);
102+
}
103+
104+
function parseFlags(argv: string[]): Record<string, string> {
105+
const flags: Record<string, string> = {};
106+
for (let i = 0; i < argv.length; i += 1) {
107+
const arg = argv[i];
108+
if (!arg.startsWith("--")) throw new Error(`Unexpected argument: ${arg}`);
109+
const [rawKey, inlineValue] = arg.slice(2).split("=", 2);
110+
const value = inlineValue ?? argv[i + 1];
111+
if (!value || value.startsWith("--")) {
112+
throw new Error(`Missing value for --${rawKey}`);
113+
}
114+
flags[rawKey] = value;
115+
if (inlineValue === undefined) i += 1;
116+
}
117+
return flags;
118+
}
119+
120+
function integerFlag(
121+
flags: Record<string, string>,
122+
name: string,
123+
defaultValue: number,
124+
): number {
125+
const raw = flags[name];
126+
if (!raw) return defaultValue;
127+
const value = Number.parseInt(raw, 10);
128+
if (!Number.isFinite(value) || String(value) !== raw) {
129+
throw new Error(`--${name} must be an integer`);
130+
}
131+
return value;
132+
}
133+
134+
function normalizePrefix(prefix: string): string {
135+
return prefix.endsWith("/") ? prefix : `${prefix}/`;
136+
}

0 commit comments

Comments
 (0)