Class: AlWorker

Inherits:
Object
  • Object
show all
Defined in:
lib/al_worker.rb

Overview

ワーカースーパークラス

Defined Under Namespace

Modules: Debug, IpcAction Classes: BroadcastMessage, Fd, Ipc, IpcClient, NumberedMessage, Program, StdoutTrap, Tcp, Timer

Constant Summary collapse

DEFAULT_WORKDIR =
"/tmp"
DEFAULT_NAME =
"al_worker"
LOG_SEVERITY =
{ :fatal=>Logger::FATAL, :error=>Logger::ERROR,
:warn=>Logger::WARN, :info=>Logger::INFO, :debug=>Logger::DEBUG }
@@log =

Returns ロガー.

Returns:

  • (Logger)

    ロガー

nil
@@mutex_sync =

Returns 同期実行用mutex.

Returns:

  • (Mutex)

    同期実行用mutex

Mutex.new

Instance Attribute Summary collapse

Class Method Summary collapse

Instance Method Summary collapse

Constructor Details

#initialize(name = nil, workdir = nil) ⇒ AlWorker

constructor

Parameters:

  • name (String) (defaults to: nil)

    識別名



212
213
214
215
216
217
218
219
220
221
222
223
# File 'lib/al_worker.rb', line 212

def initialize( name = nil, workdir = nil )
  @values = {}
  @values_rwlock = defined?(Sync) ? Sync.new : nil
  @name = name || DEFAULT_NAME
  @workdir = workdir || DEFAULT_WORKDIR
  @state = ""
  @pid_filename = File.join( @workdir, @name ) + ".pid"
  @main_queue = Queue.new

  Signal::trap( :HUP ) { signal_hup() }
  Signal::trap( :QUIT ) { signal_quit() }
end

Instance Attribute Details

#configHash, Nil

Returns 動作設定(コンフィグファイル読込結果).

Returns:

  • (Hash, Nil)

    動作設定(コンフィグファイル読込結果)



173
174
175
# File 'lib/al_worker.rb', line 173

def config
  @config
end

#log_filenameString

Returns ログファイル名(フルパス).

Returns:

  • (String)

    ログファイル名(フルパス)



188
189
190
# File 'lib/al_worker.rb', line 188

def log_filename
  @log_filename
end

#main_queueQueue

Returns メインスレッド動作依頼キュー.

Returns:

  • (Queue)

    メインスレッド動作依頼キュー



203
204
205
# File 'lib/al_worker.rb', line 203

def main_queue
  @main_queue
end

#nameString (readonly)

Returns ユニークネーム.

Returns:

  • (String)

    ユニークネーム



191
192
193
# File 'lib/al_worker.rb', line 191

def name
  @name
end

#pid_filenameString

Returns pidファイル名(フルパス).

Returns:

  • (String)

    pidファイル名(フルパス)



185
186
187
# File 'lib/al_worker.rb', line 185

def pid_filename
  @pid_filename
end

#privilegeString

Returns 実行権限ユーザ名.

Returns:

  • (String)

    実行権限ユーザ名



197
198
199
# File 'lib/al_worker.rb', line 197

def privilege
  @privilege
end

#program_nameString

Returns 現在実行中のRubyスクリプトの名前を表す文字列 $PROGRAM_NAME.

Returns:

  • (String)

    現在実行中のRubyスクリプトの名前を表す文字列 $PROGRAM_NAME



194
195
196
# File 'lib/al_worker.rb', line 194

def program_name
  @program_name
end

#stateString (readonly)

Returns ステート(ステートマシン用).

Returns:

  • (String)

    ステート(ステートマシン用)



200
201
202
# File 'lib/al_worker.rb', line 200

def state
  @state
end

#valuesHash

Returns 外部提供を目的とする値のHash IPCの関係でキーは文字列のみとする。.

Returns:

  • (Hash)

    外部提供を目的とする値のHash IPCの関係でキーは文字列のみとする。



176
177
178
# File 'lib/al_worker.rb', line 176

def values
  @values
end

#values_rwlockSync (readonly)

Returns @values の reader writer lock (require ‘sync’).

Returns:

  • (Sync)

    @values の reader writer lock (require ‘sync’)



179
180
181
# File 'lib/al_worker.rb', line 179

def values_rwlock
  @values_rwlock
end

#workdirString

Returns ワークファイルの作成場所.

Returns:

  • (String)

    ワークファイルの作成場所



182
183
184
# File 'lib/al_worker.rb', line 182

def workdir
  @workdir
end

Class Method Details

.log(*args) ⇒ Logger

ログ出力

Parameters:

  • msg (String, Object)

    エラーメッセージ

  • severity (Symbol)

    ログレベル :fatal, :error …

  • progname (String)

    プログラム名

Returns:

  • (Logger)

    Loggerオブジェクト



51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
# File 'lib/al_worker.rb', line 51

def self.log( *args )
  return nil  if ! @@log
  return @@log  if args.empty?

  msg,severity,progname = *args
  s = LOG_SEVERITY[ severity ]

  case msg
  when String
    @@log.add(s || Logger::INFO, msg, progname)

  when Exception
    @@log.add(s || Logger::ERROR, "#{msg.class} / #{msg.message}", progname)
    @@log.add(s || Logger::ERROR, "BACKTRACE: \n  " + msg.backtrace.join("\n  ") + "\n", progname)

  else
    @@log.add(s || Logger::INFO, msg.inspect, progname)
  end

  return @@log
end

.mutex_syncObject

同期実行用mutexのアクセッサ



37
38
39
# File 'lib/al_worker.rb', line 37

def self.mutex_sync()
  return @@mutex_sync
end

.na(method_name) ⇒ Object

Note:

ステートマシンで無視するイベントの記述

クラス定義中に、na :state_XXX_event_YYY の様に記述する。



123
124
125
# File 'lib/al_worker.rb', line 123

def self.na( method_name )
  define_method( method_name ) { |*args| }
end

.parse_request(req) ⇒ String, Hash

IPC定形リクエストからコマンドとパラメータを解析・取り出し

Parameters:

  • req (String)

    リクエスト

Returns:

  • (String)

    コマンド

  • (Hash)

    パラメータ



81
82
83
84
85
86
87
# File 'lib/al_worker.rb', line 81

def self.parse_request( req )
  (cmd,param) = req.split( " ", 2 )
  return cmd,{}  if param == nil
  param.strip!
  return cmd,{}  if param.empty?
  return cmd,( JSON.parse( param ) rescue { ""=>param } )
end

.realize_string(src, comment_pattern = /\s[;#]/) ⇒ String, NIl

文字列のコメントを取り必要に応じて変換する

Parameters:

  • src (String)

    source string

  • comment_pattern (Regexp) (defaults to: /\s[;#]/)

    comment pattern

Returns:

  • (String, NIl)

    result



135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
# File 'lib/al_worker.rb', line 135

def self.realize_string( src, comment_pattern = /\s[;#]/ )
  # ダブルクォートで括られている場合は文字列として取り出す
  if /^\s*"/ =~ src
    if /^\s*"(([^\\"]+|\\.)*)"/ =~ src
      s = $1
    else
      return nil
    end

  else
    # コメントの削除
    if comment_pattern && (len = comment_pattern =~ src)
      s = src[0, len].strip
    else
      s = src.strip
    end

    # 数値に変換するか?
    case s
    when /^[+-]?[\d]+$/
      return s.to_i
    when /^[+-]?([\d]+(\.[\d]*)?|\.[\d]+)([eE][+-]?[\d]+)?$/
      return s.to_f
    when /^0[xX][\h_]+$/
      return s.hex
    end
  end

  # エスケープされた文字があれば展開して返す
  return s.gsub(/\\(.)/) {
    ({"0"=>"\x00", "a"=>"\x07", "b"=>"\x08",
      "t"=>"\x09", "r"=>"\x0d", "n"=>"\x0a"}[$1]) || $1
  }
end

.reply(sock, st_code, st_msg, val = nil) ⇒ True

Note:

IPC定形リプライ

定形リプライフォーマット

 (ステータスコード) "200. Message"
 (JSONデータ)       { .... }
JSONデータは付与されない場合がある。
その判断は、ステータスコードの数字直後のピリオドの有無で行う。

Parameters:

  • sock (Socket)

    返信先ソケット

  • st_code (Integer)

    ステータスコード

  • st_msg (String)

    ステータスメッセージ

  • val (Hash) (defaults to: nil)

    リプライデータ

Returns:

  • (True)


105
106
107
108
109
110
111
112
113
114
# File 'lib/al_worker.rb', line 105

def self.reply( sock, st_code, st_msg, val = nil )
  sock.puts ("%03d" % st_code) + (val ? ". " : " ") + st_msg
  if val
    sock.puts val.to_json, ""
  end
  return true

rescue Errno::EPIPE
  Thread.exit
end

Instance Method Details

#append_default_option_to(opt) ⇒ Object

基本的なオプションの解析を、OptionParseオブジェクトへ追加

Parameters:

  • opt (OptionParse)

    OptionParseオブジェクト



259
260
261
262
263
264
265
266
267
268
269
270
271
272
# File 'lib/al_worker.rb', line 259

def append_default_option_to( opt )
  opt.on("-d", "--debug", "set debug mode.") { @flag_debug = true }
  opt.on("-k", "--kill", "kill stay process.") { @flag_kill = true }
  opt.on("-r", "--restart", "restart process.") { @flag_restart = true }
  opt.on("-p filename", "--pid=filename", "specify pid filename.") {|v|
    @pid_filename = v
  }
  opt.on("-l filename", "--log=filename", "specify log filename.") {|v|
    @log_filename = v
  }
  opt.on("-c filename", "--config=filename", "specify configfilename") {|v|
    @config_filename = v
  }
end

#daemonObject

デーモンになって実行



643
644
645
646
647
648
649
# File 'lib/al_worker.rb', line 643

def daemon()
  if @flag_debug
    run()
  else
    run( :daemon )
  end
end

#get_value(key) ⇒ Object

Note:

valueのゲッター タイムアウトなし(単一値)

値はdupして返す。

Parameters:

  • key (String)

    キー

Returns:

  • (Object)



433
434
435
436
437
438
439
440
# File 'lib/al_worker.rb', line 433

def get_value( key )
  @values_rwlock.synchronize( Sync::SH ) {
    return @values[ key.to_s ].dup rescue @values[ key.to_s ]
  }

rescue NameError =>ex
  raise "Need gem 'sync' and require 'sync' in your program."
end

#get_value_wt(key, timeout = 1) ⇒ Object, Boolean

Note:

valueのゲッター タイムアウト付き(単一値)

値はdupして返す。

Parameters:

  • key (String)

    キー

  • timeout (Numeric) (defaults to: 1)

    タイムアウト時間

Returns:

  • (Object)

  • (Boolean)

    ロック状態



476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
# File 'lib/al_worker.rb', line 476

def get_value_wt( key, timeout = 1 )
  locked = false
  (timeout * 10).times {
    locked = @values_rwlock.try_lock( Sync::SH )
    break if locked
    sleep 0.1
  }

  return (@values[ key.to_s ].dup rescue @values[ key.to_s ]), locked


rescue NameError =>ex
  raise "Need gem 'sync' and require 'sync' in your program."

ensure
  @values_rwlock.unlock( Sync::SH ) if locked
end

#get_values(keys) ⇒ Hash

Note:

valueのゲッター タイムアウトなし(複数値)

値はdupするが、簡素化のためにディープコピーは行っていない。文字列では問題ないが、配列などが格納されている場合は注意が必要。

Parameters:

  • keys (Array)

    キーの配列

Returns:

  • (Hash)



452
453
454
455
456
457
458
459
460
461
462
463
# File 'lib/al_worker.rb', line 452

def get_values( keys )
  ret = {}
  @values_rwlock.synchronize( Sync::SH ) {
    keys.each do |k|
      ret[ k.to_s ] = @values[ k.to_s ].dup rescue @values[ k.to_s ]
    end
  }
  return ret

rescue NameError =>ex
  raise "Need gem 'sync' and require 'sync' in your program."
end

#get_values_json(key = nil) ⇒ String

valueのゲッター JSON版 タイムアウトなし

Parameters:

  • key (String, Array) (defaults to: nil)

    取得する値のキー文字列

Returns:

  • (String)

    保存されている値のJSON文字列



534
535
536
537
538
539
540
541
542
543
544
545
546
# File 'lib/al_worker.rb', line 534

def get_values_json( key = nil )
  @values_rwlock.synchronize( Sync::SH ) {
    if key.class == Array
      ret = {}
      key.each { |k| ret[ k ] = @values[ k ] }
      return ret.to_json
    end
    return ( key ? { key => @values[key] } : @values ).to_json
  }

rescue NameError =>ex
  raise "Need gem 'sync' and require 'sync' in your program."
end

#get_values_json_wt(key = nil, timeout = nil) ⇒ String, Boolean

valuesのゲッター JSON版 タイムアウト付き

Parameters:

  • key (String, Array) (defaults to: nil)

    取得する値のキー文字列

  • timeout (Numeric) (defaults to: nil)

    タイムアウト時間

Returns:

  • (String)

    保存されている値のJSON文字列

  • (Boolean)

    ロック状態



557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
# File 'lib/al_worker.rb', line 557

def get_values_json_wt( key = nil, timeout = nil )
  locked = false
  timeout ||= 1       # can't change. see AlWorker::Ipc#ipc_a_get_values_wt()
  (timeout * 10).times {
    locked = @values_rwlock.try_lock( Sync::SH )
    break if locked
    sleep 0.1
  }
  if key.class == Array
    ret = {}
    key.each { |k| ret[ k ] = @values[ k ] }
    return ret.to_json, locked
  end
  return ( key ? { key => @values[key] } : @values ).to_json, locked

rescue NameError =>ex
  raise "Need gem 'sync' and require 'sync' in your program."

ensure
  @values_rwlock.unlock( Sync::SH ) if locked
end

#get_values_wt(keys, timeout = 1) ⇒ Object, Boolean

Note:

valueのゲッター タイムアウト付き(複数値)

値はdupするが、簡素化のためにディープコピーは行っていない。文字列では問題ないが、配列などが格納されている場合は注意が必要。

Parameters:

  • keys (Array)

    キーの配列

  • timeout (Numeric) (defaults to: 1)

    タイムアウト時間

Returns:

  • (Object)

  • (Boolean)

    ロック状態



506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
# File 'lib/al_worker.rb', line 506

def get_values_wt( keys, timeout = 1 )
  locked = false
  (timeout * 10).times {
    locked = @values_rwlock.try_lock( Sync::SH )
    break if locked
    sleep 0.1
  }

  ret = {}
  keys.each do |k|
    ret[ k.to_s ] = @values[ k.to_s ].dup rescue @values[ k.to_s ]
  end
  return ret, locked

rescue NameError =>ex
  raise "Need gem 'sync' and require 'sync' in your program."

ensure
  @values_rwlock.unlock( Sync::SH ) if locked
end

#initialize2Object

Note:

初期化2

常駐後に処理をさせるには、これをオーバライドする。



807
808
# File 'lib/al_worker.rb', line 807

def initialize2()
end

#load_values(filename = nil) ⇒ Object

値(@values)読み込み



608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
# File 'lib/al_worker.rb', line 608

def load_values( filename = nil )
  filename ||= File.join( @workdir, @name ) + ".values"
  digest = Digest::SHA1.file( filename ) rescue nil
  return nil if ! digest      # same as file not found.

  digestfile = File.join( File.dirname(filename), File.basename(filename,".*") ) + ".sha1"
  digestfile_value = File.read( digestfile ) rescue nil
  if digestfile_value
    return nil  if digest != digestfile_value
  end

  json = ""
  File.open( filename, "r" ) { |f|
    while txt = f.gets
      break if txt == "VALUES: \n"
    end
    if txt == "VALUES: \n"
      while txt = f.gets
        json << txt
      end
    end
  }
  return nil  if json == ""
  begin
    @values = JSON.parse( json )
    return true
  rescue
    return false
  end
end

#log(*args) ⇒ Object

ログ出力

See Also:

  • log()


816
817
818
# File 'lib/al_worker.rb', line 816

def log( *args )
  AlWorker.log( *args )
end

#no_method_error(event) ⇒ Object

メソッドエラーの場合のエラーハンドラ



871
872
873
# File 'lib/al_worker.rb', line 871

def no_method_error( event )
  raise "No action defined. state: #{@state}, event: #{event}"
end

#parse_option(argv = ARGV) ⇒ Object

基本的なオプションの解析

(OptionParseクラスを使わない場合に使用)

Parameters:

  • argv (Array<String>) (defaults to: ARGV)

    引数配列



232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
# File 'lib/al_worker.rb', line 232

def parse_option( argv = ARGV )
  i = 0
  while i < argv.size
    case argv[i]
    when "-d"                 # debug mode
      @flag_debug = true
    when "-k"                 # kill stay process.
      @flag_kill = true
    when "-r"                 # restart process.
      @flag_restart = true
    when "-p"                 # specify pid filename
      @pid_filename = argv[i += 1]
    when "-l"                 # specify log filename
      @log_filename = argv[i += 1]
    when "-c"                 # specify configfilename
      @config_filename = argv[i += 1]
    end
    i += 1
  end
end

#read_config(filename = nil) ⇒ Boolean, Nil

設定ファイルの読み込み

Parameters:

  • filename (String) (defaults to: nil)

    設定ファイル名

Returns:

  • (Boolean, Nil)

    エラー有無, 処理なしならnil



281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
# File 'lib/al_worker.rb', line 281

def read_config( filename = nil )
  # ファイル名の確定
  if filename
    @config_filename = File.expand_path( filename )
  elsif @config_filename
    @config_filename = File.expand_path( @config_filename )
  else
    fn = File.join( File.dirname( File.expand_path($0) ), @name )
    begin
      @config_filename = fn + ".ini"
      break  if File.exist?( @config_filename )

      @config_filename = fn + ".yaml"
      break  if File.exist?( @config_filename )

      @config_filename = nil
    end while false
  end

  case @config_filename
  when nil
    return nil

  when /\.ini$/
    # break case. to below.

  when /\.yaml$/
    @config = YAML.load_file( @config_filename )
    return true

  else
    warn "ERROR: Config file must be .ini or .yaml"
    return false
  end

  # iniファイルの読み込み
  config = {}
  section = nil
  flag_error = false
  file = File.open( @config_filename )
  while txt = file.gets
    case txt
    # key=value
    when /^\s*(\w+)\s*=(.*)$/
      s = AlWorker.realize_string($2)
      if !s
        warn "ERROR: #{@config_filename}:#{file.lineno}: Syntax error in value"
        flag_error = true
        next
      end
      if section
        config[section][$1.to_sym] = s
      else
        config[$1.to_sym] = s
      end

    # section
    when /^\s*\[(\w+)\]/
      config = {""=>config}  if !config.empty? && !section
      section = $1.to_sym
      config[section] ||= {}

    # comment or empty line.
    when /^\s*[;#]/, /^\s*$/
      # nothing to do.

    else
      warn "ERROR: #{@config_filename}:#{file.lineno}: Syntax error"
      flag_error = true
    end
  end
  file.close
  @config = config

  return !flag_error

rescue Psych::SyntaxError=>ex
  warn "ERROR: #{ex.message}"
  return false
end

#reply(sock, st_code, st_msg, val = nil) ⇒ Object

IPC定形リプライ

See Also:

  • reply()


826
827
828
# File 'lib/al_worker.rb', line 826

def reply( sock, st_code, st_msg, val = nil )
  AlWorker.reply( sock, st_code, st_msg, val )
end

#run(*modes) ⇒ Object

実行開始

Parameters:

  • modes (Symbol)

    動作モード nul デーモンにならずに実行:daemon デーモンで実行:nostop デーモンにならずスリープもしない:nopid プロセスIDファイルを作らない:nolog ログファイルを作らない:exit_idle_task アイドルタスクが終了したら

    プロセスも終了する
    


663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
# File 'lib/al_worker.rb', line 663

def run( *modes )
  # 実効権限変更(放棄)
  if @privilege
    uid = Etc.getpwnam( @privilege ).uid
    Process.uid = uid
    Process.euid = uid
  end

  # 停止 or 再実行?
  if @flag_kill || @flag_restart
    begin
      pid = File.read( @pid_filename ).to_i
      Process.kill( "TERM", pid )
    rescue Errno::ENOENT
      puts "Error: No pid file. '#{@pid_filename}'"
    rescue Errno::ESRCH
      puts "Error: No such pid=#{pid} process."
    rescue Errno::EPERM
      puts "Error: Operation not permitted for pid=#{pid} process."
    end

    exit(0)  if @flag_kill
    sleep 1
  end

  # 設定ファイル読み込み
  exit(1) if read_config() == false

  # ログ準備
  if !modes.include?(:nolog) && @@log == nil
    @log_filename ||= File.join( @workdir, @name ) + ".log"
    @@log = Logger.new( @log_filename, 3 )
    @@log.level = @flag_debug ? Logger::DEBUG : Logger::INFO
  end

  if ! modes.include?( :nopid )
    # 実行可/不可確認
    if File.directory?( @pid_filename )
      puts "ERROR: @pid_filename is directory."
      exit( 64 )
    end
    if File.exist?( @pid_filename )
      puts "ERROR: Still work."
      exit( 64 )
    end

    # プロセスIDファイル作成
    # (note) pid作成エラーの場合は、daemonになる前にここで検出される。
    File.open( @pid_filename, "w" ) { |file| file.write( Process.pid ) }
  end

  # 常駐処理
  if modes.include?( :daemon )
    Process.daemon()
    # プロセスIDファイル再作成
    if ! modes.include?( :nopid )
      File.open( @pid_filename, "w" ) { |file| file.write( Process.pid ) }
    end

    # stdout, stderrの差し替え
    if !modes.include?(:nolog)
      $stdout = StdoutTrap.new(:info)
      $stderr = StdoutTrap.new(:error)
    end
  end
  $PROGRAM_NAME = @program_name  if @program_name

  # 終了時処理
  at_exit {
    if ! modes.include?( :nopid )
      File.unlink( @pid_filename ) rescue 0
    end
    AlWorker.log( "finish", :info, @name )
  }

  # 初期化2
  AlWorker.log( "start", :info, @name )
  begin
    initialize2()
  rescue Exception => ex
    raise ex  if ex.class == SystemExit
    AlWorker.log( ex )
    raise ex  if STDERR.isatty
    exit( 64 )
  end

  # アイドルタスク
  if respond_to?( :idle_task, true )
    Thread.start {
      Thread.current.priority -= 1
      begin
        idle_task()
      rescue Exception => ex
        raise ex  if ex.class == SystemExit
        AlWorker.log( ex )
        if STDERR.isatty
          STDERR.puts ex.to_s
          STDERR.puts ex.backtrace.join("\n") + "\n"
        end
      end
      exit  if modes.include?( :exit_idle_task )
    }
  end

  # メインスレッド
  if modes.include?( :nostop )
    return
  else
    run_main_queue()
  end
end

#run_main_queueObject

メインスレッドでの実行依頼をキューで受けて実行する



779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
# File 'lib/al_worker.rb', line 779

def run_main_queue()
  while true
    # Rubyのデッドロック検出を避けつつ、
    # シグナルによる依頼に備えるためループを抜けないようにする。
    if Thread.list.size == 1 && @main_queue.empty?
      sleep 1
      next
    end

    req = @main_queue.pop
    case req
    when String
      log( req )
    when Proc
      req.call()
    else
      log("main_queue_process: #{req.inspect} is not a valid request.", :error)
    end
  end
end

#save_valuesObject

Note:

値(@values)保存

排他処理なし。バックアップファイルを3つまで作成する。



587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
# File 'lib/al_worker.rb', line 587

def save_values()
  filename = File.join( @workdir, @name ) + ".values"
  File.rename( filename + ".bak2", filename + ".bak3" ) rescue 0
  File.rename( filename + ".bak1", filename + ".bak2" ) rescue 0
  File.rename( filename,           filename + ".bak1" ) rescue 0

  File.open( filename, "w" ) { |f|
    f.puts "DATE: #{Time.now}"
    f.puts "NAME: #{@name}"
    f.puts "SELF: #{self.inspect}"
    f.puts "VALUES: \n#{@values.to_json}"
  }
  File.open( File.join( @workdir, @name ) + ".sha1", "w" ) { |file|
    file.write( Digest::SHA1.file( filename ) )
  }
end

#set_state(state) ⇒ Object Also known as: state=, next_state

現在のステートを宣言する

Parameters:

  • state (String)

    ステート文字列



881
882
883
884
# File 'lib/al_worker.rb', line 881

def set_state( state )
  @state = state.to_s
  AlWorker.log( "change state to #{@state}", :debug, @name )
end

#set_value(key, val) ⇒ Object

valueのセッター(単一値)

Parameters:

  • key (String)

    キー

  • val (Object)



404
405
406
407
408
409
# File 'lib/al_worker.rb', line 404

def set_value( key, val )
  @values_rwlock.synchronize( Sync::EX ) { @values[ key.to_s ] = val }

rescue NameError =>ex
  raise "Need gem 'sync' and require 'sync' in your program."
end

#set_values(values) ⇒ Object

valueのセッター(複数値)

Parameters:

  • values (Hash)

    セットする値



417
418
419
420
421
422
# File 'lib/al_worker.rb', line 417

def set_values( values )
  @values_rwlock.synchronize( Sync::EX ) { @values.merge!( values ) }

rescue NameError =>ex
  raise "Need gem 'sync' and require 'sync' in your program."
end

#signal_hupObject

Note:

シグナルハンドラ HUP

この実装はコンフィグファイルを読むだけだが、必要に応じてオーバライドして、処理変更のための仕組みを追加する。



370
371
372
373
374
375
# File 'lib/al_worker.rb', line 370

def signal_hup()
  @main_queue << Proc.new {
    log("Reload config file.")
    read_config()
  }
end

#signal_quitObject

Note:

シグナルハンドラ SIGQUIT

デバグ用

状態をファイルに書き出す。
画面があれば、表示する。


386
387
388
389
390
391
392
393
394
395
# File 'lib/al_worker.rb', line 386

def signal_quit()
  save_values()

  if STDOUT.isatty
    puts "\n===== @values ====="
    @values.keys.sort.each do |k|
      puts "#{k}=> #{@values[k]}"
    end
  end
end

#trigger_event(event, *args) ⇒ Object

ステートマシン 実行メソッド割り当て

Parameters:

  • event (String)

    イベント名

  • args (Array)

    引数



837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
# File 'lib/al_worker.rb', line 837

def trigger_event( event, *args )
  @respond_to = "from_#{@state}_event_#{event}"
  if respond_to?( @respond_to )
    AlWorker.log( "st:#{@state} ev:#{event} call:#{@respond_to}", :debug, @name )
    return __send__( @respond_to, *args )
  end

  @respond_to = "state_#{@state}_event_#{event}"
  if respond_to?( @respond_to )
    AlWorker.log( "st:#{@state} ev:#{event} call:#{@respond_to}", :debug, @name )
    return __send__( @respond_to, *args )
  end

  @respond_to = "event_#{event}"
  if respond_to?( @respond_to )
    AlWorker.log( "st:#{@state} ev:#{event} call:#{@respond_to}", :debug, @name )
    return __send__( @respond_to, *args )
  end

  @respond_to = "state_#{@state}"
  if respond_to?( @respond_to )
    AlWorker.log( "st:#{@state} ev:#{event} call:#{@respond_to}", :debug, @name )
    return __send__( @respond_to, *args )
  end

  # 実行すべきメソッドが見つからない場合
  @respond_to = ""
  no_method_error( event )
end