001package jmri.jmrix.loconet; 002 003import java.util.concurrent.DelayQueue; 004import java.util.concurrent.Delayed; 005import java.util.concurrent.TimeUnit; 006import javax.annotation.Nonnull; 007import org.slf4j.Logger; 008import org.slf4j.LoggerFactory; 009 010/** 011 * Delay LocoNet messages that need to be throttled. 012 * <p> 013 * A LocoNetThrottledTransmitter object sits in front of a LocoNetInterface 014 * (e.g. TrafficHandler) and meters out specific LocoNet messages. 015 * 016 * <p> 017 * The internal Memo class is used to hold the pending message and the time it's 018 * to be sent. Time computations are in units of milliseconds, as that's all the 019 * accuracy that's needed here. 020 * 021 * @author Bob Jacobsen Copyright (C) 2009 022 */ 023public class LocoNetThrottledTransmitter implements LocoNetInterface { 024 025 public LocoNetThrottledTransmitter(@Nonnull LocoNetInterface controller, boolean mTurnoutExtraSpace) { 026 this.controller = controller; 027 this.memo = controller.getSystemConnectionMemo(); 028 this.mTurnoutExtraSpace = mTurnoutExtraSpace; 029 030 // calculation is needed time to send on DCC: 031 // msec*nBitsInPacket*packetRepeat/bitRate*safetyFactor 032 minInterval = 1000 * (18 + 3 * 10) * 3 / 16000 * 2; 033 034 if (mTurnoutExtraSpace) { 035 minInterval = minInterval * 4; 036 } 037 038 attachServiceThread(); 039 } 040 041 /** 042 * Reference to the system connection memo. 043 */ 044 LocoNetSystemConnectionMemo memo = null; 045 046 /** 047 * Set the system connection memo associated with this traffic controller. 048 * 049 * @param m associated systemConnectionMemo object 050 */ 051 @Override 052 public void setSystemConnectionMemo(LocoNetSystemConnectionMemo m) { 053 log.debug("LnTrafficController set memo to {}", m.getUserName()); 054 memo = m; 055 } 056 057 /** 058 * Get the system connection memo associated with this traffic controller. 059 * 060 * @return the associated systemConnectionMemo object 061 */ 062 @Override 063 public LocoNetSystemConnectionMemo getSystemConnectionMemo() { 064 log.debug("getSystemConnectionMemo {} called in LnTC", memo.getUserName()); 065 return memo; 066 } 067 068 boolean mTurnoutExtraSpace; 069 070 /** 071 * Request that server thread cease operation, no more messages can be sent. 072 * Note that this returns before the thread is known to be done if it still 073 * has work pending. If you need to be sure it's done, check and wait on 074 * !running. 075 */ 076 public void dispose() { 077 disposed = true; 078 079 // put a shutdown request on the queue after any existing 080 Memo m = new Memo(null, nowMSec(), TimeUnit.MILLISECONDS, false) { 081 @Override 082 boolean requestsShutDown() { 083 return true; 084 } 085 }; 086 queue.add(m); 087 } 088 089 volatile boolean disposed = false; 090 volatile boolean running = false; 091 092 // interface being shadowed 093 LocoNetInterface controller; 094 095 // Forward methods to underlying interface 096 @Override 097 public void addLocoNetListener(int mask, LocoNetListener listener) { 098 controller.addLocoNetListener(mask, listener); 099 } 100 101 @Override 102 public void removeLocoNetListener(int mask, LocoNetListener listener) { 103 controller.removeLocoNetListener(mask, listener); 104 } 105 106 @Override 107 public boolean status() { 108 return controller.status(); 109 } 110 111 /** 112 * Accept a message to be sent after suitable delay. 113 */ 114 @Override 115 public void sendLocoNetMessage(LocoNetMessage msg) { 116 sendLocoNetMessage(msg, false); 117 } 118 119 /** 120 * Accept a message to be sent after suitable delay. 121 * 122 * @param msg Message to send; will be updated with CRC 123 * @param requestIgnoreEcho If true: Notify listeners on enqueing message, ignore echo from line. 124 * Only in effect if preference "LoconetUpdateSlotOnMessageCreation" is set. 125 */ 126 @Override 127 public void sendLocoNetMessage(LocoNetMessage msg, boolean requestIgnoreEcho) { 128 if (disposed) { 129 log.error("Message sent after queue disposed"); 130 return; 131 } 132 133 long sendTime = calcSendTimeMSec(); 134 135 Memo m = new Memo(msg, sendTime, TimeUnit.MILLISECONDS, requestIgnoreEcho); 136 queue.add(m); 137 138 } 139 140 // minimum time in msec between messages 141 long minInterval; 142 143 long lastSendTimeMSec = 0; 144 145 long calcSendTimeMSec() { 146 // next time is at least now or minInterval after latest so far 147 lastSendTimeMSec = Math.max(nowMSec(), minInterval + lastSendTimeMSec); 148 return lastSendTimeMSec; 149 } 150 151 DelayQueue<Memo> queue = new DelayQueue<Memo>(); 152 153 /** 154 * Constant for the name of the Service Thread. 155 * Requires the connection UserName prepending. 156 */ 157 public static final String SERVICE_THREAD_NAME = " LocoNetThrottledTransmitter"; 158 159 private void attachServiceThread() { 160 theServiceThread = new ServiceThread(); 161 theServiceThread.setPriority(Thread.NORM_PRIORITY); 162 theServiceThread.setName( memo.getUserName() + SERVICE_THREAD_NAME); 163 theServiceThread.setDaemon(true); 164 theServiceThread.start(); 165 } 166 167 ServiceThread theServiceThread; 168 169 class ServiceThread extends Thread { 170 171 @Override 172 public void run() { 173 running = true; 174 while (true) { 175 try { 176 Memo m = queue.take(); 177 178 // check for request to shutdown 179 if (m.requestsShutDown()) { 180 log.debug("item requests shutdown"); 181 break; 182 } 183 184 // normal request 185 if (log.isDebugEnabled()) { 186 log.debug("forwarding message: {}", m.getMessage()); 187 } 188 controller.sendLocoNetMessage(m.getMessage(), m.getRequestIgnoreEcho()); 189 // and go round again 190 } catch (InterruptedException e) { 191 // request to terminate 192 this.interrupt(); 193 break; 194 } 195 } 196 running = false; 197 } 198 } 199 200 // a separate method to ease testing by stopping clock 201 static long nowMSec() { 202 return System.currentTimeMillis(); 203 } 204 205 static class Memo implements Delayed { 206 207 public Memo(LocoNetMessage msg, long endTime, TimeUnit unit, boolean requestIgnoreEcho) { 208 this.msg = msg; 209 this.endTimeMsec = unit.toMillis(endTime); 210 this.requestIgnoreEcho = requestIgnoreEcho; 211 } 212 213 LocoNetMessage getMessage() { 214 return msg; 215 } 216 217 boolean requestsShutDown() { 218 return false; 219 } 220 221 boolean getRequestIgnoreEcho() { 222 return requestIgnoreEcho; 223 } 224 225 long endTimeMsec; 226 LocoNetMessage msg; 227 boolean requestIgnoreEcho; 228 229 @Override 230 public long getDelay(TimeUnit unit) { 231 long delay = endTimeMsec - nowMSec(); 232 return unit.convert(delay, TimeUnit.MILLISECONDS); 233 } 234 235 @Override 236 public int compareTo(Delayed d) { 237 // -1 means this is less than m 238 long delta; 239 if (d instanceof Memo) { 240 delta = this.endTimeMsec - ((Memo)d).endTimeMsec; 241 } else { 242 delta = this.getDelay(TimeUnit.MILLISECONDS) 243 - d.getDelay(TimeUnit.MILLISECONDS); 244 } 245 if (delta > 0) { 246 return 1; 247 } else if (delta < 0) { 248 return -1; 249 } else { 250 return 0; 251 } 252 } 253 254 // ensure consistent with compareTo 255 @Override 256 public boolean equals(Object o) { 257 if (o == null) { 258 return false; 259 } 260 if (o instanceof Delayed) { 261 return (compareTo((Delayed) o) == 0); 262 } else { 263 return false; 264 } 265 } 266 267 @Override 268 public int hashCode() { 269 return (int) (this.getDelay(TimeUnit.MILLISECONDS) & 0xFFFFFF); 270 } 271 } 272 273 private static final Logger log = LoggerFactory.getLogger(LocoNetThrottledTransmitter.class); 274 275}