LaserData Cloud
Laser SDK

Consumer Filters

Filter records before they cross the network. Receive only what each consumer needs.

Receive only the records you need. A consumer filter selects records on the streaming server, before they cross the network. Change data capture (CDC) records describe changes published by a producer. The 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.

98.5% less payload transfer in the complete CDC example. The reader receives 4 of 240 records and 424 of 27,953 payload bytes. Rust, Python, and TypeScript produce the same result.

At scale the picture holds. The Frostline example publishes ten million change records, 6.30 GB of payload, and reads them with four filtered groups. The groups receive 777.64 MB where reading the full feed four times would move 25.21 GB, 96.9% of payload avoided. Counting both TCP directions with strace gives 1.11 GB against 31.40 GB, 96.5% less TCP application traffic. These counts exclude TCP/IP headers and retransmissions. The same repository holds the benchmark, p50/p99/p99.9 tables, server CPU, and memory results. Traffic reduction is measured separately from latency and CPU cost. Header predicates can avoid payload decoding. Payload predicates add server work while reducing downstream bytes and application work.

Built for

Use consumer filters for CDC consumers that care about one kind of change, for alert routing, and for backfills that must skip most of a topic.

How it works

A change feed carries every change to every table. Without a filter, a consumer downloads and decodes all of it just to drop most records. With a filter, the server does that work next to the data, and the log stays the single source of truth. Savings depend on the matching fraction of payload bytes in your feed.

One topic can serve many selective consumers without sending every record to each one. Consumers can follow different subsets of the same topic and its partitions. You do not need an intermediate topic just to deliver each subset.

For example, suppose ten consumers each need 5% of a topic's payload bytes:

Payload from one 1 GiB source rangeRead everything, filter in each applicationFilter on the server
One consumer1 GiB51.2 MiB
Ten consumers combined10 GiB512 MiB
Payload transfer saved0%95%

This is an illustration of payload transfer, excluding protocol overhead. The broker still examines records for each reader. Matching records retain their full payload. Applications avoid receiving and decoding rejected payloads.

Combine payload fields, typed headers, and change evidence in one filter. A filter is declarative. It never runs your code. It can:

  • Payload predicates: Compare fields with eq, ne, lt, lte, gt, gte, in, contains, and prefix.
  • Text matching: Match a text field or header with equals, prefix, suffix, contains, glob, or regex, each optionally case-insensitive.
  • Typed headers: Compare custom headers by exact key, without decoding the payload.
  • Presence tests: Distinguish present and absent fields.
  • Composable rules: Combine tests with all, any, and not.
  • Explicit coercion: Coerce a field before comparing, so a timestamp string orders as an instant and a decimal string compares as a number.

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 require 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 filter's mismatch policy asks to see type mismatches (see Edge cases). 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 distinction matters.

A payload that cannot be decoded is a fault, which is separate from truth. The filter's fault policy decides what happens:

PolicyWhat 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 recovery position or policy
passDelivered and marked as not evaluated
dropSkipped and included in the examined count

Edge cases: foreign and mismatched records

Two record policies decide what happens to records a filter cannot judge on its own terms. Both default to reject.

PolicyCoversreject (default)pass
foreign_policyA 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 familySkipped. Never decoded and never a faultDelivered and marked as not evaluated
mismatch_policyA 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, and the filter could not decide otherwiseSkippedDelivered and marked as not evaluated

A missing or null field is not a mismatch. It is plain unknown and never delivered by the mismatch policy. A record without an agdx.ct header is decoded with the filter's codec, as before. 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.

Records delivered under pass are listed in the page's unevaluated offsets, so a consumer can quarantine them, log them, or apply its own rules.

const strict = ConsumerFilter.json(FilterExpr.pred("xyz", "eq", "Abc"))
const seeEdgeCases = ConsumerFilter.withMismatchPolicy(strict, "pass")
let see_edge_cases = ConsumerFilter::json(FilterExpr::pred("xyz", CmpOp::Eq, "Abc"))
    .with_mismatch_policy(RecordPolicy::Pass);
strict = ls.ConsumerFilter.json(ls.FilterExpr.pred("xyz", "eq", "Abc"))
see_edge_cases = strict.with_mismatch_policy("pass")

Every filter has a digest, a SHA-256 hash of its canonical form. The digest names exactly what the filter selects, so two filters with the same digest behave the same way.

Safe progress across skipped records

A fresh reader using Next starts at the first retained record, including offset 0. After offset 0 is acknowledged, Next resumes at offset 1. An unset checkpoint and a stored checkpoint of zero are different states. Do not initialize a new reader by storing zero.

Acknowledgments cover completed work, including records the server skipped. A primary page that examines records reports a safe offset, the last record it finished classifying. When you acknowledge a page, the server stores that offset, including the records it skipped. A consumer that matches one record in a million still moves forward.

Acknowledgment admission checks the source generation and current ownership, and a store that waits in the owner's queue keeps the history it was admitted under. A purge that commits meanwhile refuses it, so an old offset never lands in the new history. A revoked partition can drain accepted pages while its commit fence still permits it.

A permission change takes effect inside a page. The server checks source read permission before every read round of a page. A revocation between rounds fails the request with an authorization error and returns no partial page. The filter, its revision, and the group binding resolve 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 validates response identity and progress before exposure. 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 record-bearing pages 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. 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. A local page offers no acknowledgment offset, because a lagging replica can still hold records from before a purge, and the reader keeps no pending work for it.

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.

Checkpoints inside a page

Persist the application checkpoint before acknowledging its offset. Rust and Python expose ack_through(record). TypeScript exposes ackThrough(record). The call marks all preceding returned records on that partition as handled and stores the completed prefix. Process that full prefix first. Later records in the same page remain pending. Use this when a window receipt or transaction ends inside a page. Use ack_page or ackPage when the whole page is complete.

Quick example

Rust callers must enable the filters Cargo feature, which is also included by managed:

laser-sdk = { version = "0.5", features = ["filters"] }

Create Laser from a connection string for primary reads, so the SDK can open authenticated connections to the serving nodes. A caller-supplied raw client without connection settings supports Local diagnostics only, without acknowledgments.

This reader follows a satellite fleet. It selects a producer-reported mode change to safe, a decommission, or a telemetry report of safe mode. The complete examples publish typed records, a serde model in Rust, dataclasses in Python, and a discriminated union in TypeScript, and decode what arrives back into those types.

import { ConsumerFilter, FilterExpr } from "@laserdata/laser-sdk"

const satellites = () => FilterExpr.pred("table", "eq", "satellites")
const safeMode = ConsumerFilter.json(
  FilterExpr.any([
    FilterExpr.all([
      satellites(),
      FilterExpr.pred("op", "eq", "u"),
      FilterExpr.pred("changed", "contains", "mode"),
      FilterExpr.pred("after.mode", "eq", "safe")
    ]),
    FilterExpr.all([satellites(), FilterExpr.pred("op", "eq", "d")]),
    FilterExpr.all([
      FilterExpr.pred("event", "eq", "satellite.telemetry_changed"),
      FilterExpr.pred("fields.mode", "eq", "safe")
    ])
  ])
)

const reader = await laser
  .filters()
  .reader("orbit", "fleet_changes")
  .consumer("safe-mode-backfill")
  .inline(safeMode)
  .start({ kind: "first" })
  .build()

try {
  const record = await reader.nextRecord({ timeoutMs: 15_000 })
  // FleetChange is the record union of the complete example.
  const change = JSON.parse(new TextDecoder().decode(record.payload)) as FleetChange
  console.log(record.partitionId, record.offset, change)
  await reader.ack(record)
} finally {
  await reader.close()
}
use laser_sdk::filters::{ConsumerFilter, FilterExpr, FilterGroupRef, FilterRef, FilteredStart};
use laser_sdk::query::CmpOp;

let satellites = || FilterExpr::pred("table", CmpOp::Eq, "satellites");
let safe_mode = ConsumerFilter::json(FilterExpr::any([
    FilterExpr::all([
        satellites(),
        FilterExpr::pred("op", CmpOp::Eq, "u"),
        FilterExpr::pred("changed", CmpOp::Contains, "mode"),
        FilterExpr::pred("after.mode", CmpOp::Eq, "safe"),
    ]),
    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"),
    ]),
]));

let mut reader = laser
    .filters()
    .reader("orbit", "fleet_changes")
    .consumer("safe-mode-backfill")
    .inline(safe_mode.clone())
    .start(FilteredStart::First)
    .build()
    .await?;

let record = reader.next_record().await?;
// FleetChange is the serde model of the complete example.
let change: FleetChange = record.json()?;
println!("{} {} {change:?}", record.partition_id, record.offset);
reader.ack(&record).await?;
reader.close().await?;
import laser_sdk as ls

def satellites() -> ls.FilterExpr:
    return ls.FilterExpr.pred("table", "eq", "satellites")

safe_mode = ls.ConsumerFilter.json(
    ls.FilterExpr.any([
        ls.FilterExpr.all([
            satellites(),
            ls.FilterExpr.pred("op", "eq", "u"),
            ls.FilterExpr.pred("changed", "contains", "mode"),
            ls.FilterExpr.pred("after.mode", "eq", "safe"),
        ]),
        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"),
        ]),
    ])
)

reader = await laser.filters().reader(
    "orbit", "fleet_changes", consumer="safe-mode-backfill", filter=safe_mode, start="first"
)
try:
    record = await reader.next_record()
    print(record.partition_id, record.offset, record.message.json())
    await reader.ack(record)
finally:
    await reader.close()

Complete examples: Rust, Python, and TypeScript.

Records, pages, and ordinary consumers

Use next_record() for a record loop, or next_page() for batch processing. TypeScript uses nextRecord() and nextPage(). Both use the same filtered reader. The record API buffers poll results internally, so reading one record does not require a separate network poll.

A page is one bounded filtered poll result, with matching records and scan progress. It is not a separate transport or a CDC-only storage format. The optional AGDX command uses standard Iggy transport and leaves stored records unchanged.

A saved binding applies only to filtered reads. Join the bound group through filters().reader(...).group(...) in Rust and TypeScript, or filters().reader(..., group=...) in Python. A regular topic.consumer_group(...) continues standard Iggy polling even when that group has a saved filter binding.

For batch handling, process every returned record before calling ack_page (ackPage in TypeScript). For record handling, call ack(record) after processing. The SDK advances stored progress only through the completed prefix.

What a change record proves

Filters evaluate each record independently. 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 snapshot branch requires 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. The separate values-only filter selects it. Choose that 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 independently of mode. Keep these producer semantics explicit when adapting the example to your feed.

Test and preview

Check a filter before you run it. Neither call joins a group or stores an offset.

  • test judges one sample record you supply and explains every predicate.
  • preview judges a range of stored records in one partition and lists each verdict.
const tested = await laser.filters().test({ kind: "inline", filter: safeMode }, sample)
console.log(tested.explanation.verdict)

const preview = await laser
  .filters()
  .preview("orbit", "fleet_changes", 0, { kind: "inline", filter: safeMode }, { maxRecords: 10 })
console.log(preview.examined, preview.matched, preview.stop)
let tested = laser
    .filters()
    .test(FilterRef::Inline(safe_mode.clone()), sample, Vec::new())
    .await?;
println!("{}", tested.explanation.verdict);

let preview = laser
    .filters()
    .preview("orbit", "fleet_changes", 0, FilterRef::Inline(safe_mode.clone()))
    .max_records(10)
    .send()
    .await?;
println!("{} {} {}", preview.examined, preview.matched, preview.stop);
tested = await laser.filters().test(sample, filter=safe_mode)
print(tested["explanation"]["verdict"])

preview = await laser.filters().preview(
    "orbit", "fleet_changes", 0, filter=safe_mode, max_records=10
)
print(preview["examined"], preview["matched"], preview["stop"])

Filter on headers only

Filter opaque payloads without decoding them. 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.

import { HeaderValue } from "@laserdata/laser-sdk"

const critical = ConsumerFilter.headersOnly(FilterExpr.header("priority", "eq", 2))
const producer = laser.stream("orbit").topic("fleet_alerts").producer()
await producer.send(frame, { headers: { priority: HeaderValue.uint8(2) } })
use laser_sdk::iggy::prelude::{HeaderKey, HeaderValue};
use laser_sdk::stream::ProducerMessage;

let critical = ConsumerFilter::headers_only(FilterExpr::header("priority", CmpOp::Eq, 2_i32));
let producer = laser.stream("orbit").topic("fleet_alerts").producer().build().await?;
let message = ProducerMessage::new(frame)
    .header(HeaderKey::try_from("priority")?, HeaderValue::from(2_u8));
producer.send_message(message).await?;
critical = ls.ConsumerFilter.headers_only(ls.FilterExpr.header("priority", "eq", 2))
producer = laser.stream("orbit").topic("fleet_alerts").producer(partition=0)
await producer.send(frame, headers={"priority": ("uint8", 2)})

Header keys are custom text names matched exactly, including dots. Values can be booleans, signed or unsigned integers through 64 bits, floating-point numbers, or strings. Choose an exact-width Iggy header when space matters: uint8 uses one value byte. Numeric 2 and text "2" are different values. Use the streaming producer for typed headers. The publish().header(...) convenience accepts strings.

frame above is your encoded payload. Header-only selection works with JSON, Avro, Protobuf, or other bytes today, because it does not inspect the payload. Non-text keys and 128-bit numeric comparisons are outside the current filter contract.

One log, many kinds of events

Put every producer's events in one ordered log and let each service select its subdomain 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:

  1. Stamp a routing header on every record, such as event.type = "billing.invoice.v1.created", and the content type through the codec helpers, which set agdx.ct.
  2. Select the subdomain with a header text match, such as prefix billing., contains .v1., or a glob or regex. A headers_only filter never decodes a payload, so every format mixes safely, and it is the cheapest filter to run.
  3. Add payload conditions only where needed, inside all([header match, payload predicates]). Header children run first and all stops at the first no-match, so records of other kinds are rejected before any decode.
  4. Keep foreign_policy at reject. A record in another codec is then skipped instead of stalling the partition, even when its routing header is missing.
const billingV1 = ConsumerFilter.headersOnly(
  FilterExpr.headerText("event.type", "glob", "billing.*.v1.*")
)
const largeInvoices = ConsumerFilter.json(
  FilterExpr.all([
    FilterExpr.headerText("event.type", "prefix", "billing.invoice."),
    FilterExpr.pred("amount_cents", "gte", 100_000)
  ])
)
let billing_v1 = ConsumerFilter::headers_only(FilterExpr::header_text(
    "event.type",
    TextMatch::Glob,
    "billing.*.v1.*",
));
let large_invoices = ConsumerFilter::json(FilterExpr::all([
    FilterExpr::header_text("event.type", TextMatch::Prefix, "billing.invoice."),
    FilterExpr::pred("amount_cents", CmpOp::Gte, 100_000_i64),
]));
billing_v1 = ls.ConsumerFilter.headers_only(
    ls.FilterExpr.header_text("event.type", "glob", "billing.*.v1.*")
)
large_invoices = ls.ConsumerFilter.json(ls.FilterExpr.all([
    ls.FilterExpr.header_text("event.type", "prefix", "billing.invoice."),
    ls.FilterExpr.pred("amount_cents", "gte", 100_000),
]))

Text matching

KindMatches whenRelative cost
equalsthe whole value equals the patternlowest
prefixthe value starts with the patternlowest
suffixthe value ends with the patternlowest
containsthe pattern occurs anywherelow, one scan of the value
globthe whole value matches, * is any run, ? is one character, \ escapeslow to moderate with many stars
regexthe value contains a matchmoderate, linear in the value

Add case-insensitive matching to any kind. It lowercases the value first, which allocates a copy. The value must be text. A number, object, or array is a type mismatch and follows the mismatch policy.

On a JSON change record of about 600 bytes, one evaluation took about 1.3 to 1.5 microseconds for every text kind, the same as a plain equality test, because parsing the payload dominates. A header-only test took about 8 nanoseconds. Measure your own records with the SDK's filter_evaluation benchmark.

Regular expressions run in linear time on the server. A pattern cannot backtrack exponentially. Lookaround, backreferences, and inline flags such as (?i) are refused. The server's Rust engine decides pattern syntax. TypeScript performs structural validation and sends patterns to that engine. Use the case-insensitive option instead of (?i). Character classes such as \d and \w are Unicode-aware on the server. The TypeScript local guard rejects filters that contain a regex, because JavaScript's classes differ. Disable localGuard for these filters, or use Rust or Python for local verification. The Rust and Python guards run the server's evaluator. Case-insensitive regex uses the engine's Unicode folding. Other text predicates lowercase the pattern and value.

Saved filters and consumer groups

Save a policy once and share it across a consumer group. With laser-plane, you can save a filter in the catalog. A saved filter has a name, a description, and immutable revisions. An edit makes a new revision. Old revisions never change.

Bind a consumer group to one revision, and every member that uses filtered reads runs that revision. A member that brings a different filter is refused. The group reader needs no filter of its own, because the server applies the bound revision.

Filtering applies only to explicit filtered reads. Standard Apache Iggy polling retains its existing behavior. A filter is a delivery policy, not an access-control boundary. To change a group policy, create a new group and choose its starting position explicitly.

A saved filter has three states:

StateMeaning
activeCan be revised and bound
archivedHidden from new bindings. Groups already bound keep running it
droppedA tombstone. The id, the name, and every revision are never reused

A bound filter cannot be dropped. Release its groups first. Keep the binding returned by bind, or fetch it from the catalog. Rust and Python provide unbind_binding(binding), and TypeScript provides unbindBinding(binding). They check both its digest and its exact group identity. This also lets you clean up a stale binding after the group was deleted or recreated. The name-based unbind resolves the current group and falls back to stored names after deletion. Use the exact binding form when cleaning up historical groups. Unbinding does not permit that same group identity to adopt a different policy.

Every change carries an operation id. A retry by the same caller with the same id and mutation returns the first retained outcome. Reusing an id with another mutation or caller is rejected. The SDK waits for the outcome and reports a still-pending change with its operation id, and apply_as in Rust (applyAs in TypeScript) sends a change under an id you recorded first. There is no lifetime limit on filter-management operations. Results stay in durable, indexed plane storage, even after control-topic messages expire. Durable lookup preserves retry protection, also across 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 filter yet. Catalog lists page newest first. Pass the last id of a page as before_id to read the next page stably while filters are registered or dropped.

const filters = laser.filters()
const saved = await filters.register("sats-safe-mode", safeMode, {
  description: "Satellite mode changes, decommissions, and reported telemetry"
})

const group = { stream: "orbit", topic: "fleet_changes", group: "anomaly-desk" }
const binding = await filters.createConsumerGroup(group, saved.filterId, saved.revision)

const desk = await filters
  .reader("orbit", "fleet_changes")
  .groupId(binding.identity.groupId)
  .start({ kind: "first" })
  .build()
let filters = laser.filters();
let saved = filters
    .register("sats-safe-mode", safe_mode.clone(), "Satellite mode changes, decommissions, and reported telemetry")
    .await?;

let group = FilterGroupRef {
    stream: "orbit".to_owned(),
    topic: "fleet_changes".to_owned(),
    group: "anomaly-desk".to_owned(),
};
let binding = filters.create_consumer_group(group, saved.filter_id, saved.revision).await?;

let mut desk = filters
    .reader("orbit", "fleet_changes")
    .group_id(binding.identity.group_id)
    .start(FilteredStart::First)
    .build()
    .await?;
filters = laser.filters()
saved = await filters.register(
    "sats-safe-mode", safe_mode, description="Satellite mode changes, decommissions, and reported telemetry"
)

binding = await filters.create_consumer_group("orbit", "fleet_changes", "anomaly-desk", saved["filter_id"], saved["revision"])

desk = await filters.reader("orbit", "fleet_changes", group_id=binding["identity"]["group_id"], start="first")

A/B testing and revision controls

Create and bind once, then consume by group ID without repeating the filter. The setup helper creates a missing native group and binds it to a saved revision. These are two operations. If binding fails, the group can remain unbound. Retrying the same setup preserves an existing matching binding.

Use separate groups for A/B variants. Bind group A to revision 1 and group B to revision 2 of the same definition. Each group has independent offsets. Members of one group must use its bound revision. Supplying a different .revision(id, revision) returns conflict, because sharing offsets across different predicates can skip another variant's records.

set_revision_enabled(filter_id, revision, false) pauses new reads of a saved revision. TypeScript uses setRevisionEnabled. The response is revision_disabled. Existing delivered work can still be acknowledged, and tests and previews can still inspect the revision. Set the flag to true to resume. This flag applies to every group using that revision. It changes neither executable content nor its digest.

The complete CDC examples run two revisions through separate numeric group IDs and demonstrate pause, acknowledgment of in-flight work, and resume. A group ID is scoped to its topic. Deleting and recreating the same group name does not make a reader of the old ID join the replacement.

Stream UI lists definitions, revisions, and bound groups. Its workbench provides visual conditions and a JSON editor, explicit coercion, a producer-change requirement, typed-header sample tests, and bounded previews. Results are invalidated when their inputs change. Revision browsing includes pagination and copyable Rust, Python, TypeScript, and HTTP snippets. Catalog mutations retry a lost reply under the same operation ID.

Errors

Filter failures carry a typed reason. The ones a reader meets most often:

ReasonWhat it means
conflictA precondition failed. A different digest is bound, a name is taken, or the expected revision moved
source_changedThe source history or group identity changed. The reader restarts that partition from its stored offset once and raises this error so you know
revision_disabledThe saved revision is paused. New server reads stop until enabled. Already delivered records can still be acknowledged
membership_staleThe consumer no longer owns the partition. The SDK refreshes its assignment
not_primaryThe node is no longer the partition primary. The SDK resolves the route again
unsupportedThe server does not serve this operation, such as the catalog without laser-plane, or a filter its evaluator version or codec does not match
capacity_exhaustedA configured catalog resource limit is reached, such as definitions, revisions, group identities, or stored revision bytes. Stored operation results have no lifetime count cap
version_skewThe client and the server speak different filter op versions

A stop fault or a record larger than the page cap blocks its partition. The reader delivers every earlier match first, then raises a filter fault with the partition and offset.

Key operations

VerbWhat it does
filters().reader(stream, topic)Build a filtered reader
.consumer(name) / .group(name) / .group_id(id)Follow an independent consumer or a group by name or native ID
.inline(filter) / .revision(id, revision)Run a filter from the request or a saved revision
.partition(id) / .start(..)Pick partitions and the start position
next_page() / next_record()Read the next page or record
ack(record) / ack_page(page)Mark records done, which stores progress
filters().validate(filter)Check and compile a filter, returning its digest
filters().test(..) / preview(..)Judge a sample or stored records without storing state
create_consumer_groupCreate a missing group and bind its saved revision, returning the binding and group ID
set_revision_enabledPause or resume a saved revision without changing its content
register, revise, describeSave a filter, add a revision, change its description
bind, unbind, unbind_bindingPin a consumer group to a revision, then release the current group or an exact returned binding
bindings(filter, stream, topic, ..)List bindings, narrowed to one filter, stream, or topic
archive, deleteHide a filter from new bindings, then drop it

A local guard also checks records marked unevaluated. The server includes its decoder bounds so the guard can reproduce a size or depth fault. A marker without the applicable pass policy is rejected. Older servers without these bounds cannot verify pass-through decode faults locally.

Limits

LimitValue
Encoded filter size8 KiB
Sample header value1 KiB
Expression nodes128
Nesting depth8
Items in one in list64
Text pattern1 KiB
Compiled glob or regex programs per filter4
Budget per compiled program256 KiB
Records in one page1000
Unacknowledged record-bearing pages per partition1024 by default, configurable per reader
Bytes in one page8 MiB
Records one preview examines1,000 by default, configurable up to 10,000
Records one preview returns100
Live saved definitions1,024
Dropped definitions kept as tombstones4,096
Revisions per definition256
Retained group identities4,096
Stored revision bytes, all filters together64 MiB
Results cached in memory by default256
Concurrent mutations per plane by default32
Concurrent mutations per caller by default8

Catalog storage, memory use, and request admission are configured separately. The definition, tombstone, revision, and group-identity values above are defaults. Their 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 settings accept 0 for no limit. LD_PLANE_FILTER_OUTCOME_CACHE_ENTRIES=0 disables the cache while keeping durable retries.

LD_PLANE_FILTER_MAX_CONCURRENT_MUTATIONS and LD_PLANE_FILTER_MAX_CONCURRENT_MUTATIONS_PER_ACTOR set positive admission limits. When busy, the plane returns a retryable refusal. These limits count requests currently in progress, not historical operations.

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 topic retention, recovery uses the durable plane snapshot and remaining log. Expiring topic messages does not delete saved 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.

IGGY_PLANE_FILTERS_ENABLED turns filtering on for a deployment. IGGY_PLANE_FILTERS_MAX_CONCURRENT allows 16 concurrent scans per shard by default, and a scan past it gets a retryable refusal. IGGY_PLANE_FILTERS_MAX_EXAMINED_BYTES bounds the bytes one page scan may examine at 64 MiB. IGGY_PLANE_FILTERS_PREVIEW_MAX_EXAMINED sets how many records a preview may examine, 1,000 by default with a ceiling of 10,000.

IGGY_PLANE_FILTERS_MAX_ROUND_RECORDS caps one internal owner read at 1,000 records by default, the same ceiling an ordinary poll has, with a configurable range of 1 to 1,000. Fewer owner reads per page cost less server CPU: on a one-million-record trial the cap of 1,000 used about 19% less server CPU than 128 with the same results. After the first read of a page, the next read shrinks to what the remaining examined-byte budget covers at the average record size seen so far. The first owner read probes at most 128 records. A sudden increase in frame size can still exceed the byte budget by one owner read. The scan yields the server 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. Consumer-group reads also respect the remaining desired match count to preserve the native polling fence. IGGY_PLANE_FILTERS_MAX_PAGE_ROUNDS bounds owner reads per page (256 by default). IGGY_PLANE_FILTERS_MAX_ROUNDS separately bounds previews. A page also stops at its configured round and byte budgets, so a 1,000-record request can return fewer matches. Byte accounting includes record frames and is checked between owner reads, so a page can pass it by one read. Previews size their owner reads the same way and also check it before each record. These bounds limit work per page. They do not establish a throughput or latency guarantee.

Saved filters and group bindings reuse cached definitions. Each broker shard caches resolved policies and watches its own plane for catalog changes. Before reusing an entry, it confirms the current catalog version through the sidecar. A delayed watch cannot keep an old binding active after a completed local catalog change. Changes made on another node first reach the local plane through control-log replay. A plane restart invalidates prior version tokens. The server also checks current reader grants before reuse.

IGGY_PLANE_FILTERS_POLICY_CONFIRM_TIMEOUT_MS bounds each version confirmation at 1,000 ms by default. IGGY_PLANE_FILTERS_POLICY_CACHE_ENTRIES bounds entries per shard, with a default of 1,024 and a maximum of 65,536. IGGY_PLANE_FILTERS_POLICY_CACHE_TTL_MS sets the maximum entry lifetime, with a default of 60 seconds and a maximum of 10 minutes. Set either to zero to resolve every request. Entries hold weak references to compiled filters. Cache hits avoid repeated definition transfer and decoding, while a small version confirmation still runs for each hit. Watch failures and confirmation failures disable reuse. The cache confirmation counter separates those requests from full catalog lookups.

Payload formats and future codecs

Filter fields inside JSON, CBOR, Avro, and Protobuf without downloading every record. 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.

FormatFilter constructorValue rules
JSONConsumerFilter.json(expression)Exact 64-bit integers, finite numbers, arrays, objects, strings, booleans, and null
CBORConsumerFilter.cbor(expression)Definite-length values, text map keys, byte strings as integer arrays. No tags, duplicate keys, or non-finite numbers
AvroConsumerFilter.avro(expression, schema_ids)Registered writer schemas. Logical types use their physical values, such as integer timestamps or decimal bytes
ProtobufConsumerFilter.protobuf(expression, schema_ids)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 onlyConsumerFilter.headers_only(expression)Typed user headers, with no payload decoding

The table uses Python constructor names. Rust uses ConsumerFilter::avro(expression, schema_ids) and TypeScript uses ConsumerFilter.avro(expression, schemaIds). The CDC examples show complete typed producers and readers 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/Python or schemaId in TypeScript. This sets the existing agdx.sid header. A filter can allow several immutable writer schemas of the same codec. The schema IDs participate in its digest. A missing, unlisted, or incompatible schema ID follows the filter's foreign_policy when the predicate needs the payload: the record is skipped by default, or delivered unevaluated under pass.

The broker resolves schemas through the existing plane schema catalog before scanning records. Rust, Python, and TypeScript local guards use the same references. CBOR needs no registered schema. Protobuf groups are outside this profile. Empty repeated and map fields are absent. Use explicit timestamp coercions for Avro integer timestamps.

Byte and nesting limits bound decoding. The server announces its evaluator version and supported codecs. SDKs reject unsupported capabilities explicitly. Header-only routing works with opaque payloads in any format. Additional codecs can reuse the decoder boundary and predicate evaluator without changing Iggy transport or storage. Each new codec must define value semantics and pass the shared conformance corpus.

Running it

Filtered reads, tests, and previews run on the LaserData Iggy fork, in Laser Stack and LaserData Cloud. Saved filters and consumer group bindings also need laser-plane. Check capabilities().filters before you use them.

Binding lists can narrow by stream, topic, or both. A topic-only query matches that topic name across readable streams. In TypeScript, acknowledge the original reader-owned record or page. Copying its payload to another task does not transfer acknowledgment ownership.

On this page