Skip to content
Merged
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
42 changes: 42 additions & 0 deletions doclib/msgpack/factory.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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 #<the_object>.
#
# @param obj [Object] object to serialize
# @return [String] serialized object
#
def dump(object)
end
end
end
end
10 changes: 9 additions & 1 deletion ext/java/org/msgpack/jruby/Factory.java
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand All @@ -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);
Expand Down
2 changes: 1 addition & 1 deletion ext/java/org/msgpack/jruby/Packer.java
Original file line number Diff line number Diff line change
Expand Up @@ -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);
}
Expand Down
37 changes: 37 additions & 0 deletions ext/msgpack/factory_class.c
Original file line number Diff line number Diff line change
Expand Up @@ -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);

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I squeezed an extra optimization here. When we freeze the factory, the types registry hash gets frozen as well. This way all packers can share it safely.

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);
Expand Down Expand Up @@ -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);
Expand Down
10 changes: 8 additions & 2 deletions ext/msgpack/packer_ext_registry.c
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
103 changes: 103 additions & 0 deletions lib/msgpack/factory.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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
52 changes: 52 additions & 0 deletions spec/factory_spec.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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