Filters in depth
Edge cases, batch semantics, progress, revisions, errors, permissions, limits, and payload formats for consumer group filters
The full reference for consumer group filters. For a short introduction, start with Filters.
How a filter works
A consumer filter selects records on the streaming server, before they cross the network. A reader asks for a page of one partition. The server scans it and returns only the records the filter selects. Each record keeps its original offset, ID, timestamp, headers, and payload bytes. Everything else stays on the server, and the log stays the single source of truth.
A filter belongs to a consumer group: Laser, then stream, topic, consumer group, filter. You create the group once and give it a filter. Every Laser SDK consumer of that group then runs it, and a consumer names only the group. A group with no filter receives every record. One topic can serve many selective groups without an intermediate topic per subset.
Raw Iggy polling ignores the filter. Rust's topic.iggy_consumer_group(..), the Iggy CLI, and any native Iggy client poll the same group with standard semantics and receive every record. The filter is a delivery policy for Laser SDK consumers, not an access control boundary, and a native source-read grant still permits native consumption. Do not mix the two on one group. A native poll stores the group's offset past records that the filtered consumers never saw, and they resume after it.
Savings depend on the share of payload bytes that match. Suppose ten consumers each need 5% of a topic's payload bytes:
| Payload from one 1 GiB source range | Read everything, filter in each application | Filter on the server |
|---|---|---|
| One consumer | 1 GiB | 51.2 MiB |
| Ten consumers combined | 10 GiB | 512 MiB |
| Payload transfer saved | 0% | 95% |
This illustrates payload transfer only, without protocol overhead. The server still examines records for each group.
Measured results:
- In the CDC example, the reader receives 4 of 240 records and 424 of 27,953 payload bytes, 98.5% less payload transfer. Rust, Python, and TypeScript produce the same result.
- The Frostline example publishes ten million change records, 6.30 GB of payload, and reads them with four filtered groups. In its recorded run, made before SDK 0.5.1, the groups receive 777.64 MB, while reading the full feed four times moves 25.21 GB, 96.9% of payload avoided. Counting both TCP directions with
stracegives 1.11 GB against 31.40 GB, 96.5% less application traffic, without TCP/IP headers and retransmissions. The repository holds the benchmark, latency tables, server CPU, and memory results.
Header predicates can skip payload decoding. Payload predicates add server work while they cut downstream bytes and application work.
What a filter can test
A filter is declarative. It never runs your code.
| Test | Operators |
|---|---|
| Payload field comparison | eq, ne, lt, lte, gt, gte, in, contains, prefix |
| Text matching on a field or a header | equals, prefix, suffix, contains, glob, regex, each optionally case-insensitive |
| Typed header comparison | The same comparisons, by exact header key, without decoding the payload |
| Field presence | present, absent |
| Composition | all, any, not (built with negate) |
| Explicit coercion | pred_as (TypeScript predAs) compares a timestamp (RFC 3339, epoch seconds, milliseconds, or microseconds) as an instant, or a decimal string as a number |
A payload field is named by a path. Dots separate object keys, and [n] selects an array item, as in after.ground_stations[0]. Escape a literal ., [, ], or \ inside a key with a backslash. A digit-only key such as codes.200 stays a key. A path holds up to 16 segments and 256 bytes.
The payload codec is part of the filter. JSON, CBOR, Avro, and Protobuf support field filtering. headers_only never decodes the payload, so it works with any format. Avro and Protobuf need registered writer schemas.
Truth has three values. Most comparisons with a missing field or a type mismatch return unknown, and not keeps unknown. An unknown result selects nothing, unless the mismatch policy asks to see type mismatches. Null comparisons are the exception: eq null matches a missing or explicit-null field, and ne null excludes both. Use present and absent when the difference matters.
Faults
A payload that cannot be decoded is a fault, which is separate from truth. The filter's fault policy decides what happens:
| Policy | What happens to the record |
|---|---|
stop (default) | The page stops before it. The reader raises the fault and reads that partition again after its idle interval, so the block clears once the record ages out of the partition. To move past it sooner, create a reader with an explicit start position or use another policy |
pass | Delivered and marked as not evaluated |
drop | Skipped |
Set it with with_fault_policy (TypeScript ConsumerFilter.withFaultPolicy(filter, policy)). The Python and TypeScript constructors also take it as their last argument.
Foreign and mismatched records
Two record policies decide what happens to records a filter cannot judge on its own terms. Both default to reject.
| Policy | Covers | reject (default) | pass |
|---|---|---|---|
foreign_policy | A record whose agdx.ct header names another codec than the filter's, or whose writer schema is missing, not listed by the filter, or of another schema family | Skipped. Never decoded and never a fault | Delivered and marked as not evaluated |
mismatch_policy | A record where a predicate's value exists but has a type the predicate cannot compare, for example xyz == "Abc" against "xyz": 1, an object, or an array | Skipped | Delivered and marked as not evaluated |
A missing or null field is not a mismatch. It is plain unknown and the mismatch policy never delivers it. A record without an agdx.ct header is decoded with the filter's codec. Content type any and codes a server does not know also decode as usual. A payload that is broken in the filter's own codec is still a fault and follows the fault policy.
A group reader marks a record delivered under pass with evaluated set to false, so a consumer can quarantine it, log it, or apply its own rules. Every record of an unbound group also has evaluated set to false, because nothing was evaluated. Records from the normal consumer carry no such flag, so use the group reader when you need to tell them apart.
import { ConsumerFilter, FilterExpr } from "@laserdata/laser-sdk"
const strict = ConsumerFilter.json(FilterExpr.pred("xyz", "eq", "Abc"))
const seeEdgeCases = ConsumerFilter.withMismatchPolicy(strict, "pass")use laser_sdk::filters::{ConsumerFilter, FilterExpr, RecordPolicy};
use laser_sdk::query::CmpOp;
let see_edge_cases = ConsumerFilter::json(FilterExpr::pred("xyz", CmpOp::Eq, "Abc"))
.with_mismatch_policy(RecordPolicy::Pass);import laser_sdk as ls
strict = ls.ConsumerFilter.json(ls.FilterExpr.pred("xyz", "eq", "Abc"))
see_edge_cases = strict.with_mismatch_policy("pass")Digest
Every filter has a digest, a SHA-256 hash over a fixed domain tag and the filter's JSON encoding. Two filters with the same digest behave the same way. Equivalent expressions written differently get different digests. A group remembers the digest it was configured with. Read it with digest() in Rust, the digest property in Python, or ConsumerFilter.digest(filter) in TypeScript.
Define a filter
Build the expression, then wrap it in the codec it reads. This filter follows a satellite fleet from the CDC example. It selects a producer-reported mode change to safe, a decommission, or a telemetry report of safe mode.
import { ConsumerFilter, FilterExpr } from "@laserdata/laser-sdk"
const satellites = () => FilterExpr.pred("table", "eq", "satellites")
const transitionsOnly = FilterExpr.all([
satellites(),
FilterExpr.pred("op", "eq", "u"),
FilterExpr.pred("changed", "contains", "mode"),
FilterExpr.pred("after.mode", "eq", "safe")
])
const safeMode = ConsumerFilter.json(
FilterExpr.any([
transitionsOnly,
FilterExpr.all([satellites(), FilterExpr.pred("op", "eq", "d")]),
FilterExpr.all([
FilterExpr.pred("event", "eq", "satellite.telemetry_changed"),
FilterExpr.pred("fields.mode", "eq", "safe")
])
])
)use laser_sdk::filters::{ConsumerFilter, FilterExpr};
use laser_sdk::query::CmpOp;
let satellites = || FilterExpr::pred("table", CmpOp::Eq, "satellites");
let transitions_only = FilterExpr::all([
satellites(),
FilterExpr::pred("op", CmpOp::Eq, "u"),
FilterExpr::pred("changed", CmpOp::Contains, "mode"),
FilterExpr::pred("after.mode", CmpOp::Eq, "safe"),
]);
let safe_mode = ConsumerFilter::json(FilterExpr::any([
transitions_only.clone(),
FilterExpr::all([satellites(), FilterExpr::pred("op", CmpOp::Eq, "d")]),
FilterExpr::all([
FilterExpr::pred("event", CmpOp::Eq, "satellite.telemetry_changed"),
FilterExpr::pred("fields.mode", CmpOp::Eq, "safe"),
]),
]));import laser_sdk as ls
def satellites() -> ls.FilterExpr:
return ls.FilterExpr.pred("table", "eq", "satellites")
transitions_only = ls.FilterExpr.all([
satellites(),
ls.FilterExpr.pred("op", "eq", "u"),
ls.FilterExpr.pred("changed", "contains", "mode"),
ls.FilterExpr.pred("after.mode", "eq", "safe"),
])
safe_mode = ls.ConsumerFilter.json(
ls.FilterExpr.any([
transitions_only,
ls.FilterExpr.all([satellites(), ls.FilterExpr.pred("op", "eq", "d")]),
ls.FilterExpr.all([
ls.FilterExpr.pred("event", "eq", "satellite.telemetry_changed"),
ls.FilterExpr.pred("fields.mode", "eq", "safe"),
]),
])
)Group reads, tests, previews, and filter administration are part of the default Rust streaming feature. The filters Cargo feature, also included by managed, adds only the local evaluator and the reader's local guard:
laser-sdk = { version = "0.7", features = ["filters"] }Create Laser from a connection string, so the SDK can open authenticated connections to the serving nodes. A raw client you supply without connection settings supports local diagnostics only, without acknowledgments.
Configure once, consume by group ID
Create the group with its filter once. Consumers then name only the group, by name or by native ID. The server runs the group's filter, or none, on every read. Several instances of the same group share partition assignments and progress. Different groups keep separate progress and can each receive their own copy of the same matching records.
const topic = laser.stream("orbit").topic("fleet_changes")
const desk = topic.consumerGroup("anomaly-desk")
// One-time setup. Omit filter to create a group that reads everything.
const info = await desk.create({ filter: safeMode })
const groupId = info.id
// Consumer instances only need the topic and group ID.
const consumer = await topic.consumerGroupId(groupId).consumer({
commitPolicy: { kind: "disabled" },
batchLength: 100
})
try {
for await (const record of consumer) {
// FleetChange is the record union of the complete example.
const change = record.json<FleetChange>()
console.log(record.position.partitionId, record.position.offset, change)
await consumer.commit(record)
}
} finally {
await consumer.shutdown()
}use laser_sdk::prelude::CommitPolicy;
let topic = laser.stream("orbit").topic("fleet_changes");
let desk = topic.consumer_group("anomaly-desk");
// One-time setup. Omit .filter(..) to create a group that reads everything.
let info = desk.create().filter(safe_mode.clone()).build().await?;
let group_id = info.id;
// Consumer instances only need the topic and group ID.
let mut consumer = topic
.consumer_group_id(u64::from(group_id))
.consumer()
.batch_length(100)
.commit_policy(CommitPolicy::Disabled)
.build()
.await?;
while let Some(record) = consumer.next().await {
let record = record?;
// FleetChange is the serde model of the complete example.
let change: FleetChange = record.json()?;
println!("{} {} {change:?}", record.position.partition_id, record.position.offset);
consumer.commit(&record).await?;
}
consumer.shutdown().await?;topic = laser.stream("orbit").topic("fleet_changes")
desk = topic.consumer_group("anomaly-desk")
# One-time setup. Omit filter to create a group that reads everything.
info = await desk.create(filter=safe_mode)
group_id = info.id
# Consumer instances only need the topic and group ID.
consumer = topic.consumer_group_id(group_id).consumer(
auto_commit="disabled", batch_length=100
)
try:
async for record in consumer:
print(record.position.partition_id, record.position.offset, record.json())
await consumer.commit(record)
finally:
await consumer.shutdown()create is idempotent. An existing group is kept, and the same definition keeps its binding. A different definition on a group that already runs one fails with conflict. Create another group instead, so two selections never share one set of offsets. You can also create a plain group first and call group.filter().configure(filter) later. The definition is saved as the group's own filter in the catalog and the group is bound to it in one step. The first definition becomes revision 1. A later matching definition reuses its revision.
A released group, or a group whose own filter was deleted, can be configured again only with the digest it ran before. Only a group that never had a binding takes any filter.
Native group creation and catalog configuration are two operations. If configuration fails after the group exists, Rust returns LaserError::ConsumerGroupSetup, TypeScript throws ConsumerGroupSetupError, and Python raises ConsumerGroupSetupError, a FilterError. Each carries the group's ID, name, and identity (group_id, group_name, and identity in Python) and the reason the catalog gave, such as conflict. The group stays as it was. In Rust the reason sits on the error's source, in TypeScript read it with filterReason(error.cause), and in Python it is reason. Pass an operation_id (TypeScript operationId) you recorded first to resume the same configuration after a crash and read its first outcome.
What a group consumer needs
On a server that serves filters, a group consumer needs the group_policy_reads capability. There it runs whatever the group holds when it reads, a filter or none. On original Apache Iggy, the same call builds the native group consumer. A server that serves filters but not group reads refuses the build with unsupported and asks for an upgrade. A capability probe that never answered fails the build with a retryable error. The SDK never falls back to a native read that ignores the group's policy. Python builds the consumer without awaiting it, so call await consumer.init() when startup must fail on any of these errors.
What a batch of 100 means
A normal consumer with batch length 100 examines at most 100 source records per partition request and delivers zero through 100 records. If seven match, you receive seven. If none match, the SDK keeps the scan position and continues through the backlog. An empty poll does not mean the topic is empty.
On the normal consumer, batch_length is both the page size and the scan budget, so one poll never costs the server more than one batch of source records. A selective filter over a long backlog therefore takes many polls to reach the first match. On a group-aware server, the batch length runs from 1 to 1000. When you want the server to scan further per request, use the group reader.
Group reader
The group reader separates the two limits. count bounds delivered matches and max_examined bounds scanned source records, within the server's own limits. count defaults to 100 and runs from 1 to 1000. Without max_examined, the server's own scan budget applies. The reader returns pages, acknowledges explicitly, and reads every record when the group is unbound. Use next_record for a record loop or next_page for batch handling (TypeScript nextRecord and nextPage). Both wait until a record arrives, so bound the wait yourself. TypeScript takes timeoutMs, and in Rust and Python wrap the call in a timeout.
const reader = await desk
.reader()
.start({ kind: "first" })
.count(100)
.maxExamined(1000)
.build()
try {
const page = await reader.nextPage({ timeoutMs: 15_000 })
for (const record of page.records) {
console.log(record.partitionId, record.offset, record.json())
}
await reader.ackPage(page)
} finally {
await reader.close()
}use laser_sdk::filters::FilteredStart;
let mut reader = desk
.reader()?
.start(FilteredStart::First)
.count(100)
.max_examined(1000)
.build()
.await?;
let page = reader.next_page().await?;
for record in &page.records {
let change: FleetChange = record.json()?;
println!("{} {} {change:?}", record.partition_id, record.offset);
}
reader.ack_page(&page).await?;
reader.close().await?;reader = await desk.reader(start="first", count=100, max_examined=1000)
try:
page = await reader.next_page()
for record in page.records:
print(record.partition_id, record.offset, record.message.json())
await reader.ack_page(page)
finally:
await reader.close()A page is one bounded filtered poll result, with matching records and scan progress. It uses standard Iggy transport and leaves stored records unchanged. For batch handling, process every returned record before ack_page (TypeScript ackPage). For record handling, call ack(record) after processing. The SDK advances stored progress only through the completed prefix. Acknowledge the record or page object that the reader returned. In TypeScript, a copy of its payload in another task does not carry the acknowledgment. try_next_page() returns at once when nothing is new.
Safe progress and policy changes
A fresh group using Next starts at the first retained record, including offset 0. After offset 0 is acknowledged, Next resumes at offset 1. An unset offset and a stored offset of zero are different states, so do not initialize a new group by storing zero. Next is the default start. Pick another with start_at(ConsumerStart::First) in Rust, polling="first" in Python, and startAt: { kind: "first" } in TypeScript.
Acknowledgments cover finished work, including records the server skipped. A page that examines records reports a safe offset, the last record it finished classifying. When you acknowledge the page, the server stores that offset, skipped records included. A group that matches one record in a million still moves forward. A page that examines offsets 0 through 99 and returns matches at 10 and 20 can store progress through 99 once both are done. It cannot cross an earlier delivery that is not yet handled.
On the normal consumer, the commit policy decides when the handled prefix is stored. Polling (the default) stores it at the next poll and at shutdown, Interval on a timer, Each, Every(n), and All after deliveries, and Disabled only when you call commit(record). Python names them with auto_commit ("polling", "interval", "each", "every", "all", or "disabled") plus commit_interval_ms and commit_every. TypeScript takes commitPolicy, such as { kind: "disabled" }. Nothing is stored before a record reaches your code, so a crash delivers the batch you were handling again instead of skipping it. Use Disabled plus commit when processing must finish before progress is stored.
The server fences every acknowledgment by source incarnation, group identity, execution mode, and policy generation. An old unfiltered page cannot commit after a filter is configured. An old filtered page cannot commit after the group is released. A purge or a recreated topic invalidates old delivery handles. A revoked partition can still drain accepted pages while its commit fence permits it.
A permission change takes effect inside a page. The server checks source read permission before every read round. A revocation between rounds fails the request and returns no partial page. The group's policy resolves once when the request starts, so a catalog change applies from the next request. An acknowledgment checks offset store permission when it runs.
The SDK rejects a page or record from another reader and checks response identity and progress before it shows them to you. It reads ahead and lets you finish records in any order. It stores only the offset up to which every returned record is done. Pages with no matches store their scanned range on their own.
A reader keeps at most 1024 unacknowledged pages with records per partition and stops reading that partition until you acknowledge, so read-ahead memory stays bounded. Change the bound with max_unacked_pages in Rust and Python or maxUnackedPages in TypeScript. Zero is refused: TypeScript throws when you set it, Rust fails at build(), and Python fails when you await reader(..). Empty pages do not count.
Reads go to the partition primary by default. That is the only mode that can acknowledge. The local read mode serves diagnostics from whatever replica the connection reaches. Set it with read_mode(ReadMode::Local) in Rust, read_mode="local" in Python, and readMode("local") in TypeScript. A local page offers no acknowledgment offset, because a lagging replica can still hold records from before a purge.
A group reader joins over its own connection, so two readers of one Laser are two members, and a reader dropped without close leaves the group with its connection. A partition the group hands a reader later resumes after the group's stored offset, whatever start the reader was built with, so the previous owner's unacknowledged backlog is read, not skipped. A partition that fails waits one idle interval while the other partitions keep reading.
A read never falls back to unfiltered delivery when a filter should apply. A failed catalog lookup, a node that cannot prove its catalog is current, and a paused revision all fail the read. With the plane disabled, a group read is unfiltered only when the server holds no catalog history and the read carries no required catalog position. Otherwise it fails with catalog_unavailable.
When the policy changes
A group keeps running when its policy changes. Its consumers and readers stay joined and keep their place. When the server reports that the policy moved, the partition restarts from its stored offset and the next read runs whatever the group holds now.
| Change | What running consumers do |
|---|---|
| Configure a filter on a group that never had one | The next read runs the filter |
| Release or delete the filter | The next read returns every record |
| Pause the active revision | New reads fail with revision_disabled. Records already delivered can still be acknowledged. Enable the revision again to resume |
| Draft a revision | Nothing changes. The group keeps running the revision it is bound to |
Records delivered under the old policy but not yet stored arrive again from the stored offset, so nothing is skipped. Only an explicit commit or ack of a record read under the old policy fails, with conflict, because the server never stored that offset.
Checkpoints inside a page
Persist your application checkpoint before you acknowledge its offset. ack_through(record) (TypeScript ackThrough(record)) marks every earlier returned record on that partition as handled and stores the completed prefix. Process that full prefix first. Later records in the same page stay pending. Use it when a window receipt or a transaction ends inside a page. Use ack_page when the whole page is done.
Cluster freshness
Configuration outcomes carry a durable operation receipt and a control log position. The group handle remembers its last configuration, also across a cloned handle, and every read built from it carries that position. The serving plane confirms that operation before it answers, so an application never sees its own configuration as missing through another node.
A read from a group that nobody configured recently can hit a cached unbound answer. IGGY_PLANE_FILTERS_UNBOUND_CONFIRM_MS bounds how long a server shard reuses that answer before it asks the plane for a lookup confirmed against the control log head. The default is 1000 ms, the maximum is 60,000 ms, and zero confirms every read. A filter configured through another node takes effect within that window at the latest. It is not an instant cross-node guarantee.
Acknowledgments confirm the current policy again. A node that cannot prove its catalog is current returns the retryable catalog_unavailable. Cluster membership, primary routing, and replication keep their native Iggy behavior.
What a change record proves
Filters judge each record on its own. They do not compare it with an earlier version of the entity. In the example, the producer supplies changed, a list of columns that changed. The transition branch needs both changed containing mode and after.mode equal to safe.
A battery update from a satellite already in safe mode does not match that branch. A filter on values only selects it. Choose that kind of filter when you need current state rather than evidence of a transition.
The telemetry branch selects a report of safe mode. Repeated or forced telemetry reports can match even when the mode did not change. The decommission branch selects deletions whatever the mode. Keep these producer semantics in mind when you adapt the example to your feed.
Test and preview
Check the group's filter before you depend on it. Neither call joins the group or stores an offset.
testjudges one sample payload you supply and explains every predicate.previewjudges stored records in one partition and lists each verdict. Its first argument is the partition ID. It takesfrom_offset,max_examined,max_records, andexplain(TypeScriptfromOffset,maxExamined,maxRecords,explain).max_recordsdefaults to 20. Passexplainto list rejected records too.
const tested = await desk.filter().test(JSON.stringify(sample))
console.log(tested.explanation.verdict)
const preview = await desk.filter().preview(0, { maxRecords: 10 })
console.log(preview.examined, preview.matched, preview.stop)let payload = serde_json::to_string(&sample).expect("the sample serializes");
let tested = desk.filter().test(payload, Vec::new()).await?;
println!("{:?}", tested.explanation.verdict);
let preview = desk.filter().preview(0).await?.max_records(10).send().await?;
println!("{} {} {:?}", preview.examined, preview.matched, preview.stop);tested = await desk.filter().test(json.dumps(sample))
print(tested["explanation"]["verdict"])
preview = await desk.filter().preview(0, max_records=10)
print(preview["examined"], preview["matched"], preview["stop"])stop says why the scan ended: filled, budget, end_of_visible, fault, or oversized_record. Python returns catalog and preview replies as dicts.
test also takes typed headers, as FilterHeader values in Rust and TypeScript and a headers= dict in Python, so a headers-only filter can be checked the same way.
Local evaluator
The local evaluator gives the same verdict without a server. Compile the filter once with CompiledFilter.compile(filter), then call evaluate and explain on each record. Avro and Protobuf need the schemas at compile time: a schemas argument in Python and TypeScript, and compile_with_schemas in Rust. Python takes evaluate(payload, headers). TypeScript takes a { payload, headers } record. Rust uses CompiledFilter from laser_sdk::filters, which needs the filters feature, and its evaluate and explain also take &DecodeLimits. evaluate_with_fault (TypeScript evaluateWithFault) also returns the fault reason when the evaluator cannot judge a record.
Local guard
The local guard checks the server's work. A reader built with the guard evaluates every returned record again with the shared evaluator and fails the page on any disagreement, before your code sees it. Rust calls local_guard(true), Python passes local_guard=True, and TypeScript calls localGuard(true). The server sends its decoder bounds so the guard can reproduce a size or depth fault. A record marked as not evaluated without the matching pass policy is rejected. Older servers without these bounds cannot verify pass-through decode faults locally.
Filter on headers only
A headers_only filter reads typed user headers and never touches the payload. Use it when producers stamp a routing header and the payload is Protobuf, Avro, or any other format. The pager group below receives only critical alert frames.
import { ConsumerFilter, FilterExpr, HeaderValue } from "@laserdata/laser-sdk"
const alerts = laser.stream("orbit").topic("fleet_alerts")
await alerts.producer().send(frame, { headers: { priority: HeaderValue.uint8(2) } })
const pager = alerts.consumerGroup("pager")
await pager.create({
filter: ConsumerFilter.headersOnly(FilterExpr.header("priority", "eq", 2))
})
const reader = await pager.reader().start({ kind: "first" }).build()
const record = await reader.nextRecord({ timeoutMs: 15_000 })
console.log(record.offset, record.message.payload.byteLength)
await reader.ack(record)
await reader.close()use laser_sdk::filters::{ConsumerFilter, FilterExpr, FilteredStart};
use laser_sdk::query::CmpOp;
use laser_sdk::stream::{HeaderKey, HeaderValue, ProducerMessage};
let alerts = laser.stream("orbit").topic("fleet_alerts");
let producer = alerts.producer().build().await?;
let message = ProducerMessage::new(frame)
.header(HeaderKey::try_from("priority")?, HeaderValue::from(2_u8));
producer.send_message(message).await?;
let pager = alerts.consumer_group("pager");
pager
.create()
.filter(ConsumerFilter::headers_only(FilterExpr::header("priority", CmpOp::Eq, 2_i32)))
.build()
.await?;
let mut reader = pager.reader()?.start(FilteredStart::First).build().await?;
let record = reader.next_record().await?;
println!("{} {}", record.offset, record.message.payload.len());
reader.ack(&record).await?;
reader.close().await?;alerts = laser.stream("orbit").topic("fleet_alerts")
producer = alerts.producer()
await producer.send(frame, headers={"priority": ("uint8", 2)})
pager = alerts.consumer_group("pager")
await pager.create(
filter=ls.ConsumerFilter.headers_only(ls.FilterExpr.header("priority", "eq", 2))
)
reader = await pager.reader(start="first")
record = await reader.next_record()
print(record.offset, len(record.message.payload))
await reader.ack(record)
await reader.close()frame is your encoded payload. Header keys are text names matched exactly, including dots. Values can be booleans, signed or unsigned integers up to 64 bits, floating-point numbers, or strings. Integers compare by value across widths, so a uint8 header matches the integer 2. Numeric 2 and text "2" are different values. Pick an exact-width Iggy header when space matters: uint8 uses one value byte. Use the streaming producer for typed headers. The publish().header(..) shortcut takes strings. Non-text keys and 128-bit numeric comparisons are outside the filter contract.
Agent groups
Agents use this kind of filter without any setup. Each agent ID has its own consumer group. On agent.sessions and agent.control, the agent runtime binds that group to the headers-only filter agdx.to In [<agent id>, "*"]. It does so when the server serves filtered reads, group policies, and the catalog, and the client was built from a connection string. On a server without filters, the group stays unbound and the agent sorts records on the client. A server that serves filters but not group reads makes the agent fail to start with unsupported, a capability probe that never answered fails it with a retryable timeout, and a group already bound to another filter is refused. See Agents in depth.
One log, many kinds of events
Put every producer's events in one ordered log and let each group select its part on the server. One partition keeps a single global order. More partitions add parallel group workers. Events can differ in type, payload shape, envelope, and codec.
The pattern that holds up:
- Stamp a routing header on every record, such as
event.type = "metrics.cpu.v1.reported", and the content type through the codec helpers, which setagdx.ct. - Select the part with a header text match, such as
prefixmetrics.,contains.v1., or agloborregex. Aheaders_onlyfilter never decodes a payload, so every format mixes safely, and it is the cheapest filter to run. - Add payload conditions only where needed, inside
all([header match, payload predicates]). Header children run first andallstops at the first child that does not match, so records of other kinds are rejected before any decode. - Keep
foreign_policyatreject. A record in another codec is then skipped instead of stalling the partition, even when its routing header is missing.
const metricsV1 = ConsumerFilter.headersOnly(
FilterExpr.headerText("event.type", "glob", "metrics.*.v1.*")
)
const hotHosts = ConsumerFilter.json(
FilterExpr.all([
FilterExpr.headerText("event.type", "prefix", "metrics.cpu."),
FilterExpr.pred("cpu", "gte", 90)
])
)use laser_sdk::filters::TextMatch;
let metrics_v1 = ConsumerFilter::headers_only(FilterExpr::header_text(
"event.type",
TextMatch::Glob,
"metrics.*.v1.*",
));
let hot_hosts = ConsumerFilter::json(FilterExpr::all([
FilterExpr::header_text("event.type", TextMatch::Prefix, "metrics.cpu."),
FilterExpr::pred("cpu", CmpOp::Gte, 90_i64),
]));metrics_v1 = ls.ConsumerFilter.headers_only(
ls.FilterExpr.header_text("event.type", "glob", "metrics.*.v1.*")
)
hot_hosts = ls.ConsumerFilter.json(ls.FilterExpr.all([
ls.FilterExpr.header_text("event.type", "prefix", "metrics.cpu."),
ls.FilterExpr.pred("cpu", "gte", 90),
]))Text matching
FilterExpr.text matches a payload field and header_text (TypeScript headerText) matches a header. Both take one of these kinds:
| Kind | Matches when | Relative cost |
|---|---|---|
equals | The whole value equals the pattern | Lowest |
prefix | The value starts with the pattern | Lowest |
suffix | The value ends with the pattern | Lowest |
contains | The pattern occurs anywhere | Low, one scan of the value |
glob | The whole value matches. * is any run, ? is one character, \ escapes | Low to moderate with many stars |
regex | The value contains a match | Moderate, linear in the value |
Any kind can be case-insensitive. Rust and Python call .case_insensitive() on the expression, such as FilterExpr.text("host", "prefix", "node-").case_insensitive() in Python. TypeScript passes true as the last argument of text or headerText, or wraps the expression in filterExprCaseInsensitive(expr). Every kind except regex lowercases the pattern and the value first, which copies the value. Case-insensitive regex uses the engine's Unicode folding. The value must be text. A number, object, or array is a type mismatch and follows the mismatch policy.
On a JSON payload, parsing costs more than the text test, so the text kinds cost about the same as a plain equality test. A header-only test does not parse the payload and costs much less. Measure your own records with the SDK's filter_evaluation and filter_paths benchmarks.
Regular expressions run in linear time on the server, so a pattern cannot backtrack without bound. Lookaround, backreferences, and inline flags such as (?i) are refused. Use the case-insensitive option instead of (?i). The server's Rust engine decides pattern syntax. Character classes such as \d and \w are Unicode-aware. Rust and Python run that engine in their local evaluator and guard. TypeScript runs a port of the same syntax in linear time, so its evaluator and local guard also handle regex. It refuses the few Unicode properties that Node cannot express, such as Age and the segmentation break properties, and a pattern close to the size limit can pass locally and still be refused by the server.
Revisions, pause and resume, A/B
A group's filter has immutable revisions. revise drafts a new revision of the group's own definition. Readers keep running the active one. A draft never replaces what a bound group runs. To run a stricter variant, create a second group with that definition. Each group has its own offsets, so one variant never skips the other's records.
set_revision_enabled(revision, false) pauses the active revision. New reads then fail with revision_disabled. Records already delivered can still be acknowledged, and test and preview still work. Set it to true to resume. The flag changes neither the definition nor its digest.
release() removes the group's policy and returns the released binding. Its consumers then receive every record, and its filter stays saved in the catalog. delete() removes the group's own definition and every revision, and releases its binding in the same catalog operation. It returns whether the group had one. After either call the group keeps its digest restriction, so configure it again with the same digest or create another group for a different policy. A named filter dropped from the console or over HTTP also releases its groups and leaves a tombstone that reserves the name. Past the tombstone limit the oldest tombstones are deleted. An ID is never reused.
import { filterReason } from "@laserdata/laser-sdk"
const binding = await desk.filter().get()
if (binding === undefined) throw new Error("the desk is unbound")
// Draft on the desk's own filter. The desk keeps running its revision.
const draft = await desk.filter().revise(binding.revision, ConsumerFilter.json(transitionsOnly))
const revisions = await desk.filter().revisions(0, 10)
// Run the variant in its own group, with its own offsets.
const variant = topic.consumerGroup("anomaly-desk-transitions")
const created = await variant.create({ filter: ConsumerFilter.json(transitionsOnly) })
// Pause and resume the variant.
if (created.filter === undefined) throw new Error("the variant was created unbound")
await variant.filter().setRevisionEnabled(created.filter.revision, false)
await variant.filter().setRevisionEnabled(created.filter.revision, true)
// Another definition on a running group is refused.
try {
await desk.filter().configure(ConsumerFilter.json(transitionsOnly))
} catch (error) {
if (filterReason(error) !== "conflict") throw error
}
// Release: the desk receives every record again.
const released = await desk.filter().release()
// Delete the desk's own filter with every revision. False when it had none.
const deleted = await desk.filter().delete()use laser_sdk::filters::FilterErrorReason;
let binding = desk.filter().get().await?.expect("the desk is bound");
// Draft on the desk's own filter. The desk keeps running its revision.
let draft = desk
.filter()
.revise(binding.revision, ConsumerFilter::json(transitions_only.clone()))
.await?;
let revisions = desk.filter().revisions(0, 10).await?;
// Run the variant in its own group, with its own offsets.
let variant = topic.consumer_group("anomaly-desk-transitions");
let created = variant
.create()
.filter(ConsumerFilter::json(transitions_only.clone()))
.build()
.await?;
let variant_revision = created.filter.expect("bound").revision;
// Pause and resume the variant.
variant.filter().set_revision_enabled(variant_revision, false).await?;
variant.filter().set_revision_enabled(variant_revision, true).await?;
// Another definition on a running group is refused.
match desk.filter().configure(ConsumerFilter::json(transitions_only)).await {
Err(error) if error.filter_reason() == Some(FilterErrorReason::Conflict) => {}
other => panic!("expected conflict, got {other:?}"),
}
// Release: the desk receives every record again.
let released = desk.filter().release().await?;
// Delete the desk's own filter with every revision. False when it had none.
let deleted = desk.filter().delete().await?;binding = await desk.filter().get()
if binding is None:
raise RuntimeError("the desk is unbound")
# Draft on the desk's own filter. The desk keeps running its revision.
draft = await desk.filter().revise(binding["revision"], ls.ConsumerFilter.json(transitions_only))
revisions = await desk.filter().revisions(page=0, page_size=10)
# Run the variant in its own group, with its own offsets.
variant = topic.consumer_group("anomaly-desk-transitions")
created = await variant.create(filter=ls.ConsumerFilter.json(transitions_only))
# Pause and resume the variant.
await variant.filter().set_revision_enabled(created.filter["revision"], False)
await variant.filter().set_revision_enabled(created.filter["revision"], True)
# Another definition on a running group is refused.
try:
await desk.filter().configure(ls.ConsumerFilter.json(transitions_only))
except ls.FilterError as error:
assert error.reason == "conflict"
# Release: the desk receives every record again.
released = await desk.filter().release()
# Delete the desk's own filter with every revision. False when it had none.
deleted = await desk.filter().delete()revisions takes a page and a page size. Rust needs both. TypeScript makes both optional. Python takes them as keywords, page=0 and page_size=50 by default.
Every catalog change carries an operation ID. A retry by the same caller with the same ID and change returns the first retained outcome. Reusing an ID with another change or caller is refused. configure_as(operation_id, ..) (TypeScript configureAs) sends a configuration under an ID you recorded first. configure_with (TypeScript configureWith) takes a definition or one of the group's own revisions. Python has the same three verbs: configure(filter), configure_with(filter=, filter_id=, revision=), and configure_as(operation_id, filter=, filter_id=, revision=). Results stay in durable plane storage, even after control topic messages expire, so retry protection survives a recreated control topic. Recent results stay in memory up to the cache limit, and a miss reads durable storage.
Each change is authorized again when the catalog applies it, against the catalog as it stands at that point of its log, so a scoped deny holds even on a node that had not seen the group yet.
A group ID belongs to one topic incarnation. Deleting and recreating a group, a topic, or a stream makes a new identity that inherits no policy. A reader of the old ID does not join the replacement.
Errors
Filter failures carry a typed reason. Rust exposes it through error.filter_reason(), Python through FilterError.reason, and TypeScript through the filterReason(error) helper or FilterExecutionError.detail.reason.
| Reason | What it means |
|---|---|
conflict | A precondition failed. The group runs another definition, the expected revision moved, or an acknowledgment names a policy the group no longer holds. A reader that meets it on a read restarts that partition from its stored offset and keeps reading |
source_changed | The source history or the group's identity changed. The reader restarts that partition from its stored offset once and raises this error so you know |
revision_disabled | The active revision is paused. New server reads stop until it is enabled. Records already delivered can still be acknowledged |
catalog_unavailable | The plane cannot prove its catalog is current, or the needed configuration has not reached this node yet. Retryable |
membership_stale | The consumer no longer owns the partition. The SDK refreshes its assignment |
not_primary | The node is no longer the partition primary. The SDK resolves the route again |
not_found | The group does not exist, or an operation that needs a policy met an unbound group |
forbidden | The caller lacks the filter grant on the group path or read permission on the source |
unsupported | The server does not serve this operation, such as the catalog without laser-plane, or a filter whose evaluator version or codec it does not match |
capacity_exhausted | A configured catalog limit is reached, such as definitions, revisions, group identities, or stored revision bytes |
version_skew | The client and the server speak different filter operation versions |
invalid_request | The filter or request is malformed or out of range |
too_large | The request is too large |
unauthenticated | The connection is not signed in |
unavailable | A transient failure. The same request can succeed later |
backend | An unexpected failure that a retry cannot fix |
unknown | A reason this SDK version does not know. The error's result code still classifies it |
Faults, oversized records, and group setup failures are separate error types. A stop fault or a matching record larger than the page's reply cap blocks its partition. The reader delivers every earlier match first, then raises the error with the partition and offset. Rust returns LaserError::FilterFault or LaserError::FilterOversizedRecord, and filter_reason() returns None for them and for ConsumerGroupSetup. TypeScript throws FilterFaultError, with the fault in reason, or FilterOversizedRecordError, each with partitionId and offset. filterReason() returns undefined for them and for ConsumerGroupSetupError. Python raises FilterFaultError or FilterOversizedRecordError, both FilterError subclasses with reason set to fault or oversized_record, plus fault_reason, partition_id, and offset. A Python ConsumerGroupSetupError has reason setup_failed when the catalog gave no other reason.
Key operations
| Operation | What it does |
|---|---|
topic.consumer_group(name) / consumer_group_id(id) | The group handle, by name or native ID. Free to construct |
group.create().filter(filter).build() | Create the group, optionally with its filter. Idempotent |
group.info() | ID, name, exact identity, and the active binding |
group.consumer() | The normal consumer. Runs the group's policy. batch_length bounds page and scan |
group.reader() | The match-oriented reader. count, max_examined, max_reply_bytes, start, partition, read_mode, local_guard, idle_interval (Python: idle_interval_ms), max_unacked_pages |
next_page() / next_record() | Read the next page or record. try_next_page() returns at once when nothing is new |
ack(record) / ack_through(record) / ack_page(page) | Mark records done, which stores progress |
group.filter().configure(filter) | Give an existing group its policy. configure_with takes a definition or one of its own revisions, configure_as adds an operation ID |
group.filter().get() | The active binding, or none |
group.filter().revisions(page, size) / revise(expected, filter) | List revisions, draft a new one |
group.filter().set_revision_enabled(revision, enabled) | Pause or resume without changing content |
group.filter().release() | Remove the policy. Consumers then receive every record |
group.filter().delete() | Delete the group's own filter with every revision, releasing the group first. Returns whether it had one |
group.filter().test(payload, headers) / preview(partition) | Judge a sample or stored records without storing state |
TypeScript spells these consumerGroup, consumerGroupId, nextPage, ackThrough, configureWith, setRevisionEnabled, and so on. Python passes options as keyword arguments and returns catalog replies as dicts. The Python reader takes partitions where Rust and TypeScript take partition. Rust takes idle_interval as a Duration, and Python (idle_interval_ms) and TypeScript take milliseconds.
Permissions
Grants for a group's filter use the group path stream/topic/group. filter:read inspects and lists, filter:write drafts a revision, filter:admin configures, pauses, resumes, and releases, and filter:delete deletes the group's own filter. Native source read permission still controls consumption, and a consumer does not need to inspect the definition to run its group's policy.
In the LaserData Cloud Console, open a topic and its consumer group. The group's filter shows above its members, with its revision, digest, and policy generation. You can configure an unbound group there with the visual, schema-aware editor, pause or resume the active revision, release the policy, and preview it. Policy information that cannot be read shows as unavailable, never as "no filter".
Limits
| Limit | Value |
|---|---|
| Encoded filter size | 8 KiB |
| Field path | 256 bytes, and 16 path segments |
| Header key in a filter | 256 bytes. Iggy itself stores header keys of at most 255 bytes, and a test sample header key is limited to 255 |
| Test sample payload | 1 MiB |
| Test sample headers | 64, each value up to 1 KiB |
| Expression nodes | 128 |
| Nesting depth | 8 |
Items in one in list | 64 |
| Text pattern | 1 KiB |
| Compiled glob or regex programs per filter | 4 |
| Budget per compiled program | 256 KiB |
| Records in one page | 1000 |
| Source records one request can ask a page to examine | 100,000, lowered to the server's own budget (10,000 by default) |
| Unacknowledged pages with records per partition | 1024 by default, configurable per reader |
| Bytes in one page | 8 MiB, and 1 MiB by default for a reader |
| Records one preview examines | 1,000 by default, configurable up to 10,000 |
| Records one preview returns | 100 |
| Live saved definitions | 1,024 |
| Revisions per definition | 256 |
| Retained group identities | 4,096 |
| Stored revision bytes, all filters together | 64 MiB |
| Dropped filters kept as tombstones, oldest deleted past the limit | 4,096 |
| Results cached in memory by default | 256 |
| Concurrent changes per plane by default | 32 |
| Concurrent changes per caller by default | 8 |
| Changes per caller per second by default | 32 |
Plane settings
The definition, tombstone, revision, and group-identity values above are defaults. LD_PLANE_FILTER_MAX_DEFINITIONS, LD_PLANE_FILTER_MAX_TOMBSTONES, LD_PLANE_FILTER_MAX_REVISIONS, LD_PLANE_FILTER_MAX_GROUP_IDENTITIES, and LD_PLANE_FILTER_MAX_CATALOG_BYTES accept 0 for no limit. LD_PLANE_FILTER_OUTCOME_CACHE_ENTRIES=0 turns off the cache and keeps durable retries.
LD_PLANE_FILTER_MAX_CONCURRENT_MUTATIONS and LD_PLANE_FILTER_MAX_CONCURRENT_MUTATIONS_PER_ACTOR set positive limits on changes in progress at once. LD_PLANE_FILTER_MAX_MUTATIONS_PER_ACTOR_PER_SECOND limits each caller's change rate, with a minimum of 1. Past any of them, the plane returns a retryable unavailable.
The system control topic defaults to no expiry. Set LD_PLANE_CONTROL_EXPIRY_SECS to 604800 for seven days, 2592000 for thirty days, or 0 for no expiry. Startup applies the configured value to existing topics. With finite retention, recovery uses the durable plane snapshot and the remaining log. Expiring topic messages does not delete group filters, bindings, or operation results.
LD_PLANE_DLQ_EXPIRY_SECS and LD_PLANE_CHANGES_EXPIRY_SECS keep the plane's dead-letter and change topics for one day by default.
Server settings
IGGY_PLANE_FILTERS_ENABLED turns the filter evaluator on or off for a deployment. It is on by default. Group reads of unbound groups work without it. A bound group is refused until it is on.
| Setting | Default | What it bounds |
|---|---|---|
IGGY_PLANE_FILTERS_MAX_CONCURRENT | 16 per shard | Concurrent scans. A scan past it gets a retryable refusal |
IGGY_PLANE_FILTERS_MAX_EXAMINED_RECORDS | 10,000 | Source records one page examines |
IGGY_PLANE_FILTERS_MAX_EXAMINED_BYTES | 64 MiB | Bytes one page scan examines |
IGGY_PLANE_FILTERS_PREVIEW_MAX_EXAMINED | 1,000, at most 10,000 | Records one preview examines |
IGGY_PLANE_FILTERS_MAX_ROUND_RECORDS | 1,000, from 1 to 1,000 | Records in one internal owner read |
IGGY_PLANE_FILTERS_MAX_PAGE_ROUNDS | 256 | Owner reads per page |
IGGY_PLANE_FILTERS_MAX_ROUNDS | Separate value | Owner reads per preview |
IGGY_PLANE_FILTERS_UNBOUND_CONFIRM_MS | 1,000, at most 60,000, 0 confirms every read | How long a cached unbound answer is reused |
IGGY_PLANE_FILTERS_POLICY_CONFIRM_TIMEOUT_MS | 1,000, from 1 to 60,000 | Each policy version confirmation |
IGGY_PLANE_FILTERS_POLICY_CACHE_ENTRIES | 1,024, at most 65,536, 0 resolves every request | Cached policies per shard |
IGGY_PLANE_FILTERS_POLICY_CACHE_TTL_MS | 60 seconds, at most 10 minutes, 0 resolves every request | Lifetime of a cached policy |
Fewer owner reads per page cost less server CPU. On a one-million-record trial, a round cap of 1,000 used about 19% less server CPU than 128 with the same results. The first owner read of a page probes at most 128 records. Each later read shrinks to what the remaining byte budget covers at the average record size seen so far, so a sudden jump in record size can pass the byte budget by one owner read. The scan yields every 1 MiB of records and sizes the next read by how selective the filter was so far, so a sparse filter can scan further within the page budget. A request's max_examined caps the budget below the server's. Consumer group reads also respect the remaining wanted match count to keep the native polling fence. A page also stops at its round and byte budgets, so a 1,000-record request can return fewer matches. Previews size their owner reads the same way and also check bytes before each record. These bounds limit work per page. They do not promise throughput or latency.
Each server shard caches resolved policies and watches its own plane for catalog changes. Before it reuses an entry, it confirms the current catalog version through the sidecar, so a delayed watch cannot keep an old binding active after a finished local change. Changes made on another node reach the local plane through control log replay. A plane restart invalidates earlier version tokens. The server also checks current reader grants before reuse. Cache hits skip repeated definition transfer and decoding, while a small version confirmation still runs for each hit. Watch failures and confirmation failures turn off reuse.
Payload formats
All formats use the same predicates, fault policies, previews, and acknowledgment rules. The server returns the original payload bytes, so consumers keep their existing decoders.
| Format | Filter constructor | Value rules |
|---|---|---|
| JSON | ConsumerFilter.json(expression) | Exact 64-bit integers, finite numbers, arrays, objects, strings, booleans, and null |
| CBOR | ConsumerFilter.cbor(expression) | Definite-length values, text map keys, byte strings as integer arrays. No tags, duplicate keys, or non-finite numbers |
| Avro | ConsumerFilter.avro(expression, schema_refs) | Registered writer schemas. Logical types use their physical values, such as integer timestamps or decimal bytes. Enums are their symbol strings |
| Protobuf | ConsumerFilter.protobuf(expression, schema_refs) | Registered descriptor sets and message names. Field names follow the schema, enums are numbers, and bytes are integer arrays. A field with presence tracking (proto2, optional, messages) is absent when unset. A proto3 scalar or repeated field without it always exists, at its default when unset, so counter == 0 matches |
| Headers only | ConsumerFilter.headers_only(expression) | Typed user headers, with no payload decoding |
The table uses Python names. Rust uses ConsumerFilter::avro(expression, schema_refs) and TypeScript uses ConsumerFilter.avro(expression, schemaRefs). The CDC examples show typed producers and group readers for every format in all three languages.
Register Avro and Protobuf schemas once, then reference their IDs. Put the allowed writer IDs in the filter's schema_refs. Stamp each published record with its writer ID through schema_id in Rust and Python or schemaId in TypeScript, which sets the agdx.sid header. A filter can allow several immutable writer schemas of the same codec. The schema IDs are part of its digest. A schema ID that is missing, not listed, or registered for another schema family follows the filter's foreign_policy when a predicate needs the payload: the record is skipped by default, or delivered as not evaluated under pass.
The server resolves schemas through the plane schema catalog before it scans records. Without a schema stream, schema_refs resolve in the deployment-wide registry. with_schema_stream(stream) in Rust and Python, and ConsumerFilter.withSchemaStream(filter, stream) in TypeScript, name the stream whose registry holds them, and change the digest. Rust, Python, and TypeScript local guards use the same references. CBOR needs no registered schema. Protobuf groups are not supported. Empty repeated and map fields are present with their default empty values. Use explicit timestamp coercions for Avro integer timestamps.
Byte and nesting limits bound decoding. The server announces its evaluator version and supported codecs, and the SDKs refuse unsupported capabilities explicitly.
Availability and cost
Group reads, tests, and previews run on the LaserData Iggy fork, in Laser Stack and LaserData Cloud. Configuring a group's filter needs laser-plane. create with a filter is refused before the group is created when catalog support is unavailable. Original Apache Iggy keeps working for native streaming.
Check the filter capabilities before you use them. native covers filtered reads, tests, and previews. catalog covers configuration. group_policy_reads covers group-aware consumers.
const { filters } = await laser.capabilities()
console.log(filters.native, filters.catalog, filters.groupPolicyReads)let filters = laser.capabilities().await.filters;
println!("{} {} {}", filters.native, filters.catalog, filters.group_policy_reads);filters = (await laser.capabilities()).filters
print(filters.native, filters.catalog, filters.group_policy_reads)Measure bandwidth savings and server cost together. The server still scans source records per group. Header-only policies avoid decoder work. Payload policies add decoding and evaluation while they cut bytes sent to consumers. Use the Frostline benchmark documentation to measure CPU, memory, sample counts, and p50, p99, and p99.9 latency on your workload. The figures on this page describe recorded workloads. They are not a latency guarantee or a replacement for benchmarks on your deployment.