Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions index.d.ts
Original file line number Diff line number Diff line change
Expand Up @@ -194,6 +194,11 @@ export class Message {
getProducerName(): string;
isReplicated(): boolean;
getReplicatedFrom(): string;
/**
* The version of the schema the producer used to serialize this message,
* or `null` if the message carries no schema version.
*/
getSchemaVersion(): number | null;
getEncryptionContext(): EncryptionContext | null;
}

Expand Down
14 changes: 14 additions & 0 deletions src/Message.cc
Original file line number Diff line number Diff line change
Expand Up @@ -65,6 +65,7 @@ Napi::Object Message::Init(Napi::Env env, Napi::Object exports) {
InstanceMethod("getProducerName", &Message::GetProducerName),
InstanceMethod("isReplicated", &Message::IsReplicated),
InstanceMethod("getReplicatedFrom", &Message::GetReplicatedFrom),
InstanceMethod("getSchemaVersion", &Message::GetSchemaVersion),
InstanceMethod("getEncryptionContext", &Message::GetEncryptionContext)});

constructor = Napi::Persistent(func);
Expand Down Expand Up @@ -197,6 +198,19 @@ Napi::Value Message::GetReplicatedFrom(const Napi::CallbackInfo &info) {
return Napi::String::New(env, replicatedFrom);
}

Napi::Value Message::GetSchemaVersion(const Napi::CallbackInfo &info) {
Napi::Env env = info.Env();
if (!ValidateCMessage(env)) {
return env.Null();
}

const pulsar::Message &message = this->cMessage.get()->message;
if (!message.hasSchemaVersion()) {
return env.Null();
}
return Napi::Number::New(env, message.getLongSchemaVersion());
}

Napi::Value Message::GetEncryptionContext(const Napi::CallbackInfo &info) {
Napi::Env env = info.Env();
if (!ValidateCMessage(env)) {
Expand Down
1 change: 1 addition & 0 deletions src/Message.h
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,7 @@ class Message : public Napi::ObjectWrap<Message> {
Napi::Value GetProducerName(const Napi::CallbackInfo &info);
Napi::Value IsReplicated(const Napi::CallbackInfo &info);
Napi::Value GetReplicatedFrom(const Napi::CallbackInfo &info);
Napi::Value GetSchemaVersion(const Napi::CallbackInfo &info);
Napi::Value GetRedeliveryCount(const Napi::CallbackInfo &info);
Napi::Value GetEncryptionContext(const Napi::CallbackInfo &info);
bool ValidateCMessage(Napi::Env env);
Expand Down
64 changes: 64 additions & 0 deletions tests/end_to_end.test.js
Original file line number Diff line number Diff line change
Expand Up @@ -178,6 +178,70 @@ const Pulsar = require('../index');
await client.close();
});

test('getSchemaVersion', async () => {
const client = new Pulsar.Client({
serviceUrl: 'pulsar://localhost:6650',
operationTimeoutSeconds: 30,
});

const schemaV1 = {
schemaType: 'Json',
schema: JSON.stringify({
type: 'record',
name: 'User',
fields: [{ name: 'name', type: 'string' }],
}),
};
// Compatible evolution of schemaV1: adds a field with a default.
const schemaV2 = {
schemaType: 'Json',
schema: JSON.stringify({
type: 'record',
name: 'User',
fields: [
{ name: 'name', type: 'string' },
{ name: 'age', type: 'int', default: 0 },
],
}),
};

const topic = `persistent://public/default/schema-version-${Date.now()}`;
const consumer = await client.subscribe({
topic,
subscription: 'sub1',
schema: schemaV1,
});

const producerV1 = await client.createProducer({ topic, schema: schemaV1 });
await producerV1.send({ data: Buffer.from(JSON.stringify({ name: 'alice' })) });
const msgV1 = await consumer.receive();
consumer.acknowledge(msgV1);

const producerV2 = await client.createProducer({ topic, schema: schemaV2 });
await producerV2.send({ data: Buffer.from(JSON.stringify({ name: 'bob', age: 30 })) });
const msgV2 = await consumer.receive();
consumer.acknowledge(msgV2);

expect(msgV1.getSchemaVersion()).toBe(0);
expect(msgV2.getSchemaVersion()).toBe(1);

const bytesTopic = `persistent://public/default/schema-version-bytes-${Date.now()}`;
const bytesConsumer = await client.subscribe({ topic: bytesTopic, subscription: 'sub1' });
const bytesProducer = await client.createProducer({ topic: bytesTopic });
await bytesProducer.send({ data: Buffer.from('no-schema') });
const bytesMsg = await bytesConsumer.receive();
bytesConsumer.acknowledge(bytesMsg);

expect(bytesMsg.getSchemaVersion()).toBeNull();

await producerV1.close();
await producerV2.close();
await consumer.close();
await bytesProducer.close();
await bytesConsumer.close();
await client.close();
});

test('Produce/Consume Listener', async () => {
const client = new Pulsar.Client({
serviceUrl: 'pulsar://localhost:6650',
Expand Down
1 change: 1 addition & 0 deletions tstest.ts
Original file line number Diff line number Diff line change
Expand Up @@ -309,6 +309,7 @@ import Pulsar = require('./index');
const redeliveryCount: number = message1.getRedeliveryCount();
const isReplicated: boolean = message1.isReplicated();
const replicatedFrom: string = message1.getReplicatedFrom();
const schemaVersion: number | null = message1.getSchemaVersion();
const partitionKey: string = message1.getPartitionKey();

const message3: Pulsar.Message = await reader1.readNext();
Expand Down