From e18334a4bbf3cfd42f3506a8fea8cc3ff03e712d Mon Sep 17 00:00:00 2001 From: wj32 Date: Sat, 6 Jun 2009 03:37:24 +0000 Subject: [PATCH] added new wait manager with support for more than 63 objects git-svn-id: svn://svn.code.sf.net/p/processhacker/code@1378 21ef857c-d57f-4fe0-8362-d861dc6d29cd --- trunk/ProcessHacker.Common/WorkQueue.cs | 2 +- .../ProcessHacker.Native.csproj | 2 +- .../ProcessHacker.Native/Threading/Waiter.cs | 352 ++++++++++++++++++ .../Threading/WaiterThread.cs | 226 ----------- 4 files changed, 354 insertions(+), 228 deletions(-) create mode 100644 trunk/ProcessHacker.Native/Threading/Waiter.cs delete mode 100644 trunk/ProcessHacker.Native/Threading/WaiterThread.cs diff --git a/trunk/ProcessHacker.Common/WorkQueue.cs b/trunk/ProcessHacker.Common/WorkQueue.cs index 0e5216a90..e8792ff2c 100644 --- a/trunk/ProcessHacker.Common/WorkQueue.cs +++ b/trunk/ProcessHacker.Common/WorkQueue.cs @@ -565,7 +565,7 @@ namespace ProcessHacker.Common } } - // Check for work + // Check for work. if (_workQueue.Count > 0) { WorkItem workItem = null; diff --git a/trunk/ProcessHacker.Native/ProcessHacker.Native.csproj b/trunk/ProcessHacker.Native/ProcessHacker.Native.csproj index 3dba18a32..ba3acab81 100644 --- a/trunk/ProcessHacker.Native/ProcessHacker.Native.csproj +++ b/trunk/ProcessHacker.Native/ProcessHacker.Native.csproj @@ -144,7 +144,7 @@ - + Form diff --git a/trunk/ProcessHacker.Native/Threading/Waiter.cs b/trunk/ProcessHacker.Native/Threading/Waiter.cs new file mode 100644 index 000000000..9b713de67 --- /dev/null +++ b/trunk/ProcessHacker.Native/Threading/Waiter.cs @@ -0,0 +1,352 @@ +/* + * Process Hacker - + * wait manager + * + * Copyright (C) 2009 wj32 + * + * This file is part of Process Hacker. + * + * Process Hacker is free software; you can redistribute it and/or modify + * it under the terms of the GNU General Public License as published by + * the Free Software Foundation, either version 3 of the License, or + * (at your option) any later version. + * + * Process Hacker is distributed in the hope that it will be useful, + * but WITHOUT ANY WARRANTY; without even the implied warranty of + * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the + * GNU General Public License for more details. + * + * You should have received a copy of the GNU General Public License + * along with Process Hacker. If not, see . + */ + +using System; +using System.Collections.Generic; +using System.Text; +using System.Threading; +using ProcessHacker.Native.Api; +using ProcessHacker.Native.Objects; +using ProcessHacker.Native.Security; + +namespace ProcessHacker.Native.Threading +{ + public delegate void ObjectSignaledDelegate(ISynchronizable obj); + + public class Waiter : IDisposable + { + private class WaiterThread : IDisposable + { + public event ObjectSignaledDelegate ObjectSignaled; + + private bool _disposed = false; + private object _disposeLock = new object(); + private bool _terminating = false; + private Thread _thread; + private List _waitObjects = new List(); + private Event _waitChangedEvent = new Event(true, false); + + public WaiterThread() + { + _thread = new Thread(this.WaiterThreadStart); + _thread.IsBackground = true; + _thread.SetApartmentState(ApartmentState.STA); + _thread.Start(); + } + + ~WaiterThread() + { + this.Dispose(false); + } + + public void Dispose() + { + this.Dispose(true); + GC.SuppressFinalize(this); + } + + private void Dispose(bool disposing) + { + if (disposing) + Monitor.Enter(_disposeLock); + + try + { + if (!_disposed) + { + this.Terminate(); + _waitChangedEvent.Dispose(); + _disposed = true; + } + } + finally + { + Monitor.Exit(_disposeLock); + } + } + + public int Count + { + get + { + lock (_waitObjects) + return _waitObjects.Count; + } + } + + public ISynchronizable[] Objects + { + get + { + lock (_waitObjects) + return _waitObjects.ToArray(); + } + } + + public bool Add(ISynchronizable obj) + { + lock (_waitObjects) + { + // Check if we already have the maximum number of wait objects. + if (_waitObjects.Count >= Win32.MaximumWaitObjects - 1) + return false; + + _waitObjects.Add(obj); + this.NotifyChange(); + return true; + } + } + + public void NotifyChange() + { + _waitChangedEvent.Set(); + } + + private void OnObjectSignaled(ISynchronizable obj) + { + if (this.ObjectSignaled != null) + this.ObjectSignaled(obj); + } + + public bool Remove(ISynchronizable obj) + { + lock (_waitObjects) + { + if (!_waitObjects.Contains(obj)) + return false; + + _waitObjects.Remove(obj); + this.NotifyChange(); + return true; + } + } + + public void Terminate() + { + _terminating = true; + this.NotifyChange(); + } + + private void WaiterThreadStart() + { + ISynchronizable[] waitObjects; + + while (!_terminating) + { + lock (_waitObjects) + { + waitObjects = new ISynchronizable[_waitObjects.Count + 1]; + waitObjects[0] = _waitChangedEvent.Handle; + Array.Copy(_waitObjects.ToArray(), 0, waitObjects, 1, _waitObjects.Count); + } + + NtStatus waitStatus = NativeHandle.WaitAny(waitObjects); + + if (waitStatus == NtStatus.Wait0) + { + // The wait was changed. Go back to refresh the wait objects array. + // The event is also signaled to notify that the thread should terminate. + continue; + } + else if (waitStatus > NtStatus.Wait0 && waitStatus <= NtStatus.Wait63) + { + // One of the objects was signaled. + ISynchronizable signaledObject = waitObjects[(int)(waitStatus - NtStatus.Wait0)]; + // Remove the object now that it is signaled. + _waitObjects.Remove(signaledObject); + // Call the object-signaled event. + OnObjectSignaled(signaledObject); + } + } + } + } + + /// + /// Raised when an object is signaled. + /// + public event ObjectSignaledDelegate ObjectSignaled; + + private bool _disposed = false; + private object _disposeLock = new object(); + private List _waiterThreads = new List(); + private List _waitObjects = new List(); + + /// + /// Creates a waiter. + /// + public Waiter() + { + + } + + ~Waiter() + { + this.Dispose(false); + } + + /// + /// Releases resources used by the waiter. + /// + public void Dispose() + { + this.Dispose(true); + GC.SuppressFinalize(this); + } + + private void Dispose(bool disposing) + { + if (disposing) + Monitor.Enter(_disposeLock); + + try + { + if (!_disposed) + { + // Tell the waiter threads to terminate. + foreach (var waiterThread in _waiterThreads) + waiterThread.Terminate(); + _waiterThreads.Clear(); + + _disposed = true; + } + } + finally + { + Monitor.Exit(_disposeLock); + } + } + + public int Count + { + get + { + lock (_waitObjects) + return _waitObjects.Count; + } + } + + public ISynchronizable[] Objects + { + get + { + lock (_waitObjects) + return _waitObjects.ToArray(); + } + } + + /// + /// Adds an object for the waiter to wait on. + /// + /// The object to wait for. + public void Add(ISynchronizable obj) + { + lock (_waitObjects) + _waitObjects.Add(obj); + + foreach (var waiterThread in this.GetWaiterThreads()) + { + if (waiterThread.Add(obj)) + return; + } + + // We couldn't add the object to any existing waiter thread. + // Create a new waiter thread and add the object to that. + this.CreateWaiterThread(obj); + } + + private void BalanceWaiterThreads() + { + lock (_waitObjects) + { + // Eliminate waiter threads with no objects. + foreach (var waiterThread in this.GetWaiterThreads()) + { + if (waiterThread.Count == 0) + this.DeleteWaiterThread(waiterThread); + } + } + } + + private WaiterThread CreateWaiterThread() + { + return this.CreateWaiterThread(null); + } + + private WaiterThread CreateWaiterThread(ISynchronizable obj) + { + WaiterThread waiterThread = new WaiterThread(); + + waiterThread.ObjectSignaled += this.OnObjectSignaled; + + if (obj != null) + waiterThread.Add(obj); + + lock (_waiterThreads) + _waiterThreads.Add(waiterThread); + + return waiterThread; + } + + private void DeleteWaiterThread(WaiterThread waiterThread) + { + lock (_waiterThreads) + { + _waiterThreads.Remove(waiterThread); + waiterThread.Terminate(); + } + } + + private WaiterThread[] GetWaiterThreads() + { + lock (_waiterThreads) + return _waiterThreads.ToArray(); + } + + private void OnObjectSignaled(ISynchronizable obj) + { + if (ObjectSignaled != null) + ObjectSignaled(obj); + } + + /// + /// Removes an object the waiter is waiting on. + /// + /// An object which is currently being waited on. + public void Remove(ISynchronizable obj) + { + foreach (var waiterThread in this.GetWaiterThreads()) + { + if (waiterThread.Remove(obj)) + { + lock (_waitObjects) + _waitObjects.Remove(obj); + + this.BalanceWaiterThreads(); + return; + } + } + + // We couldn't remove the object. + throw new ArgumentException("The object is not being waited on."); + } + } +} diff --git a/trunk/ProcessHacker.Native/Threading/WaiterThread.cs b/trunk/ProcessHacker.Native/Threading/WaiterThread.cs deleted file mode 100644 index 0b86efa8b..000000000 --- a/trunk/ProcessHacker.Native/Threading/WaiterThread.cs +++ /dev/null @@ -1,226 +0,0 @@ -/* - * Process Hacker - - * waiter thread - * - * Copyright (C) 2009 wj32 - * - * This file is part of Process Hacker. - * - * Process Hacker is free software; you can redistribute it and/or modify - * it under the terms of the GNU General Public License as published by - * the Free Software Foundation, either version 3 of the License, or - * (at your option) any later version. - * - * Process Hacker is distributed in the hope that it will be useful, - * but WITHOUT ANY WARRANTY; without even the implied warranty of - * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the - * GNU General Public License for more details. - * - * You should have received a copy of the GNU General Public License - * along with Process Hacker. If not, see . - */ - -using System; -using System.Collections.Generic; -using System.Text; -using System.Threading; -using ProcessHacker.Native.Api; -using ProcessHacker.Native.Objects; -using ProcessHacker.Native.Security; - -namespace ProcessHacker.Native.Threading -{ - public delegate void ObjectSignaledDelegate(ISynchronizable obj); - - public class WaiterThread : IDisposable - { - private enum WaiterThreadMessageType - { - AddObject, - RemoveObject - } - - private class WaiterThreadMessage - { - private WaiterThreadMessageType _type; - private object _param; - - public WaiterThreadMessage(WaiterThreadMessageType type) - : this(type, null) - { } - - public WaiterThreadMessage(WaiterThreadMessageType type, object param) - { - _type = type; - _param = param; - } - - public WaiterThreadMessageType Type { get { return _type; } } - public object Param { get { return _param; } } - } - - /// - /// Raised when an object is signaled. - /// - public event ObjectSignaledDelegate ObjectSignaled; - - private bool _disposed = false; - private bool _terminating = false; - private object _disposeLock = new object(); - - private Thread _waiterThread; - private List _waitObjects = new List(); - - private Queue _waitMessageQueue = new Queue(); - private Event _waitMessageEvent = new Event(true, false); - - /// - /// Creates a waiter thread. - /// - public WaiterThread() - { - _waiterThread = new Thread(this.WaiterThreadStart); - _waiterThread.IsBackground = true; - _waiterThread.SetApartmentState(ApartmentState.STA); - _waiterThread.Start(); - } - - ~WaiterThread() - { - this.Dispose(false); - } - - /// - /// Releases resources used by the waiter thread. - /// - public void Dispose() - { - this.Dispose(true); - GC.SuppressFinalize(this); - } - - private void Dispose(bool disposing) - { - if (disposing) - Monitor.Enter(_disposeLock); - - try - { - if (!_disposed) - { - // Tell the waiter thread to terminate. - _terminating = true; - _waitMessageEvent.Set(); - - _disposed = true; - } - } - finally - { - Monitor.Exit(_disposeLock); - } - } - - /// - /// Adds an object for the waiter thread to wait on. - /// - /// The object to wait for. - public void Add(ISynchronizable obj) - { - lock (_waitMessageQueue) - { - if (_waitObjects.Count >= Win32.MaximumWaitObjects - 1) - throw new TooManyWaitObjectsException(); - - this.SendWaiterThreadMessage(new WaiterThreadMessage(WaiterThreadMessageType.AddObject, obj)); - } - } - - private void OnObjectSignaled(ISynchronizable obj) - { - if (ObjectSignaled != null) - ObjectSignaled(obj); - } - - /// - /// Removes an object the waiter thread is waiting on. - /// - /// An object which is currently being waited on. - public void Remove(ISynchronizable obj) - { - this.SendWaiterThreadMessage(new WaiterThreadMessage(WaiterThreadMessageType.RemoveObject, obj)); - } - - private void SendWaiterThreadMessage(WaiterThreadMessage message) - { - lock (_waitMessageQueue) - { - _waitMessageQueue.Enqueue(message); - _waitMessageEvent.Set(); - } - } - - private void WaiterThreadStart() - { - ISynchronizable[] waitObjects; - - while (!_terminating) - { - waitObjects = new ISynchronizable[_waitObjects.Count + 1]; - waitObjects[0] = _waitMessageEvent.Handle; - Array.Copy(_waitObjects.ToArray(), 0, waitObjects, 1, _waitObjects.Count); - - NtStatus waitStatus = NativeHandle.WaitAny(waitObjects); - - if (waitStatus == NtStatus.Wait0) - { - // We have a message in the message queue (or we are terminating). - WaiterThreadMessage message; - - // Lock and retrieve the message. - lock (_waitMessageQueue) - { - if (_waitMessageQueue.Count > 0) - message = _waitMessageQueue.Dequeue(); - else - continue; - } - - switch (message.Type) - { - case WaiterThreadMessageType.AddObject: - { - _waitObjects.Add((ISynchronizable)message.Param); - } - break; - case WaiterThreadMessageType.RemoveObject: - { - _waitObjects.Remove((ISynchronizable)message.Param); - } - break; - } - } - else if (waitStatus > NtStatus.Wait0 && waitStatus <= NtStatus.Wait63) - { - // One of the objects was signaled. - ISynchronizable signaledObject = waitObjects[(int)(waitStatus - NtStatus.Wait0)]; - // Remove the object now that it is signaled. - _waitObjects.Remove(signaledObject); - // Call the object-signaled event. - OnObjectSignaled(signaledObject); - } - } - } - } - - public class TooManyWaitObjectsException : Exception - { - public override string Message - { - get - { - return "An attempt was made to add too many objects to the wait list."; - } - } - } -}