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
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
|
/*****************************************************************************
$Id$
File: em.h
Date: 06Apr06
Copyright (C) 2006-07 by Francis Cianfrocca. All Rights Reserved.
Gmail: blackhedd
This program is free software; you can redistribute it and/or modify
it under the terms of either: 1) the GNU General Public License
as published by the Free Software Foundation; either version 2 of the
License, or (at your option) any later version; or 2) Ruby's License.
See the file COPYING for complete licensing information.
*****************************************************************************/
#ifndef __EventMachine__H_
#define __EventMachine__H_
#ifdef BUILD_FOR_RUBY
#include <ruby.h>
#ifdef HAVE_RB_THREAD_FD_SELECT
#define EmSelect rb_thread_fd_select
#else
// ruby 1.9.1 and below
#define EmSelect rb_thread_select
#endif
#ifdef HAVE_RB_THREAD_CALL_WITHOUT_GVL
#include <ruby/thread.h>
#endif
#ifdef HAVE_RB_WAIT_FOR_SINGLE_FD
#include <ruby/io.h>
#endif
#if defined(HAVE_RB_TRAP_IMMEDIATE)
#include <rubysig.h>
#elif defined(HAVE_RB_ENABLE_INTERRUPT)
extern "C" {
void rb_enable_interrupt(void);
void rb_disable_interrupt(void);
}
#define TRAP_BEG rb_enable_interrupt()
#define TRAP_END do { rb_disable_interrupt(); rb_thread_check_ints(); } while(0)
#else
#define TRAP_BEG
#define TRAP_END
#endif
// 1.9.0 compat
#ifndef RUBY_UBF_IO
#define RUBY_UBF_IO RB_UBF_DFL
#endif
#ifndef RSTRING_PTR
#define RSTRING_PTR(str) RSTRING(str)->ptr
#endif
#ifndef RSTRING_LEN
#define RSTRING_LEN(str) RSTRING(str)->len
#endif
#ifndef RSTRING_LENINT
#define RSTRING_LENINT(str) RSTRING_LEN(str)
#endif
#else
#define EmSelect select
#endif
#if !defined(HAVE_TYPE_RB_FDSET_T)
#define fd_check(n) (((n) < FD_SETSIZE) ? 1 : 0*fprintf(stderr, "fd %d too large for select\n", (n)))
// These definitions are cribbed from include/ruby/intern.h in Ruby 1.9.3,
// with this change: any macros that read or write the nth element of an
// fdset first call fd_check to make sure n is in bounds.
typedef fd_set rb_fdset_t;
#define rb_fd_zero(f) FD_ZERO(f)
#define rb_fd_set(n, f) do { if (fd_check(n)) FD_SET((n), (f)); } while(0)
#define rb_fd_clr(n, f) do { if (fd_check(n)) FD_CLR((n), (f)); } while(0)
#define rb_fd_isset(n, f) (fd_check(n) ? FD_ISSET((n), (f)) : 0)
#define rb_fd_copy(d, s, n) (*(d) = *(s))
#define rb_fd_dup(d, s) (*(d) = *(s))
#define rb_fd_resize(n, f) ((void)(f))
#define rb_fd_ptr(f) (f)
#define rb_fd_init(f) FD_ZERO(f)
#define rb_fd_init_copy(d, s) (*(d) = *(s))
#define rb_fd_term(f) ((void)(f))
#define rb_fd_max(f) FD_SETSIZE
#define rb_fd_select(n, rfds, wfds, efds, timeout) \
select(fd_check((n)-1) ? (n) : FD_SETSIZE, (rfds), (wfds), (efds), (timeout))
#define rb_thread_fd_select(n, rfds, wfds, efds, timeout) \
rb_thread_select(fd_check((n)-1) ? (n) : FD_SETSIZE, (rfds), (wfds), (efds), (timeout))
#endif
// This Solaris fix is adapted from eval_intern.h in Ruby 1.9.3:
// Solaris sys/select.h switches select to select_large_fdset to support larger
// file descriptors if FD_SETSIZE is larger than 1024 on 32bit environment.
// But Ruby doesn't change FD_SETSIZE because fd_set is allocated dynamically.
// So following definition is required to use select_large_fdset.
#ifdef HAVE_SELECT_LARGE_FDSET
#define select(n, r, w, e, t) select_large_fdset((n), (r), (w), (e), (t))
extern "C" {
int select_large_fdset(int, fd_set *, fd_set *, fd_set *, struct timeval *);
}
#endif
class EventableDescriptor;
class InotifyDescriptor;
struct SelectData_t;
/*************
enum Poller_t
*************/
enum Poller_t {
Poller_Default, // typically Select
Poller_Epoll,
Poller_Kqueue
};
/********************
class EventMachine_t
********************/
class EventMachine_t
{
public:
static int GetMaxTimerCount();
static void SetMaxTimerCount (int);
static int GetSimultaneousAcceptCount();
static void SetSimultaneousAcceptCount (int);
public:
EventMachine_t (EMCallback, Poller_t);
virtual ~EventMachine_t();
bool RunOnce();
void Run();
void ScheduleHalt();
bool Stopping();
void SignalLoopBreaker();
const uintptr_t InstallOneshotTimer (uint64_t);
const uintptr_t ConnectToServer (const char *, int, const char *, int);
const uintptr_t ConnectToUnixServer (const char *);
const uintptr_t CreateTcpServer (const char *, int);
const uintptr_t OpenDatagramSocket (const char *, int);
const uintptr_t CreateUnixDomainServer (const char*);
const uintptr_t AttachSD (SOCKET);
const uintptr_t OpenKeyboard();
//const char *Popen (const char*, const char*);
const uintptr_t Socketpair (char* const*);
void Add (EventableDescriptor*);
void Modify (EventableDescriptor*);
void Deregister (EventableDescriptor*);
const uintptr_t AttachFD (SOCKET, bool);
int DetachFD (EventableDescriptor*);
void ArmKqueueWriter (EventableDescriptor*);
void ArmKqueueReader (EventableDescriptor*);
void SetTimerQuantum (int);
static void SetuidString (const char*);
static int SetRlimitNofile (int);
pid_t SubprocessPid;
int SubprocessExitStatus;
int GetConnectionCount();
float GetHeartbeatInterval();
int SetHeartbeatInterval(float);
const uintptr_t WatchFile (const char*);
void UnwatchFile (int);
void UnwatchFile (const uintptr_t);
#ifdef HAVE_KQUEUE
void _HandleKqueueFileEvent (struct kevent*);
void _RegisterKqueueFileEvent(int);
#endif
const uintptr_t WatchPid (int);
void UnwatchPid (int);
void UnwatchPid (const uintptr_t);
#ifdef HAVE_KQUEUE
void _HandleKqueuePidEvent (struct kevent*);
#endif
uint64_t GetCurrentLoopTime() { return MyCurrentLoopTime; }
void QueueHeartbeat(EventableDescriptor*);
void ClearHeartbeat(uint64_t, EventableDescriptor*);
uint64_t GetRealTime();
Poller_t GetPoller() { return Poller; }
static int name2address (const char *server, int port, int socktype, struct sockaddr *addr, size_t *addr_len);
private:
void _RunTimers();
void _UpdateTime();
void _AddNewDescriptors();
void _ModifyDescriptors();
void _InitializeLoopBreaker();
void _CleanupSockets();
void _RunSelectOnce();
void _RunEpollOnce();
void _RunKqueueOnce();
void _ModifyEpollEvent (EventableDescriptor*);
void _DispatchHeartbeats();
timeval _TimeTilNextEvent();
void _CleanBadDescriptors();
public:
void _ReadLoopBreaker();
void _ReadInotifyEvents();
int NumCloseScheduled;
private:
enum {
MaxEpollDescriptors = 64*1024,
MaxEvents = 4096
};
int HeartbeatInterval;
EMCallback EventCallback;
class Timer_t: public Bindable_t {
};
std::multimap<uint64_t, Timer_t> Timers;
std::multimap<uint64_t, EventableDescriptor*> Heartbeats;
std::map<int, Bindable_t*> Files;
std::map<int, Bindable_t*> Pids;
std::vector<EventableDescriptor*> Descriptors;
std::vector<EventableDescriptor*> NewDescriptors;
std::set<EventableDescriptor*> ModifiedDescriptors;
SOCKET LoopBreakerReader;
SOCKET LoopBreakerWriter;
#ifdef OS_WIN32
struct sockaddr_in LoopBreakerTarget;
#endif
timeval Quantum;
uint64_t MyCurrentLoopTime;
#ifdef OS_WIN32
unsigned TickCountTickover;
unsigned LastTickCount;
#endif
#ifdef OS_DARWIN
mach_timebase_info_data_t mach_timebase;
#endif
private:
bool bTerminateSignalReceived;
SelectData_t *SelectData;
Poller_t Poller;
int epfd; // Epoll file-descriptor
#ifdef HAVE_EPOLL
struct epoll_event epoll_events [MaxEvents];
#endif
int kqfd; // Kqueue file-descriptor
#ifdef HAVE_KQUEUE
struct kevent Karray [MaxEvents];
#endif
#ifdef HAVE_INOTIFY
InotifyDescriptor *inotify; // pollable descriptor for our inotify instance
#endif
};
/*******************
struct SelectData_t
*******************/
struct SelectData_t
{
SelectData_t();
~SelectData_t();
int _Select();
void _Clear();
SOCKET maxsocket;
rb_fdset_t fdreads;
rb_fdset_t fdwrites;
rb_fdset_t fderrors;
timeval tv;
int nSockets;
};
#endif // __EventMachine__H_
|