summaryrefslogtreecommitdiff
path: root/qpid/dotnet/Qpid.Client.Transport.Socket.Blocking/BlockingSocketTransport.cs
blob: 17f911fb6dca2d44ec5af8011e3f7fefaa504ac6 (plain)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
/*
 *
 * Licensed to the Apache Software Foundation (ASF) under one
 * or more contributor license agreements.  See the NOTICE file
 * distributed with this work for additional information
 * regarding copyright ownership.  The ASF licenses this file
 * to you under the Apache License, Version 2.0 (the
 * "License"); you may not use this file except in compliance
 * with the License.  You may obtain a copy of the License at
 * 
 *   http://www.apache.org/licenses/LICENSE-2.0
 * 
 * Unless required by applicable law or agreed to in writing,
 * software distributed under the License is distributed on an
 * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
 * KIND, either express or implied.  See the License for the
 * specific language governing permissions and limitations
 * under the License.
 *
 */
using System;
using System.Collections;
using System.Threading;
using log4net;
using Apache.Qpid.Client.Protocol;
using Apache.Qpid.Framing;

namespace Qpid.Client.Transport.Socket.Blocking
{
    public class BlockingSocketTransport : ITransport
    {
//        static readonly ILog _log = LogManager.GetLogger(typeof(BlockingSocketTransport));

        // Configuration variables.
        string _host;
        int _port;
        IProtocolListener _protocolListener;

        // Runtime variables.
        private BlockingSocketProcessor _socketProcessor;
        private AmqpChannel _amqpChannel;
        
        private ReaderRunner _readerRunner;
        private Thread _readerThread;       

        public BlockingSocketTransport(string host, int port, AMQConnection connection)
        {
            _host = host;
            _port = port;
            _protocolListener = connection.ProtocolListener;
        }
        
        public void Open()
        {
            _socketProcessor = new BlockingSocketProcessor(_host, _port, _protocolListener);
            _socketProcessor.Connect();
            _amqpChannel = new AmqpChannel(_socketProcessor.ByteChannel);
            _readerRunner = new ReaderRunner(this);
            _readerThread = new Thread(new ThreadStart(_readerRunner.Run));  
            _readerThread.Start();
        }

        public string getLocalEndPoint()
        {
            return _socketProcessor.getLocalEndPoint();
        }

        public void Close()
        {
            StopReaderThread();
            _socketProcessor.Disconnect();
        }

        public IProtocolChannel ProtocolChannel { get { return _amqpChannel;  } }
        public IProtocolWriter ProtocolWriter { get { return _amqpChannel; } }

        public void StopReaderThread()
        {
            _readerRunner.Stop();
        }

        class ReaderRunner
        {
            BlockingSocketTransport _transport;
            bool _running = true;

            public ReaderRunner(BlockingSocketTransport transport)
            {
                _transport = transport;
            }

            public void Run()
            {
                try
                {
                    while (_running)
                    {
                        Queue frames = _transport.ProtocolChannel.Read();

                        foreach (IDataBlock dataBlock in frames)
                        {
                            _transport._protocolListener.OnMessage(dataBlock);
                        }
                    }
                }
                catch (Exception e)
                {
                    _transport._protocolListener.OnException(e);
                }
            }

            public void Stop()
            {
                // TODO: Check if this is thread safe. running is not volitile....
                _running = false;
            }
        }
    }
}