diff --git a/index.d.ts b/index.d.ts index 160fd52..c6311a2 100644 --- a/index.d.ts +++ b/index.d.ts @@ -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; } diff --git a/src/Message.cc b/src/Message.cc index 55fa7e3..07dfe04 100644 --- a/src/Message.cc +++ b/src/Message.cc @@ -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); @@ -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)) { diff --git a/src/Message.h b/src/Message.h index c84126c..fadba41 100644 --- a/src/Message.h +++ b/src/Message.h @@ -48,6 +48,7 @@ class Message : public Napi::ObjectWrap { 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); diff --git a/tests/end_to_end.test.js b/tests/end_to_end.test.js index d60be82..565d4fd 100644 --- a/tests/end_to_end.test.js +++ b/tests/end_to_end.test.js @@ -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', diff --git a/tstest.ts b/tstest.ts index c58f660..e8acda5 100644 --- a/tstest.ts +++ b/tstest.ts @@ -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();