-
Notifications
You must be signed in to change notification settings - Fork 390
/
tcp_session.h
291 lines (229 loc) · 8.97 KB
/
tcp_session.h
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
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
/*
* Copyright (c) 2013 Juniper Networks, Inc. All rights reserved.
*/
#ifndef __TCP_SESSION_H__
#define __TCP_SESSION_H__
#include <list>
#include <deque>
#include <boost/asio/buffer.hpp>
#include <boost/asio/io_service.hpp>
#include <boost/asio/ip/tcp.hpp>
#include <boost/asio/strand.hpp>
#include <boost/intrusive_ptr.hpp>
#include <boost/function.hpp>
#include <boost/scoped_ptr.hpp>
#include <tbb/mutex.h>
#include <tbb/task.h>
#ifndef _LIBCPP_VERSION
#include <tbb/compat/condition_variable>
#endif
#include "base/util.h"
#include "base/task.h"
#include "io/tcp_server.h"
class EventManager;
class TcpServer;
class TcpSession;
class TcpMessageWriter;
// TcpSession
//
// Concurrency: the session is created by the event manager thread, which
// also invokes the AsyncHandlers. ReleaseBuffer and Send will typically be
// invoked by a different thread.
class TcpSession {
public:
static const int kDefaultBufferSize = 4 * 1024;
enum Event {
EVENT_NONE,
ACCEPT,
CONNECT_COMPLETE,
CONNECT_FAILED,
CLOSE
};
enum Direction {
ACTIVE,
PASSIVE
};
typedef boost::asio::ip::tcp::socket Socket;
typedef boost::asio::ip::tcp::endpoint Endpoint;
typedef boost::function<void(TcpSession *, Event)> EventObserver;
typedef boost::asio::const_buffer Buffer;
// TcpSession constructor takes ownership of socket.
TcpSession(TcpServer *server, Socket *socket,
bool async_read_ready = true);
// Performs a non-blocking send operation.
virtual bool Send(const u_int8_t *data, size_t size, size_t *sent);
// Called by TcpServer to trigger async read.
virtual bool Connected(Endpoint remote);
// Called by TcpServer to trigger async read.
virtual void Accepted();
void ConnectFailed();
void Close();
virtual std::string ToString() const { return name_; }
void SetBufferSize(int buffer_size);
// Getters and setters
virtual Socket *socket() const { return socket_.get(); }
int sock_descriptor() { return socket_->native_handle(); }
TcpServer *server() { return server_.get(); }
int32_t local_port() const;
int32_t remote_port() const;
// Concurrency: changing the observer guarantees mutual exclusion with
// the observer invocation. e.g. if the caller sets the observer to NULL
// it is guaranteed that the observer will not get invoked after this
// method returns.
void set_observer(EventObserver observer);
// Buffers must be freed in arrival order.
virtual void ReleaseBuffer(Buffer buffer);
// This function returns the instance to run SessionTask.
// Returning Task::kTaskInstanceAny would allow multiple session tasks to
// run in parallel.
// Derived class may override implementation if it expects the all the
// Tasks of the session to run in specific instance
// Note: Two tasks of same task ID and task instance can't run
// at in parallel
// E.g. BgpSession is created per BgpPeer and to ensure that
// there is one SessionTask per peer, PeerIndex is returned
// from this function
virtual int GetSessionInstance() const;
static const uint8_t *BufferData(const Buffer &buffer) {
return boost::asio::buffer_cast<const uint8_t *>(buffer);
}
static size_t BufferSize(const Buffer &buffer) {
return boost::asio::buffer_size(buffer);
}
bool IsEstablished() const {
tbb::mutex::scoped_lock lock(mutex_);
return established_;
}
bool IsClosed() const {
tbb::mutex::scoped_lock lock(mutex_);
return closed_;
}
bool IsServerSession() {
if (direction_ == PASSIVE) return true;
return false;
}
Endpoint remote_endpoint() const {
return remote_;
}
Endpoint local_endpoint() const;
virtual boost::system::error_code SetSocketOptions();
static bool IsSocketErrorHard(const boost::system::error_code &ec);
void set_read_on_connect(bool read) { read_on_connect_ = read; }
void SessionEstablished(Endpoint remote, Direction direction);
virtual void AsyncReadStart();
const io::SocketStats &GetSocketStats() const { return stats_; }
void GetRxSocketStats(SocketIOStats &socket_stats) const;
void GetTxSocketStats(SocketIOStats &socket_stats) const;
int SetMd5SocketOption(uint32_t peer_ip, const std::string &md5_password);
int ClearMd5SocketOption(uint32_t peer_ip);
protected:
typedef boost::intrusive_ptr<TcpSession> TcpSessionPtr;
static void AsyncReadHandler(TcpSessionPtr session,
boost::asio::mutable_buffer buffer,
const boost::system::error_code &error,
size_t size);
static void AsyncWriteHandler(TcpSessionPtr session,
const boost::system::error_code &error);
void AsyncReadStartInternal(TcpSessionPtr session);
virtual Task* CreateReaderTask(boost::asio::mutable_buffer, size_t);
virtual ~TcpSession();
// Read handler. Called from a TBB task.
virtual void OnRead(Buffer buffer) = 0;
// Callback after socket is ready for write.
virtual void WriteReady(const boost::system::error_code &error);
virtual void AsyncReadSome(boost::asio::mutable_buffer buffer);
virtual std::size_t WriteSome(const uint8_t *data, std::size_t len,
boost::system::error_code &error);
virtual void AsyncWrite(const u_int8_t *data, std::size_t size);
virtual int reader_task_id() const {
return reader_task_id_;
}
EventObserver observer() { return observer_; }
boost::system::error_code SetSocketKeepaliveOptions(int keepalive_time,
int keepalive_intvl, int keepalive_probes);
void CloseInternal(bool call_observer, bool notify_server = true);
// Protects session state and buffer queue.
mutable tbb::mutex mutex_;
private:
class Reader;
friend class TcpServer;
friend class TcpMessageWriter;
friend void intrusive_ptr_add_ref(TcpSession *session);
friend void intrusive_ptr_release(TcpSession *session);
typedef std::list<boost::asio::mutable_buffer> BufferQueue;
typedef boost::asio::strand Strand;
static void WriteReadyInternal(TcpSessionPtr session,
const boost::system::error_code &error,
uint64_t block_start_time);
void DeferWriter();
void ReleaseBufferLocked(Buffer buffer);
void SetEstablished(Endpoint remote, Direction dir);
bool IsClosedLocked() const {
return closed_;
}
void SetName();
boost::asio::mutable_buffer AllocateBuffer();
void DeleteBuffer(boost::asio::mutable_buffer buffer);
static int reader_task_id_;
TcpServerPtr server_;
boost::scoped_ptr<Socket> socket_;
boost::scoped_ptr<Strand> io_strand_;
bool read_on_connect_;
int buffer_size_;
/**************** protected by mutex_ ****************/
bool established_; // In TCP ESTABLISHED state.
bool closed_; // Close has been called.
Endpoint remote_; // Remote end-point
Direction direction_; // direction (active, passive)
BufferQueue buffer_queue_;
/**************** end protected by mutex_ ****************/
// Protects observer manipulation and invocation. When this lock is
// held the session mutex should not be held and vice-versa.
tbb::mutex obs_mutex_;
EventObserver observer_;
io::SocketStats stats_;
boost::scoped_ptr<TcpMessageWriter> writer_;
tbb::atomic<int> refcount_;
std::string name_;
DISALLOW_COPY_AND_ASSIGN(TcpSession);
};
inline void intrusive_ptr_add_ref(TcpSession *session) {
session->refcount_.fetch_and_increment();
}
inline void intrusive_ptr_release(TcpSession *session) {
int prev = session->refcount_.fetch_and_decrement();
if (prev == 1) {
delete session;
}
}
// TcpMessageReader
//
// Provides base implementation of OnRead() for TcpSession assuming
// fixed message header length
//
class TcpMessageReader {
public:
typedef boost::asio::const_buffer Buffer;
typedef boost::function<bool(const u_int8_t *, size_t)> ReceiveCallback;
TcpMessageReader(TcpSession *session, ReceiveCallback callback);
virtual ~TcpMessageReader();
virtual void OnRead(Buffer buffer);
protected:
virtual int MsgLength(Buffer buffer, int offset) = 0;
virtual const int GetHeaderLenSize() = 0;
virtual const int GetMaxMessageSize() = 0;
private:
typedef std::deque<Buffer> BufferQueue;
// Copy the queue into one contiguous buffer.
uint8_t *BufferConcat(uint8_t *data, Buffer buffer, int msglength);
int QueueByteLength() const;
Buffer PullUp(uint8_t *data, Buffer buffer, size_t size) const;
int AllocBufferSize(int length);
TcpSession *session_;
ReceiveCallback callback_;
BufferQueue queue_;
int offset_;
int remain_;
DISALLOW_COPY_AND_ASSIGN(TcpMessageReader);
};
#endif // __TCP_SESSION_H__