@openaudr/audr-sink-lago
v0.1.0
Published
Lago sink for AUDR: deliver usage records to Lago usage-based billing.
Maintainers
Readme
@openaudr/audr-sink-lago
The Lago sink for AUDR delivers record batches from
an @openaudr/audr Client to
Lago's batch event endpoint, for usage-based billing. Each record
becomes one Lago event, routed on attribution.subscription_id and de-duplicated on
record_id. ESM with full type declarations; no runtime dependencies beyond
@openaudr/audr.
Status: alpha. The record model tracks AUDR v1.0.0; until 1.0.0, a minor release may change the public API.
Setup
npm install @openaudr/audr @openaudr/audr-sink-lagoRequires Node.js 22.12 or later.
Usage
import { Client, createRecord } from '@openaudr/audr';
import { LagoSink } from '@openaudr/audr-sink-lago';
const record = createRecord({
resource: {
provider: 'anthropic',
type: 'model',
name: 'claude-sonnet-5',
operation: 'generation',
modality: 'text',
},
usage: { llm: { input_tokens: 1200, output_tokens: 340, requests: 1 } },
run: { run_id: '01J8ZQ8Y2K3M4N5P6Q7R8S9T0V', span_id: 'turn-3', run_type: 'agent_run' },
attribution: { environment: 'production', account_id: 'acct_42', subscription_id: 'sub_42' },
});
const sink = new LagoSink(); // reads LAGO_API_KEY and LAGO_METRIC_CODE
const client = new Client(sink, {
emitter: { component: 'harness', name: 'my-harness', version: '1.4.0' },
});
const result = client.record(record);
if (!result.queued) console.warn('record rejected', result.issues);
await client.shutdown(); // drains the queue, then closes the sinkThe Client stamps its emitter onto every record that arrives without one; a record that
reaches validation with no emitter is rejected with an /emitter issue and never sent.
Record construction is documented in the
core SDK README.
A runnable version against a stand-in endpoint, which
chooses the metric code per record, is
examples/deliver.ts.
[!IMPORTANT] Lago routes on
attribution.subscription_id. A record without one is never sent:deliver()returns it rejected with detailmissing_subscription_id. A record that corrects another (corrects) is never sent either: it is rejected with detailunsupported_correction, because Lago cannot restate an event.
Configuration
Configure the sink with apiKey and metricCode. Pass them explicitly or set
LAGO_API_KEY and LAGO_METRIC_CODE; explicit values take precedence, and both are always
required. Credentials belong to the sink and never appear in logs, errors, String(sink) or
util.inspect(sink). An invalid configuration throws ConfigurationError from the
constructor.
metricCode is the code of the Lago billable metric an event is billed against: one code
for every record, or a function that chooses the code for each record.
LAGO_METRIC_CODE can only name one code. It has no default, because Lago accepts an event
whose code matches no active metric and then skips it without an error.
apiUrl, or LAGO_API_URL, selects the Lago instance. It defaults to Lago Cloud in the US
region, https://api.getlago.com. Use https://api.eu.getlago.com for the EU region, or the
address of a self-hosted instance. The value is an origin, an origin
followed by /api/v1, or the full /api/v1/events/batch URL; all three name one endpoint.
Delivery
| AUDR field | Lago event |
| --- | --- |
| attribution.subscription_id | external_subscription_id |
| record_id | transaction_id |
| timing.event_time | timestamp, in Unix seconds with millisecond precision |
| metricCode, or the code it chooses for the record | code |
| Every field, flattened | properties, such as usage__llm__input_tokens |
Lago accepts only string and number property values, so every field is flattened to a
reversible name joined by __. Arrays and attribution.labels are sent as canonical JSON
under a __json suffix, such as attribution__labels__json. The separator is fixed,
because billable metric definitions refer to these names. flattenRecord(record) returns the
properties a record becomes before it is sent.
Billable metric
Create the Lago billable metric with the code passed as metricCode, and choose its
aggregation field from the flattened names. A metric that sums usage__llm__input_tokens,
for example, bills the input tokens of every record.
A Lago metric aggregates one field, so records that meter different things, such as model tokens and tool calls, belong to different metrics. Pass a function to choose the code from the record:
const METRICS: Record<string, string> = {
generation: 'llm_tokens',
embedding: 'embedding_tokens',
tool_execution: 'tool_calls',
};
const sink = new LagoSink({ metricCode: (record) => METRICS[record.resource.operation] });One sink serves every metric. The function runs for each record, so records with different codes are delivered in the same batch, and even the same request. A separate sink is needed only for another Lago instance or API key.
A record for which the function returns undefined, or an invalid code, is rejected with
detail missing_metric_code or invalid_metric_code and never sent. The function must
depend only on the record, so that a replay chooses the same code: on a ClickHouse event
store the code is part of the de-duplication key, and the same record under another code is
billed again.
One record becomes one event under one metric code. Lago de-duplicates on transaction_id,
so the same record sent under a second code is dropped as a duplicate on a Postgres event
store.
Results
A batch is sent as sequential requests of at most 100 events, and the responses combine into
one BatchResult. Lago validates a request as a unit: one invalid event fails the whole
request with a 422 that names the offending events. The sink rejects the named records and
resends the remainder. A record that Lago already holds counts as delivered. The
reference
maps each response status to an outcome.
Transient failures are retried inside deliver(), a bounded number of times, with the
identical request body. A retry therefore carries the same transaction_id and timestamp
as the first attempt, which keeps it idempotent on both Lago event stores.
A Lago 200 confirms that the event was ingested, not that it was billed. Lago ingests and
then skips an event whose subscription does not exist or has terminated, whose metric code
matches no active metric, or whose timestamp precedes the subscription start. Check the
metric code and the subscription identifiers against Lago when usage does not appear on an
invoice.
This sink forwards every field of the record, including attribution.labels and any
x_* extension, to your Lago instance verbatim. attribution.labels MUST NOT contain PII and
resource.key_name is a label, never key material; the
security policy holds the full
rules. Enforce them where the record is built, the last point at which they can be enforced.
Self-hosted Lago
Set apiUrl to the origin of the instance, for example https://lago.internal:8443. The
sink requires HTTPS. An instance on a trusted network that serves plain HTTP is reachable
only with allowInsecureHttp: true, which sends the API key unencrypted. The sink rejects
a URL that carries userinfo, a query string, a fragment or an unrelated path, and it never
follows a redirect.
An instance that sets LAGO_EVENTS_BATCH_MAX_LENGTH below 100 refuses larger requests with
detail too_many_events. Set the client's batchMaxSize to that limit or lower.
Shutdown
client.shutdown() drains the queue and closes the sink. sink.close() is idempotent and
aborts any request in flight, and later deliveries answer closed. A batch that closing
interrupts answers closed; when earlier requests of that batch had already settled
records, it answers accepted and lists the remaining records as unknown. Delivering an
unknown record again is safe, because Lago de-duplicates on record_id.
Documentation
- Reference: options, event mapping, responses, validation handling, idempotency, retries, diagnostics
- Examples: runnable against a stand-in endpoint, without network access
- Changelog
- AUDR specification, which defines every record field
- Lago batch endpoint, which this sink calls
License
Apache-2.0. Contributions follow CONTRIBUTING.md.
