Manager.cs :  » Workflows » NServiceBus » Grid » C# / CSharp Open Source

Home
C# / CSharp Open Source
1.2.6.4 mono .net core
2.2.6.4 mono core
3.Aspect Oriented Frameworks
4.Bloggers
5.Build Systems
6.Business Application
7.Charting Reporting Tools
8.Chat Servers
9.Code Coverage Tools
10.Content Management Systems CMS
11.CRM ERP
12.Database
13.Development
14.Email
15.Forum
16.Game
17.GIS
18.GUI
19.IDEs
20.Installers Generators
21.Inversion of Control Dependency Injection
22.Issue Tracking
23.Logging Tools
24.Message
25.Mobile
26.Network Clients
27.Network Servers
28.Office
29.PDF
30.Persistence Frameworks
31.Portals
32.Profilers
33.Project Management
34.RSS RDF
35.Rule Engines
36.Script
37.Search Engines
38.Sound Audio
39.Source Control
40.SQL Clients
41.Template Engines
42.Testing
43.UML
44.Web Frameworks
45.Web Service
46.Web Testing
47.Wiki Engines
48.Windows Presentation Foundation
49.Workflows
50.XML Parsers
C# / C Sharp
C# / C Sharp by API
C# / CSharp Tutorial
C# / CSharp Open Source » Workflows » NServiceBus 
NServiceBus » Grid » Manager.cs
using System;
using System.Collections.Generic;
using System.Diagnostics;
using System.Threading;
using NServiceBus.Grid.Messages;
using NServiceBus;
using System.IO;
using System.Runtime.Serialization.Formatters.Binary;
using System.Messaging;


namespace Grid{
    public static class Manager
    {
        private static readonly string path = "storage.txt";

        private static readonly int refreshInterval = 5;
        public static int RefreshInterval
        {
            get
            {
                lock (typeof(Manager))
                    return refreshInterval;
            }
        }

        private static IBus bus;

        private static List<ManagedEndpoint> endpoints = new List<ManagedEndpoint>();

        private static readonly Timer timer = new Timer(CheckNumberOfMessages);

        static Manager()
        {
            if (File.Exists(path))
            {
                try
                {
                    Stream stream = File.OpenRead(path);
                    BinaryFormatter formatter = new BinaryFormatter();

                    object result = formatter.Deserialize(stream);
                    stream.Close();

                    endpoints = result as List<ManagedEndpoint>;
                }
                catch(Exception)
                {
                    // intentionally swallow exception
                }
            }

            timer.Change(0, refreshInterval*1000);
        }

        public static void SetBus(IBus b)
        {
            bus = b;
        }

        private static void CheckNumberOfMessages(object state)
        {
            timer.Change(int.MaxValue, int.MaxValue);

            Stopwatch watch = new Stopwatch();
            watch.Start();

            List<ManagedEndpoint> myList;

            lock(typeof(Manager))
                myList = new List<ManagedEndpoint>(endpoints);

            foreach(ManagedEndpoint endpoint in myList)
            {
                MSMQ.MSMQManagementClass qMgmt = new MSMQ.MSMQManagementClass();
                object machine = Type.Missing;
                object missing = Type.Missing;
                object formatName = "DIRECT=OS:" + Environment.MachineName + "\\private$\\" + endpoint.Queue;

                try
                {
                    qMgmt.Init(ref machine, ref missing, ref formatName);
                    endpoint.SetNumberOfMessages(qMgmt.MessageCount);

                    MessageQueue q = new MessageQueue("FormatName:" + formatName as string);

                    MessagePropertyFilter mpf = new MessagePropertyFilter();
                    mpf.SetAll();

                    q.MessageReadPropertyFilter = mpf;

                    Message m = q.Peek();
                    if (m != null)
                        endpoint.AgeOfOldestMessage = DateTime.Now - m.SentTime;
                }
                catch
                {
                    //intentionally swallow bad endpoints
                }
            }

            watch.Stop();

            long due = refreshInterval*1000 - watch.ElapsedMilliseconds;
            due = (due < 0 ? 0 : due);

            timer.Change(due, refreshInterval*1000);
        }

        public static void Save()
        {
            if (!File.Exists(path))
                File.CreateText(path).Close();

            Stream stream = File.OpenWrite(path);
            BinaryFormatter formatter = new BinaryFormatter();

            formatter.Serialize(stream, endpoints);

            stream.Close();
        }

        public static List<ManagedEndpoint> GetManagedEndpoints()
        {
            lock (typeof(Manager))
                return new List<ManagedEndpoint>(endpoints);
        }

        public static void StoreManagedEndpoints(List<ManagedEndpoint> points)
        {
            lock(typeof(Manager))
                endpoints = points;
        }

        internal static void UpdateNumberOfWorkerThreads(string queue, int number)
        {
            List<ManagedEndpoint> myList;

            lock(typeof(Manager))
                myList = new List<ManagedEndpoint>(endpoints);

            foreach (ManagedEndpoint endpoint in myList)
                foreach (Worker w in endpoint.Workers)
                {
                    string[] aList = w.Queue.Split('\\');
                    string a = aList[aList.Length - 1].ToLower();

                    string[] bList = queue.Split('\\');
                    string b = bList[bList.Length - 1].ToLower();

                    if (a == b)
                        w.SetNumberOfWorkerThreads(number);
                }
        }

        public static void RefreshNumberOfWorkerThreads(string queue)
        {
            var message = new GetNumberOfWorkerThreadsMessage();
            bus.Send(queue, message).Register(
                delegate(IAsyncResult aResult)
                    {
                        var result =
                            aResult.AsyncState as CompletionResult;
                        if (result == null)
                            return;
                        if (result.Messages == null)
                            return;
                        if (result.Messages.Length != 1)
                            return;
                        var response = result.Messages[0] as GotNumberOfWorkerThreadsMessage;
                        if (response == null)
                            return;
                        var q = result.State as string;

                        UpdateNumberOfWorkerThreads(q, response.NumberOfWorkerThreads);
                    }, queue);
        }

        public static void SetNumberOfWorkerThreads(string queue, int number)
        {
            var message = new ChangeNumberOfWorkerThreadsMessage { NumberOfWorkerThreads = number };
            bus.Send(queue, message);

            RefreshNumberOfWorkerThreads(queue);
        }

    }
}
www.java2v.com | Contact Us
Copyright 2009 - 12 Demo Source and Support. All rights reserved.
All other trademarks are property of their respective owners.