00001
00002
00003
00004
00005
00006
00007
00008
00009
00010
00011
00012
00013
00014
00015
00016
00017
00018 import threading
00019
00020 import OpenRTM_aist
00021
00022
00023
00024
00025
00026
00027
00028
00029
00030
00031
00032
00033
00034
00035
00036
00037
00038
00039
00040
00041
00042
00043
00044
00045
00046 class PublisherNew(OpenRTM_aist.PublisherBase):
00047 """
00048 """
00049
00050
00051 ALL = 0
00052 FIFO = 1
00053 SKIP = 2
00054 NEW = 3
00055
00056
00057
00058
00059
00060
00061
00062
00063
00064
00065
00066
00067
00068
00069
00070 def __init__(self):
00071 self._rtcout = OpenRTM_aist.Manager.instance().getLogbuf("PublisherNew")
00072 self._consumer = None
00073 self._buffer = None
00074 self._task = None
00075 self._retcode = self.PORT_OK
00076 self._retmutex = threading.RLock()
00077 self._pushPolicy = self.NEW
00078 self._skipn = 0
00079 self._active = False
00080 self._leftskip = 0
00081 self._profile = None
00082 self._listeners = None
00083
00084
00085
00086
00087
00088
00089
00090
00091
00092
00093
00094
00095
00096 def __del__(self):
00097 self._rtcout.RTC_TRACE("~PublisherNew()")
00098 if self._task:
00099 self._task.resume()
00100 self._task.finalize()
00101
00102 OpenRTM_aist.PeriodicTaskFactory.instance().deleteObject(self._task)
00103 del self._task
00104 self._rtcout.RTC_PARANOID("task deleted.")
00105
00106
00107 self._consumer = 0
00108
00109 self._buffer = 0
00110 return
00111
00112
00113
00114
00115
00116
00117
00118
00119
00120 def setPushPolicy(self, prop):
00121 push_policy = prop.getProperty("publisher.push_policy","new")
00122 self._rtcout.RTC_DEBUG("push_policy: %s", push_policy)
00123
00124 push_policy = OpenRTM_aist.normalize([push_policy])
00125
00126 if push_policy == "all":
00127 self._pushPolicy = self.ALL
00128
00129 elif push_policy == "fifo":
00130 self._pushPolicy = self.FIFO
00131
00132 elif push_policy == "skip":
00133 self._pushPolicy = self.SKIP
00134
00135 elif push_policy == "new":
00136 self._pushPolicy = self.NEW
00137
00138 else:
00139 self._rtcout.RTC_ERROR("invalid push_policy value: %s", push_policy)
00140 self._pushPolicy = self.NEW
00141
00142 skip_count = prop.getProperty("publisher.skip_count","0")
00143 self._rtcout.RTC_DEBUG("skip_count: %s", skip_count)
00144
00145 skipn = [self._skipn]
00146 ret = OpenRTM_aist.stringTo(skipn, skip_count)
00147 if ret:
00148 self._skipn = skipn[0]
00149 else:
00150 self._rtcout.RTC_ERROR("invalid skip_count value: %s", skip_count)
00151 self._skipn = 0
00152
00153 if self._skipn < 0:
00154 self._rtcout.RTC_ERROR("invalid skip_count value: %d", self._skipn)
00155 self._skipn = 0
00156
00157 return
00158
00159
00160
00161
00162
00163
00164
00165
00166
00167 def createTask(self, prop):
00168 factory = OpenRTM_aist.PeriodicTaskFactory.instance()
00169
00170 th = factory.getIdentifiers()
00171 self._rtcout.RTC_DEBUG("available task types: %s", OpenRTM_aist.flatten(th))
00172
00173 self._task = factory.createObject(prop.getProperty("thread_type", "default"))
00174
00175 if not self._task:
00176 self._rtcout.RTC_ERROR("Task creation failed: %s",
00177 prop.getProperty("thread_type", "default"))
00178 return self.INVALID_ARGS
00179
00180 self._rtcout.RTC_PARANOID("Task creation succeeded.")
00181
00182 mprop = prop.getNode("measurement")
00183
00184
00185 self._task.setTask(self.svc)
00186 self._task.setPeriod(0.0)
00187 self._task.executionMeasure(OpenRTM_aist.toBool(mprop.getProperty("exec_time"),
00188 "enable", "disable", True))
00189 ecount = [0]
00190 if OpenRTM_aist.stringTo(ecount, mprop.getProperty("exec_count")):
00191 self._task.executionMeasureCount(ecount[0])
00192
00193 self._task.periodicMeasure(OpenRTM_aist.toBool(mprop.getProperty("period_time"),
00194 "enable", "disable", True))
00195 pcount = [0]
00196 if OpenRTM_aist.stringTo(pcount, mprop.getProperty("period_count")):
00197 self._task.periodicMeasureCount(pcount[0])
00198
00199 self._task.suspend()
00200 self._task.activate()
00201 self._task.suspend()
00202
00203 return self.PORT_OK
00204
00205
00206
00207
00208
00209
00210
00211
00212
00213
00214
00215
00216
00217
00218
00219
00220
00221
00222
00223
00224
00225
00226
00227
00228
00229
00230
00231
00232
00233
00234
00235
00236
00237
00238
00239
00240
00241
00242
00243
00244
00245
00246
00247
00248
00249
00250
00251
00252
00253
00254
00255
00256
00257 def init(self, prop):
00258 self._rtcout.RTC_TRACE("init()")
00259 self.setPushPolicy(prop)
00260 return self.createTask(prop)
00261
00262
00263
00264
00265
00266
00267
00268
00269
00270
00271
00272
00273
00274
00275
00276
00277
00278
00279
00280
00281
00282
00283
00284
00285
00286
00287
00288 def setConsumer(self, consumer):
00289 self._rtcout.RTC_TRACE("setConsumer()")
00290
00291 if not consumer:
00292 self._rtcout.RTC_ERROR("setConsumer(consumer = 0): invalid argument.")
00293 return self.INVALID_ARGS
00294
00295 self._consumer = consumer
00296 return self.PORT_OK
00297
00298
00299
00300
00301
00302
00303
00304
00305
00306
00307
00308
00309
00310
00311
00312
00313
00314
00315
00316
00317
00318
00319
00320
00321
00322
00323
00324 def setBuffer(self, buffer):
00325 self._rtcout.RTC_TRACE("setBuffer()")
00326
00327 if not buffer:
00328 self._rtcout.RTC_ERROR("setBuffer(buffer == 0): invalid argument")
00329 return self.INVALID_ARGS
00330
00331 self._buffer = buffer
00332 return self.PORT_OK
00333
00334
00335
00336
00337
00338
00339
00340
00341
00342
00343
00344
00345
00346
00347
00348
00349
00350
00351
00352
00353
00354
00355
00356
00357
00358
00359
00360
00361
00362
00363
00364
00365
00366
00367
00368
00369 def setListener(self, info, listeners):
00370 self._rtcout.RTC_TRACE("setListener()")
00371
00372 if not listeners:
00373 self._rtcout.RTC_ERROR("setListeners(listeners == 0): invalid argument")
00374 return self.INVALID_ARGS
00375
00376 self._profile = info
00377 self._listeners = listeners
00378
00379 return self.PORT_OK
00380
00381
00382
00383
00384
00385
00386
00387
00388
00389
00390
00391
00392
00393
00394
00395
00396
00397
00398
00399
00400
00401
00402
00403
00404
00405
00406
00407
00408
00409
00410
00411
00412
00413
00414
00415
00416
00417
00418
00419
00420
00421
00422
00423
00424
00425
00426
00427
00428
00429
00430
00431
00432
00433
00434
00435
00436
00437
00438
00439
00440
00441
00442
00443
00444
00445
00446
00447
00448
00449
00450
00451
00452
00453
00454
00455
00456
00457
00458
00459 def write(self, data, sec, usec):
00460 self._rtcout.RTC_PARANOID("write()")
00461
00462 if not self._consumer or not self._buffer or not self._listeners:
00463 return self.PRECONDITION_NOT_MET
00464
00465 if self._retcode == self.CONNECTION_LOST:
00466 self._rtcout.RTC_DEBUG("write(): connection lost.")
00467 return self._retcode
00468
00469 if self._retcode == self.SEND_FULL:
00470 self._rtcout.RTC_DEBUG("write(): InPort buffer is full.")
00471 ret = self._buffer.write(data, sec, usec)
00472 self._task.signal()
00473 return self.BUFFER_FULL
00474
00475
00476 assert(self._buffer != 0)
00477
00478 self.onBufferWrite(data)
00479 ret = self._buffer.write(data, sec, usec)
00480
00481 self._task.signal()
00482 self._rtcout.RTC_DEBUG("%s = write()", OpenRTM_aist.DataPortStatus.toString(ret))
00483
00484 return self.convertReturn(ret, data)
00485
00486
00487
00488
00489
00490
00491
00492
00493
00494
00495
00496
00497
00498
00499
00500
00501
00502
00503
00504
00505
00506
00507
00508
00509
00510
00511
00512
00513
00514 def isActive(self):
00515 return self._active
00516
00517
00518
00519
00520
00521
00522
00523
00524
00525
00526
00527
00528
00529
00530
00531
00532
00533
00534
00535
00536
00537
00538
00539
00540
00541
00542
00543 def activate(self):
00544 self._active = True
00545 return self.PORT_OK
00546
00547
00548
00549
00550
00551
00552
00553
00554
00555
00556
00557
00558
00559
00560
00561
00562
00563
00564
00565
00566
00567
00568
00569
00570
00571
00572
00573 def deactivate(self):
00574 self._active = False;
00575 return self.PORT_OK
00576
00577
00578
00579
00580
00581
00582
00583
00584
00585
00586
00587
00588
00589
00590
00591 def svc(self):
00592 guard = OpenRTM_aist.ScopedLock(self._retmutex)
00593
00594 if self._pushPolicy == self.ALL:
00595 self._retcode = self.pushAll()
00596 return 0
00597 elif self._pushPolicy == self.FIFO:
00598 self._retcode = self.pushFifo()
00599 return 0
00600 elif self._pushPolicy == self.SKIP:
00601 self._retcode = self.pushSkip()
00602 return 0
00603 elif self._pushPolicy == self.NEW:
00604 self._retcode = self.pushNew()
00605 return 0
00606 else:
00607 self._retcode = self.pushNew()
00608
00609 return 0
00610
00611
00612
00613
00614
00615 def pushAll(self):
00616 self._rtcout.RTC_TRACE("pushAll()")
00617 try:
00618
00619 while self._buffer.readable() > 0:
00620 cdr = self._buffer.get()
00621 self.onBufferRead(cdr)
00622
00623 self.onSend(cdr)
00624 ret = self._consumer.put(cdr)
00625
00626 if ret != self.PORT_OK:
00627 self._rtcout.RTC_DEBUG("%s = consumer.put()", OpenRTM_aist.DataPortStatus.toString(ret))
00628 return self.invokeListener(ret, cdr)
00629 self.onReceived(cdr)
00630
00631 self._buffer.advanceRptr()
00632
00633 return self.PORT_OK
00634 except:
00635 self._rtcout.RTC_ERROR(OpenRTM_aist.Logger.print_exception())
00636 return self.CONNECTION_LOST
00637
00638 return self.PORT_ERROR
00639
00640
00641
00642
00643
00644 def pushFifo(self):
00645 self._rtcout.RTC_TRACE("pushFifo()")
00646
00647 try:
00648 cdr = self._buffer.get()
00649 self.onBufferRead(cdr)
00650
00651 self.onSend(cdr)
00652 ret = self._consumer.put(cdr)
00653
00654 if ret != self.PORT_OK:
00655 self._rtcout.RTC_DEBUG("%s = consumer.put()", OpenRTM_aist.DataPortStatus.toString(ret))
00656 return self.invokeListener(ret, cdr)
00657 self.onReceived(cdr)
00658
00659 self._buffer.advanceRptr()
00660
00661 return self.PORT_OK
00662 except:
00663 self._rtcout.RTC_ERROR(OpenRTM_aist.Logger.print_exception())
00664 return self.CONNECTION_LOST
00665
00666 return self.PORT_ERROR
00667
00668
00669
00670
00671
00672 def pushSkip(self):
00673 self._rtcout.RTC_TRACE("pushSkip()")
00674 try:
00675 ret = self.PORT_OK
00676 preskip = self._buffer.readable() + self._leftskip
00677 loopcnt = preskip/(self._skipn+1)
00678 postskip = self._skipn - self._leftskip
00679
00680 for i in range(loopcnt):
00681 self._buffer.advanceRptr(postskip)
00682 cdr = self._buffer.get()
00683 self.onBufferRead(cdr)
00684
00685 self.onSend(cdr)
00686 ret = self._consumer.put(cdr)
00687 if ret != self.PORT_OK:
00688 self._buffer.advanceRptr(-postskip)
00689 self._rtcout.RTC_DEBUG("%s = consumer.put()", OpenRTM_aist.DataPortStatus.toString(ret))
00690 return self.invokeListener(ret, cdr)
00691
00692 self.onReceived(cdr)
00693 postskip = self._skipn + 1
00694
00695 self._buffer.advanceRptr(self._buffer.readable())
00696
00697 if loopcnt == 0:
00698
00699 self._leftskip = preskip % (self._skipn + 1)
00700 else:
00701 if self._retcode != self.PORT_OK:
00702
00703 self._leftskip = 0
00704 else:
00705
00706 self._leftskip = preskip % (self._skipn + 1)
00707
00708 return ret
00709
00710 except:
00711 self._rtcout.RTC_ERROR(OpenRTM_aist.Logger.print_exception())
00712 return self.CONNECTION_LOST
00713
00714 return self.PORT_ERROR
00715
00716
00717
00718
00719
00720 def pushNew(self):
00721 self._rtcout.RTC_TRACE("pushNew()")
00722 try:
00723 self._buffer.advanceRptr(self._buffer.readable() - 1)
00724
00725 cdr = self._buffer.get()
00726 self.onBufferRead(cdr)
00727
00728 self.onSend(cdr)
00729 ret = self._consumer.put(cdr)
00730
00731 if ret != self.PORT_OK:
00732 self._rtcout.RTC_DEBUG("%s = consumer.put()", OpenRTM_aist.DataPortStatus.toString(ret))
00733 return self.invokeListener(ret, cdr)
00734
00735 self.onReceived(cdr)
00736 self._buffer.advanceRptr()
00737
00738 return self.PORT_OK
00739
00740 except:
00741 self._rtcout.RTC_ERROR(OpenRTM_aist.Logger.print_exception())
00742 return self.CONNECTION_LOST
00743
00744 return self.PORT_ERROR
00745
00746
00747
00748
00749
00750
00751
00752
00753
00754
00755
00756
00757
00758
00759
00760
00761
00762
00763
00764
00765
00766
00767
00768
00769
00770
00771
00772
00773
00774
00775
00776
00777
00778
00779
00780
00781
00782
00783
00784
00785
00786
00787
00788
00789
00790
00791
00792
00793
00794
00795
00796
00797
00798
00799
00800
00801
00802
00803
00804 def convertReturn(self, status, data):
00805
00806
00807
00808
00809
00810
00811
00812
00813
00814
00815 if status == OpenRTM_aist.BufferStatus.BUFFER_OK:
00816 return self.PORT_OK
00817
00818 elif status == OpenRTM_aist.BufferStatus.BUFFER_ERROR:
00819 return self.BUFFER_ERROR
00820
00821 elif status == OpenRTM_aist.BufferStatus.BUFFER_FULL:
00822 self.onBufferFull(data)
00823 return self.BUFFER_FULL
00824
00825 elif status == OpenRTM_aist.BufferStatus.NOT_SUPPORTED:
00826 return self.PORT_ERROR
00827
00828 elif status == OpenRTM_aist.BufferStatus.TIMEOUT:
00829 self.onBufferWriteTimeout(data)
00830 return self.BUFFER_TIMEOUT
00831
00832 elif status == OpenRTM_aist.BufferStatus.PRECONDITION_NOT_MET:
00833 return self.PRECONDITION_NOT_MET
00834
00835 else:
00836 return self.PORT_ERROR
00837
00838
00839
00840
00841
00842
00843
00844
00845
00846
00847
00848
00849
00850
00851
00852
00853
00854
00855
00856
00857
00858 def invokeListener(self, status, data):
00859
00860
00861
00862 if status == self.PORT_ERROR:
00863 self.onReceiverError(data)
00864 return self.PORT_ERROR
00865
00866 elif status == self.SEND_FULL:
00867 self.onReceiverFull(data)
00868 return self.SEND_FULL
00869
00870 elif status == self.SEND_TIMEOUT:
00871 self.onReceiverTimeout(data)
00872 return self.SEND_TIMEOUT
00873
00874 elif status == self.CONNECTION_LOST:
00875 self.onReceiverError(data)
00876 return self.CONNECTION_LOST
00877
00878 elif status == self.UNKNOWN_ERROR:
00879 self.onReceiverError(data)
00880 return self.UNKNOWN_ERROR
00881
00882 else:
00883 self.onReceiverError(data)
00884 return self.PORT_ERROR
00885
00886
00887
00888
00889
00890
00891
00892
00893
00894
00895
00896 def onBufferWrite(self, data):
00897 if self._listeners is not None and self._profile is not None:
00898 self._listeners.connectorData_[OpenRTM_aist.ConnectorDataListenerType.ON_BUFFER_WRITE].notify(self._profile, data)
00899 return
00900
00901
00902
00903
00904
00905
00906
00907
00908
00909
00910
00911 def onBufferFull(self, data):
00912 if self._listeners is not None and self._profile is not None:
00913 self._listeners.connectorData_[OpenRTM_aist.ConnectorDataListenerType.ON_BUFFER_FULL].notify(self._profile, data)
00914 return
00915
00916
00917
00918
00919
00920
00921
00922
00923
00924
00925
00926 def onBufferWriteTimeout(self, data):
00927 if self._listeners is not None and self._profile is not None:
00928 self._listeners.connectorData_[OpenRTM_aist.ConnectorDataListenerType.ON_BUFFER_WRITE_TIMEOUT].notify(self._profile, data)
00929 return
00930
00931
00932
00933
00934
00935
00936
00937
00938
00939
00940
00941 def onBufferWriteOverwrite(self, data):
00942 if self._listeners is not None and self._profile is not None:
00943 self._listeners.connectorData_[OpenRTM_aist.ConnectorDataListenerType.ON_BUFFER_OVERWRITE].notify(self._profile, data)
00944 return
00945
00946
00947
00948
00949
00950
00951
00952
00953
00954
00955
00956 def onBufferRead(self, data):
00957 if self._listeners is not None and self._profile is not None:
00958 self._listeners.connectorData_[OpenRTM_aist.ConnectorDataListenerType.ON_BUFFER_READ].notify(self._profile, data)
00959 return
00960
00961
00962
00963
00964
00965
00966
00967
00968
00969
00970
00971 def onSend(self, data):
00972 if self._listeners is not None and self._profile is not None:
00973 self._listeners.connectorData_[OpenRTM_aist.ConnectorDataListenerType.ON_SEND].notify(self._profile, data)
00974 return
00975
00976
00977
00978
00979
00980
00981
00982
00983
00984
00985
00986 def onReceived(self, data):
00987 if self._listeners is not None and self._profile is not None:
00988 self._listeners.connectorData_[OpenRTM_aist.ConnectorDataListenerType.ON_RECEIVED].notify(self._profile, data)
00989 return
00990
00991
00992
00993
00994
00995
00996
00997
00998
00999
01000
01001 def onReceiverFull(self, data):
01002 if self._listeners is not None and self._profile is not None:
01003 self._listeners.connectorData_[OpenRTM_aist.ConnectorDataListenerType.ON_RECEIVER_FULL].notify(self._profile, data)
01004 return
01005
01006
01007
01008
01009
01010
01011
01012
01013
01014
01015
01016 def onReceiverTimeout(self, data):
01017 if self._listeners is not None and self._profile is not None:
01018 self._listeners.connectorData_[OpenRTM_aist.ConnectorDataListenerType.ON_RECEIVER_TIMEOUT].notify(self._profile, data)
01019 return
01020
01021
01022
01023
01024
01025
01026
01027
01028
01029
01030
01031 def onReceiverError(self, data):
01032 if self._listeners is not None and self._profile is not None:
01033 self._listeners.connectorData_[OpenRTM_aist.ConnectorDataListenerType.ON_RECEIVER_ERROR].notify(self._profile, data)
01034 return
01035
01036
01037
01038
01039
01040
01041
01042
01043
01044
01045
01046 def onSenderError(self):
01047 if self._listeners is not None and self._profile is not None:
01048 self._listeners.connector_[OpenRTM_aist.ConnectorListenerType.ON_SENDER_ERROR].notify(self._profile)
01049 return
01050
01051
01052
01053 def PublisherNewInit():
01054 OpenRTM_aist.PublisherFactory.instance().addFactory("new",
01055 OpenRTM_aist.PublisherNew,
01056 OpenRTM_aist.Delete)