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_listcommand, 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:
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"- 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 atest.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 latesttailOffsetsmessages 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 aKafkaReadResponse, same askafka.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)- returnthis, so calls chain.printTopics()- logs each topic at debug level; for debugging a test, not for assertions.andThen()- returns the underlyingKafkaAgent, 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 rawvaluematches exactly (case-insensitive).assertThatMessageCountIsAtLeast(minCount)- asserts against the current (possibly filtered) message count.filterMessagesHavingSubstring(substring)- narrowsmessagesto those whose value containssubstring.filterMessagesHaving<T>(predicate)- narrowsmessagesto those whose value, JSON-parsed asT, satisfiespredicate. 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 asT; unparseable ones are skipped. All other metadata (key, timestamp, offset, partition) is lost, same as the Java client'smessagesAsListOf.unfilter()- cancels anyfilterMessagesHaving*calls, restoring the original result set.printMessages()- logs each message at debug level; for debugging a test, not for assertions.andThen()- returns the underlyingKafkaAgent.
All filter/assert methods return this, so they chain freely and in any order, same as the Java client.