root/src/stratcon_datastore.c

Revision 88a71780101cbf23034aa0cb840f9f0368fda2dd, 21.8 kB (checked in by Theo Schlossnagle <jesus@omniti.com>, 5 years ago)

fixes #126

  • Property mode set to 100644
Line 
1 /*
2  * Copyright (c) 2007, OmniTI Computer Consulting, Inc.
3  * All rights reserved.
4  *
5  * Redistribution and use in source and binary forms, with or without
6  * modification, are permitted provided that the following conditions are
7  * met:
8  *
9  *     * Redistributions of source code must retain the above copyright
10  *       notice, this list of conditions and the following disclaimer.
11  *     * Redistributions in binary form must reproduce the above
12  *       copyright notice, this list of conditions and the following
13  *       disclaimer in the documentation and/or other materials provided
14  *       with the distribution.
15  *     * Neither the name OmniTI Computer Consulting, Inc. nor the names
16  *       of its contributors may be used to endorse or promote products
17  *       derived from this software without specific prior written
18  *       permission.
19  *
20  * THIS SOFTWARE IS PROVIDED BY THE COPYRIGHT HOLDERS AND CONTRIBUTORS
21  * "AS IS" AND ANY EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT
22  * LIMITED TO, THE IMPLIED WARRANTIES OF MERCHANTABILITY AND FITNESS FOR
23  * A PARTICULAR PURPOSE ARE DISCLAIMED. IN NO EVENT SHALL THE COPYRIGHT
24  * OWNER OR CONTRIBUTORS BE LIABLE FOR ANY DIRECT, INDIRECT, INCIDENTAL,
25  * SPECIAL, EXEMPLARY, OR CONSEQUENTIAL DAMAGES (INCLUDING, BUT NOT
26  * LIMITED TO, PROCUREMENT OF SUBSTITUTE GOODS OR SERVICES; LOSS OF USE,
27  * DATA, OR PROFITS; OR BUSINESS INTERRUPTION) HOWEVER CAUSED AND ON ANY
28  * THEORY OF LIABILITY, WHETHER IN CONTRACT, STRICT LIABILITY, OR TORT
29  * (INCLUDING NEGLIGENCE OR OTHERWISE) ARISING IN ANY WAY OUT OF THE USE
30  * OF THIS SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF SUCH DAMAGE.
31  */
32
33 #include "noit_defines.h"
34 #include "eventer/eventer.h"
35 #include "utils/noit_log.h"
36 #include "utils/noit_b64.h"
37 #include "stratcon_datastore.h"
38 #include "stratcon_realtime_http.h"
39 #include "stratcon_iep.h"
40 #include "noit_conf.h"
41 #include "noit_check.h"
42 #include <unistd.h>
43 #include <netinet/in.h>
44 #include <sys/un.h>
45 #include <arpa/inet.h>
46 #include <libpq-fe.h>
47 #include <zlib.h>
48
49 static char *check_loadall = NULL;
50 static const char *check_loadall_conf = "/stratcon/database/statements/allchecks";
51 static char *check_find = NULL;
52 static const char *check_find_conf = "/stratcon/database/statements/findcheck";
53 static char *check_insert = NULL;
54 static const char *check_insert_conf = "/stratcon/database/statements/check";
55 static char *status_insert = NULL;
56 static const char *status_insert_conf = "/stratcon/database/statements/status";
57 static char *metric_insert_numeric = NULL;
58 static const char *metric_insert_numeric_conf = "/stratcon/database/statements/metric_numeric";
59 static char *metric_insert_text = NULL;
60 static const char *metric_insert_text_conf = "/stratcon/database/statements/metric_text";
61 static char *config_insert = NULL;
62 static const char *config_insert_conf = "/stratcon/database/statements/config";
63
64 static struct datastore_onlooker_list {
65   void (*dispatch)(stratcon_datastore_op_t, struct sockaddr *, void *);
66   struct datastore_onlooker_list *next;
67 } *onlookers = NULL;
68
69 #define GET_QUERY(a) do { \
70   if(a == NULL) \
71     if(!noit_conf_get_string(NULL, a ## _conf, &(a))) \
72       goto bad_row; \
73 } while(0)
74
75 #define MAX_PARAMS 8
76 #define POSTGRES_PARTS \
77   PGresult *res; \
78   int rv; \
79   int nparams; \
80   int metric_type; \
81   char *paramValues[MAX_PARAMS]; \
82   int paramLengths[MAX_PARAMS]; \
83   int paramFormats[MAX_PARAMS]; \
84   int paramAllocd[MAX_PARAMS];
85
86 typedef struct ds_single_detail {
87   POSTGRES_PARTS
88 } ds_single_detail;
89 typedef struct ds_job_detail {
90   /* Postgres specific stuff */
91   POSTGRES_PARTS
92
93   char *data;  /* The raw string, NULL means the stream is done -- commit. */
94   struct realtime_tracker *rt;
95
96   int problematic;
97   eventer_t completion_event; /* This event should be registered if non NULL */
98   struct ds_job_detail *next;
99 } ds_job_detail;
100
101 typedef struct {
102   struct sockaddr *remote;
103   eventer_jobq_t  *jobq;
104   /* Postgres specific stuff */
105   PGconn          *dbh;
106   ds_job_detail   *head;
107   ds_job_detail   *tail;
108 } conn_q;
109
110 static int stratcon_database_connect(conn_q *cq);
111
112 static void
113 free_params(ds_single_detail *d) {
114   int i;
115   for(i=0; i<d->nparams; i++)
116     if(d->paramAllocd[i] && d->paramValues[i])
117       free(d->paramValues[i]);
118 }
119
120 static void
121 __append(conn_q *q, ds_job_detail *d) {
122   d->next = NULL;
123   if(!q->head) q->head = q->tail = d;
124   else {
125     q->tail->next = d;
126     q->tail = d;
127   }
128 }
129 static void
130 __remove_until(conn_q *q, ds_job_detail *d) {
131   ds_job_detail *next;
132   while(q->head && q->head != d) {
133     next = q->head;
134     q->head = q->head->next;
135     free_params((ds_single_detail *)next);
136     if(next->data) free(next->data);
137     free(next);
138   }
139   if(!q->head) q->tail = NULL;
140 }
141
142 noit_hash_table ds_conns;
143
144 conn_q *
145 __get_conn_q_for_remote(struct sockaddr *remote) {
146   void *vcq;
147   conn_q *cq;
148   char queue_name[128] = "datastore_";
149   static const char __zeros[4] = { 0 };
150   int len = 0;
151   if(remote) {
152     switch(remote->sa_family) {
153       case AF_INET:
154         len = sizeof(struct sockaddr_in);
155         inet_ntop(remote->sa_family, &((struct sockaddr_in *)remote)->sin_addr,
156                   queue_name + strlen("datastore_"), len);
157         break;
158       case AF_INET6:
159        len = sizeof(struct sockaddr_in6);
160         inet_ntop(remote->sa_family, &((struct sockaddr_in6 *)remote)->sin6_addr,
161                   queue_name + strlen("datastore_"), len);
162        break;
163       case AF_UNIX:
164         len = SUN_LEN(((struct sockaddr_un *)remote));
165         snprintf(queue_name, sizeof(queue_name), "datastore_%s", ((struct sockaddr_un *)remote)->sun_path);
166         break;
167       default: return NULL;
168     }
169   }
170   else {
171     /* This is a dummy connection */
172     remote = (struct sockaddr *)__zeros;
173     snprintf(queue_name, sizeof(queue_name), "datastore_default");
174     len = 4;
175   }
176   if(noit_hash_retrieve(&ds_conns, (const char *)remote, len, &vcq))
177     return vcq;
178   cq = calloc(1, sizeof(*cq));
179   cq->remote = malloc(len);
180   memcpy(cq->remote, remote, len);
181   cq->jobq = calloc(1, sizeof(*cq->jobq));
182   eventer_jobq_init(cq->jobq, queue_name);
183   cq->jobq->backq = eventer_default_backq();
184   /* Add one thread */
185   eventer_jobq_increase_concurrency(cq->jobq);
186   noit_hash_store(&ds_conns, (const char *)cq->remote, len, cq);
187   return cq;
188 }
189
190 typedef enum {
191   DS_EXEC_SUCCESS = 0,
192   DS_EXEC_ROW_FAILED = 1,
193   DS_EXEC_TXN_FAILED = 2,
194 } execute_outcome_t;
195
196 static char *
197 __noit__strndup(const char *src, int len) {
198   int slen;
199   char *dst;
200   for(slen = 0; slen < len; slen++)
201     if(src[slen] == '\0') break;
202   dst = malloc(slen + 1);
203   memcpy(dst, src, slen);
204   dst[slen] = '\0';
205   return dst;
206 }
207 #define DECLARE_PARAM_STR(str, len) do { \
208   d->paramValues[d->nparams] = __noit__strndup(str, len); \
209   d->paramLengths[d->nparams] = len; \
210   d->paramFormats[d->nparams] = 0; \
211   d->paramAllocd[d->nparams] = 1; \
212   if(!strcmp(d->paramValues[d->nparams], "[[null]]")) { \
213     free(d->paramValues[d->nparams]); \
214     d->paramValues[d->nparams] = NULL; \
215     d->paramLengths[d->nparams] = 0; \
216     d->paramAllocd[d->nparams] = 0; \
217   } \
218   d->nparams++; \
219 } while(0)
220 #define DECLARE_PARAM_INT(i) do { \
221   int buffer__len; \
222   char buffer__[32]; \
223   snprintf(buffer__, sizeof(buffer__), "%d", (i)); \
224   buffer__len = strlen(buffer__); \
225   DECLARE_PARAM_STR(buffer__, buffer__len); \
226 } while(0)
227
228 #define PG_GET_STR_COL(dest, row, name) do { \
229   int colnum = PQfnumber(d->res, name); \
230   dest = NULL; \
231   if (colnum >= 0) \
232     dest = PQgetisnull(d->res, row, colnum) \
233          ? NULL : PQgetvalue(d->res, row, colnum); \
234 } while(0)
235
236 #define PG_EXEC(cmd) do { \
237   d->res = PQexecParams(cq->dbh, cmd, d->nparams, NULL, \
238                         (const char * const *)d->paramValues, \
239                         d->paramLengths, d->paramFormats, 0); \
240   d->rv = PQresultStatus(d->res); \
241   if(d->rv != PGRES_COMMAND_OK && \
242      d->rv != PGRES_TUPLES_OK) { \
243     noitL(noit_error, "stratcon datasource bad (%d): %s\n", \
244           d->rv, PQresultErrorMessage(d->res)); \
245     PQclear(d->res); \
246     goto bad_row; \
247   } \
248 } while(0)
249
250 static int
251 stratcon_datastore_asynch_drive_iep(eventer_t e, int mask, void *closure,
252                                     struct timeval *now) {
253   conn_q *cq = closure;
254   ds_job_detail *d;
255   int i, row_count = 0, good = 0;
256   char buff[1024];
257
258   if(!(mask & EVENTER_ASYNCH_WORK)) return 0;
259   if(mask & EVENTER_ASYNCH_CLEANUP) return 0;
260
261   stratcon_database_connect(cq);
262   d = calloc(1, sizeof(*d));
263   GET_QUERY(check_loadall);
264   PG_EXEC(check_loadall);
265   row_count = PQntuples(d->res);
266  
267   for(i=0; i<row_count; i++) {
268     char *remote, *id, *target, *module, *name;
269     PG_GET_STR_COL(remote, i, "remote_address");
270     PG_GET_STR_COL(id, i, "id");
271     PG_GET_STR_COL(target, i, "target");
272     PG_GET_STR_COL(module, i, "module");
273     PG_GET_STR_COL(name, i, "name");
274     snprintf(buff, sizeof(buff), "C\t0.000\t%s\t%s\t%s\t%s\n", id, target, module, name);
275     stratcon_iep_line_processor(DS_OP_INSERT, NULL, buff);
276     good++;
277   }
278   noitL(noit_error, "Staged %d/%d remembered checks into IEP\n", good, row_count);
279  bad_row:
280   PQclear(d->res);
281   return 0;
282 }
283 void
284 stratcon_datastore_iep_check_preload() {
285   eventer_t e;
286   conn_q *cq;
287   cq = __get_conn_q_for_remote(NULL);
288
289   e = eventer_alloc();
290   e->mask = EVENTER_ASYNCH;
291   e->callback = stratcon_datastore_asynch_drive_iep;
292   e->closure = cq;
293   eventer_add_asynch(cq->jobq, e);
294 }
295 execute_outcome_t
296 stratcon_datastore_find(conn_q *cq, ds_job_detail *d) {
297   char *val;
298   int row_count;
299
300   if(!d->nparams) DECLARE_PARAM_INT(d->rt->sid);
301   GET_QUERY(check_find);
302   PG_EXEC(check_find);
303   row_count = PQntuples(d->res);
304   if(row_count != 1) goto bad_row;
305
306   /* Get the check uuid */
307   PG_GET_STR_COL(val, 0, "id");
308   if(!val) goto bad_row;
309   if(uuid_parse(val, d->rt->checkid)) goto bad_row;
310
311   /* Get the remote_address (which noit owns this) */
312   PG_GET_STR_COL(val, 0, "remote_address");
313   if(!val) goto bad_row;
314   d->rt->noit = strdup(val);
315
316   PQclear(d->res);
317   return DS_EXEC_SUCCESS;
318  bad_row:
319   return DS_EXEC_ROW_FAILED;
320 }
321 execute_outcome_t
322 stratcon_datastore_execute(conn_q *cq, struct sockaddr *r, ds_job_detail *d) {
323   int type, len;
324   char *final_buff;
325   uLong final_len, actual_final_len;;
326   char *token;
327
328   type = d->data[0];
329
330   /* Parse the log line, but only if we haven't already */
331   if(!d->nparams) {
332     char raddr[128];
333     char *scp, *ecp;
334
335     /* setup our remote address */
336     raddr[0] = '\0';
337     switch(r->sa_family) {
338       case AF_INET:
339         inet_ntop(AF_INET, &(((struct sockaddr_in *)r)->sin_addr),
340                   raddr, sizeof(raddr));
341         break;
342       case AF_INET6:
343         inet_ntop(AF_INET6, &(((struct sockaddr_in6 *)r)->sin6_addr),
344                   raddr, sizeof(raddr));
345         break;
346       default:
347         noitL(noit_error, "remote address of family %d\n", r->sa_family);
348     }
349  
350     scp = d->data;
351 #define PROCESS_NEXT_FIELD(t,l) do { \
352   if(!*scp) goto bad_row; \
353   ecp = strchr(scp, '\t'); \
354   if(!ecp) goto bad_row; \
355   token = scp; \
356   len = (ecp-scp); \
357   scp = ecp + 1; \
358 } while(0)
359 #define PROCESS_LAST_FIELD(t,l) do { \
360   if(!*scp) ecp = scp; \
361   else { \
362     ecp = scp + strlen(scp); /* Puts us at the '\0' */ \
363     if(*(ecp-1) == '\n') ecp--; /* We back up on letter if we ended in \n */ \
364   } \
365   t = scp; \
366   l = (ecp-scp); \
367 } while(0)
368
369     PROCESS_NEXT_FIELD(token,len); /* Skip the leader, we know what we are */
370     switch(type) {
371       /* See noit_check_log.c for log description */
372       case 'n':
373         DECLARE_PARAM_STR(raddr, strlen(raddr));
374         DECLARE_PARAM_STR("noitd",5); /* node_type */
375         PROCESS_NEXT_FIELD(token,len);
376         DECLARE_PARAM_STR(token,len); /* timestamp */
377
378         /* This is the expected uncompressed len */
379         PROCESS_NEXT_FIELD(token,len);
380         final_len = atoi(token);
381         final_buff = malloc(final_len);
382         if(!final_buff) goto bad_row;
383  
384         /* The last token is b64 endoded and compressed.
385          * we need to decode it, declare it and then free it.
386          */
387         PROCESS_LAST_FIELD(token, len);
388         /* We can in-place decode this */
389         len = noit_b64_decode((char *)token, len,
390                               (unsigned char *)token, len);
391         if(len <= 0) {
392           noitL(noit_error, "noitd config base64 decoding error.\n");
393           free(final_buff);
394           goto bad_row;
395         }
396         actual_final_len = final_len;
397         if(Z_OK != uncompress((Bytef *)final_buff, &actual_final_len,
398                               (unsigned char *)token, len)) {
399           noitL(noit_error, "noitd config decompression failure.\n");
400           free(final_buff);
401           goto bad_row;
402         }
403         if(final_len != actual_final_len) {
404           noitL(noit_error, "noitd config decompression error.\n");
405           free(final_buff);
406           goto bad_row;
407         }
408         DECLARE_PARAM_STR(final_buff, final_len);
409         free(final_buff);
410         break;
411       case 'C':
412         DECLARE_PARAM_STR(raddr, strlen(raddr));
413         PROCESS_NEXT_FIELD(token,len);
414         DECLARE_PARAM_STR(token,len); /* timestamp */
415         PROCESS_NEXT_FIELD(token, len);
416         DECLARE_PARAM_STR(token,len); /* uuid */
417         PROCESS_NEXT_FIELD(token, len);
418         DECLARE_PARAM_STR(token,len); /* target */
419         PROCESS_NEXT_FIELD(token, len);
420         DECLARE_PARAM_STR(token,len); /* module */
421         PROCESS_LAST_FIELD(token, len);
422         DECLARE_PARAM_STR(token,len); /* name */
423         break;
424       case 'M':
425         PROCESS_NEXT_FIELD(token,len);
426         DECLARE_PARAM_STR(token,len); /* timestamp */
427         PROCESS_NEXT_FIELD(token, len);
428         DECLARE_PARAM_STR(token,len); /* uuid */
429         PROCESS_NEXT_FIELD(token, len);
430         DECLARE_PARAM_STR(token,len); /* name */
431         PROCESS_NEXT_FIELD(token,len);
432         d->metric_type = *token;
433         PROCESS_LAST_FIELD(token,len);
434         DECLARE_PARAM_STR(token,len); /* value */
435         break;
436       case 'S':
437         PROCESS_NEXT_FIELD(token,len);
438         DECLARE_PARAM_STR(token,len); /* timestamp */
439         PROCESS_NEXT_FIELD(token, len);
440         DECLARE_PARAM_STR(token,len); /* uuid */
441         PROCESS_NEXT_FIELD(token, len);
442         DECLARE_PARAM_STR(token,len); /* state */
443         PROCESS_NEXT_FIELD(token, len);
444         DECLARE_PARAM_STR(token,len); /* availability */
445         PROCESS_NEXT_FIELD(token, len);
446         DECLARE_PARAM_STR(token,len); /* duration */
447         PROCESS_LAST_FIELD(token,len);
448         DECLARE_PARAM_STR(token,len); /* status */
449         break;
450       default:
451         goto bad_row;
452     }
453
454   }
455
456   /* Now execute the query */
457   switch(type) {
458     case 'n':
459       GET_QUERY(config_insert);
460       PG_EXEC(config_insert);
461       PQclear(d->res);
462       break;
463     case 'C':
464       GET_QUERY(check_insert);
465       PG_EXEC(check_insert);
466       PQclear(d->res);
467       break;
468     case 'S':
469       GET_QUERY(status_insert);
470       PG_EXEC(status_insert);
471       PQclear(d->res);
472       break;
473     case 'M':
474       switch(d->metric_type) {
475         case METRIC_INT32:
476         case METRIC_UINT32:
477         case METRIC_INT64:
478         case METRIC_UINT64:
479         case METRIC_DOUBLE:
480           GET_QUERY(metric_insert_numeric);
481           PG_EXEC(metric_insert_numeric);
482           PQclear(d->res);
483           break;
484         case METRIC_STRING:
485           GET_QUERY(metric_insert_text);
486           PG_EXEC(metric_insert_text);
487           PQclear(d->res);
488           break;
489         default:
490           goto bad_row;
491       }
492       break;
493     default:
494       /* should never get here */
495       goto bad_row;
496   }
497   return DS_EXEC_SUCCESS;
498  bad_row:
499   return DS_EXEC_ROW_FAILED;
500 }
501 static int
502 stratcon_database_connect(conn_q *cq) {
503   char dsn[512];
504   noit_hash_iter iter = NOIT_HASH_ITER_ZERO;
505   const char *k, *v;
506   int klen;
507   noit_hash_table *t;
508
509   dsn[0] = '\0';
510   t = noit_conf_get_hash(NULL, "/stratcon/database/dbconfig");
511   while(noit_hash_next_str(t, &iter, &k, &klen, &v)) {
512     if(dsn[0]) strlcat(dsn, " ", sizeof(dsn));
513     strlcat(dsn, k, sizeof(dsn));
514     strlcat(dsn, "=", sizeof(dsn));
515     strlcat(dsn, v, sizeof(dsn));
516   }
517
518   if(cq->dbh) {
519     if(PQstatus(cq->dbh) == CONNECTION_OK) return 0;
520     PQreset(cq->dbh);
521     if(PQstatus(cq->dbh) == CONNECTION_OK) return 0;
522     noitL(noit_error, "Error reconnecting to database: '%s'\nError: %s\n",
523           dsn, PQerrorMessage(cq->dbh));
524     return -1;
525   }
526
527   cq->dbh = PQconnectdb(dsn);
528   if(!cq->dbh) return -1;
529   if(PQstatus(cq->dbh) == CONNECTION_OK) return 0;
530   noitL(noit_error, "Error connection to database: '%s'\nError: %s\n",
531         dsn, PQerrorMessage(cq->dbh));
532   return -1;
533 }
534 static int
535 stratcon_datastore_savepoint_op(conn_q *cq, const char *p,
536                                 const char *name) {
537   int rv = -1;
538   PGresult *res;
539   char cmd[128];
540   strlcpy(cmd, p, sizeof(cmd));
541   strlcat(cmd, name, sizeof(cmd));
542   if((res = PQexec(cq->dbh, cmd)) == NULL) return -1;
543   if(PQresultStatus(res) == PGRES_COMMAND_OK) rv = 0;
544   PQclear(res);
545   return rv;
546 }
547 static int
548 stratcon_datastore_do(conn_q *cq, const char *cmd) {
549   PGresult *res;
550   int rv = -1;
551   if((res = PQexec(cq->dbh, cmd)) == NULL) return -1;
552   if(PQresultStatus(res) == PGRES_COMMAND_OK) rv = 0;
553   PQclear(res);
554   return rv;
555 }
556 #define BUSTED(cq) do { \
557   PQfinish((cq)->dbh); \
558   (cq)->dbh = NULL; \
559   goto full_monty; \
560 } while(0)
561 #define SAVEPOINT(name) do { \
562   if(stratcon_datastore_savepoint_op(cq, "SAVEPOINT ", name)) BUSTED(cq); \
563   last_sp = current; \
564 } while(0)
565 #define ROLLBACK_TO_SAVEPOINT(name) do { \
566   if(stratcon_datastore_savepoint_op(cq, "ROLLBACK TO SAVEPOINT ", name)) \
567     BUSTED(cq); \
568   last_sp = NULL; \
569 } while(0)
570 #define RELEASE_SAVEPOINT(name) do { \
571   if(stratcon_datastore_savepoint_op(cq, "RELEASE SAVEPOINT ", name)) \
572     BUSTED(cq); \
573   last_sp = NULL; \
574 } while(0)
575 int
576 stratcon_datastore_asynch_lookup(eventer_t e, int mask, void *closure,
577                                  struct timeval *now) {
578   conn_q *cq = closure;
579   ds_job_detail *current, *next;
580   if(!(mask & EVENTER_ASYNCH_WORK)) return 0;
581
582   if(!cq->head) return 0;
583
584   stratcon_database_connect(cq);
585
586   current = cq->head;
587   while(current) {
588     if(current->rt) {
589       next = current->next;
590       stratcon_datastore_find(cq, current);
591       current = next;
592     }
593     else if(current->completion_event) {
594       next = current->next;
595       eventer_add(current->completion_event);
596       current = next;
597       __remove_until(cq, current);
598     }
599     else current = current->next;
600   }
601   return 0;
602 }
603 int
604 stratcon_datastore_asynch_execute(eventer_t e, int mask, void *closure,
605                                   struct timeval *now) {
606   int i;
607   conn_q *cq = closure;
608   ds_job_detail *current, *last_sp;
609   if(!(mask & EVENTER_ASYNCH_WORK)) return 0;
610   if(mask & EVENTER_ASYNCH_CLEANUP) return 0;
611   if(!cq->head) return 0;
612
613  full_monty:
614   /* Make sure we have a connection */
615   i = 1;
616   while(stratcon_database_connect(cq)) {
617     noitL(noit_error, "Error connecting to database\n");
618     sleep(i);
619     i *= 2;
620     i = MIN(i, 16);
621   }
622
623   current = cq->head;
624   last_sp = NULL;
625   if(stratcon_datastore_do(cq, "BEGIN")) BUSTED(cq);
626   while(current) {
627     execute_outcome_t rv;
628     if(current->data) {
629       if(!last_sp) SAVEPOINT("batch");
630  
631       if(current->problematic) {
632         noitL(noit_error, "Failed noit line: %s", current->data);
633         RELEASE_SAVEPOINT("batch");
634         current = current->next;
635         continue;
636       }
637       rv = stratcon_datastore_execute(cq, cq->remote, current);
638       switch(rv) {
639         case DS_EXEC_SUCCESS:
640           current = current->next;
641           break;
642         case DS_EXEC_ROW_FAILED:
643           /* rollback to savepoint, mark this record as bad and start again */
644           current->problematic = 1;
645           current = last_sp;
646           ROLLBACK_TO_SAVEPOINT("batch");
647           break;
648         case DS_EXEC_TXN_FAILED:
649           BUSTED(cq);
650       }
651     }
652     if(current->completion_event) {
653       if(last_sp) RELEASE_SAVEPOINT("batch");
654       if(stratcon_datastore_do(cq, "COMMIT")) BUSTED(cq);
655       eventer_add(current->completion_event);
656       current = current->next;
657       __remove_until(cq, current);
658     }
659   }
660   return 0;
661 }
662 void
663 stratcon_datastore_push(stratcon_datastore_op_t op,
664                         struct sockaddr *remote, void *operand) {
665   conn_q *cq;
666   eventer_t e;
667   ds_job_detail *dsjd;
668   struct datastore_onlooker_list *nnode;
669
670   for(nnode = onlookers; nnode; nnode = nnode->next)
671     nnode->dispatch(op,remote,operand);
672
673   cq = __get_conn_q_for_remote(remote);
674   dsjd = calloc(1, sizeof(*dsjd));
675   switch(op) {
676     case DS_OP_FIND:
677       dsjd->rt = operand;
678       __append(cq, dsjd);
679       break;
680     case DS_OP_INSERT:
681       dsjd->data = operand;
682       __append(cq, dsjd);
683       break;
684     case DS_OP_FIND_COMPLETE:
685     case DS_OP_CHKPT:
686       dsjd->completion_event = operand;
687       __append(cq,dsjd);
688       e = eventer_alloc();
689       e->mask = EVENTER_ASYNCH;
690       if(op == DS_OP_FIND_COMPLETE)
691         e->callback = stratcon_datastore_asynch_lookup;
692       else if(op == DS_OP_CHKPT)
693         e->callback = stratcon_datastore_asynch_execute;
694       e->closure = cq;
695       eventer_add_asynch(cq->jobq, e);
696       break;
697   }
698 }
699
700 int
701 stratcon_datastore_saveconfig(void *unused) {
702   int rv = -1;
703   conn_q _cq = { 0 }, *cq = &_cq;
704   char *buff;
705   ds_single_detail _d = { 0 }, *d = &_d;
706
707   if(stratcon_database_connect(cq) == 0) {
708     char time_as_str[20];
709     size_t len;
710     buff = noit_conf_xml_in_mem(&len);
711     if(!buff) goto bad_row;
712
713     snprintf(time_as_str, sizeof(time_as_str), "%lu", (long unsigned int)time(NULL));
714     DECLARE_PARAM_STR("0.0.0.0", 7);
715     DECLARE_PARAM_STR("stratcond", 9);
716     DECLARE_PARAM_STR(time_as_str, strlen(time_as_str));
717     DECLARE_PARAM_STR(buff, len);
718     free(buff);
719
720     GET_QUERY(config_insert);
721     PG_EXEC(config_insert);
722     PQclear(d->res);
723     rv = 0;
724
725     bad_row:
726       free_params(d);
727   }
728   if(cq->dbh) PQfinish(cq->dbh);
729   return rv;
730 }
731
732 void
733 stratcon_datastore_register_onlooker(void (*f)(stratcon_datastore_op_t,
734                                                struct sockaddr *, void *)) {
735   struct datastore_onlooker_list *nnode;
736   nnode = calloc(1, sizeof(*nnode));
737   nnode->dispatch = f;
738   nnode->next = onlookers;
739   while(noit_atomic_casptr((void **)&onlookers, nnode, nnode->next) != (void *)nnode->next)
740     nnode->next = onlookers;
741 }
Note: See TracBrowser for help on using the browser.