diff --git a/OpenSim/Region/CoreModules/Avatar/InstantMessage/MessageTransferModule.cs b/OpenSim/Region/CoreModules/Avatar/InstantMessage/MessageTransferModule.cs index efb9421153..f20294210b 100644 --- a/OpenSim/Region/CoreModules/Avatar/InstantMessage/MessageTransferModule.cs +++ b/OpenSim/Region/CoreModules/Avatar/InstantMessage/MessageTransferModule.cs @@ -56,6 +56,8 @@ namespace OpenSim.Region.CoreModules.Avatar.InstantMessage public event UndeliveredMessage OnUndeliveredMessage; + public static ObjectJobEngine IMXMLRPCSendWorkers; + private IPresenceService m_PresenceService; protected IPresenceService PresenceService { @@ -72,8 +74,7 @@ namespace OpenSim.Region.CoreModules.Avatar.InstantMessage IConfig cnf = config.Configs["Messaging"]; if (cnf != null) { - if (cnf.GetString("MessageTransferModule", - "MessageTransferModule") != "MessageTransferModule") + if (cnf.GetString("MessageTransferModule", "MessageTransferModule") != "MessageTransferModule") { return; } @@ -82,6 +83,8 @@ namespace OpenSim.Region.CoreModules.Avatar.InstantMessage } m_log.Debug("[MESSAGE TRANSFER]: Module enabled"); m_Enabled = true; + + IMXMLRPCSendWorkers = new ObjectJobEngine(DoSendIMviaXMLRPC, "IMXMLRPCSendWorkers", 3); } public virtual void AddRegion(Scene scene) @@ -102,8 +105,7 @@ namespace OpenSim.Region.CoreModules.Avatar.InstantMessage if (!m_Enabled) return; - MainServer.Instance.AddXmlRPCHandler( - "grid_instant_message", processXMLRPCGridInstantMessage); + MainServer.Instance.AddXmlRPCHandler("grid_instant_message", processXMLRPCGridInstantMessage); } public virtual void RegionLoaded(Scene scene) @@ -432,44 +434,40 @@ namespace OpenSim.Region.CoreModules.Avatar.InstantMessage return resp; } - - /// - /// delegate for sending a grid instant message asynchronously - /// - private delegate void GridInstantMessageDelegate(GridInstantMessage im, MessageResultNotification result); - - private class GIM { + private class GIMData { public GridInstantMessage im; public MessageResultNotification result; }; - private Queue pendingInstantMessages = new Queue(); - private int numInstantMessageThreads = 0; + //private Queue pendingInstantMessages = new Queue(); + //private int numInstantMessageThreads = 0; private void SendGridInstantMessageViaXMLRPC(GridInstantMessage im, MessageResultNotification result) { - lock (pendingInstantMessages) { - if (numInstantMessageThreads >= 4) { - GIM gim = new GIM(); - gim.im = im; - gim.result = result; - pendingInstantMessages.Enqueue(gim); - } else { - ++ numInstantMessageThreads; - //m_log.DebugFormat("[SendGridInstantMessageViaXMLRPC]: ++numInstantMessageThreads={0}", numInstantMessageThreads); - GridInstantMessageDelegate d = SendGridInstantMessageViaXMLRPCAsyncMain; - d.BeginInvoke(im, result, GridInstantMessageCompleted, d); - } + GIMData gim = new GIMData() + { + im = im, + result = result + }; + + IMXMLRPCSendWorkers.Enqueue(gim); + } + + private void DoSendIMviaXMLRPC(object o) + { + try + { + GIMData data = o as GIMData; + if (o == null) + return; + SendGridInstantMessageViaXMLRPCAsync(data.im, data.result, UUID.Zero); + } + catch (Exception e) + { + m_log.Error("[SendGridInstantMessageViaXMLRPC]: exception " + e.Message); } } - - - private void GridInstantMessageCompleted(IAsyncResult iar) - { - GridInstantMessageDelegate d = (GridInstantMessageDelegate)iar.AsyncState; - d.EndInvoke(iar); - } - + /// /// Internal SendGridInstantMessage over XMLRPC method. /// @@ -478,32 +476,8 @@ namespace OpenSim.Region.CoreModules.Avatar.InstantMessage /// Pass in 0 the first time this method is called. It will be called recursively with the last /// regionhandle tried /// - private void SendGridInstantMessageViaXMLRPCAsyncMain(GridInstantMessage im, MessageResultNotification result) + private void SendGridInstantMessageViaXMLRPCAsync(GridInstantMessage im, MessageResultNotification result, UUID prevRegionID) { - GIM gim; - do { - try { - SendGridInstantMessageViaXMLRPCAsync(im, result, UUID.Zero); - } catch (Exception e) { - m_log.Error("[SendGridInstantMessageViaXMLRPC]: exception " + e.Message); - } - lock (pendingInstantMessages) { - if (pendingInstantMessages.Count > 0) { - gim = pendingInstantMessages.Dequeue(); - im = gim.im; - result = gim.result; - } else { - gim = null; - -- numInstantMessageThreads; - //m_log.DebugFormat("[SendGridInstantMessageViaXMLRPC]: --numInstantMessageThreads={0}", numInstantMessageThreads); - } - } - } while (gim != null); - } - - private void SendGridInstantMessageViaXMLRPCAsync(GridInstantMessage im, MessageResultNotification result, UUID prevRegionID) - { - UUID toAgentID = new UUID(im.toAgentID); PresenceInfo upd = null; bool lookupAgent = false; @@ -528,7 +502,6 @@ namespace OpenSim.Region.CoreModules.Avatar.InstantMessage } } - // Are we needing to look-up an agent? if (lookupAgent) { @@ -603,8 +576,7 @@ namespace OpenSim.Region.CoreModules.Avatar.InstantMessage // The version that spawns the thread is SendGridInstantMessageViaXMLRPC // This is recursive!!!!! - SendGridInstantMessageViaXMLRPCAsync(im, result, - upd.RegionID); + SendGridInstantMessageViaXMLRPCAsync(im, result, upd.RegionID); } } else @@ -632,11 +604,8 @@ namespace OpenSim.Region.CoreModules.Avatar.InstantMessage XmlRpcRequest GridReq = new XmlRpcRequest("grid_instant_message", SendParams); try { - XmlRpcResponse GridResp = GridReq.Send(reginfo.ServerURI, 3000); - Hashtable responseData = (Hashtable)GridResp.Value; - if (responseData.ContainsKey("success")) { if ((string)responseData["success"] == "TRUE")