AnyEvent-FastPing

 view release on metacpan or  search on metacpan

FastPing.xs  view on Meta::CPAN

#include "XSUB.h"

#include <pthread.h>

#include <stdio.h>
#include <stdlib.h>
#include <string.h>

#include <time.h>
#include <poll.h>
#include <unistd.h>
#include <inttypes.h>
#include <fcntl.h>
#include <errno.h>
#include <limits.h>

#include <sys/types.h>
#include <sys/time.h>
#include <sys/socket.h>

#include <netinet/in.h>
#include <arpa/inet.h>

#ifdef __linux
# include <linux/icmp.h>
#endif
#if ENABLE_IPV6 && !defined (__CYGWIN__)
# include <netinet/icmp6.h>
#endif

#define ICMP4_ECHO       8
#define ICMP4_ECHO_REPLY 0
#define ICMP6_ECHO       128
#define ICMP6_ECHO_REPLY 129

#define DRAIN_INTERVAL 1e-6 // how long to wait when sendto returns ENOBUFS, in seconds
#define MIN_INTERVAL   1e-6 // minimum packet send interval, in seconds

#define HDR_SIZE_IP4  20
#define HDR_SIZE_IP6  48

static int thr_res[2]; // worker thread finished status
static int icmp4_fd = -1;
static int icmp6_fd = -1;

/*****************************************************************************/

typedef double tstamp;

static tstamp
NOW (void)
{
  struct timeval tv;

  gettimeofday (&tv, 0);

  return tv.tv_sec + tv.tv_usec * 1e-6;
}

static void
ssleep (tstamp wait)
{
#if defined (__SVR4) && defined (__sun)
  struct timeval tv;

  tv.tv_sec  = wait;
  tv.tv_usec = (wait - tv.tv_sec) * 1e6;

  select (0, 0, 0, 0, &tv);
#elif defined(_WIN32)
  Sleep ((unsigned long)(delay * 1e3));
#else
  struct timespec ts;

  ts.tv_sec  = wait;
  ts.tv_nsec = (wait - ts.tv_sec) * 1e9;

  nanosleep (&ts, 0);
#endif
}

/*****************************************************************************/

typedef struct
{
  uint8_t version_ihl;
  uint8_t tos;
  uint16_t tot_len;

  uint16_t id;
  uint16_t flags;

  uint8_t ttl;
  uint8_t protocol;
  uint16_t cksum;

  uint32_t src;
  uint32_t dst;
} IP4HDR;

/*****************************************************************************/

typedef uint8_t addr_tt[16];

typedef struct
{
  tstamp next;
  tstamp interval;
  int addrlen;

  addr_tt lo, hi; /* only if !addrcnt */

  int addrcnt;
  /* addrcnt addresses follow */
} RANGE;

typedef struct
{
  RANGE **ranges;
  int rangecnt, rangemax;

FastPing.xs  view on Meta::CPAN

{
  return pkt->id    == pinger->magic1
      && pkt->seq   == pinger->magic2
      && pkt->magic == pinger->magic3;
}

static void
ts_to_pkt (PKT *pkt, tstamp ts)
{
  /* move 12 bits of seconds into the 32 bit fractional part */
  /* leaving 20 bits subsecond resolution and 44 bits of integers */
  /* (of which 32 are typically usable) */
  ts *= 1. / 4096.;

  pkt->stamp_hi = ts;
  pkt->stamp_lo = (ts - pkt->stamp_hi) * 4294967296.;
}

static tstamp
pkt_to_ts (PKT *pkt)
{
  return pkt->stamp_hi *  4096.
       + pkt->stamp_lo * (4096. / 4294967296.);
}

static void
pkt_cksum (PKT *pkt)
{
  uint_fast32_t sum = 0;
  uint32_t *wp = (uint32_t *)pkt;
  int len = sizeof (*pkt) / 4;

  do
    {
      uint_fast32_t w = *(volatile uint32_t *)wp++;
      sum += (w & 0xffff) + (w >> 16);
    }
  while (--len);

  sum = (sum >> 16) + (sum & 0xffff);   /* add high 16 to low 16 */
  sum += sum >> 16;                     /* add carry */

  pkt->cksum = ~sum;
}

/*****************************************************************************/

static void
range_free (RANGE *self)
{
  free (self);
}

/* like sendto, but retries on failure */
static void
xsendto (int fd, void *buf, size_t len, int flags, void *sa, int salen)
{
  tstamp wait = DRAIN_INTERVAL / 2.;

  while (sendto (fd, buf, len, flags, sa, salen) < 0 && errno == ENOBUFS)
    ssleep (wait *= 2.);
}

// ping current address, return true and increment if more to ping
static int
range_send_ping (RANGE *self, PKT *pkt)
{
  // send ping
  uint8_t *addr;
  int addrlen;

  if (self->addrcnt)
    addr = (self->addrcnt - 1) * self->addrlen + (uint8_t *)(self + 1);
  else
    addr = sizeof (addr_tt) - self->addrlen + self->lo;

  addrlen = self->addrlen;

  /* convert ipv4 mapped addresses - this only works for host lists */
  /* this tries to match 0000:0000:0000:0000:0000:ffff:a.b.c.d */
  /* efficiently but also with few insns */
  if (addrlen == 16 && !addr [0] && icmp4_fd >= 0
      && !(              addr [ 1]
           | addr [ 2] | addr [ 3]
           | addr [ 4] | addr [ 5]
           | addr [ 6] | addr [ 7]
           | addr [ 8] | addr [ 9]
           | (255-addr [10]) | (255-addr [11])))
    {
      addr += 12;
      addrlen -= 12;
    }

  pkt->cksum = 0;

  if (addrlen == 4)
    {
      struct sockaddr_in sa;

      pkt->type = ICMP4_ECHO;
      pkt_cksum (pkt);

      sa.sin_family = AF_INET;
      sa.sin_port   = 0;

      memcpy (&sa.sin_addr, addr, sizeof (sa.sin_addr));

      xsendto (icmp4_fd, pkt, sizeof (*pkt), 0, &sa, sizeof (sa));
    }
  else
    {
#if ENABLE_IPV6
      struct sockaddr_in6 sa;

      pkt->type = ICMP6_ECHO;

      sa.sin6_family   = AF_INET6;
      sa.sin6_port     = 0;
      sa.sin6_flowinfo = 0;
      sa.sin6_scope_id = 0;

FastPing.xs  view on Meta::CPAN

  self->ranges [k] = elem;
}

static void
upheap (PINGER *self, int k)
{
  RANGE *elem = self->ranges [k];

  while (k)
    {
      int j = (k - 1) >> 1;

      if (self->ranges [j]->next <= elem->next)
        break;

      self->ranges [k] = self->ranges [j];

      k = j;
    }

  self->ranges [k] = elem;
}

static void *
ping_proc (void *self_)
{
  PINGER *self = (PINGER *)self_;
  PKT pkt;

  memset (&pkt, 0, sizeof (pkt));

  tstamp now = NOW ();

  pkt.code   = 0;
  pkt.id     = self->magic1;
  pkt.seq    = self->magic2;
  pkt.magic  = self->magic3;
  pkt.pinger = self->id;

  if (self->next < now)
    self->next = now;

  while (self->rangecnt)
    {
      RANGE *range = self->ranges [0];

      // ranges [0] is always the next range to ping
      tstamp wait = range->next - now;

      // compare with the global frequency limit
      {
        tstamp diff = self->next - now;

        if (wait < diff)
          wait = diff; // global rate limit overrides
        else
          self->next = range->next; // fast forward
      }

      if (wait > 0.)
        ssleep (wait);

      now = NOW ();

      ts_to_pkt (&pkt, now);

      if (!range_send_ping (range, &pkt))
        {
          self->ranges [0] = self->ranges [--self->rangecnt];
          range_free (range);
        }
      else
        range->next = self->next + range->interval;

      downheap (self);

      self->next += self->interval;
      now = NOW ();
    }

  ssleep (self->maxrtt);

  {
    uint16_t id = self->id;

    write (thr_res [1], &id, sizeof (id));
  }

  return 0;
}

/*****************************************************************************/

/* NetBSD, Solaris... */
#ifndef PTHREAD_STACK_MIN
# define PTHREAD_STACK_MIN 0
#endif

static void
pinger_start (PINGER *self)
{
  sigset_t fullsigset, oldsigset;
  pthread_attr_t attr;

  if (self->running)
    return;

  sigfillset (&fullsigset);

  pthread_attr_init (&attr);
  pthread_attr_setstacksize (&attr, PTHREAD_STACK_MIN < sizeof (long) * 2048 ? sizeof (long) * 2048 : PTHREAD_STACK_MIN);

  pthread_sigmask (SIG_SETMASK, &fullsigset, &oldsigset);

  if (pthread_create (&self->thrid, &attr, ping_proc, (void *)self))
    croak ("AnyEvent::FastPing: unable to create pinger thread");

  pthread_sigmask (SIG_SETMASK, &oldsigset, 0);

  self->running = 1;
}

static void
pinger_stop (PINGER *self)
{
  if (!self->running)
    return;

  self->running = 0;
  pthread_cancel (self->thrid);
  pthread_join (self->thrid, 0);
}

static void
pinger_init (PINGER *self)
{
  memset (self, 0, sizeof (PINGER));

  if (firstfree >= 0)
    {
      self->id = firstfree;



( run in 0.866 second using v1.01-cache-2.11-cpan-acf6aa7dc9e )