Skip to main content

Kafka

Reaches Driver8's Kafka capability the same way the Java client does: a WebDriver session with browserName: "pn5-driver8", then plain REST commands relative to that session (no CDP, no browser involved).

Status: the wire protocol has been cross-checked against a real, currently-passing Java client CI run (identical session flow and kafka_list command, real topics returned), and the NodeJS fixture/client wiring has been exercised locally end to end (up to the point of reaching the farm). It has not yet been run to completion against a live farm from this package's own CI or a tester's machine - if you're among the first to do so, please report back the result on DQA-122.

Usage​

import { test, expect } from '@pumpo5.dev/core'

test('lists Kafka topics', async ({ kafka }) => {
const result = await kafka.listTopics()
expect(result.topics.length).toBeGreaterThan(0)
})

The kafka fixture creates a Driver8 session before the test and closes it after, whether the test passes or fails.

listTopics() and read() return fluent response objects (KafkaListResponse/KafkaReadResponse), mirroring the Java client's own KafkaListResponse/KafkaReadResponse - same method names, same behavior. result.topics / result.messages remain plain, directly accessible arrays; the assertion/filtering methods below are additive. See Response assertions.

Connection configuration​

The kafka fixture resolves its connection details from a named profile, defaulting to KAFKA. For the default profile, values are looked up as, in order of precedence:

  1. config.conf, under the dotted key derived from the profile name (KAFKA → kafka):
    kafka.url="my-kafka-broker:9092"
    kafka.username="myUsername"
    kafka.password="myPassword"
  2. Environment variables named <PROFILE>_<SUFFIX>, e.g. KAFKA_URL, KAFKA_USERNAME, KAFKA_PASSWORD.

url, username and password are all required; the fixture throws with a descriptive error naming both possible sources if any is missing. These map directly onto the Java client's pn5:kafkaUrl, pn5:kafkaUsername and pn5:kafkaPassword capabilities (see Kafka capabilities reference).

Using a different profile​

Override the profile with test.use(), e.g. to target a second Kafka broker in the same test file:

test.describe('secondary broker', () => {
test.use({ kafkaProfile: 'KAFKA_ANALYTICS' })

test('lists topics on the analytics broker', async ({ kafka }) => {
const result = await kafka.listTopics()
expect(result.topics).toContain('analytics.events')
})
})

This reads kafka.analytics.* from config.conf / KAFKA_ANALYTICS_* environment variables instead of the KAFKA defaults. This is the closest NodeJS equivalent to how the Java client picks a different @Capability-annotated interface per test.

Note: a file-level test.skip(...) guard (used to skip cleanly when unconfigured) only evaluates the default profile's variables. If you add a test.use({ kafkaProfile: ... }) override in the same file, make sure that profile is configured too, or gate that specific test/describe block separately.

Method reference​

kafka.listTopics()​

Returns a KafkaListResponse (topics: string[], plus the assertion methods below). Lists topics visible to the configured user.

kafka.createTopic(topic, partitions?)​

Creates a topic. partitions defaults to the broker's own default if omitted.

kafka.read(query)​

interface KafkaReadQuery {
topic: string
partition?: number
offset?: number
consumerGroupId?: string
maxCount?: number
tailOffsets?: number // most recent N messages per partition
dateStart?: number // epoch millis
dateEnd?: number // epoch millis
headers?: Record<string, string>
}

Returns a KafkaReadResponse (messages: KafkaMessage[], plus the assertion/filtering methods below), where each KafkaMessage has partition, offset, key, value, timestamp and headers. This is the low-level form - pass a plain query object directly. For age/timestamp filters, kafka.queryTopic() (below) is usually more convenient, since it computes the epoch millis for you and validates the filter combination.

Example - read the last 10 messages per partition on orders.created and assert on their contents:

test('recently created orders contain the expected status', async ({ kafka }) => {
const { messages } = await kafka.read({ topic: 'orders.created', tailOffsets: 10 })

const orderIds = messages.map((m) => JSON.parse(m.value).orderId)
expect(orderIds).toContain(expectedOrderId)
})

To read only messages produced after the test started, pass dateStart (epoch millis) instead of tailOffsets:

const testStartedAt = Date.now()
// ... perform the action under test that should publish a message ...
const { messages } = await kafka.read({ topic: 'orders.created', dateStart: testStartedAt })

kafka.queryTopic(topic)​

Returns a KafkaReadQueryBuilder, mirroring the Java client's ReadQuery - same method names and behavior, including its age/timestamp computation and filter-combination validation. This is the fluent alternative to kafka.read(query) above; both end up sending the same request.

test('order-created events for a cancelled order look right', async ({ kafka }) => {
const result = await kafka
.queryTopic('orders.created')
.forAgedBetween(0, 60) // messages produced in the last 60s
.limitedToMaxMessageCountOf(50)
.getMessages()

result.assertThatMessageCountIsAtLeast(1)
})
  • onPartition(partition) - limits the query to a single partition.
  • atOffset(offset) - limits the query to a specific offset. Cannot be combined with any other filter (forLatestPerPartitionCount, forTimestampAfter/Before/Between, forYoungerThan, forOlderThan, forAgedBetween) - calling one after the other throws, in either order, same as the Java client.
  • limitedToMaxMessageCountOf(maxCount) - caps the total number of messages returned.
  • forLatestPerPartitionCount(tailOffsets) - returns only the latest tailOffsets messages per partition.
  • forTimestampAfter(millisStart) / forTimestampBefore(millisEnd) / forTimestampBetween(start, end) - filters by an explicit epoch-millis range.
  • forYoungerThan(seconds) / forOlderThan(ageInSeconds) / forAgedBetween(minAgeInSeconds, maxAgeInSeconds) - filters by age, computed from the current instant at call time (Date.now() - seconds * 1000) - you don't need to compute epoch millis yourself.
  • getMessages() - sends the built query and returns a KafkaReadResponse, same as kafka.read(query).

kafka.writeMessage(topic, message, options?)​

options is { key?, headers?, partitionId? }. Returns { offset, timestamp, serializedKeySize, serializedValueSize, partition }.

test('publishes a test order event', async ({ kafka }) => {
const result = await kafka.writeMessage(
'orders.created',
JSON.stringify({ orderId: '123', status: 'CREATED' }),
{ key: '123', headers: { 'content-type': 'application/json' } }
)

expect(result.partition).toBeGreaterThanOrEqual(0)
})

kafka.close()​

Closes the Driver8 session. Called automatically by the fixture after the test; only call this yourself if you created a KafkaAgent directly instead of using the fixture.

Response assertions​

listTopics() and read() return fluent wrapper objects with the same methods (names and behavior) as the Java client's KafkaListResponse/KafkaReadResponse. These are entirely optional - using Playwright's own expect() against .topics/.messages directly, as in the examples above, works just as well. The wrappers exist for parity with existing Java-authored test suites being ported, and because they read naturally in a fluent chain.

Failed assertions throw a Node AssertionError (via node:assert), which Playwright reports as a normal test failure with the assertion's message.

KafkaListResponse​

test('the expected topics exist, and a decommissioned one does not', async ({ kafka }) => {
const result = await kafka.listTopics()

result
.assertThatContainsTopicNamed('orders.created')
.assertThatContainsTopicNamed('orders.updated')
.assertThatDoesNotContainTopicNamed('orders.legacy')
})
  • assertThatContainsTopicNamed(topic) / assertThatDoesNotContainTopicNamed(topic) - return this, so calls chain.
  • printTopics() - logs each topic at debug level; for debugging a test, not for assertions.
  • andThen() - returns the underlying KafkaAgent, to continue the flow without breaking out of the fluent chain.

KafkaReadResponse​

test('order-created events for a cancelled order look right', async ({ kafka }) => {
const result = await kafka.read({ topic: 'orders.created', tailOffsets: 50 })

result
.assertThatMessageCountIsAtLeast(1)
.filterMessagesHaving<{ orderId: string; status: string }>((order) => order.orderId === expectedOrderId)
.assertThatMessageCountIsAtLeast(1)
.assertThatContainsMessageValue(JSON.stringify({ orderId: expectedOrderId, status: 'CREATED' }))
})
  • assertThatContainsMessageValue(value) - asserts at least one message's raw value matches exactly (case-insensitive).
  • assertThatMessageCountIsAtLeast(minCount) - asserts against the current (possibly filtered) message count.
  • filterMessagesHavingSubstring(substring) - narrows messages to those whose value contains substring.
  • filterMessagesHaving<T>(predicate) - narrows messages to those whose value, JSON-parsed as T, satisfies predicate. Messages that fail to parse as JSON are treated as not matching (dropped), the same way the Java client silently skips messages it can't deserialize into the given PoJo class.
  • messagesAsListOf<T>() - returns the (possibly filtered) messages' values, each JSON-parsed as T; unparseable ones are skipped. All other metadata (key, timestamp, offset, partition) is lost, same as the Java client's messagesAsListOf.
  • unfilter() - cancels any filterMessagesHaving* calls, restoring the original result set.
  • printMessages() - logs each message at debug level; for debugging a test, not for assertions.
  • andThen() - returns the underlying KafkaAgent.

All filter/assert methods return this, so they chain freely and in any order, same as the Java client.