From: Eric Hodel Date: 2005-09-01T11:01:11+09:00 Subject: Re: Rogue threads found! The plot thickens.... --Apple-Mail-14--943503632 Content-Transfer-Encoding: 7bit Content-Type: text/plain; charset=US-ASCII; delsp=yes; format=flowed On 31 Aug 2005, at 14:51, Eric Hodel wrote: > On 30 Aug 2005, at 16:40, Townley, Andrew wrote: > >> I'm open to suggestions, but I can't just do it without the timer >> because I need to get control back every n seconds so I can do things >> like graceful shutdown, etc. > > If you need safe concurrent access like a Queue and timeouts you > may find rinda/tuplespace.rb useful. Somewhere around here I have > a stream implementation for it (but its not terribly difficult to > write from scratch). I think it can be modified to have timeouts > on pop. > >> Here's the full test program (not out to win any style awards with >> the >> calls to read, btw) :) > > Give me a bit and I think I can make your test program with with a > TupleSpace streams. $ cat test.rb require 'ts_stream' ts = Rinda::TupleSpace.new 1 stream = Rinda::Stream.new ts, 1 def read(stream, timeout) puts "READ: #{stream.length} elements" puts "*** #{stream.pop timeout} ***" rescue Rinda::Stream::ClosedError puts "CLOSED: #{stream.length} elements" rescue Rinda::Stream::TimeoutError puts "TIMEOUT: #{stream.length} elements" end 6.times { read stream, 1 } stream.push "one" 6.times { read stream, 1 } p Thread.list $ ruby test.rb READ: 0 elements TIMEOUT: 0 elements READ: 0 elements TIMEOUT: 0 elements READ: 0 elements TIMEOUT: 0 elements READ: 0 elements TIMEOUT: 0 elements READ: 0 elements TIMEOUT: 0 elements READ: 0 elements TIMEOUT: 0 elements READ: 1 elements *** one *** READ: 0 elements TIMEOUT: 0 elements READ: 0 elements TIMEOUT: 0 elements READ: 0 elements TIMEOUT: 0 elements READ: 0 elements TIMEOUT: 0 elements READ: 0 elements TIMEOUT: 0 elements [#, #] --Apple-Mail-14--943503632 Content-Transfer-Encoding: 7bit Content-Type: text/x-ruby-script; x-unix-mode=0644; name="ts_stream.rb" Content-Disposition: attachment; filename=ts_stream.rb require 'rinda/tuplespace' ## # streams will be stored as [:stream, stream_index, position, value] class Rinda::Stream ## # Raised when attempting to pop or push onto a closed Stream. class ClosedError < RuntimeError; end ## # Raised when a Stream operation times out. class TimeoutError < RuntimeError; end ## # The number of this stream. attr_reader :stream_id ## # Creates a new stream on +ts+ def initialize(ts, stream_id) @ts = ts @stream_id = stream_id begin head_tup = @ts.read [:stream, @stream_id, :head, nil], true rescue Rinda::RequestExpiredError @ts.write [:stream, @stream_id, :head, 0] @ts.write [:stream, @stream_id, :tail, 0] end end def length tail_index = @ts.read([:stream, @stream_id, :tail, nil]).last head_index = @ts.read([:stream, @stream_id, :head, nil]).last length = tail_index - head_index return length < 0 ? 0 : length end def push(value) index = @ts.take([:stream, @stream_id, :tail, nil]).last if index.nil? then @ts.write [:stream, @stream_id, :tail, nil] raise ClosedError end @ts.write [:stream, @stream_id, :tail, index + 1] @ts.write [:stream, @stream_id, index, value] return value end ## # Removes the oldest value from the stream. +sec+ can be a time-to-wait in # seconds or a Renewer object (see Rinda documentation). #pop can wait up # to twice the value of +sec+. def pop(sec = nil) # Find the head begin index = @ts.take([:stream, @stream_id, :head, nil], sec).last rescue Rinda::RequestExpiredError raise TimeoutError end # Check if the stream is closed begin last = @ts.read [:stream, @stream_id, :closed, nil], true rescue Rinda::RequestExpiredError last = nil else last = last.last end if index == last then @ts.write [:stream, @stream_id, :head, index] raise ClosedError end # Grab our value begin value = @ts.take([:stream, @stream_id, index, nil], sec).last rescue Rinda::RequestExpiredError @ts.write [:stream, @stream_id, :head, index] raise TimeoutError else @ts.write [:stream, @stream_id, :head, index + 1] end return value end ## # Closes the stream for any further writing. Reading may continue until the # stream is empty. An already blocket #get will remain blocked. def close index = @ts.take([:stream, @stream_id, :tail, nil]).last @ts.write [:stream, @stream_id, :tail, nil] @ts.write [:stream, @stream_id, :closed, index] return nil end end if __FILE__ == $0 then require 'test/unit' class Rinda::StreamTest < Test::Unit::TestCase def setup @ts = Rinda::TupleSpace.new 1 @stream = Rinda::Stream.new @ts, 0 end def test_create assert @stream, "Stream must not be nil" end def test_put_get @stream.push 0 @stream.push 1 assert 0, @stream.pop assert 1, @stream.pop end def test_length assert_equal 0, @stream.length @stream.push 0 assert_equal 1, @stream.length @stream.push 1 assert_equal 2, @stream.length assert 0, @stream.pop assert_equal 1, @stream.length assert 1, @stream.pop assert_equal 0, @stream.length begin @stream.pop true rescue Rinda::Stream::TimeoutError end assert_equal 0, @stream.length end def test_close @stream.push 0 @stream.close assert_raises Rinda::Stream::ClosedError do @stream.push 1 end assert_equal 0, @stream.pop, "Finish reading from the stream" assert_raises Rinda::Stream::ClosedError do @stream.pop end end def test_pop_timeout assert_raises Rinda::Stream::TimeoutError do @stream.pop true end assert_raises Rinda::Stream::TimeoutError do @stream.pop 1 end @stream.push :value assert_nothing_raised do assert_equal :value, @stream.pop(1) end end def test_concurrency_but_not_really push_threads = [] pop_threads = [] value_one = nil value_two = nil push_threads << Thread.start do @stream.push 0 end push_threads << Thread.start do @stream.push 1 end pop_threads << Thread.start do value_one = @stream.pop end pop_threads << Thread.start do value_two = @stream.pop end push_threads.each do |t| t.join end @stream.close pop_threads.each do |t| t.join end assert_equal [0, 1], [value_one, value_two].sort end end end --Apple-Mail-14--943503632 Content-Transfer-Encoding: 7bit Content-Type: text/plain; charset=US-ASCII; format=flowed -- Eric Hodel - drbrain@segment7.net - http://segment7.net FEC2 57F1 D465 EB15 5D6E 7C11 332A 551C 796C 9F04 --Apple-Mail-14--943503632--