git migration: Refactor the ASTERISK_FILE_VERSION macro
[asterisk/asterisk.git] / res / res_corosync.c
1 /*
2  * Asterisk -- An open source telephony toolkit.
3  *
4  * Copyright (C) 2007, Digium, Inc.
5  * Copyright (C) 2012, Russell Bryant
6  *
7  * Russell Bryant <russell@russellbryant.net>
8  *
9  * See http://www.asterisk.org for more information about
10  * the Asterisk project. Please do not directly contact
11  * any of the maintainers of this project for assistance;
12  * the project provides a web site, mailing lists and IRC
13  * channels for your use.
14  *
15  * This program is free software, distributed under the terms of
16  * the GNU General Public License Version 2. See the LICENSE file
17  * at the top of the source tree.
18  */
19
20 /*!
21  * \file
22  * \author Russell Bryant <russell@russellbryant.net>
23  *
24  * This module is based on and replaces the previous res_ais module.
25  */
26
27 /*** MODULEINFO
28         <depend>corosync</depend>
29         <defaultenabled>no</defaultenabled>
30         <support_level>extended</support_level>
31  ***/
32
33 #include "asterisk.h"
34
35 ASTERISK_REGISTER_FILE();
36
37 #include <corosync/cpg.h>
38 #include <corosync/cfg.h>
39
40 #include "asterisk/module.h"
41 #include "asterisk/logger.h"
42 #include "asterisk/poll-compat.h"
43 #include "asterisk/config.h"
44 #include "asterisk/event.h"
45 #include "asterisk/cli.h"
46 #include "asterisk/devicestate.h"
47 #include "asterisk/app.h"
48 #include "asterisk/stasis.h"
49 #include "asterisk/stasis_message_router.h"
50
51 AST_RWLOCK_DEFINE_STATIC(event_types_lock);
52
53 static void publish_mwi_to_stasis(struct ast_event *event);
54 static void publish_device_state_to_stasis(struct ast_event *event);
55
56 /*! \brief The internal topic used for message forwarding and pings */
57 static struct stasis_topic *corosync_aggregate_topic;
58
59 /*! \brief Our \ref stasis message router */
60 static struct stasis_message_router *stasis_router;
61
62 /*! \brief Internal accessor for our topic */
63 static struct stasis_topic *corosync_topic(void)
64 {
65         return corosync_aggregate_topic;
66 }
67
68 /*! \brief A payload wrapper around a corosync ping event */
69 struct corosync_ping_payload {
70         /*! The corosync ping event being passed over \ref stasis */
71         struct ast_event *event;
72 };
73
74 /*! \brief Destructor for the \ref corosync_ping_payload wrapper object */
75 static void corosync_ping_payload_dtor(void *obj)
76 {
77         struct corosync_ping_payload *payload = obj;
78
79         ast_free(payload->event);
80 }
81
82 /*! \brief Convert a Corosync PING to a \ref ast_event */
83 static struct ast_event *corosync_ping_to_event(struct stasis_message *message)
84 {
85         struct corosync_ping_payload *payload;
86         struct ast_event *event;
87         struct ast_eid *event_eid;
88
89         if (!message) {
90                 return NULL;
91         }
92
93         payload = stasis_message_data(message);
94
95         if (!payload->event) {
96                 return NULL;
97         }
98
99         event_eid = (struct ast_eid *)ast_event_get_ie_raw(payload->event, AST_EVENT_IE_EID);
100
101         event = ast_event_new(AST_EVENT_PING,
102                                 AST_EVENT_IE_EID, AST_EVENT_IE_PLTYPE_RAW, event_eid, sizeof(*event_eid),
103                                 AST_EVENT_IE_END);
104
105         return event;
106 }
107
108 STASIS_MESSAGE_TYPE_DEFN_LOCAL(corosync_ping_message_type,
109         .to_event = corosync_ping_to_event, );
110
111 /*! \brief Publish a Corosync ping to \ref stasis */
112 static void publish_corosync_ping_to_stasis(struct ast_event *event)
113 {
114         struct corosync_ping_payload *payload;
115         struct stasis_message *message;
116
117         ast_assert(ast_event_get_type(event) == AST_EVENT_PING);
118         ast_assert(event != NULL);
119
120         if (!corosync_ping_message_type()) {
121                 return;
122         }
123
124         payload = ao2_t_alloc(sizeof(*payload), corosync_ping_payload_dtor, "Create ping payload");
125         if (!payload) {
126                 return;
127         }
128         payload->event = event;
129
130         message = stasis_message_create(corosync_ping_message_type(), payload);
131         if (!message) {
132                 ao2_t_ref(payload, -1, "Destroy payload on off nominal");
133                 return;
134         }
135
136         stasis_publish(corosync_topic(), message);
137
138         ao2_t_ref(payload, -1, "Hand ref to stasis");
139         ao2_t_ref(message, -1, "Hand ref to stasis");
140 }
141
142 static struct {
143         const char *name;
144         struct stasis_forward *sub;
145         unsigned char publish;
146         unsigned char publish_default;
147         unsigned char subscribe;
148         unsigned char subscribe_default;
149         struct stasis_topic *(* topic_fn)(void);
150         struct stasis_cache *(* cache_fn)(void);
151         struct stasis_message_type *(* message_type_fn)(void);
152         void (* publish_to_stasis)(struct ast_event *);
153 } event_types[] = {
154         [AST_EVENT_MWI] = { .name = "mwi",
155                             .topic_fn = ast_mwi_topic_all,
156                             .cache_fn = ast_mwi_state_cache,
157                             .message_type_fn = ast_mwi_state_type,
158                             .publish_to_stasis = publish_mwi_to_stasis, },
159         [AST_EVENT_DEVICE_STATE_CHANGE] = { .name = "device_state",
160                                             .topic_fn = ast_device_state_topic_all,
161                                             .cache_fn = ast_device_state_cache,
162                                             .message_type_fn = ast_device_state_message_type,
163                                             .publish_to_stasis = publish_device_state_to_stasis, },
164         [AST_EVENT_PING] = { .name = "ping",
165                              .publish_default = 1,
166                              .subscribe_default = 1,
167                              .topic_fn = corosync_topic,
168                              .message_type_fn = corosync_ping_message_type,
169                              .publish_to_stasis = publish_corosync_ping_to_stasis, },
170 };
171
172 static struct {
173         pthread_t id;
174         int alert_pipe[2];
175         unsigned int stop:1;
176 } dispatch_thread = {
177         .id = AST_PTHREADT_NULL,
178         .alert_pipe = { -1, -1 },
179 };
180
181 static cpg_handle_t cpg_handle;
182 static corosync_cfg_handle_t cfg_handle;
183
184 #ifdef HAVE_COROSYNC_CFG_STATE_TRACK
185 static void cfg_state_track_cb(
186                 corosync_cfg_state_notification_buffer_t *notification_buffer,
187                 cs_error_t error);
188 #endif /* HAVE_COROSYNC_CFG_STATE_TRACK */
189
190 static void cfg_shutdown_cb(corosync_cfg_handle_t cfg_handle,
191                 corosync_cfg_shutdown_flags_t flags);
192
193 static corosync_cfg_callbacks_t cfg_callbacks = {
194 #ifdef HAVE_COROSYNC_CFG_STATE_TRACK
195         .corosync_cfg_state_track_callback = cfg_state_track_cb,
196 #endif /* HAVE_COROSYNC_CFG_STATE_TRACK */
197         .corosync_cfg_shutdown_callback = cfg_shutdown_cb,
198 };
199
200 /*! \brief Publish a received MWI \ref ast_event to \ref stasis */
201 static void publish_mwi_to_stasis(struct ast_event *event)
202 {
203         const char *mailbox;
204         const char *context;
205         unsigned int new_msgs;
206         unsigned int old_msgs;
207         struct ast_eid *event_eid;
208
209         ast_assert(ast_event_get_type(event) == AST_EVENT_MWI);
210
211         mailbox = ast_event_get_ie_str(event, AST_EVENT_IE_MAILBOX);
212         context = ast_event_get_ie_str(event, AST_EVENT_IE_CONTEXT);
213         new_msgs = ast_event_get_ie_uint(event, AST_EVENT_IE_NEWMSGS);
214         old_msgs = ast_event_get_ie_uint(event, AST_EVENT_IE_OLDMSGS);
215         event_eid = (struct ast_eid *)ast_event_get_ie_raw(event, AST_EVENT_IE_EID);
216
217         if (ast_strlen_zero(mailbox) || ast_strlen_zero(context)) {
218                 return;
219         }
220
221         if (new_msgs > INT_MAX) {
222                 new_msgs = INT_MAX;
223         }
224
225         if (old_msgs > INT_MAX) {
226                 old_msgs = INT_MAX;
227         }
228
229         if (ast_publish_mwi_state_full(mailbox, context, (int)new_msgs,
230                                        (int)old_msgs, NULL, event_eid)) {
231                 char eid[16];
232                 ast_eid_to_str(eid, sizeof(eid), event_eid);
233                 ast_log(LOG_WARNING, "Failed to publish MWI message for %s@%s from %s\n",
234                         mailbox, context, eid);
235         }
236 }
237
238 /*! \brief Publish a received device state \ref ast_event to \ref stasis */
239 static void publish_device_state_to_stasis(struct ast_event *event)
240 {
241         const char *device;
242         enum ast_device_state state;
243         unsigned int cachable;
244         struct ast_eid *event_eid;
245
246         ast_assert(ast_event_get_type(event) == AST_EVENT_DEVICE_STATE_CHANGE);
247
248         device = ast_event_get_ie_str(event, AST_EVENT_IE_DEVICE);
249         state = ast_event_get_ie_uint(event, AST_EVENT_IE_STATE);
250         cachable = ast_event_get_ie_uint(event, AST_EVENT_IE_CACHABLE);
251         event_eid = (struct ast_eid *)ast_event_get_ie_raw(event, AST_EVENT_IE_EID);
252
253         if (ast_strlen_zero(device)) {
254                 return;
255         }
256
257         if (ast_publish_device_state_full(device, state, cachable, event_eid)) {
258                 char eid[16];
259                 ast_eid_to_str(eid, sizeof(eid), event_eid);
260                 ast_log(LOG_WARNING, "Failed to publish device state message for %s from %s\n",
261                         device, eid);
262         }
263 }
264
265 static void cpg_deliver_cb(cpg_handle_t handle, const struct cpg_name *group_name,
266                 uint32_t nodeid, uint32_t pid, void *msg, size_t msg_len);
267
268 static void cpg_confchg_cb(cpg_handle_t handle, const struct cpg_name *group_name,
269                 const struct cpg_address *member_list, size_t member_list_entries,
270                 const struct cpg_address *left_list, size_t left_list_entries,
271                 const struct cpg_address *joined_list, size_t joined_list_entries);
272
273 static cpg_callbacks_t cpg_callbacks = {
274         .cpg_deliver_fn = cpg_deliver_cb,
275         .cpg_confchg_fn = cpg_confchg_cb,
276 };
277
278 #ifdef HAVE_COROSYNC_CFG_STATE_TRACK
279 static void cfg_state_track_cb(
280                 corosync_cfg_state_notification_buffer_t *notification_buffer,
281                 cs_error_t error)
282 {
283 }
284 #endif /* HAVE_COROSYNC_CFG_STATE_TRACK */
285
286 static void cfg_shutdown_cb(corosync_cfg_handle_t cfg_handle,
287                 corosync_cfg_shutdown_flags_t flags)
288 {
289 }
290
291 static void cpg_deliver_cb(cpg_handle_t handle, const struct cpg_name *group_name,
292                 uint32_t nodeid, uint32_t pid, void *msg, size_t msg_len)
293 {
294         struct ast_event *event;
295         void (*publish_handler)(struct ast_event *) = NULL;
296         enum ast_event_type event_type;
297
298         if (msg_len < ast_event_minimum_length()) {
299                 ast_debug(1, "Ignoring event that's too small. %u < %u\n",
300                         (unsigned int) msg_len,
301                         (unsigned int) ast_event_minimum_length());
302                 return;
303         }
304
305         if (!ast_eid_cmp(&ast_eid_default, ast_event_get_ie_raw(msg, AST_EVENT_IE_EID))) {
306                 /* Don't feed events back in that originated locally. */
307                 return;
308         }
309
310         event_type = ast_event_get_type(msg);
311         if (event_type > AST_EVENT_TOTAL) {
312                 /* Egads, we don't support this */
313                 return;
314         }
315
316         ast_rwlock_rdlock(&event_types_lock);
317         publish_handler = event_types[event_type].publish_to_stasis;
318         if (!event_types[event_type].subscribe || !publish_handler) {
319                 /* We are not configured to subscribe to these events or
320                    we have no way to publish it internally. */
321                 ast_rwlock_unlock(&event_types_lock);
322                 return;
323         }
324         ast_rwlock_unlock(&event_types_lock);
325
326         if (!(event = ast_malloc(msg_len))) {
327                 return;
328         }
329
330         memcpy(event, msg, msg_len);
331
332         if (event_type == AST_EVENT_PING) {
333                 const struct ast_eid *eid;
334                 char buf[128] = "";
335
336                 eid = ast_event_get_ie_raw(event, AST_EVENT_IE_EID);
337                 ast_eid_to_str(buf, sizeof(buf), (struct ast_eid *) eid);
338                 ast_log(LOG_NOTICE, "Got event PING from server with EID: '%s'\n", buf);
339         }
340         ast_debug(5, "Publishing event %s (%u) to stasis\n",
341                 ast_event_get_type_name(event), event_type);
342         publish_handler(event);
343 }
344
345 static void publish_to_corosync(struct stasis_message *message)
346 {
347         cs_error_t cs_err;
348         struct iovec iov;
349         struct ast_event *event;
350
351         event = stasis_message_to_event(message);
352         if (!event) {
353                 return;
354         }
355
356         if (ast_eid_cmp(&ast_eid_default, ast_event_get_ie_raw(event, AST_EVENT_IE_EID))) {
357                 /* If the event didn't originate from this server, don't send it back out. */
358                 ast_event_destroy(event);
359                 return;
360         }
361
362         if (ast_event_get_type(event) == AST_EVENT_PING) {
363                 const struct ast_eid *eid;
364                 char buf[128] = "";
365
366                 eid = ast_event_get_ie_raw(event, AST_EVENT_IE_EID);
367                 ast_eid_to_str(buf, sizeof(buf), (struct ast_eid *) eid);
368                 ast_log(LOG_NOTICE, "Sending event PING from this server with EID: '%s'\n", buf);
369         }
370
371         iov.iov_base = (void *)event;
372         iov.iov_len = ast_event_get_size(event);
373
374         ast_debug(5, "Publishing event %s (%u) to corosync\n",
375                 ast_event_get_type_name(event), ast_event_get_type(event));
376
377         /* The stasis subscription will only exist if we are configured to publish
378          * these events, so just send away. */
379         if ((cs_err = cpg_mcast_joined(cpg_handle, CPG_TYPE_FIFO, &iov, 1)) != CS_OK) {
380                 ast_log(LOG_WARNING, "CPG mcast failed (%u)\n", cs_err);
381         }
382 }
383
384 static void stasis_message_cb(void *data, struct stasis_subscription *sub, struct stasis_message *message)
385 {
386         if (!message) {
387                 return;
388         }
389
390         publish_to_corosync(message);
391 }
392
393 static int dump_cache_cb(void *obj, void *arg, int flags)
394 {
395         struct stasis_message *message = obj;
396
397         if (!message) {
398                 return 0;
399         }
400
401         publish_to_corosync(message);
402
403         return 0;
404 }
405
406 static void cpg_confchg_cb(cpg_handle_t handle, const struct cpg_name *group_name,
407                 const struct cpg_address *member_list, size_t member_list_entries,
408                 const struct cpg_address *left_list, size_t left_list_entries,
409                 const struct cpg_address *joined_list, size_t joined_list_entries)
410 {
411         unsigned int i;
412
413         /* If any new nodes have joined, dump our cache of events we are publishing
414          * that originated from this server. */
415
416         if (!joined_list_entries) {
417                 return;
418         }
419
420         for (i = 0; i < ARRAY_LEN(event_types); i++) {
421                 struct ao2_container *messages;
422
423                 ast_rwlock_rdlock(&event_types_lock);
424                 if (!event_types[i].publish) {
425                         ast_rwlock_unlock(&event_types_lock);
426                         continue;
427                 }
428
429                 if (!event_types[i].cache_fn || !event_types[i].message_type_fn) {
430                         ast_rwlock_unlock(&event_types_lock);
431                         continue;
432                 }
433
434                 messages = stasis_cache_dump_by_eid(event_types[i].cache_fn(),
435                         event_types[i].message_type_fn(),
436                         &ast_eid_default);
437                 ast_rwlock_unlock(&event_types_lock);
438
439                 ao2_callback(messages, OBJ_NODATA, dump_cache_cb, NULL);
440
441                 ao2_t_ref(messages, -1, "Dispose of dumped cache");
442         }
443 }
444
445 static void *dispatch_thread_handler(void *data)
446 {
447         cs_error_t cs_err;
448         struct pollfd pfd[3] = {
449                 { .events = POLLIN, },
450                 { .events = POLLIN, },
451                 { .events = POLLIN, },
452         };
453
454         if ((cs_err = cpg_fd_get(cpg_handle, &pfd[0].fd)) != CS_OK) {
455                 ast_log(LOG_ERROR, "Failed to get CPG fd.  This module is now broken.\n");
456                 return NULL;
457         }
458
459         if ((cs_err = corosync_cfg_fd_get(cfg_handle, &pfd[1].fd)) != CS_OK) {
460                 ast_log(LOG_ERROR, "Failed to get CFG fd.  This module is now broken.\n");
461                 return NULL;
462         }
463
464         pfd[2].fd = dispatch_thread.alert_pipe[0];
465
466         while (!dispatch_thread.stop) {
467                 int res;
468
469                 cs_err = CS_OK;
470
471                 pfd[0].revents = 0;
472                 pfd[1].revents = 0;
473                 pfd[2].revents = 0;
474
475                 res = ast_poll(pfd, ARRAY_LEN(pfd), -1);
476                 if (res == -1 && errno != EINTR && errno != EAGAIN) {
477                         ast_log(LOG_ERROR, "poll() error: %s (%d)\n", strerror(errno), errno);
478                         continue;
479                 }
480
481                 if (pfd[0].revents & POLLIN) {
482                         if ((cs_err = cpg_dispatch(cpg_handle, CS_DISPATCH_ALL)) != CS_OK) {
483                                 ast_log(LOG_WARNING, "Failed CPG dispatch: %u\n", cs_err);
484                         }
485                 }
486
487                 if (pfd[1].revents & POLLIN) {
488                         if ((cs_err = corosync_cfg_dispatch(cfg_handle, CS_DISPATCH_ALL)) != CS_OK) {
489                                 ast_log(LOG_WARNING, "Failed CFG dispatch: %u\n", cs_err);
490                         }
491                 }
492
493                 if (cs_err == CS_ERR_LIBRARY || cs_err == CS_ERR_BAD_HANDLE) {
494                         struct cpg_name name;
495
496                         /* If corosync gets restarted out from under Asterisk, try to recover. */
497
498                         ast_log(LOG_NOTICE, "Attempting to recover from corosync failure.\n");
499
500                         if ((cs_err = corosync_cfg_initialize(&cfg_handle, &cfg_callbacks)) != CS_OK) {
501                                 ast_log(LOG_ERROR, "Failed to initialize cfg (%d)\n", (int) cs_err);
502                                 sleep(5);
503                                 continue;
504                         }
505
506                         if ((cs_err = cpg_initialize(&cpg_handle, &cpg_callbacks) != CS_OK)) {
507                                 ast_log(LOG_ERROR, "Failed to initialize cpg (%d)\n", (int) cs_err);
508                                 sleep(5);
509                                 continue;
510                         }
511
512                         if ((cs_err = cpg_fd_get(cpg_handle, &pfd[0].fd)) != CS_OK) {
513                                 ast_log(LOG_ERROR, "Failed to get CPG fd.\n");
514                                 sleep(5);
515                                 continue;
516                         }
517
518                         if ((cs_err = corosync_cfg_fd_get(cfg_handle, &pfd[1].fd)) != CS_OK) {
519                                 ast_log(LOG_ERROR, "Failed to get CFG fd.\n");
520                                 sleep(5);
521                                 continue;
522                         }
523
524                         ast_copy_string(name.value, "asterisk", sizeof(name.value));
525                         name.length = strlen(name.value);
526                         if ((cs_err = cpg_join(cpg_handle, &name)) != CS_OK) {
527                                 ast_log(LOG_ERROR, "Failed to join cpg (%d)\n", (int) cs_err);
528                                 sleep(5);
529                                 continue;
530                         }
531
532                         ast_log(LOG_NOTICE, "Corosync recovery complete.\n");
533                 }
534         }
535
536         return NULL;
537 }
538
539 static char *corosync_show_members(struct ast_cli_entry *e, int cmd, struct ast_cli_args *a)
540 {
541         cs_error_t cs_err;
542         cpg_iteration_handle_t cpg_iter;
543         struct cpg_iteration_description_t cpg_desc;
544         unsigned int i;
545
546         switch (cmd) {
547         case CLI_INIT:
548                 e->command = "corosync show members";
549                 e->usage =
550                         "Usage: corosync show members\n"
551                         "       Show corosync cluster members\n";
552                 return NULL;
553
554         case CLI_GENERATE:
555                 return NULL;    /* no completion */
556         }
557
558         if (a->argc != e->args) {
559                 return CLI_SHOWUSAGE;
560         }
561
562         cs_err = cpg_iteration_initialize(cpg_handle, CPG_ITERATION_ALL, NULL, &cpg_iter);
563
564         if (cs_err != CS_OK) {
565                 ast_cli(a->fd, "Failed to initialize CPG iterator.\n");
566                 return CLI_FAILURE;
567         }
568
569         ast_cli(a->fd, "\n"
570                     "=============================================================\n"
571                     "=== Cluster members =========================================\n"
572                     "=============================================================\n"
573                     "===\n");
574
575         for (i = 1, cs_err = cpg_iteration_next(cpg_iter, &cpg_desc);
576                         cs_err == CS_OK;
577                         cs_err = cpg_iteration_next(cpg_iter, &cpg_desc), i++) {
578                 corosync_cfg_node_address_t addrs[8];
579                 int num_addrs = 0;
580                 unsigned int j;
581
582                 cs_err = corosync_cfg_get_node_addrs(cfg_handle, cpg_desc.nodeid,
583                                 ARRAY_LEN(addrs), &num_addrs, addrs);
584                 if (cs_err != CS_OK) {
585                         ast_log(LOG_WARNING, "Failed to get node addresses\n");
586                         continue;
587                 }
588
589                 ast_cli(a->fd, "=== Node %u\n", i);
590                 ast_cli(a->fd, "=== --> Group: %s\n", cpg_desc.group.value);
591
592                 for (j = 0; j < num_addrs; j++) {
593                         struct sockaddr *sa = (struct sockaddr *) addrs[j].address;
594                         size_t sa_len = (size_t) addrs[j].address_length;
595                         char buf[128];
596
597                         getnameinfo(sa, sa_len, buf, sizeof(buf), NULL, 0, NI_NUMERICHOST);
598
599                         ast_cli(a->fd, "=== --> Address %u: %s\n", j + 1, buf);
600                 }
601
602         }
603
604         ast_cli(a->fd, "===\n"
605                        "=============================================================\n"
606                        "\n");
607
608         cpg_iteration_finalize(cpg_iter);
609
610         return CLI_SUCCESS;
611 }
612
613 static char *corosync_ping(struct ast_cli_entry *e, int cmd, struct ast_cli_args *a)
614 {
615         struct ast_event *event;
616
617         switch (cmd) {
618         case CLI_INIT:
619                 e->command = "corosync ping";
620                 e->usage =
621                         "Usage: corosync ping\n"
622                         "       Send a test ping to the cluster.\n"
623                         "A NOTICE will be in the log for every ping received\n"
624                         "on a server.\n  If you send a ping, you should see a NOTICE\n"
625                         "in the log for every server in the cluster.\n";
626                 return NULL;
627
628         case CLI_GENERATE:
629                 return NULL;    /* no completion */
630         }
631
632         if (a->argc != e->args) {
633                 return CLI_SHOWUSAGE;
634         }
635
636         event = ast_event_new(AST_EVENT_PING, AST_EVENT_IE_END);
637
638         if (!event) {
639                 return CLI_FAILURE;
640         }
641
642         ast_rwlock_rdlock(&event_types_lock);
643         event_types[AST_EVENT_PING].publish_to_stasis(event);
644         ast_rwlock_unlock(&event_types_lock);
645
646         return CLI_SUCCESS;
647 }
648
649 static char *corosync_show_config(struct ast_cli_entry *e, int cmd, struct ast_cli_args *a)
650 {
651         unsigned int i;
652
653         switch (cmd) {
654         case CLI_INIT:
655                 e->command = "corosync show config";
656                 e->usage =
657                         "Usage: corosync show config\n"
658                         "       Show configuration loaded from res_corosync.conf\n";
659                 return NULL;
660
661         case CLI_GENERATE:
662                 return NULL;    /* no completion */
663         }
664
665         if (a->argc != e->args) {
666                 return CLI_SHOWUSAGE;
667         }
668
669         ast_cli(a->fd, "\n"
670                     "=============================================================\n"
671                     "=== res_corosync config =====================================\n"
672                     "=============================================================\n"
673                     "===\n");
674
675         ast_rwlock_rdlock(&event_types_lock);
676         for (i = 0; i < ARRAY_LEN(event_types); i++) {
677                 if (event_types[i].publish) {
678                         ast_cli(a->fd, "=== ==> Publishing Event Type: %s\n",
679                                         event_types[i].name);
680                 }
681                 if (event_types[i].subscribe) {
682                         ast_cli(a->fd, "=== ==> Subscribing to Event Type: %s\n",
683                                         event_types[i].name);
684                 }
685         }
686         ast_rwlock_unlock(&event_types_lock);
687
688         ast_cli(a->fd, "===\n"
689                        "=============================================================\n"
690                        "\n");
691
692         return CLI_SUCCESS;
693 }
694
695 static struct ast_cli_entry corosync_cli[] = {
696         AST_CLI_DEFINE(corosync_show_config, "Show configuration"),
697         AST_CLI_DEFINE(corosync_show_members, "Show cluster members"),
698         AST_CLI_DEFINE(corosync_ping, "Send a test ping to the cluster"),
699 };
700
701 enum {
702         PUBLISH,
703         SUBSCRIBE,
704 };
705
706 static int set_event(const char *event_type, int pubsub)
707 {
708         unsigned int i;
709
710         for (i = 0; i < ARRAY_LEN(event_types); i++) {
711                 if (!event_types[i].name || strcasecmp(event_type, event_types[i].name)) {
712                         continue;
713                 }
714
715                 switch (pubsub) {
716                 case PUBLISH:
717                         event_types[i].publish = 1;
718                         break;
719                 case SUBSCRIBE:
720                         event_types[i].subscribe = 1;
721                         break;
722                 }
723
724                 break;
725         }
726
727         return (i == ARRAY_LEN(event_types)) ? -1 : 0;
728 }
729
730 static int load_general_config(struct ast_config *cfg)
731 {
732         struct ast_variable *v;
733         int res = 0;
734         unsigned int i;
735
736         ast_rwlock_wrlock(&event_types_lock);
737
738         for (i = 0; i < ARRAY_LEN(event_types); i++) {
739                 event_types[i].publish = event_types[i].publish_default;
740                 event_types[i].subscribe = event_types[i].subscribe_default;
741         }
742
743         for (v = ast_variable_browse(cfg, "general"); v && !res; v = v->next) {
744                 if (!strcasecmp(v->name, "publish_event")) {
745                         res = set_event(v->value, PUBLISH);
746                 } else if (!strcasecmp(v->name, "subscribe_event")) {
747                         res = set_event(v->value, SUBSCRIBE);
748                 } else {
749                         ast_log(LOG_WARNING, "Unknown option '%s'\n", v->name);
750                 }
751         }
752
753         for (i = 0; i < ARRAY_LEN(event_types); i++) {
754                 if (event_types[i].publish && !event_types[i].sub) {
755                         event_types[i].sub = stasis_forward_all(event_types[i].topic_fn(),
756                                                                                                         corosync_topic());
757                         stasis_message_router_add(stasis_router,
758                                                   event_types[i].message_type_fn(),
759                                                   stasis_message_cb,
760                                                   NULL);
761                 } else if (!event_types[i].publish && event_types[i].sub) {
762                         event_types[i].sub = stasis_forward_cancel(event_types[i].sub);
763                         stasis_message_router_remove(stasis_router,
764                                                      event_types[i].message_type_fn());
765                 }
766         }
767
768         ast_rwlock_unlock(&event_types_lock);
769
770         return res;
771 }
772
773 static int load_config(unsigned int reload)
774 {
775         static const char filename[] = "res_corosync.conf";
776         struct ast_config *cfg;
777         const char *cat = NULL;
778         struct ast_flags config_flags = { 0 };
779         int res = 0;
780
781         cfg = ast_config_load(filename, config_flags);
782
783         if (cfg == CONFIG_STATUS_FILEMISSING || cfg == CONFIG_STATUS_FILEINVALID) {
784                 return -1;
785         }
786
787         while ((cat = ast_category_browse(cfg, cat))) {
788                 if (!strcasecmp(cat, "general")) {
789                         res = load_general_config(cfg);
790                 } else {
791                         ast_log(LOG_WARNING, "Unknown configuration section '%s'\n", cat);
792                 }
793         }
794
795         ast_config_destroy(cfg);
796
797         return res;
798 }
799
800 static void cleanup_module(void)
801 {
802         cs_error_t cs_err;
803         unsigned int i;
804
805         if (stasis_router) {
806
807                 /* Unsubscribe all topic forwards and cancel all message routes */
808                 ast_rwlock_wrlock(&event_types_lock);
809                 for (i = 0; i < ARRAY_LEN(event_types); i++) {
810                         if (event_types[i].sub) {
811                                 event_types[i].sub = stasis_forward_cancel(event_types[i].sub);
812                                 stasis_message_router_remove(stasis_router,
813                                                              event_types[i].message_type_fn());
814                         }
815                         event_types[i].publish = 0;
816                         event_types[i].subscribe = 0;
817                 }
818                 ast_rwlock_unlock(&event_types_lock);
819
820                 stasis_message_router_unsubscribe_and_join(stasis_router);
821                 stasis_router = NULL;
822         }
823
824         if (corosync_aggregate_topic) {
825                 ao2_t_ref(corosync_aggregate_topic, -1, "Dispose of topic on cleanup");
826                 corosync_aggregate_topic = NULL;
827         }
828
829         STASIS_MESSAGE_TYPE_CLEANUP(corosync_ping_message_type);
830
831         if (dispatch_thread.id != AST_PTHREADT_NULL) {
832                 char meepmeep = 'x';
833                 dispatch_thread.stop = 1;
834                 if (ast_carefulwrite(dispatch_thread.alert_pipe[1], &meepmeep, 1,
835                                         5000) == -1) {
836                         ast_log(LOG_ERROR, "Failed to write to pipe: %s (%d)\n",
837                                         strerror(errno), errno);
838                 }
839                 pthread_join(dispatch_thread.id, NULL);
840         }
841
842         if (dispatch_thread.alert_pipe[0] != -1) {
843                 close(dispatch_thread.alert_pipe[0]);
844                 dispatch_thread.alert_pipe[0] = -1;
845         }
846
847         if (dispatch_thread.alert_pipe[1] != -1) {
848                 close(dispatch_thread.alert_pipe[1]);
849                 dispatch_thread.alert_pipe[1] = -1;
850         }
851
852         if (cpg_handle && (cs_err = cpg_finalize(cpg_handle)) != CS_OK) {
853                 ast_log(LOG_ERROR, "Failed to finalize cpg (%d)\n", (int) cs_err);
854         }
855         cpg_handle = 0;
856
857         if (cfg_handle && (cs_err = corosync_cfg_finalize(cfg_handle)) != CS_OK) {
858                 ast_log(LOG_ERROR, "Failed to finalize cfg (%d)\n", (int) cs_err);
859         }
860         cfg_handle = 0;
861 }
862
863 static int load_module(void)
864 {
865         cs_error_t cs_err;
866         enum ast_module_load_result res = AST_MODULE_LOAD_FAILURE;
867         struct cpg_name name;
868
869         corosync_aggregate_topic = stasis_topic_create("corosync_aggregate_topic");
870         if (!corosync_aggregate_topic) {
871                 ast_log(AST_LOG_ERROR, "Failed to create stasis topic for corosync\n");
872                 goto failed;
873         }
874
875         stasis_router = stasis_message_router_create(corosync_aggregate_topic);
876         if (!stasis_router) {
877                 ast_log(AST_LOG_ERROR, "Failed to create message router for corosync topic\n");
878                 goto failed;
879         }
880
881         if (STASIS_MESSAGE_TYPE_INIT(corosync_ping_message_type) != 0) {
882                 ast_log(AST_LOG_ERROR, "Failed to initialize corosync ping message type\n");
883                 goto failed;
884         }
885
886         if ((cs_err = corosync_cfg_initialize(&cfg_handle, &cfg_callbacks)) != CS_OK) {
887                 ast_log(LOG_ERROR, "Failed to initialize cfg: (%d)\n", (int) cs_err);
888                 goto failed;
889         }
890
891         if ((cs_err = cpg_initialize(&cpg_handle, &cpg_callbacks)) != CS_OK) {
892                 ast_log(LOG_ERROR, "Failed to initialize cpg: (%d)\n", (int) cs_err);
893                 goto failed;
894         }
895
896         ast_copy_string(name.value, "asterisk", sizeof(name.value));
897         name.length = strlen(name.value);
898
899         if ((cs_err = cpg_join(cpg_handle, &name)) != CS_OK) {
900                 ast_log(LOG_ERROR, "Failed to join: (%d)\n", (int) cs_err);
901                 goto failed;
902         }
903
904         if (pipe(dispatch_thread.alert_pipe) == -1) {
905                 ast_log(LOG_ERROR, "Failed to create alert pipe: %s (%d)\n",
906                                 strerror(errno), errno);
907                 goto failed;
908         }
909
910         if (ast_pthread_create_background(&dispatch_thread.id, NULL,
911                         dispatch_thread_handler, NULL)) {
912                 ast_log(LOG_ERROR, "Error starting CPG dispatch thread.\n");
913                 goto failed;
914         }
915
916         if (load_config(0)) {
917                 /* simply not configured is not a fatal error */
918                 res = AST_MODULE_LOAD_DECLINE;
919                 goto failed;
920         }
921
922         ast_cli_register_multiple(corosync_cli, ARRAY_LEN(corosync_cli));
923
924         return AST_MODULE_LOAD_SUCCESS;
925
926 failed:
927         cleanup_module();
928
929         return res;
930 }
931
932 static int unload_module(void)
933 {
934         ast_cli_unregister_multiple(corosync_cli, ARRAY_LEN(corosync_cli));
935
936         cleanup_module();
937
938         return 0;
939 }
940
941 AST_MODULE_INFO_STANDARD_EXTENDED(ASTERISK_GPL_KEY, "Corosync");
942