4 Copyright (C) Andrew Tridgell 2006
6 This program is free software; you can redistribute it and/or modify
7 it under the terms of the GNU General Public License as published by
8 the Free Software Foundation; either version 3 of the License, or
9 (at your option) any later version.
11 This program is distributed in the hope that it will be useful,
12 but WITHOUT ANY WARRANTY; without even the implied warranty of
13 MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
14 GNU General Public License for more details.
16 You should have received a copy of the GNU General Public License
17 along with this program; if not, see <http://www.gnu.org/licenses/>.
22 #include "lib/tdb/include/tdb.h"
23 #include "lib/events/events.h"
24 #include "lib/util/dlinklist.h"
25 #include "system/network.h"
26 #include "system/filesys.h"
27 #include "system/wait.h"
28 #include "../include/ctdb.h"
29 #include "../include/ctdb_private.h"
30 #include <sys/socket.h>
32 static void daemon_incoming_packet(void *, struct ctdb_req_header *);
35 static void print_exit_message(void)
37 DEBUG(DEBUG_NOTICE,("CTDB daemon shutting down\n"));
41 /* called when the "startup" event script has finished */
42 static void ctdb_start_transport(struct ctdb_context *ctdb)
44 if (ctdb->methods == NULL) {
45 DEBUG(DEBUG_ALERT,(__location__ " startup event finished but transport is DOWN.\n"));
46 ctdb_fatal(ctdb, "transport is not initialized but startup completed");
49 /* start the transport running */
50 if (ctdb->methods->start(ctdb) != 0) {
51 DEBUG(DEBUG_ALERT,("transport failed to start!\n"));
52 ctdb_fatal(ctdb, "transport failed to start");
55 /* start the recovery daemon process */
56 if (ctdb_start_recoverd(ctdb) != 0) {
57 DEBUG(DEBUG_ALERT,("Failed to start recovery daemon\n"));
61 /* Make sure we log something when the daemon terminates */
62 atexit(print_exit_message);
64 /* start monitoring for connected/disconnected nodes */
65 ctdb_start_keepalive(ctdb);
67 /* start monitoring for node health */
68 ctdb_start_monitoring(ctdb);
70 /* start periodic update of tcp tickle lists */
71 ctdb_start_tcp_tickle_update(ctdb);
73 /* start listening for recovery daemon pings */
74 ctdb_control_recd_ping(ctdb);
77 static void block_signal(int signum)
81 memset(&act, 0, sizeof(act));
83 act.sa_handler = SIG_IGN;
84 sigemptyset(&act.sa_mask);
85 sigaddset(&act.sa_mask, signum);
86 sigaction(signum, &act, NULL);
91 send a packet to a client
93 static int daemon_queue_send(struct ctdb_client *client, struct ctdb_req_header *hdr)
95 client->ctdb->statistics.client_packets_sent++;
96 return ctdb_queue_send(client->queue, (uint8_t *)hdr, hdr->length);
100 message handler for when we are in daemon mode. This redirects the message
103 static void daemon_message_handler(struct ctdb_context *ctdb, uint64_t srvid,
104 TDB_DATA data, void *private_data)
106 struct ctdb_client *client = talloc_get_type(private_data, struct ctdb_client);
107 struct ctdb_req_message *r;
110 /* construct a message to send to the client containing the data */
111 len = offsetof(struct ctdb_req_message, data) + data.dsize;
112 r = ctdbd_allocate_pkt(ctdb, ctdb, CTDB_REQ_MESSAGE,
113 len, struct ctdb_req_message);
114 CTDB_NO_MEMORY_VOID(ctdb, r);
116 talloc_set_name_const(r, "req_message packet");
119 r->datalen = data.dsize;
120 memcpy(&r->data[0], data.dptr, data.dsize);
122 daemon_queue_send(client, &r->hdr);
129 this is called when the ctdb daemon received a ctdb request to
130 set the srvid from the client
132 int daemon_register_message_handler(struct ctdb_context *ctdb, uint32_t client_id, uint64_t srvid)
134 struct ctdb_client *client = ctdb_reqid_find(ctdb, client_id, struct ctdb_client);
136 if (client == NULL) {
137 DEBUG(DEBUG_ERR,("Bad client_id in daemon_request_register_message_handler\n"));
140 res = ctdb_register_message_handler(ctdb, client, srvid, daemon_message_handler, client);
142 DEBUG(DEBUG_ERR,(__location__ " Failed to register handler %llu in daemon\n",
143 (unsigned long long)srvid));
145 DEBUG(DEBUG_INFO,(__location__ " Registered message handler for srvid=%llu\n",
146 (unsigned long long)srvid));
149 /* this is a hack for Samba - we now know the pid of the Samba client */
150 if ((srvid & 0xFFFFFFFF) == srvid &&
151 kill(srvid, 0) == 0) {
153 DEBUG(DEBUG_INFO,(__location__ " Registered PID %u for client %u\n",
154 (unsigned)client->pid, client_id));
160 this is called when the ctdb daemon received a ctdb request to
161 remove a srvid from the client
163 int daemon_deregister_message_handler(struct ctdb_context *ctdb, uint32_t client_id, uint64_t srvid)
165 struct ctdb_client *client = ctdb_reqid_find(ctdb, client_id, struct ctdb_client);
166 if (client == NULL) {
167 DEBUG(DEBUG_ERR,("Bad client_id in daemon_request_deregister_message_handler\n"));
170 return ctdb_deregister_message_handler(ctdb, srvid, client);
175 destroy a ctdb_client
177 static int ctdb_client_destructor(struct ctdb_client *client)
179 struct ctdb_db_context *ctdb_db;
181 ctdb_takeover_client_destructor_hook(client);
182 ctdb_reqid_remove(client->ctdb, client->client_id);
183 if (client->ctdb->statistics.num_clients) {
184 client->ctdb->statistics.num_clients--;
187 if (client->num_persistent_updates != 0) {
188 DEBUG(DEBUG_ERR,(__location__ " Client disconnecting with %u persistent updates in flight. Starting recovery\n", client->num_persistent_updates));
189 client->ctdb->recovery_mode = CTDB_RECOVERY_ACTIVE;
191 ctdb_db = find_ctdb_db(client->ctdb, client->db_id);
193 DEBUG(DEBUG_ERR, (__location__ " client exit while transaction "
194 "commit active. Forcing recovery.\n"));
195 client->ctdb->recovery_mode = CTDB_RECOVERY_ACTIVE;
196 ctdb_db->transaction_active = false;
204 this is called when the ctdb daemon received a ctdb request message
205 from a local client over the unix domain socket
207 static void daemon_request_message_from_client(struct ctdb_client *client,
208 struct ctdb_req_message *c)
213 /* maybe the message is for another client on this node */
214 if (ctdb_get_pnn(client->ctdb)==c->hdr.destnode) {
215 ctdb_request_message(client->ctdb, (struct ctdb_req_header *)c);
219 /* its for a remote node */
220 data.dptr = &c->data[0];
221 data.dsize = c->datalen;
222 res = ctdb_daemon_send_message(client->ctdb, c->hdr.destnode,
225 DEBUG(DEBUG_ERR,(__location__ " Failed to send message to remote node %u\n",
231 struct daemon_call_state {
232 struct ctdb_client *client;
234 struct ctdb_call *call;
235 struct timeval start_time;
239 complete a call from a client
241 static void daemon_call_from_client_callback(struct ctdb_call_state *state)
243 struct daemon_call_state *dstate = talloc_get_type(state->async.private_data,
244 struct daemon_call_state);
245 struct ctdb_reply_call *r;
248 struct ctdb_client *client = dstate->client;
249 struct ctdb_db_context *ctdb_db = state->ctdb_db;
251 talloc_steal(client, dstate);
252 talloc_steal(dstate, dstate->call);
254 res = ctdb_daemon_call_recv(state, dstate->call);
256 DEBUG(DEBUG_ERR, (__location__ " ctdbd_call_recv() returned error\n"));
257 if (client->ctdb->statistics.pending_calls > 0) {
258 client->ctdb->statistics.pending_calls--;
260 ctdb_latency(ctdb_db, "call_from_client_cb 1", &client->ctdb->statistics.max_call_latency, dstate->start_time);
264 length = offsetof(struct ctdb_reply_call, data) + dstate->call->reply_data.dsize;
265 r = ctdbd_allocate_pkt(client->ctdb, dstate, CTDB_REPLY_CALL,
266 length, struct ctdb_reply_call);
268 DEBUG(DEBUG_ERR, (__location__ " Failed to allocate reply_call in ctdb daemon\n"));
269 if (client->ctdb->statistics.pending_calls > 0) {
270 client->ctdb->statistics.pending_calls--;
272 ctdb_latency(ctdb_db, "call_from_client_cb 2", &client->ctdb->statistics.max_call_latency, dstate->start_time);
275 r->hdr.reqid = dstate->reqid;
276 r->datalen = dstate->call->reply_data.dsize;
277 memcpy(&r->data[0], dstate->call->reply_data.dptr, r->datalen);
279 res = daemon_queue_send(client, &r->hdr);
281 DEBUG(DEBUG_ERR, (__location__ " Failed to queue packet from daemon to client\n"));
283 ctdb_latency(ctdb_db, "call_from_client_cb 3", &client->ctdb->statistics.max_call_latency, dstate->start_time);
285 if (client->ctdb->statistics.pending_calls > 0) {
286 client->ctdb->statistics.pending_calls--;
290 struct ctdb_daemon_packet_wrap {
291 struct ctdb_context *ctdb;
296 a wrapper to catch disconnected clients
298 static void daemon_incoming_packet_wrap(void *p, struct ctdb_req_header *hdr)
300 struct ctdb_client *client;
301 struct ctdb_daemon_packet_wrap *w = talloc_get_type(p,
302 struct ctdb_daemon_packet_wrap);
304 DEBUG(DEBUG_CRIT,(__location__ " Bad packet type '%s'\n", talloc_get_name(p)));
308 client = ctdb_reqid_find(w->ctdb, w->client_id, struct ctdb_client);
309 if (client == NULL) {
310 DEBUG(DEBUG_ERR,(__location__ " Packet for disconnected client %u\n",
318 daemon_incoming_packet(client, hdr);
323 this is called when the ctdb daemon received a ctdb request call
324 from a local client over the unix domain socket
326 static void daemon_request_call_from_client(struct ctdb_client *client,
327 struct ctdb_req_call *c)
329 struct ctdb_call_state *state;
330 struct ctdb_db_context *ctdb_db;
331 struct daemon_call_state *dstate;
332 struct ctdb_call *call;
333 struct ctdb_ltdb_header header;
336 struct ctdb_context *ctdb = client->ctdb;
337 struct ctdb_daemon_packet_wrap *w;
339 ctdb->statistics.total_calls++;
340 if (client->ctdb->statistics.pending_calls > 0) {
341 ctdb->statistics.pending_calls++;
344 ctdb_db = find_ctdb_db(client->ctdb, c->db_id);
346 DEBUG(DEBUG_ERR, (__location__ " Unknown database in request. db_id==0x%08x",
348 if (client->ctdb->statistics.pending_calls > 0) {
349 ctdb->statistics.pending_calls--;
355 key.dsize = c->keylen;
357 w = talloc(ctdb, struct ctdb_daemon_packet_wrap);
358 CTDB_NO_MEMORY_VOID(ctdb, w);
361 w->client_id = client->client_id;
363 ret = ctdb_ltdb_lock_fetch_requeue(ctdb_db, key, &header,
364 (struct ctdb_req_header *)c, &data,
365 daemon_incoming_packet_wrap, w, True);
367 /* will retry later */
368 if (client->ctdb->statistics.pending_calls > 0) {
369 ctdb->statistics.pending_calls--;
377 DEBUG(DEBUG_ERR,(__location__ " Unable to fetch record\n"));
378 if (client->ctdb->statistics.pending_calls > 0) {
379 ctdb->statistics.pending_calls--;
384 dstate = talloc(client, struct daemon_call_state);
385 if (dstate == NULL) {
386 ctdb_ltdb_unlock(ctdb_db, key);
387 DEBUG(DEBUG_ERR,(__location__ " Unable to allocate dstate\n"));
388 if (client->ctdb->statistics.pending_calls > 0) {
389 ctdb->statistics.pending_calls--;
393 dstate->start_time = timeval_current();
394 dstate->client = client;
395 dstate->reqid = c->hdr.reqid;
396 talloc_steal(dstate, data.dptr);
398 call = dstate->call = talloc_zero(dstate, struct ctdb_call);
400 ctdb_ltdb_unlock(ctdb_db, key);
401 DEBUG(DEBUG_ERR,(__location__ " Unable to allocate call\n"));
402 if (client->ctdb->statistics.pending_calls > 0) {
403 ctdb->statistics.pending_calls--;
405 ctdb_latency(ctdb_db, "call_from_client 1", &ctdb->statistics.max_call_latency, dstate->start_time);
409 call->call_id = c->callid;
411 call->call_data.dptr = c->data + c->keylen;
412 call->call_data.dsize = c->calldatalen;
413 call->flags = c->flags;
415 if (header.dmaster == ctdb->pnn) {
416 state = ctdb_call_local_send(ctdb_db, call, &header, &data);
418 state = ctdb_daemon_call_send_remote(ctdb_db, call, &header);
421 ctdb_ltdb_unlock(ctdb_db, key);
424 DEBUG(DEBUG_ERR,(__location__ " Unable to setup call send\n"));
425 if (client->ctdb->statistics.pending_calls > 0) {
426 ctdb->statistics.pending_calls--;
428 ctdb_latency(ctdb_db, "call_from_client 2", &ctdb->statistics.max_call_latency, dstate->start_time);
431 talloc_steal(state, dstate);
432 talloc_steal(client, state);
434 state->async.fn = daemon_call_from_client_callback;
435 state->async.private_data = dstate;
439 static void daemon_request_control_from_client(struct ctdb_client *client,
440 struct ctdb_req_control *c);
442 /* data contains a packet from the client */
443 static void daemon_incoming_packet(void *p, struct ctdb_req_header *hdr)
445 struct ctdb_client *client = talloc_get_type(p, struct ctdb_client);
447 struct ctdb_context *ctdb = client->ctdb;
449 /* place the packet as a child of a tmp_ctx. We then use
450 talloc_free() below to free it. If any of the calls want
451 to keep it, then they will steal it somewhere else, and the
452 talloc_free() will be a no-op */
453 tmp_ctx = talloc_new(client);
454 talloc_steal(tmp_ctx, hdr);
456 if (hdr->ctdb_magic != CTDB_MAGIC) {
457 ctdb_set_error(client->ctdb, "Non CTDB packet rejected in daemon\n");
461 if (hdr->ctdb_version != CTDB_VERSION) {
462 ctdb_set_error(client->ctdb, "Bad CTDB version 0x%x rejected in daemon\n", hdr->ctdb_version);
466 switch (hdr->operation) {
468 ctdb->statistics.client.req_call++;
469 daemon_request_call_from_client(client, (struct ctdb_req_call *)hdr);
472 case CTDB_REQ_MESSAGE:
473 ctdb->statistics.client.req_message++;
474 daemon_request_message_from_client(client, (struct ctdb_req_message *)hdr);
477 case CTDB_REQ_CONTROL:
478 ctdb->statistics.client.req_control++;
479 daemon_request_control_from_client(client, (struct ctdb_req_control *)hdr);
483 DEBUG(DEBUG_CRIT,(__location__ " daemon: unrecognized operation %u\n",
488 talloc_free(tmp_ctx);
492 called when the daemon gets a incoming packet
494 static void ctdb_daemon_read_cb(uint8_t *data, size_t cnt, void *args)
496 struct ctdb_client *client = talloc_get_type(args, struct ctdb_client);
497 struct ctdb_req_header *hdr;
504 client->ctdb->statistics.client_packets_recv++;
506 if (cnt < sizeof(*hdr)) {
507 ctdb_set_error(client->ctdb, "Bad packet length %u in daemon\n",
511 hdr = (struct ctdb_req_header *)data;
512 if (cnt != hdr->length) {
513 ctdb_set_error(client->ctdb, "Bad header length %u expected %u\n in daemon",
514 (unsigned)hdr->length, (unsigned)cnt);
518 if (hdr->ctdb_magic != CTDB_MAGIC) {
519 ctdb_set_error(client->ctdb, "Non CTDB packet rejected\n");
523 if (hdr->ctdb_version != CTDB_VERSION) {
524 ctdb_set_error(client->ctdb, "Bad CTDB version 0x%x rejected in daemon\n", hdr->ctdb_version);
528 DEBUG(DEBUG_DEBUG,(__location__ " client request %u of type %u length %u from "
529 "node %u to %u\n", hdr->reqid, hdr->operation, hdr->length,
530 hdr->srcnode, hdr->destnode));
532 /* it is the responsibility of the incoming packet function to free 'data' */
533 daemon_incoming_packet(client, hdr);
536 static void ctdb_accept_client(struct event_context *ev, struct fd_event *fde,
537 uint16_t flags, void *private_data)
539 struct sockaddr_un addr;
542 struct ctdb_context *ctdb = talloc_get_type(private_data, struct ctdb_context);
543 struct ctdb_client *client;
545 struct peercred_struct cr;
546 socklen_t crl = sizeof(struct peercred_struct);
549 socklen_t crl = sizeof(struct ucred);
552 memset(&addr, 0, sizeof(addr));
554 fd = accept(ctdb->daemon.sd, (struct sockaddr *)&addr, &len);
560 set_close_on_exec(fd);
562 client = talloc_zero(ctdb, struct ctdb_client);
564 if (getsockopt(fd, SOL_SOCKET, SO_PEERID, &cr, &crl) == 0) {
566 if (getsockopt(fd, SOL_SOCKET, SO_PEERCRED, &cr, &crl) == 0) {
568 talloc_asprintf(client, "struct ctdb_client: pid:%u", (unsigned)cr.pid);
573 client->client_id = ctdb_reqid_new(ctdb, client);
574 ctdb->statistics.num_clients++;
576 client->queue = ctdb_queue_setup(ctdb, client, fd, CTDB_DS_ALIGNMENT,
577 ctdb_daemon_read_cb, client);
579 talloc_set_destructor(client, ctdb_client_destructor);
585 create a unix domain socket and bind it
586 return a file descriptor open on the socket
588 static int ux_socket_bind(struct ctdb_context *ctdb)
590 struct sockaddr_un addr;
592 ctdb->daemon.sd = socket(AF_UNIX, SOCK_STREAM, 0);
593 if (ctdb->daemon.sd == -1) {
597 set_close_on_exec(ctdb->daemon.sd);
598 set_nonblocking(ctdb->daemon.sd);
600 memset(&addr, 0, sizeof(addr));
601 addr.sun_family = AF_UNIX;
602 strncpy(addr.sun_path, ctdb->daemon.name, sizeof(addr.sun_path));
604 if (bind(ctdb->daemon.sd, (struct sockaddr *)&addr, sizeof(addr)) == -1) {
605 DEBUG(DEBUG_CRIT,("Unable to bind on ctdb socket '%s'\n", ctdb->daemon.name));
609 if (chown(ctdb->daemon.name, geteuid(), getegid()) != 0 ||
610 chmod(ctdb->daemon.name, 0700) != 0) {
611 DEBUG(DEBUG_CRIT,("Unable to secure ctdb socket '%s', ctdb->daemon.name\n", ctdb->daemon.name));
616 if (listen(ctdb->daemon.sd, 100) != 0) {
617 DEBUG(DEBUG_CRIT,("Unable to listen on ctdb socket '%s'\n", ctdb->daemon.name));
624 close(ctdb->daemon.sd);
625 ctdb->daemon.sd = -1;
629 static void sig_child_handler(struct event_context *ev,
630 struct signal_event *se, int signum, int count,
634 // struct ctdb_context *ctdb = talloc_get_type(private_data, struct ctdb_context);
639 pid = waitpid(-1, &status, WNOHANG);
641 DEBUG(DEBUG_ERR, (__location__ " waitpid() returned error. errno:%d\n", errno));
645 DEBUG(DEBUG_DEBUG, ("SIGCHLD from %d\n", (int)pid));
651 start the protocol going as a daemon
653 int ctdb_start_daemon(struct ctdb_context *ctdb, bool do_fork)
656 struct fd_event *fde;
657 const char *domain_socket_name;
658 struct signal_event *se;
660 /* get rid of any old sockets */
661 unlink(ctdb->daemon.name);
663 /* create a unix domain stream socket to listen to */
664 res = ux_socket_bind(ctdb);
666 DEBUG(DEBUG_ALERT,(__location__ " Failed to open CTDB unix domain socket\n"));
670 if (do_fork && fork()) {
674 tdb_reopen_all(False);
679 if (open("/dev/null", O_RDONLY) != 0) {
680 DEBUG(DEBUG_ALERT,(__location__ " Failed to setup stdin on /dev/null\n"));
684 block_signal(SIGPIPE);
686 if (ctdb->do_setsched) {
687 /* try to set us up as realtime */
688 ctdb_set_scheduler(ctdb);
691 /* ensure the socket is deleted on exit of the daemon */
692 domain_socket_name = talloc_strdup(talloc_autofree_context(), ctdb->daemon.name);
693 if (domain_socket_name == NULL) {
694 DEBUG(DEBUG_ALERT,(__location__ " talloc_strdup failed.\n"));
698 ctdb->ev = event_context_init(NULL);
700 ctdb_set_child_logging(ctdb);
702 /* force initial recovery for election */
703 ctdb->recovery_mode = CTDB_RECOVERY_ACTIVE;
705 if (strcmp(ctdb->transport, "tcp") == 0) {
706 int ctdb_tcp_init(struct ctdb_context *);
707 ret = ctdb_tcp_init(ctdb);
709 #ifdef USE_INFINIBAND
710 if (strcmp(ctdb->transport, "ib") == 0) {
711 int ctdb_ibw_init(struct ctdb_context *);
712 ret = ctdb_ibw_init(ctdb);
716 DEBUG(DEBUG_ERR,("Failed to initialise transport '%s'\n", ctdb->transport));
720 if (ctdb->methods == NULL) {
721 DEBUG(DEBUG_ALERT,(__location__ " Can not initialize transport. ctdb->methods is NULL\n"));
722 ctdb_fatal(ctdb, "transport is unavailable. can not initialize.");
725 /* initialise the transport */
726 if (ctdb->methods->initialise(ctdb) != 0) {
727 ctdb_fatal(ctdb, "transport failed to initialise");
730 /* attach to any existing persistent databases */
731 if (ctdb_attach_persistent(ctdb) != 0) {
732 ctdb_fatal(ctdb, "Failed to attach to persistent databases\n");
735 /* start frozen, then let the first election sort things out */
736 if (ctdb_blocking_freeze(ctdb)) {
737 ctdb_fatal(ctdb, "Failed to get initial freeze\n");
740 /* now start accepting clients, only can do this once frozen */
741 fde = event_add_fd(ctdb->ev, ctdb, ctdb->daemon.sd,
742 EVENT_FD_READ|EVENT_FD_AUTOCLOSE,
743 ctdb_accept_client, ctdb);
745 /* tell all other nodes we've just started up */
746 ctdb_daemon_send_control(ctdb, CTDB_BROADCAST_ALL,
747 0, CTDB_CONTROL_STARTUP, 0,
748 CTDB_CTRL_FLAG_NOREPLY,
749 tdb_null, NULL, NULL);
751 /* release any IPs we hold from previous runs of the daemon */
752 ctdb_release_all_ips(ctdb);
754 /* start the transport going */
755 ctdb_start_transport(ctdb);
757 /* set up a handler to pick up sigchld */
758 se = event_add_signal(ctdb->ev, ctdb,
763 DEBUG(DEBUG_CRIT,("Failed to set up signal handler for SIGCHLD\n"));
767 /* go into a wait loop to allow other nodes to complete */
768 event_loop_wait(ctdb->ev);
770 DEBUG(DEBUG_CRIT,("event_loop_wait() returned. this should not happen\n"));
775 allocate a packet for use in daemon<->daemon communication
777 struct ctdb_req_header *_ctdb_transport_allocate(struct ctdb_context *ctdb,
779 enum ctdb_operation operation,
780 size_t length, size_t slength,
784 struct ctdb_req_header *hdr;
786 length = MAX(length, slength);
787 size = (length+(CTDB_DS_ALIGNMENT-1)) & ~(CTDB_DS_ALIGNMENT-1);
789 if (ctdb->methods == NULL) {
790 DEBUG(DEBUG_ERR,(__location__ " Unable to allocate transport packet for operation %u of length %u. Transport is DOWN.\n",
791 operation, (unsigned)length));
795 hdr = (struct ctdb_req_header *)ctdb->methods->allocate_pkt(mem_ctx, size);
797 DEBUG(DEBUG_ERR,("Unable to allocate transport packet for operation %u of length %u\n",
798 operation, (unsigned)length));
801 talloc_set_name_const(hdr, type);
802 memset(hdr, 0, slength);
803 hdr->length = length;
804 hdr->operation = operation;
805 hdr->ctdb_magic = CTDB_MAGIC;
806 hdr->ctdb_version = CTDB_VERSION;
807 hdr->generation = ctdb->vnn_map->generation;
808 hdr->srcnode = ctdb->pnn;
813 struct daemon_control_state {
814 struct daemon_control_state *next, *prev;
815 struct ctdb_client *client;
816 struct ctdb_req_control *c;
818 struct ctdb_node *node;
822 callback when a control reply comes in
824 static void daemon_control_callback(struct ctdb_context *ctdb,
825 int32_t status, TDB_DATA data,
826 const char *errormsg,
829 struct daemon_control_state *state = talloc_get_type(private_data,
830 struct daemon_control_state);
831 struct ctdb_client *client = state->client;
832 struct ctdb_reply_control *r;
835 /* construct a message to send to the client containing the data */
836 len = offsetof(struct ctdb_reply_control, data) + data.dsize;
838 len += strlen(errormsg);
840 r = ctdbd_allocate_pkt(ctdb, state, CTDB_REPLY_CONTROL, len,
841 struct ctdb_reply_control);
842 CTDB_NO_MEMORY_VOID(ctdb, r);
844 r->hdr.reqid = state->reqid;
846 r->datalen = data.dsize;
848 memcpy(&r->data[0], data.dptr, data.dsize);
850 r->errorlen = strlen(errormsg);
851 memcpy(&r->data[r->datalen], errormsg, r->errorlen);
854 daemon_queue_send(client, &r->hdr);
860 fail all pending controls to a disconnected node
862 void ctdb_daemon_cancel_controls(struct ctdb_context *ctdb, struct ctdb_node *node)
864 struct daemon_control_state *state;
865 while ((state = node->pending_controls)) {
866 DLIST_REMOVE(node->pending_controls, state);
867 daemon_control_callback(ctdb, (uint32_t)-1, tdb_null,
868 "node is disconnected", state);
873 destroy a daemon_control_state
875 static int daemon_control_destructor(struct daemon_control_state *state)
878 DLIST_REMOVE(state->node->pending_controls, state);
884 this is called when the ctdb daemon received a ctdb request control
885 from a local client over the unix domain socket
887 static void daemon_request_control_from_client(struct ctdb_client *client,
888 struct ctdb_req_control *c)
892 struct daemon_control_state *state;
893 TALLOC_CTX *tmp_ctx = talloc_new(client);
895 if (c->hdr.destnode == CTDB_CURRENT_NODE) {
896 c->hdr.destnode = client->ctdb->pnn;
899 state = talloc(client, struct daemon_control_state);
900 CTDB_NO_MEMORY_VOID(client->ctdb, state);
902 state->client = client;
903 state->c = talloc_steal(state, c);
904 state->reqid = c->hdr.reqid;
905 if (ctdb_validate_pnn(client->ctdb, c->hdr.destnode)) {
906 state->node = client->ctdb->nodes[c->hdr.destnode];
907 DLIST_ADD(state->node->pending_controls, state);
912 talloc_set_destructor(state, daemon_control_destructor);
914 if (c->flags & CTDB_CTRL_FLAG_NOREPLY) {
915 talloc_steal(tmp_ctx, state);
918 data.dptr = &c->data[0];
919 data.dsize = c->datalen;
920 res = ctdb_daemon_send_control(client->ctdb, c->hdr.destnode,
921 c->srvid, c->opcode, client->client_id,
923 data, daemon_control_callback,
926 DEBUG(DEBUG_ERR,(__location__ " Failed to send control to remote node %u\n",
930 talloc_free(tmp_ctx);
934 register a call function
936 int ctdb_daemon_set_call(struct ctdb_context *ctdb, uint32_t db_id,
937 ctdb_fn_t fn, int id)
939 struct ctdb_registered_call *call;
940 struct ctdb_db_context *ctdb_db;
942 ctdb_db = find_ctdb_db(ctdb, db_id);
943 if (ctdb_db == NULL) {
947 call = talloc(ctdb_db, struct ctdb_registered_call);
951 DLIST_ADD(ctdb_db->calls, call);
958 this local messaging handler is ugly, but is needed to prevent
959 recursion in ctdb_send_message() when the destination node is the
960 same as the source node
962 struct ctdb_local_message {
963 struct ctdb_context *ctdb;
968 static void ctdb_local_message_trigger(struct event_context *ev, struct timed_event *te,
969 struct timeval t, void *private_data)
971 struct ctdb_local_message *m = talloc_get_type(private_data,
972 struct ctdb_local_message);
975 res = ctdb_dispatch_message(m->ctdb, m->srvid, m->data);
977 DEBUG(DEBUG_ERR, (__location__ " Failed to dispatch message for srvid=%llu\n",
978 (unsigned long long)m->srvid));
983 static int ctdb_local_message(struct ctdb_context *ctdb, uint64_t srvid, TDB_DATA data)
985 struct ctdb_local_message *m;
986 m = talloc(ctdb, struct ctdb_local_message);
987 CTDB_NO_MEMORY(ctdb, m);
992 m->data.dptr = talloc_memdup(m, m->data.dptr, m->data.dsize);
993 if (m->data.dptr == NULL) {
998 /* this needs to be done as an event to prevent recursion */
999 event_add_timed(ctdb->ev, m, timeval_zero(), ctdb_local_message_trigger, m);
1006 int ctdb_daemon_send_message(struct ctdb_context *ctdb, uint32_t pnn,
1007 uint64_t srvid, TDB_DATA data)
1009 struct ctdb_req_message *r;
1012 if (ctdb->methods == NULL) {
1013 DEBUG(DEBUG_ERR,(__location__ " Failed to send message. Transport is DOWN\n"));
1017 /* see if this is a message to ourselves */
1018 if (pnn == ctdb->pnn) {
1019 return ctdb_local_message(ctdb, srvid, data);
1022 len = offsetof(struct ctdb_req_message, data) + data.dsize;
1023 r = ctdb_transport_allocate(ctdb, ctdb, CTDB_REQ_MESSAGE, len,
1024 struct ctdb_req_message);
1025 CTDB_NO_MEMORY(ctdb, r);
1027 r->hdr.destnode = pnn;
1029 r->datalen = data.dsize;
1030 memcpy(&r->data[0], data.dptr, data.dsize);
1032 ctdb_queue_packet(ctdb, &r->hdr);