diff --git a/Makefile.in b/Makefile.in index 2a088c4..bb48a59 100644 --- a/Makefile.in +++ b/Makefile.in @@ -30,6 +30,8 @@ acct.c \ addr.c \ bsd.c \ cap.c \ +cache.c \ +pf.c \ conv.c \ darkstat.c \ daylog.c \ @@ -159,10 +161,13 @@ addr.o: addr.c addr.h bsd.o: bsd.c bsd.h config.h cdefs.h cap.o: cap.c acct.h cdefs.h cap.h config.h conv.h decode.h addr.h err.h \ hosts_db.h linktypes.h localip.h now.h opt.h queue.h str.h +cache.o: cache.c cache.h +pf.o: pf.c pf.h acct.h cdefs.h config.h conv.h decode.h addr.h err.h \ + hosts_db.h localip.h now.h opt.h queue.h str.h cache.h conv.o: conv.c conv.h err.h cdefs.h darkstat.o: darkstat.c acct.h cap.h cdefs.h config.h conv.h daylog.h \ - graph_db.h db.h dns.h err.h hosts_db.h addr.h http.h localip.h ncache.h \ - now.h pidfile.h str.h + graph_db.h db.h dns.h err.h hosts_db.h addr.h http.h localip.h \ + ncache.h now.h pidfile.h str.h pf.h daylog.o: daylog.c cdefs.h err.h daylog.h graph_db.h str.h now.h db.o: db.c err.h cdefs.h hosts_db.h addr.h graph_db.h db.h decode.o: decode.c cdefs.h decode.h addr.h err.h opt.h diff --git a/acct.c b/acct.c index 540bf35..4e9abd6 100644 --- a/acct.c +++ b/acct.c @@ -174,7 +174,7 @@ void acct_for(const struct pktsummary * const sm, #endif /* Totals. */ - acct_total_packets++; + acct_total_packets += sm->packets; acct_total_bytes += sm->len; /* Graphs. */ diff --git a/cache.c b/cache.c new file mode 100644 index 0000000..d12317b --- /dev/null +++ b/cache.c @@ -0,0 +1,208 @@ +/* $OpenBSD: cache.c,v 1.8 2019/01/17 05:56:29 tedu Exp $ */ +/* + * Copyright (c) 2001, 2007 Can Erkin Acar + * + * Permission to use, copy, modify, and distribute this software for any + * purpose with or without fee is hereby granted, provided that the above + * copyright notice and this permission notice appear in all copies. + * + * THE SOFTWARE IS PROVIDED "AS IS" AND THE AUTHOR DISCLAIMS ALL WARRANTIES + * WITH REGARD TO THIS SOFTWARE INCLUDING ALL IMPLIED WARRANTIES OF + * MERCHANTABILITY AND FITNESS. IN NO EVENT SHALL THE AUTHOR BE LIABLE FOR + * ANY SPECIAL, DIRECT, INDIRECT, OR CONSEQUENTIAL DAMAGES OR ANY DAMAGES + * WHATSOEVER RESULTING FROM LOSS OF USE, DATA OR PROFITS, WHETHER IN AN + * ACTION OF CONTRACT, NEGLIGENCE OR OTHER TORTIOUS ACTION, ARISING OUT OF + * OR IN CONNECTION WITH THE USE OR PERFORMANCE OF THIS SOFTWARE. + */ + +#ifdef __OpenBSD__ + +#include +#include +#include + +#include +#include + +#include +#ifdef TEST_COMPAT +#include "pfvarmux.h" +#else +#include +#endif +#include + +#include +#include +#include + +#include + +#include "cache.h" + +/* prototypes */ +void update_state(struct sc_ent *, struct pfsync_state *); +struct sc_ent *cache_state(struct pfsync_state *); +static __inline int sc_cmp(struct sc_ent *s1, struct sc_ent *s2); + +/* initialize the tree and queue */ +RB_HEAD(sc_tree, sc_ent) sctree; +TAILQ_HEAD(sc_queue, sc_ent) scq1, scq2, scq_free; +RB_GENERATE(sc_tree, sc_ent, tlink, sc_cmp) + +struct sc_queue *scq_act = NULL; +struct sc_queue *scq_exp = NULL; + +int cache_max = 0; +int cache_size = 0; + +struct sc_ent *sc_store = NULL; + +/* preallocate the cache and insert into the 'free' queue */ +int +cache_init(int max) +{ + int n; + static int initialized = 0; + + if (max < 0 || initialized) + return (1); + + if (max == 0) { + sc_store = NULL; + } else { + sc_store = reallocarray(NULL, max, sizeof(struct sc_ent)); + if (sc_store == NULL) + return (1); + } + + RB_INIT(&sctree); + TAILQ_INIT(&scq1); + TAILQ_INIT(&scq2); + TAILQ_INIT(&scq_free); + + scq_act = &scq1; + scq_exp = &scq2; + + for (n = 0; n < max; n++) + TAILQ_INSERT_HEAD(&scq_free, sc_store + n, qlink); + + cache_size = cache_max = max; + initialized++; + + return (0); +} + +void +update_state(struct sc_ent *prev, struct pfsync_state *new) +{ + assert (prev != NULL && new != NULL); + + u_int64_t bytes[2] = { COUNTER(new->bytes[0]), COUNTER(new->bytes[1]) }; + prev->bytes_delta[0] = bytes[0] - prev->bytes[0]; + prev->bytes_delta[1] = bytes[1] - prev->bytes[1]; + prev->bytes[0] = bytes[0]; + prev->bytes[1] = bytes[1]; + + u_int64_t packets[2] = { COUNTER(new->packets[0]), COUNTER(new->packets[1]) }; + prev->packets_delta[0] = packets[0] - prev->packets[0]; + prev->packets_delta[1] = packets[1] - prev->packets[1]; + prev->packets[0] = packets[0]; + prev->packets[1] = packets[1]; +} + +void +add_state(struct pfsync_state *st) +{ + struct sc_ent *ent; + assert(st != NULL); + + if (cache_max == 0) + return; + + if (TAILQ_EMPTY(&scq_free)) + return; + + ent = TAILQ_FIRST(&scq_free); + TAILQ_REMOVE(&scq_free, ent, qlink); + + cache_size--; + + ent->id = st->id; + ent->creatorid = st->creatorid; + ent->bytes[0] = COUNTER(st->bytes[0]); + ent->bytes[1] = COUNTER(st->bytes[1]); + ent->packets[0] = COUNTER(st->packets[0]); + ent->packets[1] = COUNTER(st->packets[1]); + + RB_INSERT(sc_tree, &sctree, ent); + TAILQ_INSERT_HEAD(scq_act, ent, qlink); +} + +/* must be called only once for each state before cache_endupdate */ +struct sc_ent * +cache_state(struct pfsync_state *st) +{ + struct sc_ent ent, *old; + + if (cache_max == 0) + return (NULL); + + ent.id = st->id; + ent.creatorid = st->creatorid; + old = RB_FIND(sc_tree, &sctree, &ent); + + if (old == NULL) { + add_state(st); + return (NULL); + } + + if (COUNTER(st->bytes[0]) < old->bytes[0] || COUNTER(st->bytes[1]) < old->bytes[1]) + return (NULL); + + update_state(old, st); + + /* move to active queue */ + TAILQ_REMOVE(scq_exp, old, qlink); + TAILQ_INSERT_HEAD(scq_act, old, qlink); + + return (old); +} + +/* remove the states that are not updated in this cycle */ +void +cache_endupdate(void) +{ + struct sc_queue *tmp; + struct sc_ent *ent; + + while (! TAILQ_EMPTY(scq_exp)) { + ent = TAILQ_FIRST(scq_exp); + TAILQ_REMOVE(scq_exp, ent, qlink); + RB_REMOVE(sc_tree, &sctree, ent); + TAILQ_INSERT_HEAD(&scq_free, ent, qlink); + cache_size++; + } + + tmp = scq_act; + scq_act = scq_exp; + scq_exp = tmp; +} + +static __inline int +sc_cmp(struct sc_ent *a, struct sc_ent *b) +{ + if (a->id > b->id) + return (1); + if (a->id < b->id) + return (-1); + if (a->creatorid > b->creatorid) + return (1); + if (a->creatorid < b->creatorid) + return (-1); + return (0); +} + +#endif + +/* vim:set ts=3 sw=3 tw=78 expandtab: */ diff --git a/cache.h b/cache.h new file mode 100644 index 0000000..db08271 --- /dev/null +++ b/cache.h @@ -0,0 +1,49 @@ +/* $OpenBSD: cache.h,v 1.6 2019/01/17 05:56:29 tedu Exp $ */ +/* + * Copyright (c) 2001, 2007 Can Erkin Acar + * + * Permission to use, copy, modify, and distribute this software for any + * purpose with or without fee is hereby granted, provided that the above + * copyright notice and this permission notice appear in all copies. + * + * THE SOFTWARE IS PROVIDED "AS IS" AND THE AUTHOR DISCLAIMS ALL WARRANTIES + * WITH REGARD TO THIS SOFTWARE INCLUDING ALL IMPLIED WARRANTIES OF + * MERCHANTABILITY AND FITNESS. IN NO EVENT SHALL THE AUTHOR BE LIABLE FOR + * ANY SPECIAL, DIRECT, INDIRECT, OR CONSEQUENTIAL DAMAGES OR ANY DAMAGES + * WHATSOEVER RESULTING FROM LOSS OF USE, DATA OR PROFITS, WHETHER IN AN + * ACTION OF CONTRACT, NEGLIGENCE OR OTHER TORTIOUS ACTION, ARISING OUT OF + * OR IN CONNECTION WITH THE USE OR PERFORMANCE OF THIS SOFTWARE. + */ + +#ifndef _CACHE_H_ +#define _CACHE_H_ + +#include +#include +#include +#include +#include +#include +#include +#include +#include + +struct sc_ent { + RB_ENTRY(sc_ent) tlink; + TAILQ_ENTRY(sc_ent) qlink; + u_int64_t id; + u_int32_t creatorid; + u_int64_t bytes[2]; + u_int64_t bytes_delta[2]; + u_int64_t packets[2]; + u_int64_t packets_delta[2]; +}; + +int cache_init(int); +void cache_endupdate(void); +struct sc_ent *cache_state(struct pfsync_state *); +extern int cache_max, cache_size; + +#define COUNTER(c) ((((u_int64_t) ntohl(c[0]))<<32) + ntohl(c[1])) + +#endif diff --git a/cap.c b/cap.c index e2bb4e2..893809e 100644 --- a/cap.c +++ b/cap.c @@ -391,6 +391,7 @@ static void callback(u_char *user, if (opt_want_hexdump) hexdump(pdata, pheader->caplen, iface->linkhdr); memset(&sm, 0, sizeof(sm)); + sm.packets = 1; if (iface->linkhdr->decoder(pheader, pdata, &sm)) acct_for(&sm, &iface->local_ips); } diff --git a/darkstat.c b/darkstat.c index e3282ec..a20dd84 100644 --- a/darkstat.c +++ b/darkstat.c @@ -23,6 +23,9 @@ #include "now.h" #include "pidfile.h" #include "str.h" +#ifdef __OpenBSD__ +#include "pf.h" +#endif #include #include @@ -68,6 +71,13 @@ static unsigned long parsenum(const char *str, return n; } +static int opt_pf_seen = 0; +#ifdef __OpenBSD__ +static void cb_pf(const char *arg) { + opt_pf_seen = 1; +} +#endif + static int opt_iface_seen = 0; static void cb_interface(const char *arg) { cap_add_ifname(arg); @@ -220,6 +230,9 @@ static struct cmdline_arg cmdline_args[] = { {"--hexdump", NULL, cb_hexdump, 0}, {"--version", NULL, cb_version, 0}, {"--help", NULL, cb_help, 0}, +#ifdef __OpenBSD__ + {"--pf", NULL, cb_pf, 0}, +#endif {NULL, NULL, NULL, 0} }; @@ -306,15 +319,17 @@ static void parse_cmdline(const int argc, char * const *argv) { if (opt_want_syslog) openlog("darkstat", LOG_NDELAY | LOG_PID, LOG_DAEMON); - /* default value */ + /* some default values */ + if (opt_chroot_dir == NULL) + opt_chroot_dir = CHROOT_DIR; if (opt_privdrop_user == NULL) opt_privdrop_user = PRIVDROP_USER; /* sanity check args */ - if (!opt_iface_seen && opt_capfile == NULL) + if (!opt_pf_seen && !opt_iface_seen && opt_capfile == NULL) errx(1, "must specify either interface (-i) or capture file (-r)"); - if (opt_iface_seen && opt_capfile != NULL) + if (opt_pf_seen && opt_iface_seen && opt_capfile != NULL) errx(1, "can't specify both interface (-i) and capture file (-r)"); if ((opt_hosts_max != 0) && (opt_hosts_keep >= opt_hosts_max)) { @@ -385,7 +400,13 @@ main(int argc, char **argv) /* do this first as it forks - minimize memory use */ if (opt_want_dns) dns_init(opt_privdrop_user); - cap_start(opt_want_promisc); /* needs root */ + if (opt_pf_seen) { +#ifdef __OpenBSD__ + pfsync_start(); +#endif + } else { + cap_start(opt_want_promisc); + } http_init_base(opt_base); http_listen(opt_bindport); ncache_init(); /* must do before chroot() */ @@ -422,7 +443,13 @@ main(int argc, char **argv) FD_ZERO(&rs); FD_ZERO(&ws); - cap_fd_set(&rs, &max_fd, &timeout, &use_timeout); + if (opt_pf_seen) { +#ifdef __OpenBSD__ + pfsync_fd_set(&rs, &max_fd, &timeout, &use_timeout); +#endif + } else { + cap_fd_set(&rs, &max_fd, &timeout, &use_timeout); + } http_fd_set(&rs, &ws, &max_fd, &timeout, &use_timeout); select_ret = select(max_fd+1, &rs, &ws, NULL, @@ -454,7 +481,13 @@ main(int argc, char **argv) } graph_rotate(); - cap_ret = cap_poll(&rs); + if (opt_pf_seen) { +#ifdef __OpenBSD__ + cap_ret = pfsync_poll(); +#endif + } else { + cap_ret = cap_poll(&rs); + } dns_poll(); http_poll(&rs, &ws); timer_stop(&t, 1000000000, "event processing took longer than a second"); diff --git a/decode.h b/decode.h index 827c0ad..17aca48 100644 --- a/decode.h +++ b/decode.h @@ -27,7 +27,8 @@ struct pktsummary { /* Fields are in host byte order (except IPs) */ struct addr src, dst; - uint16_t len; + uint64_t packets; + uint64_t len; uint8_t proto; /* IPPROTO_INVALID means don't do proto accounting */ uint8_t tcp_flags; /* only for TCP */ uint16_t src_port, dst_port; /* only for TCP, UDP */ diff --git a/pf.c b/pf.c new file mode 100644 index 0000000..2f0347b --- /dev/null +++ b/pf.c @@ -0,0 +1,225 @@ +/* darkstat 3 + * copyright (c) 2001-2014 Emil Mikulic. + * + * pf.c: read pf states, and hand them off to decode and acct. + * + * You may use, modify and redistribute this file under the terms of the + * GNU General Public License version 2. (see COPYING.GPL) + */ + +#ifdef __OpenBSD__ + +#include "pf.h" +#include "acct.h" +#include "cdefs.h" +#include "config.h" +#include "conv.h" +#include "decode.h" +#include "err.h" +#include "hosts_db.h" +#include "localip.h" +#include "now.h" +#include "opt.h" +#include "queue.h" +#include "str.h" +#include "cache.h" + +#include +#include +#include +#include +#include +#include + +#define MIN_NUM_STATES 1024 +#define NUM_STATE_INC 1024 + +#define DEFAULT_CACHE_SIZE 10000 + +struct pfsync_state *state_buf = NULL; +size_t state_buf_len = 0; +size_t num_states = 0; + +int pf_dev = -1; + +struct local_ips local_ips; + +void pfsync_start(void) { + struct str *ifs = str_make(); + str_appendn(ifs, "", 1); /* NUL terminate */ + { + size_t _; + str_extract(ifs, &_, &title_interfaces); + } + + pf_dev = open("/dev/pf", O_RDONLY); + if (pf_dev == -1) { + errx(1, "pfsync_start"); + } + + if (cache_init(1<<16)) { + errx(1, "pfsync_start: cache_init"); + } + + // FIXME: initialize ip addresses based on egress, by default? + localip_init(&local_ips); +} + +void pfsync_fd_set(fd_set *read_set, int *max_fd, + struct timeval *timeout, int *need_timeout) { + assert(*need_timeout == 0); /* we're first to get a shot at the fd_set */ + + *need_timeout = 1; + timeout->tv_sec = 0; + timeout->tv_usec = 1000; +} + +void +alloc_buf(size_t ns) +{ + size_t len; + + if (ns < MIN_NUM_STATES) + ns = MIN_NUM_STATES; + + len = ns; + + if (len >= state_buf_len) { + len += NUM_STATE_INC; + state_buf = reallocarray(state_buf, len, + sizeof(struct pfsync_state)); + if (state_buf == NULL) + errx(1, "realloc"); + state_buf_len = len; + } +} + +int state_simple_idx(struct pfsync_state *s, int idx) { + struct pfsync_state_key *sk, *nk; + + sk = &s->key[PF_SK_STACK]; + nk = &s->key[PF_SK_WIRE]; + + if (sk->af != nk->af || PF_ANEQ(&sk->addr[idx], &nk->addr[idx], sk->af) || + sk->port[idx] != nk->port[idx] || sk->rdomain != nk->rdomain) { + return 0; + } + return 1; +} + +void pfsync_handle_state(struct pfsync_state * s, struct sc_ent * ent) { + int afto, dir; + +#undef v4 + + // FIXME: support ipv6 + if (s->key[0].af != AF_INET || s->key[1].af != AF_INET) { + return; + } + if (!state_simple_idx(s, 0) || !state_simple_idx(s, 1)) { + return; + } + + u_int64_t bytes_delta[2]; + u_int64_t packets_delta[2]; + if (ent != NULL) { + bytes_delta[0] = ent->bytes_delta[0]; + bytes_delta[1] = ent->bytes_delta[1]; + packets_delta[0] = ent->packets_delta[0]; + packets_delta[1] = ent->packets_delta[1]; + } else { + bytes_delta[0] = COUNTER(s->bytes[0]); + bytes_delta[1] = COUNTER(s->bytes[1]); + packets_delta[0] = COUNTER(s->packets[0]); + packets_delta[1] = COUNTER(s->packets[1]); + } + + afto = s->key[PF_SK_STACK].af == s->key[PF_SK_WIRE].af ? 0 : 1; + dir = afto ? PF_OUT : s->direction; + + struct pktsummary sm; + memset(&sm, 0, sizeof(sm)); + + sm.src.family = IPv4; + sm.dst.family = IPv4; + sm.proto = s->proto; + + sm.len = (dir == PF_OUT) ? bytes_delta[0] : bytes_delta[1]; + sm.packets = (dir == PF_OUT) ? packets_delta[0] : packets_delta[1]; + + if (sm.len > 0) { + struct pfsync_state_key *ks = &s->key[afto ? PF_SK_STACK : PF_SK_WIRE]; + sm.src.ip.v4 = ks->addr[1].pfa.v4.s_addr; + sm.src_port = ntohs(ks->port[1]); + + sm.dst.ip.v4 = ks->addr[0].pfa.v4.s_addr; + sm.dst_port = ntohs(ks->port[0]); + + acct_for(&sm, &local_ips); + } + + sm.len = (dir == PF_OUT) ? bytes_delta[1] : bytes_delta[0]; + sm.packets = (dir == PF_OUT) ? packets_delta[1] : packets_delta[0]; + if (sm.len > 0) { + struct pfsync_state_key *ks = &s->key[PF_SK_STACK]; + sm.src.ip.v4 = ks->addr[0].pfa.v4.s_addr; + sm.src_port = ntohs(ks->port[0]); + + sm.dst.ip.v4 = ks->addr[1].pfa.v4.s_addr; + sm.dst_port = ntohs(ks->port[1]); + + acct_for(&sm, &local_ips); + } +} + +int pfsync_poll(void) { + int n; + struct pfioc_states ps; + + ps.ps_len = 0; + ps.ps_states = NULL; + + if (ioctl(pf_dev, DIOCGETSTATES, &ps) == -1) { + errx(1, "DIOCGETSTATES"); + } + + for (;;) { + size_t sbytes = state_buf_len * sizeof(struct pfsync_state); + + ps.ps_len = sbytes; + ps.ps_states = state_buf; + + if (ioctl(pf_dev, DIOCGETSTATES, &ps) == -1) { + errx(1, "DIOCGETSTATES"); + } + num_states = ps.ps_len / sizeof(struct pfsync_state); + + if (ps.ps_len < sbytes) + break; + + alloc_buf(num_states); + } + + for (n = 0; n < num_states; n++) { + struct sc_ent *ent = cache_state(state_buf + n); + if (ent != NULL) { + pfsync_handle_state(state_buf + n, ent); + } + } + cache_endupdate(); + + return 1; +} + +void pfsync_stop(void) { + int ret = close(pf_dev); + if (ret == -1) { + errx(1, "pfsync_stop: close"); + } + + localip_free(&local_ips); +} + +#endif + +/* vim:set ts=3 sw=3 tw=78 expandtab: */ diff --git a/pf.h b/pf.h new file mode 100644 index 0000000..8e6d5fd --- /dev/null +++ b/pf.h @@ -0,0 +1,17 @@ +/* darkstat 3 + * copyright (c) 2001-2014 Emil Mikulic. + * + * pf.h: interface to OpenBSD pf. + */ + +#include /* OpenBSD needs this before select */ +#include /* FreeBSD 4 needs this for struct timeval */ +#include + +void pfsync_start(void); +void pfsync_fd_set(fd_set *read_set, int *max_fd, + struct timeval *timeout, int *need_timeout); +int pfsync_poll(void); +void pfsync_stop(void); + +/* vim:set ts=3 sw=3 tw=78 expandtab: */