From 0dc9fd7ecff82504fbeb3a94f8f0ae3db019bdbc Mon Sep 17 00:00:00 2001 From: Robert Sutton Date: Mon, 21 Nov 2022 11:27:10 +1100 Subject: [PATCH 1/9] make AsyncEventPump a Daemon --- .../manager/internal/AsyncEventPump.java | 396 ++++++++++-------- 1 file changed, 222 insertions(+), 174 deletions(-) diff --git a/src/main/java/org/asteriskjava/manager/internal/AsyncEventPump.java b/src/main/java/org/asteriskjava/manager/internal/AsyncEventPump.java index d17b6ae4c..47f5309fe 100644 --- a/src/main/java/org/asteriskjava/manager/internal/AsyncEventPump.java +++ b/src/main/java/org/asteriskjava/manager/internal/AsyncEventPump.java @@ -1,6 +1,10 @@ package org.asteriskjava.manager.internal; -import com.google.common.util.concurrent.RateLimiter; +import java.lang.ref.WeakReference; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.LinkedBlockingQueue; +import java.util.concurrent.TimeUnit; + import org.asteriskjava.lock.Locker; import org.asteriskjava.manager.event.ManagerEvent; import org.asteriskjava.manager.response.ManagerResponse; @@ -8,10 +12,7 @@ import org.asteriskjava.util.Log; import org.asteriskjava.util.LogFactory; -import java.lang.ref.WeakReference; -import java.util.concurrent.CountDownLatch; -import java.util.concurrent.LinkedBlockingQueue; -import java.util.concurrent.TimeUnit; +import com.google.common.util.concurrent.RateLimiter; /** * AsyncEventPump delivers events and responses to a Dispatcher without blocking @@ -20,174 +21,221 @@ * * @author rsutton */ -public class AsyncEventPump implements Dispatcher, Runnable { - private final Log logger = LogFactory.getLog(AsyncEventPump.class); - - private static final long MAX_SAFE_EVENT_AGE = 500; - - private final LinkedBlockingQueue queue = new LinkedBlockingQueue<>(20000); - private final Dispatcher dispatcher; - private volatile boolean stop = false; - private final WeakReference owner; - - private final Thread thread; - private volatile boolean terminated = false; - - private final String name; - - /** - * @param owner: A weak reference to the owner is created, should it be - * garbage collected then AsyncEventPump will shutdown. - * @param dispatcher: The dispatcher that AsyncEventPump should deliver - * events to. - * @param threadName: The AsyncEventPump's thread will be named with a - * variant of threadName - */ - AsyncEventPump(Object owner, Dispatcher dispatcher, String threadName) { - this.dispatcher = dispatcher; - this.owner = new WeakReference<>(owner); - name = threadName + ":AsyncEventPump"; - thread = new Thread(this, name); - thread.start(); - } - - @Override - public void run() { - try { - logger.info("starting"); - RateLimiter rateLimiter = RateLimiter.create(2); - while (!stop || !queue.isEmpty()) { - try { - EventWrapper wrapper = queue.poll(1, TimeUnit.MINUTES); - if (wrapper != null) { - if (wrapper.timer.timeTaken() > MAX_SAFE_EVENT_AGE && rateLimiter.tryAcquire()) { - logger.warn("The following message will only appear once per second!\n" + "Event dispatched " - + wrapper.timer.timeTaken() - + " MS after arriving, your ManagerEvent handlers are too slow!\n" - + "You should also check for Garbage Collection issues.\n" + "There are " + queue.size() - + " events waiting to be processed in the queue.\n" + "Event was " - + wrapper.getPayloadAsString()); - - } - // assume we need to process all queued events in - // MAX_SAFE_EVENT_AGE - int requiredHandlingTime = (int) (MAX_SAFE_EVENT_AGE / Math.max(1, queue.size())); - - if (wrapper.response != null) { - dispatcher.dispatchResponse(wrapper.response, requiredHandlingTime); - } else if (wrapper.event != null) { - dispatcher.dispatchEvent(wrapper.event, requiredHandlingTime); - } else if (wrapper.poison != null) { - wrapper.poison.countDown(); - } - } else if (owner.get() == null) { - stop = true; - logger.error("The owner has been garbage collected!"); - } - - } catch (InterruptedException e) { - logger.error(e); - } catch (Exception e) { - logger.error(e, e); - } - } - } finally { - terminated = true; - logger.warn("AsyncEventPump has exited"); - } - - } - - /** - * call stop() to cause the AsyncEventPump to stop, it will first empty the - * queue. - */ - public void stop() { - logger.info(name + " Requesting AsyncEventPump to stop"); - if (terminated) { - logger.warn(name + " AsyncEventPump is already stopped"); - if (!queue.isEmpty()) { - logger.error(name + " There are unprocessed events in the queue"); - } - - return; - } - EventWrapper poisonWrapper = new EventWrapper(); - queue.add(poisonWrapper); - stop = true; - LogTime timer = new LogTime(); - try { - int queueSize = queue.size(); - while (!poisonWrapper.poison.await(5, TimeUnit.SECONDS)) { - // still waiting for the poison to be consumed. - if (queueSize == queue.size()) { - if (!terminated) { - Locker.dumpThread(thread, name + " AsyncEventPump thread is blocked here..."); - } - throw new RuntimeException(name + " Failed to shutdown AsyncEventPump cleanly!"); - - } - queueSize = queue.size(); - logger.info(name + " Waiting for AsyncEventPump to Stop... "); - - if (timer.timeTaken() > 60_000) { - throw new RuntimeException(name + " Failed to shutdown AsyncEventPump cleanly!"); - } - } - } catch (InterruptedException e1) { - logger.error(name + e1.getMessage()); - } - } - - /** - * add a ManagerResponse to the queue, only if the queue is not full - */ - @Override - public void dispatchResponse(ManagerResponse response, Integer requiredHandlingTime) { - if (!queue.offer(new EventWrapper(response))) { - logger.error(name + " Event queue is full, not processing ManagerResponse " + response); - } - } - - /** - * add a ManagerEvent to the queue, only if the queue is not full - */ - @Override - public void dispatchEvent(ManagerEvent event, Integer requiredHandlingTime) { - if (!queue.offer(new EventWrapper(event))) { - logger.error(name + " Event queue is full, not processing ManagerEvent " + event); - } - } - - private static class EventWrapper { - LogTime timer = new LogTime(); - ManagerResponse response; - ManagerEvent event; - CountDownLatch poison; - - EventWrapper() { - // poison - poison = new CountDownLatch(1); - } - - public String getPayloadAsString() { - if (response != null) { - return response.toString(); - } else if (event != null) { - return event.toString(); - } - return "Poison"; - - } - - EventWrapper(ManagerResponse response) { - this.response = response; - } - - EventWrapper(ManagerEvent event) { - this.event = event; - } - - } +public class AsyncEventPump implements Dispatcher, Runnable +{ + private final Log logger = LogFactory.getLog(AsyncEventPump.class); + + private static final long MAX_SAFE_EVENT_AGE = 500; + + private final LinkedBlockingQueue queue = new LinkedBlockingQueue<>(20000); + private final Dispatcher dispatcher; + private volatile boolean stop = false; + private final WeakReference owner; + + private final Thread thread; + private volatile boolean terminated = false; + + private final String name; + + /** + * @param owner: + * A weak reference to the owner is created, should it be garbage + * collected then AsyncEventPump will shutdown. + * @param dispatcher: + * The dispatcher that AsyncEventPump should deliver events to. + * @param threadName: + * The AsyncEventPump's thread will be named with a variant of + * threadName + */ + AsyncEventPump(Object owner, Dispatcher dispatcher, String threadName) + { + this.dispatcher = dispatcher; + this.owner = new WeakReference<>(owner); + name = threadName + ":AsyncEventPump"; + thread = new Thread(this, name); + thread.setDaemon(true); + thread.start(); + } + + @Override + public void run() + { + try + { + logger.info("starting"); + RateLimiter rateLimiter = RateLimiter.create(2); + while (!stop || !queue.isEmpty()) + { + try + { + EventWrapper wrapper = queue.poll(1, TimeUnit.MINUTES); + if (wrapper != null) + { + if (wrapper.timer.timeTaken() > MAX_SAFE_EVENT_AGE && rateLimiter.tryAcquire()) + { + logger.warn("The following message will only appear once per second!\n" + + "Event dispatched " + wrapper.timer.timeTaken() + + " MS after arriving, your ManagerEvent handlers are too slow!\n" + + "You should also check for Garbage Collection issues.\n" + "There are " + + queue.size() + " events waiting to be processed in the queue.\n" + "Event was " + + wrapper.getPayloadAsString()); + + } + // assume we need to process all queued events in + // MAX_SAFE_EVENT_AGE + int requiredHandlingTime = (int) (MAX_SAFE_EVENT_AGE / Math.max(1, queue.size())); + + if (wrapper.response != null) + { + dispatcher.dispatchResponse(wrapper.response, requiredHandlingTime); + } + else if (wrapper.event != null) + { + dispatcher.dispatchEvent(wrapper.event, requiredHandlingTime); + } + else if (wrapper.poison != null) + { + wrapper.poison.countDown(); + } + } + else if (owner.get() == null) + { + stop = true; + logger.error("The owner has been garbage collected!"); + } + + } + catch (InterruptedException e) + { + logger.error(e); + } + catch (Exception e) + { + logger.error(e, e); + } + } + } + finally + { + terminated = true; + logger.warn("AsyncEventPump has exited"); + } + + } + + /** + * call stop() to cause the AsyncEventPump to stop, it will first empty the + * queue. + */ + @Override + public void stop() + { + logger.info(name + " Requesting AsyncEventPump to stop"); + if (terminated) + { + logger.warn(name + " AsyncEventPump is already stopped"); + if (!queue.isEmpty()) + { + logger.error(name + " There are unprocessed events in the queue"); + } + + return; + } + EventWrapper poisonWrapper = new EventWrapper(); + queue.add(poisonWrapper); + stop = true; + LogTime timer = new LogTime(); + try + { + int queueSize = queue.size(); + while (!poisonWrapper.poison.await(5, TimeUnit.SECONDS)) + { + // still waiting for the poison to be consumed. + if (queueSize == queue.size()) + { + if (!terminated) + { + Locker.dumpThread(thread, name + " AsyncEventPump thread is blocked here..."); + } + throw new RuntimeException(name + " Failed to shutdown AsyncEventPump cleanly!"); + + } + queueSize = queue.size(); + logger.info(name + " Waiting for AsyncEventPump to Stop... "); + + if (timer.timeTaken() > 60_000) + { + throw new RuntimeException(name + " Failed to shutdown AsyncEventPump cleanly!"); + } + } + } + catch (InterruptedException e1) + { + logger.error(name + e1.getMessage()); + } + } + + /** + * add a ManagerResponse to the queue, only if the queue is not full + */ + @Override + public void dispatchResponse(ManagerResponse response, Integer requiredHandlingTime) + { + if (!queue.offer(new EventWrapper(response))) + { + logger.error(name + " Event queue is full, not processing ManagerResponse " + response); + } + } + + /** + * add a ManagerEvent to the queue, only if the queue is not full + */ + @Override + public void dispatchEvent(ManagerEvent event, Integer requiredHandlingTime) + { + if (!queue.offer(new EventWrapper(event))) + { + logger.error(name + " Event queue is full, not processing ManagerEvent " + event); + } + } + + private static class EventWrapper + { + LogTime timer = new LogTime(); + ManagerResponse response; + ManagerEvent event; + CountDownLatch poison; + + EventWrapper() + { + // poison + poison = new CountDownLatch(1); + } + + public String getPayloadAsString() + { + if (response != null) + { + return response.toString(); + } + else if (event != null) + { + return event.toString(); + } + return "Poison"; + + } + + EventWrapper(ManagerResponse response) + { + this.response = response; + } + + EventWrapper(ManagerEvent event) + { + this.event = event; + } + + } } From a644fad15fc6e59d99e36fb7dc4062b7d9ee7e79 Mon Sep 17 00:00:00 2001 From: Robert Sutton Date: Mon, 21 Nov 2022 11:29:06 +1100 Subject: [PATCH 2/9] update version number --- pom.xml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pom.xml b/pom.xml index cae2d34bd..cd799fdcc 100644 --- a/pom.xml +++ b/pom.xml @@ -79,7 +79,7 @@ - 3.34.0 + 3.34.1.nj UTF-8 From 632b79607a49c76213895c65a75b52a450227259 Mon Sep 17 00:00:00 2001 From: Robert Sutton Date: Fri, 6 Jan 2023 14:19:13 +1100 Subject: [PATCH 3/9] update dep --- pom.xml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pom.xml b/pom.xml index cd799fdcc..215cb5815 100644 --- a/pom.xml +++ b/pom.xml @@ -123,7 +123,7 @@ org.mockito mockito-core - 4.8.0 + 4.11.0 test From 774699775b5b4f8554bb18286d3f972953811f05 Mon Sep 17 00:00:00 2001 From: Robert Sutton Date: Fri, 6 Jan 2023 14:34:02 +1100 Subject: [PATCH 4/9] for version 3.35.0-nj --- README.md | 2 +- pom.xml | 2 +- triggerRelease.dart | 12 ++++-------- 3 files changed, 6 insertions(+), 10 deletions(-) diff --git a/README.md b/README.md index 0a4d46f1b..2e86e27f7 100644 --- a/README.md +++ b/README.md @@ -49,7 +49,7 @@ Asterisk-Java 3.x (Java 1.8 and Asterisk Version 10 thru 18) (master) org.asteriskjava asterisk-java - 3.34.0 + 3.35.0-nj INSTALLATION FROM SOURCE diff --git a/pom.xml b/pom.xml index 215cb5815..968f57bec 100644 --- a/pom.xml +++ b/pom.xml @@ -79,7 +79,7 @@ - 3.34.1.nj + 3.35.0-nj UTF-8 diff --git a/triggerRelease.dart b/triggerRelease.dart index 8399effad..7e11d0759 100755 --- a/triggerRelease.dart +++ b/triggerRelease.dart @@ -15,6 +15,7 @@ void main() { read(join(dir, "pom.xml")).forEach((line) { if (line.contains("")) { line = line.replaceFirst("-SNAPSHOT", ""); + line = line.replaceFirst("-nj", ""); var parts = line.split("."); if (parts.length != 3) { exit(1); @@ -25,14 +26,9 @@ void main() { } }); - var postFix = ""; - if (confirm("Is this a SNAPSHOT release (y/n)?")) { - postFix = "-SNAPSHOT"; - rev++; - } else { - rev = 0; - minor++; - } + var postFix = "-nj"; + rev = 0; + minor++; String version = "$major.$minor.$rev$postFix"; From 585e449a28fd71ee5618e266510dd575d884db88 Mon Sep 17 00:00:00 2001 From: Robert Sutton Date: Tue, 7 Feb 2023 13:29:18 +1100 Subject: [PATCH 5/9] suppress property warning --- .../manager/event/QueueMemberStatusEvent.java | 95 +++++++++++-------- 1 file changed, 58 insertions(+), 37 deletions(-) diff --git a/src/main/java/org/asteriskjava/manager/event/QueueMemberStatusEvent.java b/src/main/java/org/asteriskjava/manager/event/QueueMemberStatusEvent.java index 1529654c1..147f8afda 100644 --- a/src/main/java/org/asteriskjava/manager/event/QueueMemberStatusEvent.java +++ b/src/main/java/org/asteriskjava/manager/event/QueueMemberStatusEvent.java @@ -22,41 +22,62 @@ * @author Asteria Solutions Group, Inc. http://www.asteriasgi.com/ * @version $Id$ */ -public class QueueMemberStatusEvent extends QueueMemberEvent { - /** - * Serializable version identifier - */ - private static final long serialVersionUID = -2293926744791895763L; - - private String ringinuse; - private Integer wrapuptime; - - /** - * @param source - */ - public QueueMemberStatusEvent(Object source) { - super(source); - } - - /** - * @return the ringinuse - */ - public String getRinginuse() { - return ringinuse; - } - - /** - * @param ringinuse the ringinuse to set - */ - public void setRinginuse(String ringinuse) { - this.ringinuse = ringinuse; - } - - public Integer getWrapuptime() { - return wrapuptime; - } - - public void setWrapuptime(Integer wrapuptime) { - this.wrapuptime = wrapuptime; - } +public class QueueMemberStatusEvent extends QueueMemberEvent +{ + /** + * Serializable version identifier + */ + private static final long serialVersionUID = -2293926744791895763L; + + private String ringinuse; + private Integer wrapuptime; + + private String logintime; + + /** + * @param source + */ + public QueueMemberStatusEvent(Object source) + { + super(source); + } + + /** + * @return the ringinuse + */ + public String getRinginuse() + { + return ringinuse; + } + + /** + * @param ringinuse + * the ringinuse to set + */ + public void setRinginuse(String ringinuse) + { + this.ringinuse = ringinuse; + } + + @Override + public Integer getWrapuptime() + { + return wrapuptime; + } + + @Override + public void setWrapuptime(Integer wrapuptime) + { + this.wrapuptime = wrapuptime; + } + + public String getLogintime() + { + return logintime; + } + + public void setLogintime(String logintime) + { + this.logintime = logintime; + } } From 71b8192aa283e3a16f5b4437c40363bbd6655e52 Mon Sep 17 00:00:00 2001 From: Robert Sutton Date: Tue, 7 Feb 2023 13:34:58 +1100 Subject: [PATCH 6/9] for version 3.36.0-nj --- README.md | 2 +- pom.xml | 2 +- pubspec.lock | 119 ++++++++++++++++++++++++++++++++++----------------- 3 files changed, 81 insertions(+), 42 deletions(-) diff --git a/README.md b/README.md index 2e86e27f7..579041b95 100644 --- a/README.md +++ b/README.md @@ -49,7 +49,7 @@ Asterisk-Java 3.x (Java 1.8 and Asterisk Version 10 thru 18) (master) org.asteriskjava asterisk-java - 3.35.0-nj + 3.36.0-nj INSTALLATION FROM SOURCE diff --git a/pom.xml b/pom.xml index 968f57bec..bfe2dbb86 100644 --- a/pom.xml +++ b/pom.xml @@ -79,7 +79,7 @@ - 3.35.0-nj + 3.36.0-nj UTF-8 diff --git a/pubspec.lock b/pubspec.lock index c0143ffb5..16800d779 100644 --- a/pubspec.lock +++ b/pubspec.lock @@ -5,273 +5,312 @@ packages: dependency: transitive description: name: archive - url: "https://pub.dartlang.org" + sha256: a92e39b291073bb840a72cf43d96d2a63c74e9a485d227833e8ea0054d16ad16 + url: "https://pub.dev" source: hosted version: "3.1.2" args: dependency: "direct main" description: name: args - url: "https://pub.dartlang.org" + sha256: "37a4264b0b7fb930e94c0c47558f3b6c4f4e9cb7e655a3ea373131d79b2dc0cc" + url: "https://pub.dev" source: hosted version: "2.0.0" async: dependency: transitive description: name: async - url: "https://pub.dartlang.org" + sha256: "6eda8392a48ae1de7ea438c91a4ba3e77205f043e7013102a424863aa6db368f" + url: "https://pub.dev" source: hosted version: "2.5.0" basic_utils: dependency: transitive description: name: basic_utils - url: "https://pub.dartlang.org" + sha256: "82976655f5fa4af987afbff7ab16a6ffe2ea5283b06aff6a2bc6cc4889005a59" + url: "https://pub.dev" source: hosted version: "3.0.0-nullsafety.1" charcode: dependency: transitive description: name: charcode - url: "https://pub.dartlang.org" + sha256: "8e36feea6de5ea69f2199f29cf42a450a855738c498b57c0b980e2d3cca9c362" + url: "https://pub.dev" source: hosted version: "1.2.0" collection: dependency: transitive description: name: collection - url: "https://pub.dartlang.org" + sha256: "6d4193120997ecfd09acf0e313f13dc122b119e5eca87ef57a7d065ec9183762" + url: "https://pub.dev" source: hosted version: "1.15.0" convert: dependency: transitive description: name: convert - url: "https://pub.dartlang.org" + sha256: df567b950053d83b4dba3e8c5799c411895d146f82b2147114b666a4fd9a80dd + url: "https://pub.dev" source: hosted version: "3.0.0" crypto: dependency: transitive description: name: crypto - url: "https://pub.dartlang.org" + sha256: "8be10341257b613566fdc9fd073c46f7c032ed329b1c732bda17aca29f2366c8" + url: "https://pub.dev" source: hosted version: "3.0.0" csv: dependency: transitive description: name: csv - url: "https://pub.dartlang.org" + sha256: "5772e255feafae00699c5c781dc3c426cbb12243c4e1a55a3ccb311b49444967" + url: "https://pub.dev" source: hosted version: "5.0.0" dcli: dependency: "direct main" description: name: dcli - url: "https://pub.dartlang.org" + sha256: d8c763628b15e8c14a90b2dcbfb22c55f0c621b5a6885e363dcdbf089e3b74ab + url: "https://pub.dev" source: hosted version: "0.51.5" equatable: dependency: transitive description: name: equatable - url: "https://pub.dartlang.org" + sha256: "8867afc0f09dd9aff976817fbd94e9a35a3a2e58bef3b39a5376918321ec97a4" + url: "https://pub.dev" source: hosted version: "2.0.0" ffi: dependency: transitive description: name: ffi - url: "https://pub.dartlang.org" + sha256: d97fffd9d86f3dccc7a9059128b468a99320c69007cc9d41a3a1bda07d4e86dc + url: "https://pub.dev" source: hosted version: "1.0.0" file: dependency: transitive description: name: file - url: "https://pub.dartlang.org" + sha256: "1b92bec4fc2a72f59a8e15af5f52cd441e4a7860b49499d69dfa817af20e925d" + url: "https://pub.dev" source: hosted - version: "6.1.0" + version: "6.1.4" glob: dependency: transitive description: name: glob - url: "https://pub.dartlang.org" + sha256: "36a6ea2cac1f93742ecee02250fb498122c0993eb948120ff0dbef6cd694beb8" + url: "https://pub.dev" source: hosted version: "2.0.0" http: dependency: transitive description: name: http - url: "https://pub.dartlang.org" + sha256: "0a48a4e44ec1b6a52eb93b12d129f5b74ee6dbb27703439c965f1bd86f7be59f" + url: "https://pub.dev" source: hosted version: "0.13.0" http_parser: dependency: transitive description: name: http_parser - url: "https://pub.dartlang.org" + sha256: e362d639ba3bc07d5a71faebb98cde68c05bfbcfbbb444b60b6f60bb67719185 + url: "https://pub.dev" source: hosted version: "4.0.0" ini: dependency: transitive description: name: ini - url: "https://pub.dartlang.org" + sha256: "578497bf9e1480bd3f7975c41e851e2f9d02eba4c158bf1af8772b5d278e2df1" + url: "https://pub.dev" source: hosted version: "2.1.0-beta" json_annotation: dependency: transitive description: name: json_annotation - url: "https://pub.dartlang.org" + sha256: "0aab0aad3dde662cdc02d3a6fa6c9f20687a91089900bd25c64f3badcc90bf21" + url: "https://pub.dev" source: hosted version: "4.0.0" logging: dependency: transitive description: name: logging - url: "https://pub.dartlang.org" + sha256: "3730d4c02b0c2d1db80ef9904e27fa796d75474f572a70011e0e616ee6bfc0ff" + url: "https://pub.dev" source: hosted version: "1.0.0" matcher: dependency: transitive description: name: matcher - url: "https://pub.dartlang.org" + sha256: "38c7be344ac5057e10161a5ecb00c9d9d67ed2f150001278601dd27d9fe64206" + url: "https://pub.dev" source: hosted version: "0.12.10" meta: dependency: transitive description: name: meta - url: "https://pub.dartlang.org" + sha256: "98a7492d10d7049ea129fd4e50f7cdd2d5008522b1dfa1148bbbc542b9dd21f7" + url: "https://pub.dev" source: hosted version: "1.3.0" path: dependency: "direct main" description: name: path - url: "https://pub.dartlang.org" + sha256: "2ad4cddff7f5cc0e2d13069f2a3f7a73ca18f66abd6f5ecf215219cdb3638edb" + url: "https://pub.dev" source: hosted version: "1.8.0" pedantic: dependency: transitive description: name: pedantic - url: "https://pub.dartlang.org" + sha256: "8f6460c77a98ad2807cd3b98c67096db4286f56166852d0ce5951bb600a63594" + url: "https://pub.dev" source: hosted version: "1.11.0" pointycastle: dependency: transitive description: name: pointycastle - url: "https://pub.dartlang.org" + sha256: "666796a2673d20ad06c6b6641994bd54f94d12857c0b93a64e1a42252c2f7e4a" + url: "https://pub.dev" source: hosted version: "3.0.0-nullsafety.2" posix: dependency: transitive description: name: posix - url: "https://pub.dartlang.org" + sha256: "64c1b8e4448812aa564ae5970e77dce69080a26edaebcd11c77df10e1ab4c9b3" + url: "https://pub.dev" source: hosted version: "2.0.0" pub_semver: dependency: transitive description: name: pub_semver - url: "https://pub.dartlang.org" + sha256: "59ed538734419e81f7fc18c98249ae72c3c7188bdd9dceff2840585227f79843" + url: "https://pub.dev" source: hosted version: "2.0.0" pubspec2: dependency: transitive description: name: pubspec2 - url: "https://pub.dartlang.org" + sha256: ae0e0f960e8308c90adf21f1ea42e2d42c930642c4d51c18d3253f03e66b8b46 + url: "https://pub.dev" source: hosted version: "1.0.0" quiver: dependency: transitive description: name: quiver - url: "https://pub.dartlang.org" + sha256: ca3ae1a9b6c2a7712f6818c8ce6bbcc6cf999115f4fb51b1fa8b7bfe1037d70e + url: "https://pub.dev" source: hosted version: "3.0.0" random_string: dependency: transitive description: name: random_string - url: "https://pub.dartlang.org" + sha256: ed4b83bc132a34d5a87298fdefa427eaf9fa9c9e95af996bc98295e5b1a37477 + url: "https://pub.dev" source: hosted version: "2.2.0-nullsafety" source_span: dependency: transitive description: name: source_span - url: "https://pub.dartlang.org" + sha256: d5f89a9e52b36240a80282b3dc0667dd36e53459717bb17b8fb102d30496606a + url: "https://pub.dev" source: hosted version: "1.8.1" stack_trace: dependency: transitive description: name: stack_trace - url: "https://pub.dartlang.org" + sha256: f8d9f247e2f9f90e32d1495ff32dac7e4ae34ffa7194c5ff8fcc0fd0e52df774 + url: "https://pub.dev" source: hosted version: "1.10.0" string_scanner: dependency: transitive description: name: string_scanner - url: "https://pub.dartlang.org" + sha256: dd11571b8a03f7cadcf91ec26a77e02bfbd6bbba2a512924d3116646b4198fc4 + url: "https://pub.dev" source: hosted version: "1.1.0" term_glyph: dependency: transitive description: name: term_glyph - url: "https://pub.dartlang.org" + sha256: a88162591b02c1f3a3db3af8ce1ea2b374bd75a7bb8d5e353bcfbdc79d719830 + url: "https://pub.dev" source: hosted version: "1.2.0" typed_data: dependency: transitive description: name: typed_data - url: "https://pub.dartlang.org" + sha256: "53bdf7e979cfbf3e28987552fd72f637e63f3c8724c9e56d9246942dc2fa36ee" + url: "https://pub.dev" source: hosted version: "1.3.0" uri: dependency: transitive description: name: uri - url: "https://pub.dartlang.org" + sha256: "889eea21e953187c6099802b7b4cf5219ba8f3518f604a1033064d45b1b8268a" + url: "https://pub.dev" source: hosted version: "1.0.0" uuid: dependency: transitive description: name: uuid - url: "https://pub.dartlang.org" + sha256: "29aa6310432a6fc0b8e2f1e8ef21f28da128cb65d80a335e2354b99435e99d50" + url: "https://pub.dev" source: hosted version: "3.0.1" validators2: dependency: transitive description: name: validators2 - url: "https://pub.dartlang.org" + sha256: c95746d5176ed1c01aff95fcd1ef22b6459f48a8aa26577ddf7a1e21919a81dc + url: "https://pub.dev" source: hosted version: "3.0.0" vin_decoder: dependency: transitive description: name: vin_decoder - url: "https://pub.dartlang.org" + sha256: "98bf7d4c9dfc0c636e7befd7f702c6b0f611660b1a98e9b2a66ccfc49267f861" + url: "https://pub.dev" source: hosted version: "0.2.0-nullsafety" yaml: dependency: transitive description: name: yaml - url: "https://pub.dartlang.org" + sha256: "3cee79b1715110341012d27756d9bae38e650588acd38d3f3c610822e1337ace" + url: "https://pub.dev" source: hosted version: "3.1.0" sdks: From b6094e2c07560c0e751954aa3d1f78afdbf691cd Mon Sep 17 00:00:00 2001 From: Robert Sutton Date: Sat, 20 Jan 2024 21:09:27 +1100 Subject: [PATCH 7/9] update log4j --- pom.xml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pom.xml b/pom.xml index bfe2dbb86..f1ec21289 100644 --- a/pom.xml +++ b/pom.xml @@ -99,7 +99,7 @@ org.apache.logging.log4j log4j-core - 2.19.0 + 2.22.0 provided From 2408c31635733dc2b4f02584bd5e0fa909588949 Mon Sep 17 00:00:00 2001 From: Robert Sutton Date: Sun, 23 Feb 2025 21:27:24 +1100 Subject: [PATCH 8/9] re-format --- .../pbx/internal/asterisk/MeetmeRoom.java | 101 +++-- .../internal/asterisk/MeetmeRoomControl.java | 232 ++++++---- .../pbx/internal/core/AsteriskPBX.java | 403 ++++++++++++------ 3 files changed, 498 insertions(+), 238 deletions(-) diff --git a/src/main/java/org/asteriskjava/pbx/internal/asterisk/MeetmeRoom.java b/src/main/java/org/asteriskjava/pbx/internal/asterisk/MeetmeRoom.java index f01856574..c1412ba8f 100644 --- a/src/main/java/org/asteriskjava/pbx/internal/asterisk/MeetmeRoom.java +++ b/src/main/java/org/asteriskjava/pbx/internal/asterisk/MeetmeRoom.java @@ -1,19 +1,20 @@ package org.asteriskjava.pbx.internal.asterisk; +import java.util.LinkedList; + import org.asteriskjava.lock.Lockable; import org.asteriskjava.lock.Locker.LockCloser; import org.asteriskjava.pbx.Channel; import org.asteriskjava.util.Log; import org.asteriskjava.util.LogFactory; -import java.util.LinkedList; - /* * This class tracks the status, channel names and number of participants in * a meetme room. The hangup function will hangup all known participants in * this meetme room. */ -public class MeetmeRoom extends Lockable { +public class MeetmeRoom extends Lockable +{ /** * The asterisk room number. This will be value offset from the Meetme Base. * e.g. if the base is 3750 and this is the third allocated room, then the @@ -35,7 +36,8 @@ public class MeetmeRoom extends Lockable { private RoomOwner owner = null; - public MeetmeRoom(final int number) { + public MeetmeRoom(final int number) + { this.roomNumber = number; } @@ -43,46 +45,58 @@ public MeetmeRoom(final int number) { * returns true if the channel was added to the list of channels in this * meetme. if the channel is already in the meetme, returns false */ - public boolean addChannel(final Channel channel) { - try (LockCloser closer = this.withLock()) { + public boolean addChannel(final Channel channel) + { + try (LockCloser closer = this.withLock()) + { boolean newChannel = false; - if (!this.channels.contains(channel)) { + if (!this.channels.contains(channel)) + { this.channels.add(channel); this.channelCount++; newChannel = true; - } else + } + else MeetmeRoom.logger.error("rejecting " + channel + " already in meetme."); //$NON-NLS-1$ //$NON-NLS-2$ return newChannel; } } - public int getChannelCount() { - try (LockCloser closer = this.withLock()) { + public int getChannelCount() + { + try (LockCloser closer = this.withLock()) + { return this.channelCount; } } - public Channel[] getChannels() { - try (LockCloser closer = this.withLock()) { + public Channel[] getChannels() + { + try (LockCloser closer = this.withLock()) + { final Channel list[] = new Channel[this.channels.size()]; int cnt = 0; - for (final Channel channel : this.channels) { + for (final Channel channel : this.channels) + { list[cnt++] = channel; } return list; } } - public boolean getForceClose() { + public boolean getForceClose() + { return this.forceClose; } - public Long getLastUpdated() { + public Long getLastUpdated() + { return this.lastUpdated; } - public String getRoomNumber() { + public String getRoomNumber() + { return ("" + this.roomNumber); //$NON-NLS-1$ } @@ -91,21 +105,26 @@ public String getRoomNumber() { * * @return */ - public boolean isActive() { + public boolean isActive() + { return this.active; } - public void removeChannel(final Channel channel) { - try (LockCloser closer = this.withLock()) { + public void removeChannel(final Channel channel) + { + try (LockCloser closer = this.withLock()) + { final boolean channelCountInSync = this.channelCount == this.channels.size(); final boolean removed = this.channels.remove(channel); - if (!removed) { + if (!removed) + { MeetmeRoom.logger.warn("An attempt to remove an non-existing channel " + channel + " from Meetme Room " //$NON-NLS-1$ //$NON-NLS-2$ + this.getRoomNumber()); } - if (channelCountInSync && removed) { + if (channelCountInSync && removed) + { this.channelCount--; } @@ -123,12 +142,15 @@ public void removeChannel(final Channel channel) { // so // decrementing it keeps us in sync with asterisk // and eventually we will get back in sync (hopefully). - if (!channelCountInSync && removed) { + if (!channelCountInSync && removed) + { this.channelCount--; } - if ((this.channels.size() < 2) && (this.channels.size() > 0)) { - if (!this.channels.get(0).isLocal()) { + if ((this.channels.size() < 2) && (this.channels.size() > 0)) + { + if (!this.channels.get(0).isLocal()) + { logger.warn("One channel left in the meet me room " + this.channels.get(0) + " room " + this.roomNumber); //$NON-NLS-1$ } } @@ -136,15 +158,18 @@ public void removeChannel(final Channel channel) { } - public void setActive() { + public void setActive() + { this.active = true; } - public void setForceClose(final boolean canClose) { + public void setForceClose(final boolean canClose) + { this.forceClose = canClose; } - public void setInactive() { + public void setInactive() + { this.active = false; this.channels.clear(); this.forceClose = false; @@ -153,7 +178,8 @@ public void setInactive() { } - public void setLastUpdated() { + public void setLastUpdated() + { this.lastUpdated = System.currentTimeMillis(); } @@ -164,25 +190,32 @@ public void setLastUpdated() { * * @param channelCount */ - public void resetChannelCount(final int resetChannelCount) { + public void resetChannelCount(final int resetChannelCount) + { this.channelCount = resetChannelCount; } - public RoomOwner getOwner() { + public RoomOwner getOwner() + { return owner; } - public void setOwner(RoomOwner newOwner) { + public void setOwner(RoomOwner newOwner) + { owner = newOwner; owner.setRoom(this); setActive(); } - public void removeOwner(RoomOwner toRemove) { - if (owner == toRemove || owner == null) { + public void removeOwner(RoomOwner toRemove) + { + if (owner == toRemove || owner == null) + { owner = null; - } else { + } + else + { logger.error( "Tring to remove the owner, but it's not the current owner. Owner=" + owner + " caller=" + toRemove); } diff --git a/src/main/java/org/asteriskjava/pbx/internal/asterisk/MeetmeRoomControl.java b/src/main/java/org/asteriskjava/pbx/internal/asterisk/MeetmeRoomControl.java index a6d7970a4..18be05293 100644 --- a/src/main/java/org/asteriskjava/pbx/internal/asterisk/MeetmeRoomControl.java +++ b/src/main/java/org/asteriskjava/pbx/internal/asterisk/MeetmeRoomControl.java @@ -1,12 +1,28 @@ package org.asteriskjava.pbx.internal.asterisk; +import java.util.HashMap; +import java.util.HashSet; +import java.util.Map; +import java.util.Map.Entry; +import java.util.concurrent.atomic.AtomicReference; + import org.asteriskjava.AsteriskVersion; import org.asteriskjava.live.ManagerCommunicationException; import org.asteriskjava.lock.Locker.LockCloser; -import org.asteriskjava.pbx.*; +import org.asteriskjava.pbx.AsteriskSettings; +import org.asteriskjava.pbx.Channel; +import org.asteriskjava.pbx.ListenerPriority; +import org.asteriskjava.pbx.PBX; +import org.asteriskjava.pbx.PBXException; +import org.asteriskjava.pbx.PBXFactory; import org.asteriskjava.pbx.asterisk.wrap.actions.CommandAction; import org.asteriskjava.pbx.asterisk.wrap.actions.ConfbridgeListAction; -import org.asteriskjava.pbx.asterisk.wrap.events.*; +import org.asteriskjava.pbx.asterisk.wrap.events.ConfbridgeListEvent; +import org.asteriskjava.pbx.asterisk.wrap.events.ManagerEvent; +import org.asteriskjava.pbx.asterisk.wrap.events.MeetMeJoinEvent; +import org.asteriskjava.pbx.asterisk.wrap.events.MeetMeLeaveEvent; +import org.asteriskjava.pbx.asterisk.wrap.events.ResponseEvent; +import org.asteriskjava.pbx.asterisk.wrap.events.ResponseEvents; import org.asteriskjava.pbx.asterisk.wrap.response.CommandResponse; import org.asteriskjava.pbx.asterisk.wrap.response.ManagerResponse; import org.asteriskjava.pbx.internal.core.AsteriskPBX; @@ -15,13 +31,8 @@ import org.asteriskjava.util.Log; import org.asteriskjava.util.LogFactory; -import java.util.HashMap; -import java.util.HashSet; -import java.util.Map; -import java.util.Map.Entry; -import java.util.concurrent.atomic.AtomicReference; - -public class MeetmeRoomControl extends EventListenerBaseClass implements CoherentManagerEventListener { +public class MeetmeRoomControl extends EventListenerBaseClass implements CoherentManagerEventListener +{ /* * listens for a channel entering or leaving meetme rooms. when there is * only 1 channel left in a room it sets it as inactive .It will not set to @@ -41,16 +52,22 @@ public class MeetmeRoomControl extends EventListenerBaseClass implements Coheren private final static AtomicReference self = new AtomicReference<>(); - synchronized public static void init(PBX pbx, final int roomCount) throws NoMeetmeException { - if (MeetmeRoomControl.self.get() != null) { + synchronized public static void init(PBX pbx, final int roomCount) throws NoMeetmeException + { + if (MeetmeRoomControl.self.get() != null) + { logger.warn("The MeetmeRoomControl has already been initialised."); //$NON-NLS-1$ - } else { + } + else + { MeetmeRoomControl.self.set(new MeetmeRoomControl(pbx, roomCount)); } } - public static MeetmeRoomControl getInstance() { - if (MeetmeRoomControl.self.get() == null) { + public static MeetmeRoomControl getInstance() + { + if (MeetmeRoomControl.self.get() == null) + { throw new IllegalStateException( "The MeetmeRoomControl has not been initialised. Please call MeetmeRoomControl.init()."); //$NON-NLS-1$ } @@ -59,7 +76,8 @@ public static MeetmeRoomControl getInstance() { } - private MeetmeRoomControl(PBX pbx, final int roomCount) throws NoMeetmeException { + private MeetmeRoomControl(PBX pbx, final int roomCount) throws NoMeetmeException + { super("MeetmeRoomControl", pbx); //$NON-NLS-1$ this.roomCount = roomCount; final AsteriskSettings settings = PBXFactory.getActiveProfile(); @@ -71,8 +89,9 @@ private MeetmeRoomControl(PBX pbx, final int roomCount) throws NoMeetmeException } @Override - public HashSet> requiredEvents() { - HashSet> required = new HashSet<>(); + public HashSet> requiredEvents() + { + HashSet> required = new HashSet<>(); required.add(MeetMeJoinEvent.class); required.add(MeetMeLeaveEvent.class); @@ -84,45 +103,58 @@ public HashSet> requiredEvents() { * returns the next available meetme room, or null if no rooms are * available. */ - public MeetmeRoom findAvailableRoom(RoomOwner newOwner) { - try (LockCloser closer = this.withLock()) { + public MeetmeRoom findAvailableRoom(RoomOwner newOwner) + { + try (LockCloser closer = this.withLock()) + { int count = 0; - for (final MeetmeRoom room : this.rooms) { - if (MeetmeRoomControl.logger.isDebugEnabled()) { + for (final MeetmeRoom room : this.rooms) + { + if (MeetmeRoomControl.logger.isDebugEnabled()) + { MeetmeRoomControl.logger.debug("room " + room.getRoomNumber() + " count " + count); } - if (room.getOwner() == null || !room.getOwner().isRoomStillRequired()) { + if (room.getOwner() == null || !room.getOwner().isRoomStillRequired()) + { /* * new code to attempt to recover uncleared meetme rooms * safely */ - try { + try + { final Long lastUpdated = room.getLastUpdated(); final long now = System.currentTimeMillis(); - if (lastUpdated != null) { + if (lastUpdated != null) + { final long elapsedTime = now - lastUpdated; MeetmeRoomControl.logger.error( "room: " + room.getRoomNumber() + " count: " + count + " elapsed: " + elapsedTime); - if ((elapsedTime > 1800000) && (room.getChannelCount() < 2)) { + if ((elapsedTime > 1800000) && (room.getChannelCount() < 2)) + { MeetmeRoomControl.logger.error("clearing room"); //$NON-NLS-1$ room.setInactive(); } } - } catch (final Exception e) { + } + catch (final Exception e) + { /* * attempt to make this new change safe */ MeetmeRoomControl.logger.error(e, e); } - if (room.getChannelCount() == 0) { + if (room.getChannelCount() == 0) + { room.setInactive(); room.setOwner(newOwner); MeetmeRoomControl.logger.info("Returning available room " + room.getRoomNumber()); return room; } - } else { + } + else + { logger.warn("Meetme " + room.getRoomNumber() + " is still in use by " + room.getOwner()); } count++; @@ -139,11 +171,15 @@ public MeetmeRoom findAvailableRoom(RoomOwner newOwner) { * @param roomNumber the meetme room number * @return */ - private MeetmeRoom findMeetmeRoom(final String roomNumber) { - try (LockCloser closer = this.withLock()) { + private MeetmeRoom findMeetmeRoom(final String roomNumber) + { + try (LockCloser closer = this.withLock()) + { MeetmeRoom foundRoom = null; - for (final MeetmeRoom room : this.rooms) { - if (room.getRoomNumber().compareToIgnoreCase(roomNumber) == 0) { + for (final MeetmeRoom room : this.rooms) + { + if (room.getRoomNumber().compareToIgnoreCase(roomNumber) == 0) + { foundRoom = room; break; } @@ -153,66 +189,83 @@ private MeetmeRoom findMeetmeRoom(final String roomNumber) { } - MeetmeRoom getRoom(final int room) { - try (LockCloser closer = this.withLock()) { + MeetmeRoom getRoom(final int room) + { + try (LockCloser closer = this.withLock()) + { return this.rooms[room]; } } @Override - public void onManagerEvent(final ManagerEvent event) { + public void onManagerEvent(final ManagerEvent event) + { MeetmeRoom room; - if (event instanceof MeetMeJoinEvent) { + if (event instanceof MeetMeJoinEvent) + { final MeetMeJoinEvent evt = (MeetMeJoinEvent) event; room = this.findMeetmeRoom(evt.getMeetMe()); final Channel channel = evt.getChannel(); - if (room != null) { - if (room.addChannel(channel)) { + if (room != null) + { + if (room.addChannel(channel)) + { MeetmeRoomControl.logger.debug(channel + " has joined the conference " //$NON-NLS-1$ + room.getRoomNumber() + " channelCount " + (room.getChannelCount())); //$NON-NLS-1$ room.setLastUpdated(); } } } - if (event instanceof MeetMeLeaveEvent) { + if (event instanceof MeetMeLeaveEvent) + { final MeetMeLeaveEvent evt = (MeetMeLeaveEvent) event; room = this.findMeetmeRoom(evt.getMeetMe()); final Channel channel = evt.getChannel(); - if (room != null) { + if (room != null) + { // ignore local dummy channels// && // !channel.toUpperCase().startsWith("LOCAL/")) { - if (MeetmeRoomControl.logger.isDebugEnabled()) { + if (MeetmeRoomControl.logger.isDebugEnabled()) + { MeetmeRoomControl.logger.debug(channel + " has left the conference " //$NON-NLS-1$ + room.getRoomNumber() + " channel count " + (room.getChannelCount())); //$NON-NLS-1$ } room.removeChannel(channel); room.setLastUpdated(); - if (room.getChannelCount() < 2 && room.getForceClose()) { + if (room.getChannelCount() < 2 && room.getForceClose()) + { this.hangupChannels(room); room.setInactive(); } - if (room.getChannelCount() < 1) { + if (room.getChannelCount() < 1) + { room.setInactive(); } } } } - public void hangupChannels(final MeetmeRoom room) { + public void hangupChannels(final MeetmeRoom room) + { final Channel Channels[] = room.getChannels(); - if (room.isActive()) { + if (room.isActive()) + { PBX pbx = PBXFactory.getActivePBX(); - for (final Channel channel : Channels) { + for (final Channel channel : Channels) + { room.removeChannel(channel); - try { + try + { logger.warn("Hanging up"); pbx.hangup(channel); - } catch (IllegalArgumentException | IllegalStateException | PBXException e) { + } + catch (IllegalArgumentException | IllegalStateException | PBXException e) + { logger.error(e, e); } @@ -220,66 +273,89 @@ public void hangupChannels(final MeetmeRoom room) { } } - private void configure(AsteriskPBX pbx) throws NoMeetmeException { + private void configure(AsteriskPBX pbx) throws NoMeetmeException + { final int base = this.meetmeBaseAddress; - for (int r = 0; r < this.roomCount; r++) { + for (int r = 0; r < this.roomCount; r++) + { this.rooms[r] = new MeetmeRoom(r + base); } - try { + try + { String command; - if (pbx.getVersion().isAtLeast(AsteriskVersion.ASTERISK_13)) { + if (pbx.getVersion().isAtLeast(AsteriskVersion.ASTERISK_13)) + { command = "ConfBridge list"; //$NON-NLS-1$ ConfbridgeListAction action = new ConfbridgeListAction(); final ResponseEvents response = pbx.sendEventGeneratingAction(action, 3000); Map roomChannelCount = new HashMap<>(); - for (ResponseEvent event : response.getEvents()) { + for (ResponseEvent event : response.getEvents()) + { ConfbridgeListEvent e = (ConfbridgeListEvent) event; Integer current = roomChannelCount.get(e.getConference()); - if (current == null) { + if (current == null) + { roomChannelCount.put(e.getConference(), 1); - } else { + } + else + { roomChannelCount.put(e.getConference(), current + 1); } } - for (Entry entry : roomChannelCount.entrySet()) { + for (Entry entry : roomChannelCount.entrySet()) + { setRoomCount(entry.getKey(), entry.getValue(), Integer.parseInt(entry.getKey())); } this.meetmeInstalled = true; - } else { - if (pbx.getVersion().isAtLeast(AsteriskVersion.ASTERISK_1_6)) { + } + else + { + if (pbx.getVersion().isAtLeast(AsteriskVersion.ASTERISK_1_6)) + { command = "meetme list"; //$NON-NLS-1$ - } else { + } + else + { command = "meetme"; //$NON-NLS-1$ } final CommandAction commandAction = new CommandAction(command); final ManagerResponse response = pbx.sendAction(commandAction, 3000); - if (!(response instanceof CommandResponse)) { + if (!(response instanceof CommandResponse)) + { throw new ManagerCommunicationException(response.getMessage(), null); } final CommandResponse commandResponse = (CommandResponse) response; MeetmeRoomControl.logger.debug("parsing active meetme rooms"); //$NON-NLS-1$ - for (final String line : commandResponse.getResult()) { + for (final String line : commandResponse.getResult()) + { this.parseMeetme(line); this.meetmeInstalled = true; MeetmeRoomControl.logger.debug(line); } } - } catch (final NoMeetmeException e) { + } + catch (final NoMeetmeException e) + { throw e; - } catch (final Exception e) { + } + catch (final Exception e) + { MeetmeRoomControl.logger.error(e, e); throw new NoMeetmeException(e.getLocalizedMessage()); } } - private void parseMeetme(final String line) throws NoMeetmeException { - try (LockCloser closer = this.withLock()) { - if (line != null) { + private void parseMeetme(final String line) throws NoMeetmeException + { + try (LockCloser closer = this.withLock()) + { + if (line != null) + { if (line.toLowerCase().startsWith("no such command 'meetme'")) //$NON-NLS-1$ { throw new NoMeetmeException("Asterisk is not configured correctly! Please enable the MeetMe app"); //$NON-NLS-1$ @@ -290,7 +366,8 @@ private void parseMeetme(final String line) throws NoMeetmeException { && (!line.toLowerCase().startsWith("* total number")) //$NON-NLS-1$ && (!line.toLowerCase().startsWith("no such conference")) //$NON-NLS-1$ && (!line.toLowerCase().startsWith("no such command 'meetme")) //$NON-NLS-1$ - ) { + ) + { // Update the stats on each meetme final String roomNumber = line.substring(0, 10).trim(); final String tmp = line.substring(11, 25).trim(); @@ -303,13 +380,17 @@ private void parseMeetme(final String line) throws NoMeetmeException { } } - private void setRoomCount(final String roomNumber, final int channelCount, final int roomNo) { + private void setRoomCount(final String roomNumber, final int channelCount, final int roomNo) + { Integer base = this.meetmeBaseAddress; // First check if its one of our rooms. - if ((roomNo >= base) && (roomNo < (base + this.roomCount))) { + if ((roomNo >= base) && (roomNo < (base + this.roomCount))) + { final MeetmeRoom room = this.findMeetmeRoom(roomNumber); - if (room != null) { - if (room.getChannelCount() != channelCount) { + if (room != null) + { + if (room.getChannelCount() != channelCount) + { /* * After a restart there may have been meetme rooms left up * and running with live calls. We need to identify any @@ -331,17 +412,20 @@ private void setRoomCount(final String roomNumber, final int channelCount, final } } - public void stop() { + public void stop() + { this.close(); } @Override - public ListenerPriority getPriority() { + public ListenerPriority getPriority() + { return ListenerPriority.NORMAL; } - public boolean isMeetmeInstalled() { + public boolean isMeetmeInstalled() + { return this.meetmeInstalled; } diff --git a/src/main/java/org/asteriskjava/pbx/internal/core/AsteriskPBX.java b/src/main/java/org/asteriskjava/pbx/internal/core/AsteriskPBX.java index 9147f06e7..415e6b28f 100644 --- a/src/main/java/org/asteriskjava/pbx/internal/core/AsteriskPBX.java +++ b/src/main/java/org/asteriskjava/pbx/internal/core/AsteriskPBX.java @@ -32,7 +32,8 @@ import java.util.Map; import java.util.concurrent.TimeUnit; -public enum AsteriskPBX implements PBX, ChannelHangupListener { +public enum AsteriskPBX implements PBX, ChannelHangupListener +{ SELF; @@ -46,8 +47,10 @@ public enum AsteriskPBX implements PBX, ChannelHangupListener { private LiveChannelManager liveChannels; - AsteriskPBX() { - try { + AsteriskPBX() + { + try + { CoherentManagerConnection.init(); this.muteSupported = CoherentManagerConnection.getInstance().isMuteAudioSupported(); @@ -56,13 +59,16 @@ public enum AsteriskPBX implements PBX, ChannelHangupListener { liveChannels = new LiveChannelManager(); MeetmeRoomControl.init(this, AsteriskPBX.MAX_MEETME_ROOMS); - } catch (Exception e) { + } + catch (Exception e) + { throw new RuntimeException(e); } } @Override - public void performPostCreationTasks() { + public void performPostCreationTasks() + { liveChannels.performPostCreationTasks(); } @@ -71,19 +77,22 @@ public void performPostCreationTasks() { * cleanup. */ @Override - public void shutdown() { + public void shutdown() + { MeetmeRoomControl.getInstance().stop(); CoherentManagerConnection.getInstance().shutDown(); } @Override - public boolean isBridgeSupported() { + public boolean isBridgeSupported() + { return this.bridgeSupport; } @Override public BlindTransferActivity blindTransfer(Call call, Call.OperandChannel channelToTransfer, EndPoint transferTarget, - CallerID toCallerID, boolean autoAnswer, long timeout, String dialOptions) { + CallerID toCallerID, boolean autoAnswer, long timeout, String dialOptions) + { final CompletionAdaptor completion = new CompletionAdaptor<>(); final BlindTransferActivityImpl transfer = new BlindTransferActivityImpl(call, channelToTransfer, transferTarget, @@ -97,15 +106,17 @@ public BlindTransferActivity blindTransfer(Call call, Call.OperandChannel channe @Override public void blindTransfer(Call call, Call.OperandChannel channelToTransfer, EndPoint transferTarget, CallerID toCallerID, - boolean autoAnswer, long timeout, ActivityCallback listener, String dialOptions) { + boolean autoAnswer, long timeout, ActivityCallback listener, String dialOptions) + { new BlindTransferActivityImpl(call, channelToTransfer, transferTarget, toCallerID, autoAnswer, timeout, listener, dialOptions); } public BlindTransferActivity blindTransfer(Channel agentChannel, EndPoint transferTarget, CallerID toCallerID, - boolean autoAnswer, int timeout, ActivityCallback iCallback, String dialOptions) - throws PBXException { + boolean autoAnswer, int timeout, ActivityCallback iCallback, String dialOptions) + throws PBXException + { return new BlindTransferActivityImpl(agentChannel, transferTarget, toCallerID, autoAnswer, timeout, iCallback, dialOptions); @@ -119,7 +130,8 @@ public BlindTransferActivity blindTransfer(Channel agentChannel, EndPoint transf * @param direction * @throws PBXException */ - public BridgeActivity bridge(final Channel lhsChannel, final Channel rhsChannel) throws PBXException { + public BridgeActivity bridge(final Channel lhsChannel, final Channel rhsChannel) throws PBXException + { final CompletionAdaptor completion = new CompletionAdaptor<>(); final BridgeActivityImpl bridge = new BridgeActivityImpl(lhsChannel, rhsChannel, completion); @@ -131,7 +143,8 @@ public BridgeActivity bridge(final Channel lhsChannel, final Channel rhsChannel) } @Override - public void split(final Call callToSplit) throws PBXException { + public void split(final Call callToSplit) throws PBXException + { final CompletionAdaptor completion = new CompletionAdaptor<>(); new SplitActivityImpl(callToSplit, completion); @@ -141,13 +154,15 @@ public void split(final Call callToSplit) throws PBXException { } @Override - public SplitActivity split(final Call callToSplit, final ActivityCallback listener) { + public SplitActivity split(final Call callToSplit, final ActivityCallback listener) + { return new SplitActivityImpl(callToSplit, listener); } @Override - public RedirectToActivity redirectToActivity(final Channel channel, final ActivityCallback listener) { + public RedirectToActivity redirectToActivity(final Channel channel, final ActivityCallback listener) + { return new RedirectToActivityImpl(channel, listener); } @@ -159,7 +174,8 @@ public RedirectToActivity redirectToActivity(final Channel channel, final Activi */ @Override public JoinActivity join(Call lhs, OperandChannel originatingOperand, Call rhs, OperandChannel acceptingOperand, - CallDirection direction) { + CallDirection direction) + { final CompletionAdaptor completion = new CompletionAdaptor<>(); final JoinActivityImpl join = new JoinActivityImpl(lhs, originatingOperand, rhs, acceptingOperand, direction, @@ -178,27 +194,31 @@ public JoinActivity join(Call lhs, OperandChannel originatingOperand, Call rhs, */ @Override public void join(Call lhs, OperandChannel originatingOperand, Call rhs, OperandChannel acceptingOperand, - CallDirection direction, ActivityCallback listener) { + CallDirection direction, ActivityCallback listener) + { new JoinActivityImpl(lhs, originatingOperand, rhs, acceptingOperand, direction, listener); } @Override - public void conference(final Channel channelOne, final Channel channelTwo, final Channel channelThree) { + public void conference(final Channel channelOne, final Channel channelTwo, final Channel channelThree) + { // TODO Auto-generated method stub } @Override public void conference(final Channel channelOne, final Channel channelTwo, final Channel channelThree, - final ActivityCallback callback) { + final ActivityCallback callback) + { // TODO Auto-generated method stub } @Override public DialActivity dial(final EndPoint from, final CallerID fromCallerID, final EndPoint to, final CallerID toCallerID, - String dialOptions) { + String dialOptions) + { final CompletionAdaptor completion = new CompletionAdaptor<>(); final DialActivityImpl dialer = new DialActivityImpl(from, to, toCallerID, false, completion, null, dialOptions); @@ -209,12 +229,14 @@ public DialActivity dial(final EndPoint from, final CallerID fromCallerID, final } public DialLocalToAgiActivity dialLocalToAgi(final EndPoint from, final CallerID fromCallerID, - ActivityCallback callback, Map channelVarsToSet) { + ActivityCallback callback, Map channelVarsToSet) + { return new DialLocalToAgiActivity(from, fromCallerID, callback, channelVarsToSet); } public DialActivity dial(final EndPoint from, final CallerID fromCallerID, final EndPoint to, final CallerID toCallerID, - final ActivityCallback callback, Map channelVarsToSet, String dialOptions) { + final ActivityCallback callback, Map channelVarsToSet, String dialOptions) + { final DialActivityImpl dialer = new DialActivityImpl(from, to, toCallerID, false, callback, channelVarsToSet, dialOptions); return dialer; @@ -222,7 +244,8 @@ public DialActivity dial(final EndPoint from, final CallerID fromCallerID, final @Override public void dial(final EndPoint from, final CallerID fromCallerID, final EndPoint to, final CallerID toCallerID, - final ActivityCallback callback, String dialOptions) { + final ActivityCallback callback, String dialOptions) + { new DialActivityImpl(from, to, toCallerID, false, callback, null, dialOptions); } @@ -231,13 +254,16 @@ public void dial(final EndPoint from, final CallerID fromCallerID, final EndPoin * Convenience method to hangup the call without having to extract the * channel yourself. */ - public void hangup(Call call) throws PBXException { + public void hangup(Call call) throws PBXException + { this.hangup(call.getOriginatingParty()); } @Override - public void hangup(final Channel channel) throws PBXException { - if (channel.isLive()) { + public void hangup(final Channel channel) throws PBXException + { + if (channel.isLive()) + { logger.info("Sending hangup action for channel: " + channel); //$NON-NLS-1$ PBX pbx = PBXFactory.getActivePBX(); @@ -245,33 +271,42 @@ public void hangup(final Channel channel) throws PBXException { throw new PBXException("Channel: " + channel + " cannot be retrieved as it is still in transition."); final HangupAction hangup = new HangupAction(channel); - try { + try + { channel.setCurrentActivityAction(new AgiChannelActivityHangup()); CoherentManagerConnection.sendAction(hangup, 1000); - } catch (IllegalArgumentException | IllegalStateException | IOException | TimeoutException e) { + } + catch (IllegalArgumentException | IllegalStateException | IOException | TimeoutException e) + { logger.error(e, e); throw new PBXException(e); } - } else + } + else logger.debug("Suppressed hangup for " + channel + " as it was already hungup"); //$NON-NLS-1$ //$NON-NLS-2$ } @Override - public void hangup(final Channel channel, final ActivityCallback callback) { + public void hangup(final Channel channel, final ActivityCallback callback) + { throw new UnsupportedOperationException("Not yet implemented."); //$NON-NLS-1$ } @Override - public HoldActivity hold(final Channel channel) { + public HoldActivity hold(final Channel channel) + { final CompletionAdaptor completion = new CompletionAdaptor<>(); HoldActivity activity = null; - try { + try + { activity = new HoldActivityImpl(channel, completion); completion.waitForCompletion(10, TimeUnit.SECONDS); - } catch (final Exception e) { + } + catch (final Exception e) + { logger.error(e, e); } @@ -279,12 +314,14 @@ public HoldActivity hold(final Channel channel) { } @Override - public boolean isMuteSupported() { + public boolean isMuteSupported() + { return this.muteSupported; } @Override - public ParkActivity park(final Call call, final Channel parkChannel) { + public ParkActivity park(final Call call, final Channel parkChannel) + { final CompletionAdaptor completion = new CompletionAdaptor<>(); final ParkActivity activity = new ParkActivityImpl(call, parkChannel, completion); @@ -295,58 +332,70 @@ public ParkActivity park(final Call call, final Channel parkChannel) { } @Override - public void park(final Call call, final Channel parkChannel, final ActivityCallback callback) { + public void park(final Call call, final Channel parkChannel, final ActivityCallback callback) + { new ParkActivityImpl(call, parkChannel, callback); } @Override - public void sendDTMF(final Channel channel, final DTMFTone tone) throws PBXException { - try { + public void sendDTMF(final Channel channel, final DTMFTone tone) throws PBXException + { + try + { if (!waitForChannelToQuiescent(channel, 3000)) throw new PBXException("Channel: " + channel + " cannot play dtmf as it is still in transition."); CoherentManagerConnection.sendAction(new PlayDtmfAction(channel, tone), 1000); - } catch (final Exception e) { + } + catch (final Exception e) + { logger.error(e, e); throw new PBXException(e); } } @Override - public void sendDTMF(final Channel channel, final DTMFTone tone, final ActivityCallback callback) { + public void sendDTMF(final Channel channel, final DTMFTone tone, final ActivityCallback callback) + { // TODO Auto-generated method stub } @Override - public void transferToMusicOnHold(final Channel channel) throws PBXException { + public void transferToMusicOnHold(final Channel channel) throws PBXException + { final RedirectCall transfer = new RedirectCall(); transfer.redirect(channel, new AgiChannelActivityHold()); } - public String getManagementContext() { + public String getManagementContext() + { final AsteriskSettings settings = PBXFactory.getActiveProfile(); return settings.getManagementContext(); } @Override - public Channel getChannelByEndPoint(final EndPoint endPoint) { + public Channel getChannelByEndPoint(final EndPoint endPoint) + { return this.liveChannels.getChannelByEndPoint(endPoint); } @Override - public void channelHangup(Channel channel, Integer cause, String causeText) { + public void channelHangup(Channel channel, Integer cause, String causeText) + { this.liveChannels.remove((ChannelProxy) channel); } - public DialPlanExtension getExtensionPark() { + public DialPlanExtension getExtensionPark() + { final AsteriskSettings settings = PBXFactory.getActiveProfile(); return this.buildDialPlanExtension(settings.getExtensionPark()); } @Override - public EndPoint getExtensionAgi() { + public EndPoint getExtensionAgi() + { final AsteriskSettings settings = PBXFactory.getActiveProfile(); return this.buildDialPlanExtension(settings.getAgiExtension()); } @@ -358,11 +407,15 @@ public EndPoint getExtensionAgi() { */ @Override - public EndPoint buildEndPoint(final String fullyQualifiedEndPoint) { + public EndPoint buildEndPoint(final String fullyQualifiedEndPoint) + { EndPoint endPoint = null; - try { + try + { endPoint = new EndPointImpl(fullyQualifiedEndPoint); - } catch (final IllegalArgumentException e) { + } + catch (final IllegalArgumentException e) + { logger.warn(e, e); } return endPoint; @@ -373,14 +426,18 @@ public EndPoint buildEndPoint(final String fullyQualifiedEndPoint) { * have a tech specified then the defaultTech is used. */ @Override - public EndPoint buildEndPoint(final TechType defaultTech, final String endPointName) { + public EndPoint buildEndPoint(final TechType defaultTech, final String endPointName) + { EndPoint endPoint = null; - try { + try + { if (endPointName == null || endPointName.trim().length() == 0) endPoint = new EndPointImpl(); else endPoint = new EndPointImpl(defaultTech, endPointName); - } catch (final IllegalArgumentException e) { + } + catch (final IllegalArgumentException e) + { logger.error(e, e); } return endPoint; @@ -388,7 +445,8 @@ public EndPoint buildEndPoint(final TechType defaultTech, final String endPointN } @Override - public EndPoint buildEndPoint(final TechType defaultTech, final Trunk trunk, final String endPointName) { + public EndPoint buildEndPoint(final TechType defaultTech, final Trunk trunk, final String endPointName) + { return new EndPointImpl(defaultTech, trunk, endPointName); } @@ -397,11 +455,15 @@ public EndPoint buildEndPoint(final TechType defaultTech, final Trunk trunk, fin * Builds an end point from an end point name. If the endpoint name doesn't * have a tech specified then the defaultTech is used. */ - public DialPlanExtension buildDialPlanExtension(final String extension) { + public DialPlanExtension buildDialPlanExtension(final String extension) + { DialPlanExtension dialPlan = null; - try { + try + { dialPlan = new DialPlanExtension(extension); - } catch (final IllegalArgumentException e) { + } + catch (final IllegalArgumentException e) + { logger.error(e, e); } return dialPlan; @@ -409,7 +471,8 @@ public DialPlanExtension buildDialPlanExtension(final String extension) { } @Override - public CallerID buildCallerID(final String number, final String name) { + public CallerID buildCallerID(final String number, final String name) + { return new CallerIDImpl(number, name); } @@ -418,28 +481,36 @@ public CallerID buildCallerID(final String number, final String name) { * * @param event */ - public CallerID buildCallerID(final AbstractChannelEvent event) { + public CallerID buildCallerID(final AbstractChannelEvent event) + { final String number = event.getCallerIdNum(); final String name = event.getCallerIdName(); return this.buildCallerID(number, name); } - public Channel registerChannel(final String channelName, final String uniqueIdParam) throws InvalidChannelName { + public Channel registerChannel(final String channelName, final String uniqueIdParam) throws InvalidChannelName + { String uniqueID = uniqueIdParam; - if (uniqueIdParam == null || uniqueIdParam.length() == 0) { + if (uniqueIdParam == null || uniqueIdParam.length() == 0) + { uniqueID = ChannelImpl.UNKNOWN_UNIQUE_ID; } - if (channelName == null || channelName.trim().length() == 0) { + if (channelName == null || channelName.trim().length() == 0) + { throw new IllegalArgumentException("Channel name must not be empty"); } Channel proxy = findChannel(cleanChannelName(channelName), null); - if (proxy == null) { + if (proxy == null) + { logger.info("Couldn't find the channel " + channelName + ", creating it"); proxy = internalRegisterChannel(channelName, uniqueID); - } else { - if (uniqueID != null && !uniqueID.equals(proxy.getUniqueId())) { + } + else + { + if (uniqueID != null && !uniqueID.equals(proxy.getUniqueId())) + { logger.info( "Found the channel(" + proxy.getUniqueId() + "), but with a different uniqueId (" + uniqueID + ")"); @@ -461,13 +532,16 @@ public Channel registerChannel(final String channelName, final String uniqueIdPa * @return * @throws InvalidChannelName */ - public Channel internalRegisterChannel(final String channelName, final String uniqueID) throws InvalidChannelName { + public Channel internalRegisterChannel(final String channelName, final String uniqueID) throws InvalidChannelName + { ChannelProxy proxy = null; - try (LockCloser closer = this.liveChannels.withLock()) { + try (LockCloser closer = this.liveChannels.withLock()) + { String localUniqueID = (uniqueID == null ? ChannelImpl.UNKNOWN_UNIQUE_ID : uniqueID); proxy = this.findChannel(cleanChannelName(channelName), localUniqueID); - if (proxy == null) { + if (proxy == null) + { proxy = new ChannelProxy(new ChannelImpl(channelName, localUniqueID)); logger.debug("Creating new Channel Proxy " + proxy); this.liveChannels.add(proxy); @@ -483,17 +557,21 @@ public Channel internalRegisterChannel(final String channelName, final String un * @param name * @return */ - private String cleanChannelName(final String name) { + private String cleanChannelName(final String name) + { String cleanedName = name.trim().toUpperCase(); return cleanedName; } - public Channel registerHangupChannel(String channel, String uniqueId) throws InvalidChannelName { + public Channel registerHangupChannel(String channel, String uniqueId) throws InvalidChannelName + { Channel newChannel = null; - try (LockCloser closer = this.liveChannels.withLock()) { + try (LockCloser closer = this.liveChannels.withLock()) + { newChannel = this.findChannel(channel, uniqueId); - if (newChannel == null) { + if (newChannel == null) + { // WE don't add this channel to the liveChannels as it is in the // process // of being hungup so we don't need to track it. @@ -508,20 +586,24 @@ public Channel registerHangupChannel(String channel, String uniqueId) throws Inv return newChannel; } - public ChannelProxy findChannel(final String channelName, final String uniqueID) { + public ChannelProxy findChannel(final String channelName, final String uniqueID) + { return this.liveChannels.findChannel(channelName, uniqueID); } - public MeetmeRoom acquireMeetmeRoom(RoomOwner owner) { + public MeetmeRoom acquireMeetmeRoom(RoomOwner owner) + { return MeetmeRoomControl.getInstance().findAvailableRoom(owner); } - public void addListener(FilteredManagerListener listener) { + public void addListener(FilteredManagerListener listener) + { CoherentManagerConnection connection = CoherentManagerConnection.getInstance(); connection.addListener(listener); } - public void removeListener(FilteredManagerListener listener) { + public void removeListener(FilteredManagerListener listener) + { CoherentManagerConnection connection = CoherentManagerConnection.getInstance(); connection.removeListener(listener); } @@ -537,70 +619,86 @@ public void removeListener(FilteredManagerListener listener) { * @throws TimeoutException */ public ManagerResponse sendAction(ManagerAction theAction) - throws IllegalArgumentException, IllegalStateException, IOException, TimeoutException { + throws IllegalArgumentException, IllegalStateException, IOException, TimeoutException + { return CoherentManagerConnection.sendAction(theAction, 30000); } public ManagerResponse sendAction(ManagerAction theAction, int timeout) - throws IllegalArgumentException, IllegalStateException, IOException, TimeoutException { + throws IllegalArgumentException, IllegalStateException, IOException, TimeoutException + { return CoherentManagerConnection.sendAction(theAction, timeout); } public ResponseEvents sendEventGeneratingAction(EventGeneratingAction action) - throws EventTimeoutException, IllegalArgumentException, IllegalStateException, IOException { + throws EventTimeoutException, IllegalArgumentException, IllegalStateException, IOException + { ResponseEvents events = CoherentManagerConnection.sendEventGeneratingAction(action); return events; } public ResponseEvents sendEventGeneratingAction(EventGeneratingAction action, int timeout) - throws EventTimeoutException, IllegalArgumentException, IllegalStateException, IOException { + throws EventTimeoutException, IllegalArgumentException, IllegalStateException, IOException + { return CoherentManagerConnection.sendEventGeneratingAction(action, timeout); } - public void setVariable(Channel channel, String name, String value) throws PBXException { + public void setVariable(Channel channel, String name, String value) throws PBXException + { CoherentManagerConnection.getInstance().setVariable(channel, name, value); } - public void sendActionNoWait(final ManagerAction action) { + public void sendActionNoWait(final ManagerAction action) + { CoherentManagerConnection.sendActionNoWait(action); } - public String getVariable(Channel channel, String name) { + public String getVariable(Channel channel, String name) + { return CoherentManagerConnection.getInstance().getVariable(channel, name); } - public AsteriskVersion getVersion() { + public AsteriskVersion getVersion() + { return CoherentManagerConnection.getInstance().getVersion(); } - public boolean isConnected() { + public boolean isConnected() + { return ((CoherentManagerConnection.managerConnection != null) && (CoherentManagerConnection.managerConnection.getState() == ManagerConnectionState.CONNECTED)); } - public boolean isMeetmeInstalled() { + public boolean isMeetmeInstalled() + { return MeetmeRoomControl.getInstance().isMeetmeInstalled(); } @Override - public boolean isChannel(String channelName) { + public boolean isChannel(String channelName) + { boolean isChannel = false; - try { + try + { internalRegisterChannel(channelName, ChannelImpl.UNKNOWN_UNIQUE_ID); isChannel = true; - } catch (InvalidChannelName e) { + } + catch (InvalidChannelName e) + { // if we get here then its not avalid channel name. } return isChannel; } - static public String getSIPADDHeader(final boolean inherit, final boolean targetIsSIP) { + static public String getSIPADDHeader(final boolean inherit, final boolean targetIsSIP) + { String sipHeader = "SIPADDHEADER"; //$NON-NLS-1$ - if (!targetIsSIP || inherit) { + if (!targetIsSIP || inherit) + { sipHeader = "__" + sipHeader; //$NON-NLS-1$ } return sipHeader; @@ -612,27 +710,35 @@ static public String getSIPADDHeader(final boolean inherit, final boolean target * * @param channel * @param timeout the time to wait (in milliseconds) for the channel to - * become quiescent. + * become quiescent. */ @Override - public boolean waitForChannelsToQuiescent(List channels, long timeout) { + public boolean waitForChannelsToQuiescent(List channels, long timeout) + { long elapsed = 0; - while (elapsed < timeout && !channelsAreQuiesent(channels)) { - try { + while (elapsed < timeout && !channelsAreQuiesent(channels)) + { + try + { Thread.sleep(200); - } catch (InterruptedException e) { + } + catch (InterruptedException e) + { logger.error(e, e); } logger.info("Waiting for channesl to Quiescent"); elapsed += 200; } - if (elapsed > timeout / 2) { + if (elapsed > timeout / 2) + { logger.warn("Took " + elapsed + "ms for channels to Quiescent"); } - if (elapsed >= timeout) { + if (elapsed >= timeout) + { logger.error("Channels didn't Quiescent"); - for (Channel channel : channels) { + for (Channel channel : channels) + { logger.error(channel); } @@ -640,20 +746,24 @@ public boolean waitForChannelsToQuiescent(List channels, long timeout) return timeout > elapsed; } - private boolean channelsAreQuiesent(List channels) { + private boolean channelsAreQuiesent(List channels) + { boolean ret = true; - for (Channel channel : channels) { + for (Channel channel : channels) + { ret &= channel.isQuiescent(); } return ret; } - public boolean moveChannelToAgi(Channel channel) throws PBXException { + public boolean moveChannelToAgi(Channel channel) throws PBXException + { if (!waitForChannelToQuiescent(channel, 3000)) throw new PBXException("Channel: " + channel + " cannot be transfered as it is still in transition."); boolean isInAgi = channel.isInAgi(); - if (!isInAgi) { + if (!isInAgi) + { final AsteriskSettings profile = PBXFactory.getActiveProfile(); channel.setCurrentActivityAction(new AgiChannelActivityHold()); @@ -662,21 +772,26 @@ public boolean moveChannelToAgi(Channel channel) throws PBXException { logger.error("Issuing redirect on channel " + channel + " to move it to the agi"); - try { + try + { final ManagerResponse response = sendAction(redirect, 1000); if ((response != null) && (response.getResponse().compareToIgnoreCase("success") == 0))//$NON-NLS-1$ { int limit = 50; - while (!channel.isInAgi() && limit-- > 0) { + while (!channel.isInAgi() && limit-- > 0) + { Thread.sleep(100); } isInAgi = channel.isInAgi(); - if (!isInAgi) { + if (!isInAgi) + { logger.error("Failed to move channel"); } } - } catch (final Exception e) { + } + catch (final Exception e) + { logger.error(e, e); } @@ -685,35 +800,42 @@ public boolean moveChannelToAgi(Channel channel) throws PBXException { } - public void moveChannelTo(Channel channel, String context, String exten, int prio) { + public void moveChannelTo(Channel channel, String context, String exten, int prio) + { DialPlanExtension ext = this.buildDialPlanExtension(exten); channel.setCurrentActivityAction(new AgiChannelActivityHold()); final RedirectAction redirect = new RedirectAction(channel, context, ext, prio); - try { + try + { sendAction(redirect, 1000); - } catch (final Exception e) { + } + catch (final Exception e) + { logger.error(e, e); } } @Override - public boolean waitForChannelToQuiescent(Channel channel, int timeout) { + public boolean waitForChannelToQuiescent(Channel channel, int timeout) + { List channels = new LinkedList<>(); channels.add(channel); return waitForChannelsToQuiescent(channels, timeout); } - public ChannelProxy getProxyById(String id) { + public ChannelProxy getProxyById(String id) + { return liveChannels.findProxyById(id); } public DialToAgiActivityImpl dialToAgi(EndPoint endPoint, CallerID callerID, AgiChannelActivityAction action, - ActivityCallback iCallback, Map channelVarsToSet) { + ActivityCallback iCallback, Map channelVarsToSet) + { final CompletionAdaptor completion = new CompletionAdaptor<>(); @@ -726,9 +848,12 @@ public DialToAgiActivityImpl dialToAgi(EndPoint endPoint, CallerID callerID, Agi final ActivityStatusEnum status; - if (dialer.isSuccess()) { + if (dialer.isSuccess()) + { status = ActivityStatusEnum.SUCCESS; - } else { + } + else + { status = ActivityStatusEnum.FAILURE; } @@ -738,7 +863,8 @@ public DialToAgiActivityImpl dialToAgi(EndPoint endPoint, CallerID callerID, Agi } public DialToAgiActivityImpl dialToAgiWithAbort(EndPoint endPoint, CallerID callerID, int timeout, - AgiChannelActivityAction action, ActivityCallback iCallback) { + AgiChannelActivityAction action, ActivityCallback iCallback) + { return new DialToAgiActivityImpl(endPoint, callerID, timeout, false, iCallback, null, action); @@ -756,13 +882,16 @@ public DialToAgiActivityImpl dialToAgiWithAbort(EndPoint endPoint, CallerID call * @throws AuthenticationFailedException * @throws TimeoutException */ - public boolean createAgiEntryPoint() throws IOException, AuthenticationFailedException, TimeoutException { + public boolean createAgiEntryPoint() throws IOException, AuthenticationFailedException, TimeoutException + { - try { + try + { AsteriskPBX pbx = (AsteriskPBX) PBXFactory.getActivePBX(); AsteriskSettings profile = PBXFactory.getActiveProfile(); - if (!checkDialplanExists(profile)) { + if (!checkDialplanExists(profile)) + { String host = profile.getAgiHost(); String agi = profile.getAgiExtension(); @@ -774,7 +903,9 @@ public boolean createAgiEntryPoint() throws IOException, AuthenticationFailedExc return checkDialplanExists(profile); } return true; - } catch (Exception e) { + } + catch (Exception e) + { logger.error(e); return false; @@ -782,13 +913,17 @@ public boolean createAgiEntryPoint() throws IOException, AuthenticationFailedExc } - public boolean checkDialplanExists(String dialPlan, String context) throws IOException, TimeoutException { + public boolean checkDialplanExists(String dialPlan, String context) throws IOException, TimeoutException + { String command; - if (getVersion().isAtLeast(AsteriskVersion.ASTERISK_1_6)) { + if (getVersion().isAtLeast(AsteriskVersion.ASTERISK_1_6)) + { // TODO: Use ShowDialplanAction instead of CommandAction? command = "dialplan show " + context; - } else { + } + else + { command = "show dialplan " + context; } @@ -799,12 +934,14 @@ public boolean checkDialplanExists(String dialPlan, String context) throws IOExc } public boolean checkDialplanExists(AsteriskSettings profile) - throws IllegalArgumentException, IllegalStateException, IOException, TimeoutException { + throws IllegalArgumentException, IllegalStateException, IOException, TimeoutException + { return checkDialplanExists(ACTIVITY_AGI, profile.getManagementContext()); } - public String addAsteriskExtension(String extNumber, int priority, String command) throws Exception { + public String addAsteriskExtension(String extNumber, int priority, String command) throws Exception + { String ext = "dialplan add extension " + extNumber + "," + priority + "," + command; CommandAction action = new CommandAction(ext); @@ -813,7 +950,8 @@ public String addAsteriskExtension(String extNumber, int priority, String comman List line = response.getResult(); String tmp = "Extension '" + extNumber + "," + priority + ","; - if (line.stream().anyMatch(answer -> answer.substring(0, tmp.length()).compareToIgnoreCase(tmp) == 0)) { + if (line.stream().anyMatch(answer -> answer.substring(0, tmp.length()).compareToIgnoreCase(tmp) == 0)) + { return "OK"; } @@ -821,21 +959,26 @@ public String addAsteriskExtension(String extNumber, int priority, String comman } @Override - public Trunk buildTrunk(final String trunk) { - return new Trunk() { + public Trunk buildTrunk(final String trunk) + { + return new Trunk() + { @Override - public String getTrunkAsString() { + public String getTrunkAsString() + { return trunk; } }; } - public List getChannelList() { + public List getChannelList() + { return liveChannels.getChannelList(); } - public boolean expectRenameEvents() { + public boolean expectRenameEvents() + { return expectRenameEvents; } From e5fc32b4083954f4b7d2bd4e723b8d685bc6d3eb Mon Sep 17 00:00:00 2001 From: Robert Sutton Date: Sun, 23 Feb 2025 22:52:47 +1100 Subject: [PATCH 9/9] remove recycling of meetme rooms --- .../pbx/internal/asterisk/MeetmeRoom.java | 30 ++-- .../internal/asterisk/MeetmeRoomControl.java | 150 ++---------------- 2 files changed, 27 insertions(+), 153 deletions(-) diff --git a/src/main/java/org/asteriskjava/pbx/internal/asterisk/MeetmeRoom.java b/src/main/java/org/asteriskjava/pbx/internal/asterisk/MeetmeRoom.java index c1412ba8f..fba073c47 100644 --- a/src/main/java/org/asteriskjava/pbx/internal/asterisk/MeetmeRoom.java +++ b/src/main/java/org/asteriskjava/pbx/internal/asterisk/MeetmeRoom.java @@ -1,6 +1,7 @@ package org.asteriskjava.pbx.internal.asterisk; import java.util.LinkedList; +import java.util.UUID; import org.asteriskjava.lock.Lockable; import org.asteriskjava.lock.Locker.LockCloser; @@ -15,17 +16,13 @@ */ public class MeetmeRoom extends Lockable { - /** - * The asterisk room number. This will be value offset from the Meetme Base. - * e.g. if the base is 3750 and this is the third allocated room, then the - * roomNumber will equal 3753. - */ - private final int roomNumber; private static final Log logger = LogFactory.getLog(MeetmeRoom.class); LinkedList channels = new LinkedList<>(); + private final UUID roomNumber = UUID.randomUUID(); + private int channelCount = 0; private boolean active = false; @@ -36,9 +33,11 @@ public class MeetmeRoom extends Lockable private RoomOwner owner = null; - public MeetmeRoom(final int number) + public MeetmeRoom(RoomOwner owner) { - this.roomNumber = number; + this.owner = owner; + owner.setRoom(this); + active = true; } /* @@ -158,11 +157,6 @@ public void removeChannel(final Channel channel) } - public void setActive() - { - this.active = true; - } - public void setForceClose(final boolean canClose) { this.forceClose = canClose; @@ -201,11 +195,13 @@ public RoomOwner getOwner() return owner; } - public void setOwner(RoomOwner newOwner) + public void clearOwner() { - owner = newOwner; - owner.setRoom(this); - setActive(); + if (owner != null) + { + owner.setRoom(null); + } + owner = null; } public void removeOwner(RoomOwner toRemove) diff --git a/src/main/java/org/asteriskjava/pbx/internal/asterisk/MeetmeRoomControl.java b/src/main/java/org/asteriskjava/pbx/internal/asterisk/MeetmeRoomControl.java index 18be05293..8db168ff6 100644 --- a/src/main/java/org/asteriskjava/pbx/internal/asterisk/MeetmeRoomControl.java +++ b/src/main/java/org/asteriskjava/pbx/internal/asterisk/MeetmeRoomControl.java @@ -3,11 +3,10 @@ import java.util.HashMap; import java.util.HashSet; import java.util.Map; -import java.util.Map.Entry; +import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.atomic.AtomicReference; import org.asteriskjava.AsteriskVersion; -import org.asteriskjava.live.ManagerCommunicationException; import org.asteriskjava.lock.Locker.LockCloser; import org.asteriskjava.pbx.AsteriskSettings; import org.asteriskjava.pbx.Channel; @@ -15,7 +14,6 @@ import org.asteriskjava.pbx.PBX; import org.asteriskjava.pbx.PBXException; import org.asteriskjava.pbx.PBXFactory; -import org.asteriskjava.pbx.asterisk.wrap.actions.CommandAction; import org.asteriskjava.pbx.asterisk.wrap.actions.ConfbridgeListAction; import org.asteriskjava.pbx.asterisk.wrap.events.ConfbridgeListEvent; import org.asteriskjava.pbx.asterisk.wrap.events.ManagerEvent; @@ -23,8 +21,6 @@ import org.asteriskjava.pbx.asterisk.wrap.events.MeetMeLeaveEvent; import org.asteriskjava.pbx.asterisk.wrap.events.ResponseEvent; import org.asteriskjava.pbx.asterisk.wrap.events.ResponseEvents; -import org.asteriskjava.pbx.asterisk.wrap.response.CommandResponse; -import org.asteriskjava.pbx.asterisk.wrap.response.ManagerResponse; import org.asteriskjava.pbx.internal.core.AsteriskPBX; import org.asteriskjava.pbx.internal.core.CoherentManagerEventListener; import org.asteriskjava.pbx.internal.managerAPI.EventListenerBaseClass; @@ -44,7 +40,7 @@ public class MeetmeRoomControl extends EventListenerBaseClass implements Coheren private Integer meetmeBaseAddress; - private MeetmeRoom rooms[]; + private final Map rooms = new ConcurrentHashMap<>(); private int roomCount; @@ -82,7 +78,7 @@ private MeetmeRoomControl(PBX pbx, final int roomCount) throws NoMeetmeException this.roomCount = roomCount; final AsteriskSettings settings = PBXFactory.getActiveProfile(); this.meetmeBaseAddress = settings.getMeetmeBaseAddress(); - this.rooms = new MeetmeRoom[roomCount]; + this.configure((AsteriskPBX) pbx); this.startListener(); @@ -108,7 +104,7 @@ public MeetmeRoom findAvailableRoom(RoomOwner newOwner) try (LockCloser closer = this.withLock()) { int count = 0; - for (final MeetmeRoom room : this.rooms) + for (final MeetmeRoom room : this.rooms.values()) { if (MeetmeRoomControl.logger.isDebugEnabled()) { @@ -147,9 +143,9 @@ public MeetmeRoom findAvailableRoom(RoomOwner newOwner) if (room.getChannelCount() == 0) { room.setInactive(); - room.setOwner(newOwner); - MeetmeRoomControl.logger.info("Returning available room " + room.getRoomNumber()); - return room; + room.clearOwner(); + MeetmeRoomControl.logger.info("freeing available room " + room.getRoomNumber()); + rooms.remove(room.getRoomNumber()); } } @@ -159,8 +155,12 @@ public MeetmeRoom findAvailableRoom(RoomOwner newOwner) } count++; } - MeetmeRoomControl.logger.error("no more available rooms"); - return null; + + MeetmeRoom room = new MeetmeRoom(newOwner); + rooms.put(room.getRoomNumber(), room); + MeetmeRoomControl.logger.info("Returning available room " + room.getRoomNumber()); + + return room; } } @@ -175,28 +175,11 @@ private MeetmeRoom findMeetmeRoom(final String roomNumber) { try (LockCloser closer = this.withLock()) { - MeetmeRoom foundRoom = null; - for (final MeetmeRoom room : this.rooms) - { - if (room.getRoomNumber().compareToIgnoreCase(roomNumber) == 0) - { - foundRoom = room; - break; - } - } - return foundRoom; + return rooms.get(roomNumber); } } - MeetmeRoom getRoom(final int room) - { - try (LockCloser closer = this.withLock()) - { - return this.rooms[room]; - } - } - @Override public void onManagerEvent(final ManagerEvent event) { @@ -275,18 +258,10 @@ public void hangupChannels(final MeetmeRoom room) private void configure(AsteriskPBX pbx) throws NoMeetmeException { - final int base = this.meetmeBaseAddress; - for (int r = 0; r < this.roomCount; r++) - { - this.rooms[r] = new MeetmeRoom(r + base); - } - try { - String command; if (pbx.getVersion().isAtLeast(AsteriskVersion.ASTERISK_13)) { - command = "ConfBridge list"; //$NON-NLS-1$ ConfbridgeListAction action = new ConfbridgeListAction(); final ResponseEvents response = pbx.sendEventGeneratingAction(action, 3000); Map roomChannelCount = new HashMap<>(); @@ -304,44 +279,9 @@ private void configure(AsteriskPBX pbx) throws NoMeetmeException roomChannelCount.put(e.getConference(), current + 1); } } - for (Entry entry : roomChannelCount.entrySet()) - { - setRoomCount(entry.getKey(), entry.getValue(), Integer.parseInt(entry.getKey())); - - } this.meetmeInstalled = true; } - else - { - if (pbx.getVersion().isAtLeast(AsteriskVersion.ASTERISK_1_6)) - { - command = "meetme list"; //$NON-NLS-1$ - } - else - { - command = "meetme"; //$NON-NLS-1$ - } - final CommandAction commandAction = new CommandAction(command); - final ManagerResponse response = pbx.sendAction(commandAction, 3000); - if (!(response instanceof CommandResponse)) - { - throw new ManagerCommunicationException(response.getMessage(), null); - } - - final CommandResponse commandResponse = (CommandResponse) response; - MeetmeRoomControl.logger.debug("parsing active meetme rooms"); //$NON-NLS-1$ - for (final String line : commandResponse.getResult()) - { - this.parseMeetme(line); - this.meetmeInstalled = true; - MeetmeRoomControl.logger.debug(line); - } - } - } - catch (final NoMeetmeException e) - { - throw e; } catch (final Exception e) { @@ -350,68 +290,6 @@ private void configure(AsteriskPBX pbx) throws NoMeetmeException } } - private void parseMeetme(final String line) throws NoMeetmeException - { - try (LockCloser closer = this.withLock()) - { - if (line != null) - { - if (line.toLowerCase().startsWith("no such command 'meetme'")) //$NON-NLS-1$ - { - throw new NoMeetmeException("Asterisk is not configured correctly! Please enable the MeetMe app"); //$NON-NLS-1$ - } - - if ((!line.toLowerCase().startsWith("no active meetme conferences.")) //$NON-NLS-1$ - && (!line.toLowerCase().startsWith("conf num")) //$NON-NLS-1$ - && (!line.toLowerCase().startsWith("* total number")) //$NON-NLS-1$ - && (!line.toLowerCase().startsWith("no such conference")) //$NON-NLS-1$ - && (!line.toLowerCase().startsWith("no such command 'meetme")) //$NON-NLS-1$ - ) - { - // Update the stats on each meetme - final String roomNumber = line.substring(0, 10).trim(); - final String tmp = line.substring(11, 25).trim(); - final int channelCount = Integer.parseInt(tmp); - - final int roomNo = Integer.valueOf(roomNumber); - setRoomCount(roomNumber, channelCount, roomNo); - } - } - } - } - - private void setRoomCount(final String roomNumber, final int channelCount, final int roomNo) - { - Integer base = this.meetmeBaseAddress; - // First check if its one of our rooms. - if ((roomNo >= base) && (roomNo < (base + this.roomCount))) - { - final MeetmeRoom room = this.findMeetmeRoom(roomNumber); - if (room != null) - { - if (room.getChannelCount() != channelCount) - { - /* - * After a restart there may have been meetme rooms left up - * and running with live calls. We need to identify any - * active rooms so we don't accidentally re-use an active - * room which would result in a crossed channel. - */ - MeetmeRoomControl.logger.warn("Room number: " + room.getRoomNumber() //$NON-NLS-1$ - + " has a server side channel count = " + channelCount //$NON-NLS-1$ - + " when the channel count for that room is: " + room.getChannelCount() //$NON-NLS-1$ - + " the server side channel count will be reset."); //$NON-NLS-1$ - } - room.resetChannelCount(channelCount); - room.setActive(); - } - // else "Found roomNumber:" + roomNumber + " but it was not - // in the list of rooms managed by MeetmeRoomControl."); - // //$NON-NLS-1$ //$NON-NLS-2$ - - } - } - public void stop() { this.close();