From 5f9e9f4f1e65ab558ae743905164d87e02948a51 Mon Sep 17 00:00:00 2001 From: Jean Boussier Date: Mon, 14 Feb 2022 16:03:46 +0100 Subject: [PATCH] Implement Factory#pool Fix: https://github.com/msgpack/msgpack-ruby/issues/255 When handling small payload creating the Packer and Unpacker instance can easily account for the majority of the serialization and deserialization time. For this reason it is recommanded to reuse these objects, and `Factory#pool` is a helper for that. --- doclib/msgpack/factory.rb | 42 ++++++++++ ext/java/org/msgpack/jruby/Factory.java | 10 ++- ext/java/org/msgpack/jruby/Packer.java | 2 +- ext/msgpack/factory_class.c | 37 +++++++++ ext/msgpack/packer_ext_registry.c | 10 ++- lib/msgpack/factory.rb | 103 ++++++++++++++++++++++++ spec/factory_spec.rb | 52 ++++++++++++ 7 files changed, 252 insertions(+), 4 deletions(-) diff --git a/doclib/msgpack/factory.rb b/doclib/msgpack/factory.rb index 685ef8a9..8a4e8a43 100644 --- a/doclib/msgpack/factory.rb +++ b/doclib/msgpack/factory.rb @@ -98,5 +98,47 @@ def registered_types(selector=:both) # def type_registered?(klass_or_type, selector=:both) end + + # + # Creates a MessagePack::PooledFactory instance of the given size. + # + # PooledFactory keeps Packer and Unpacker instance in a pool for improved performance. + # Note that the size defines how many instances are kept in cache, not the maximum of instances + # that can be created. If the pool limit is reached, a new instance is created anyway. + # + # @param size [Fixnum] specify how many Packer and Unpacker to keep in cache (default 1) + # @param options [Hash] Combined options for Packer and Unpacker. See Packer#initialize and Unpacker#initialize + # for supported options. + def pool(size=1, **options) + end + + class Pool + # + # Deserializes an object from the string or io and returns it. + # + # If there're not enough data to deserialize one object, this method raises EOFError. + # If data format is invalid, this method raises MessagePack::MalformedFormatError. + # If the object nests too deeply, this method raises MessagePack::StackError. + # + # @param data [String] + # @return [Object] deserialized object + # + # See Unpacker#initialize for supported options. + # + def load(data) + end + + # + # Serialize the passed value + # + # If it could not serialize the object, it raises + # NoMethodError: undefined method `to_msgpack' for #. + # + # @param obj [Object] object to serialize + # @return [String] serialized object + # + def dump(object) + end + end end end diff --git a/ext/java/org/msgpack/jruby/Factory.java b/ext/java/org/msgpack/jruby/Factory.java index bfa73227..5e2eaa61 100644 --- a/ext/java/org/msgpack/jruby/Factory.java +++ b/ext/java/org/msgpack/jruby/Factory.java @@ -25,7 +25,7 @@ public class Factory extends RubyObject { private static final long serialVersionUID = 8441284623445322492L; private final Ruby runtime; - private final ExtensionRegistry extensionRegistry; + private ExtensionRegistry extensionRegistry; private boolean hasSymbolExtType; private boolean hasBigIntExtType; @@ -52,6 +52,14 @@ public IRubyObject initialize(ThreadContext ctx) { return this; } + @JRubyMethod(name = "dup") + public IRubyObject dup() { + Factory clone = (Factory)super.dup(); + clone.extensionRegistry = extensionRegistry(); + clone.hasSymbolExtType = hasSymbolExtType; + return clone; + } + @JRubyMethod(name = "packer", optional = 2) public Packer packer(ThreadContext ctx, IRubyObject[] args) { return Packer.newPacker(ctx, extensionRegistry(), hasSymbolExtType, hasBigIntExtType, args); diff --git a/ext/java/org/msgpack/jruby/Packer.java b/ext/java/org/msgpack/jruby/Packer.java index de86f60a..7a383014 100644 --- a/ext/java/org/msgpack/jruby/Packer.java +++ b/ext/java/org/msgpack/jruby/Packer.java @@ -267,7 +267,7 @@ public IRubyObject size(ThreadContext ctx) { return buffer.size(ctx); } - @JRubyMethod(name = "clear") + @JRubyMethod(name = "clear", alias = { "reset" }) public IRubyObject clear(ThreadContext ctx) { return buffer.clear(ctx); } diff --git a/ext/msgpack/factory_class.c b/ext/msgpack/factory_class.c index ba1ef21e..612e4a4f 100644 --- a/ext/msgpack/factory_class.c +++ b/ext/msgpack/factory_class.c @@ -88,6 +88,41 @@ static VALUE Factory_initialize(int argc, VALUE* argv, VALUE self) return Qnil; } +static VALUE Factory_dup(VALUE self) +{ + VALUE clone = Factory_alloc(rb_obj_class(self)); + + FACTORY(self, fc); + FACTORY(clone, cloned_fc); + + cloned_fc->has_symbol_ext_type = fc->has_symbol_ext_type; + cloned_fc->pkrg = fc->pkrg; + msgpack_unpacker_ext_registry_borrow(fc->ukrg, &cloned_fc->ukrg); + msgpack_packer_ext_registry_dup(&fc->pkrg, &cloned_fc->pkrg); + + return clone; +} + +static VALUE Factory_freeze(VALUE self) { + if(!rb_obj_frozen_p(self)) { + FACTORY(self, fc); + + if (RTEST(fc->pkrg.hash)) { + rb_hash_freeze(fc->pkrg.hash); + if (!RTEST(fc->pkrg.cache)) { + // If the factory is frozen, we can safely share the packer cache between + // all packers. So we eagerly create it now so it's available when #packer + // is called. + fc->pkrg.cache = rb_hash_new(); + } + } + + rb_obj_freeze(self); + } + + return self; +} + VALUE MessagePack_Factory_packer(int argc, VALUE* argv, VALUE self) { FACTORY(self, fc); @@ -233,6 +268,8 @@ void MessagePack_Factory_module_init(VALUE mMessagePack) rb_define_alloc_func(cMessagePack_Factory, Factory_alloc); rb_define_method(cMessagePack_Factory, "initialize", Factory_initialize, -1); + rb_define_method(cMessagePack_Factory, "dup", Factory_dup, 0); + rb_define_method(cMessagePack_Factory, "freeze", Factory_freeze, 0); rb_define_method(cMessagePack_Factory, "packer", MessagePack_Factory_packer, -1); rb_define_method(cMessagePack_Factory, "unpacker", MessagePack_Factory_unpacker, -1); diff --git a/ext/msgpack/packer_ext_registry.c b/ext/msgpack/packer_ext_registry.c index abf1fc41..d01c0a37 100644 --- a/ext/msgpack/packer_ext_registry.c +++ b/ext/msgpack/packer_ext_registry.c @@ -43,8 +43,14 @@ void msgpack_packer_ext_registry_mark(msgpack_packer_ext_registry_t* pkrg) void msgpack_packer_ext_registry_dup(msgpack_packer_ext_registry_t* src, msgpack_packer_ext_registry_t* dst) { - dst->hash = RTEST(src->hash) ? rb_hash_dup(src->hash) : Qnil; - dst->cache = RTEST(src->cache) ? rb_hash_dup(src->cache): Qnil; + if(RTEST(src->hash) && !rb_obj_frozen_p(src->hash)) { + dst->hash = rb_hash_dup(src->hash); + dst->cache = RTEST(src->cache) ? rb_hash_dup(src->cache) : Qnil; + } else { + // If the type registry is frozen we can safely share it, and share the cache as well. + dst->hash = src->hash; + dst->cache = src->cache; + } } VALUE msgpack_packer_ext_registry_put(msgpack_packer_ext_registry_t* pkrg, diff --git a/lib/msgpack/factory.rb b/lib/msgpack/factory.rb index 4ce2243c..6606967e 100644 --- a/lib/msgpack/factory.rb +++ b/lib/msgpack/factory.rb @@ -77,5 +77,108 @@ def dump(v, *rest) packer.full_pack end alias :pack :dump + + def pool(size = 1, **options) + Pool.new( + frozen? ? self : dup.freeze, + size, + options.empty? ? nil : options, + ) + end + + class Pool + if RUBY_ENGINE == "ruby" + class AbstractPool + def initialize(size, &block) + @size = size + @new_member = block + @members = [] + end + + def checkout + @members.pop || @new_member.call + end + + def checkin(member) + # If the pool is already full, we simply drop the extra member. + # This is because contrary to a connection pool, creating an extra instance + # is extremely unlikely to cause some kind of resource exhaustion. + # + # We could cycle the members (keep the newer one) but first It's more work and second + # the older member might have been created pre-fork, so it might be at least partially + # in shared memory. + if member && @members.size < @size + member.reset + @members << member + end + end + end + else + class AbstractPool + def initialize(size, &block) + @size = size + @new_member = block + @members = [] + @mutex = Mutex.new + end + + def checkout + @mutex.synchronize { @members.pop } || @new_member.call + end + + def checkin(member) + @mutex.synchronize do + if member && @members.size < @size + member.reset + @members << member + end + end + end + end + end + + class PackerPool < AbstractPool + private + + def reset(packer) + packer.clear + end + end + + class UnpackerPool < AbstractPool + private + + def reset(unpacker) + unpacker.reset + end + end + + def initialize(factory, size, options = nil) + options = nil if !options || options.empty? + @factory = factory + @packers = PackerPool.new(size) { factory.packer(options) } + @unpackers = UnpackerPool.new(size) { factory.unpacker(options) } + end + + def load(data) + unpacker = @unpackers.checkout + begin + unpacker.feed_reference(data) + unpacker.full_unpack + ensure + @unpackers.checkin(unpacker) + end + end + + def dump(object) + packer = @packers.checkout + begin + packer.write(object) + packer.full_pack + ensure + @packers.checkin(packer) + end + end + end end end diff --git a/spec/factory_spec.rb b/spec/factory_spec.rb index d98ca557..833d8006 100644 --- a/spec/factory_spec.rb +++ b/spec/factory_spec.rb @@ -474,4 +474,56 @@ class << Symbol expect(MessagePack.unpack(MessagePack.pack(dm2))).to eq(dm2) end end + + describe '#pool' do + let(:factory) { described_class.new } + + it 'responds to serializers interface' do + pool = factory.pool(1) + expect(pool.load(pool.dump(42))).to be == 42 + end + + it 'types can be registered before the pool is created' do + factory.register_type(0x00, Symbol) + pool = factory.pool(1) + expect(pool.load(pool.dump(:foo))).to be == :foo + end + + it 'types cannot be registered after the pool is created' do + pool = factory.pool(1) + factory.register_type(0x20, ::MyType) + + expect do + pool.dump(MyType.new(1, 2)) + end.to raise_error NoMethodError + + payload = factory.dump(MyType.new(1, 2)) + expect do + pool.load(payload) + end.to raise_error MessagePack::UnknownExtTypeError + end + + it 'support symbolize_keys: true' do + pool = factory.pool(1, symbolize_keys: true) + expect(pool.load(pool.dump('foo' => 1))).to be == { foo: 1 } + end + + it 'support freeze: true' do + pool = factory.pool(1, freeze: true) + expect(pool.load(pool.dump('foo'))).to be_frozen + end + + it 'is thread safe' do + pool = factory.pool(1) + + threads = 10.times.map do + Thread.new do + 1_000.times do |i| + expect(pool.load(pool.dump(i))).to be == i + end + end + end + threads.each(&:join) + end + end end