#include "defs.h"
static int client_insane(struct server_info **, int, server_info_t);
void
client_init(void)
{
}
void
client_main(struct server_info **info_ary, int count)
{
struct server_info *best_off;
struct server_info *best_freq;
double last_freq;
double freq;
double offset;
int calc_offset_correction;
int didreconnect;
int i;
int insane;
last_freq = 0.0;
for (;;) {
calc_offset_correction = !sysntp_offset_correction_is_running();
for (i = 0; i < count; ++i)
client_poll(info_ary[i], min_sleep_opt, calc_offset_correction);
best_off = NULL;
best_freq = NULL;
for (i = 0; i < count; ++i)
client_check(&info_ary[i], &best_off, &best_freq);
if (best_off) {
insane = client_insane(info_ary, count, best_off);
if (insane > 0) {
best_off->server_insane = 0;
} else if (insane == 0) {
best_off = NULL;
} else {
best_off->server_insane = 60 * 60;
logdebuginfo(best_off, 1,
"excessive offset deviation, mapping out\n");
best_off = NULL;
}
}
if (best_off) {
offset = best_off->lin_sumoffset / best_off->lin_countoffset;
lin_resetalloffsets(info_ary, count);
if (offset < -COURSE_OFFSET_CORRECTION_LIMIT ||
offset > COURSE_OFFSET_CORRECTION_LIMIT ||
quickset_opt
) {
freq = sysntp_correct_course_offset(offset);
quickset_opt = 0;
} else {
freq = sysntp_correct_offset(offset);
}
} else {
freq = 0.0;
}
if (best_freq) {
freq += best_freq->lin_cache_freq;
if (last_freq != freq) {
sysntp_correct_freq(freq);
last_freq = freq;
}
}
didreconnect = 0;
for (i = 0; i < count; ++i)
client_manage_polling_mode(info_ary[i], &didreconnect);
if (didreconnect)
client_check_duplicate_ips(info_ary, count);
usleep(min_sleep_opt * 1000000 + random() % 500000);
}
}
void
client_poll(server_info_t info, int poll_interval, int calc_offset_correction)
{
struct timeval rtv;
struct timeval ltv;
struct timeval lbtv;
double offset;
if (info->server_insane > poll_interval)
info->server_insane -= poll_interval;
else
info->server_insane = 0;
if (info->poll_sleep > poll_interval) {
info->poll_sleep -= poll_interval;
return;
}
info->poll_sleep = 0;
if (info->fd < 0) {
if (info->poll_failed < 0x7FFFFFFF)
++info->poll_failed;
return;
}
logdebuginfo(info, 4, "poll, ");
if (udp_ntptimereq(info->fd, &rtv, <v, &lbtv) < 0) {
++info->poll_failed;
logdebug(4, "no response (%d failures in a row)\n", info->poll_failed);
if (info->poll_failed == POLL_FAIL_RESET) {
if (info->lin_count != 0) {
logdebuginfo(info, 4, "resetting regression due to failures\n");
}
lin_reset(info);
}
return;
}
++info->poll_count;
info->poll_failed = 0;
offset = tv_delta_double(&rtv, <v);
while (info) {
if (debug_level >= 4) {
struct tm *tp;
char buf[64];
time_t t;
t = rtv.tv_sec;
tp = localtime(&t);
strftime(buf, sizeof(buf), "%d-%b-%Y %H:%M:%S", tp);
logdebug(4, "%s.%03ld ", buf, rtv.tv_usec / 1000);
}
lin_regress(info, <v, &lbtv, offset, calc_offset_correction);
info = info->altinfo;
if (info && debug_level >= 4) {
logdebug(4, "%*.*s: poll, ",
(int)strlen(info->target),
(int)strlen(info->target), "(alt)");
}
}
}
void
client_check(struct server_info **checkp,
struct server_info **best_off,
struct server_info **best_freq)
{
struct server_info *check = *checkp;
struct server_info *info;
int min_samples;
if (check->lin_count >= LIN_RESTART / 2 && check->altinfo == NULL) {
info = malloc(sizeof(*info));
assert(info != NULL);
bcopy(check, info, sizeof(*info));
check->altinfo = info;
lin_reset(info);
}
if (check->lin_count >= LIN_RESTART) {
if ((info = check->altinfo) && info->lin_count >= LIN_RESTART / 2) {
double freq_diff;
freq_diff = info->lin_cache_freq - check->lin_cache_freq;
logdebuginfo(info, 4, "Switching to alternate, Frequency "
"difference is %6.3f ppm\n",
freq_diff * 1.0E+6);
*checkp = info;
free(check);
check = info;
}
}
info = *best_freq;
if ((check->lin_count >= 8 && fabs(check->lin_cache_corr) >= 0.99) ||
(check->lin_count >= 16 && fabs(check->lin_cache_corr) >= 0.96)
) {
if (info == NULL ||
fabs(check->lin_cache_corr) > fabs(info->lin_cache_corr)
) {
info = check;
*best_freq = info;
}
}
info = *best_off;
if (info && info->poll_mode == POLL_MAINTAIN)
min_samples = 8;
else
min_samples = 4;
if (check->lin_countoffset >= min_samples &&
(check->lin_cache_stddev <
fabs(check->lin_sumoffset / check->lin_countoffset / 4)) &&
check->server_insane == 0
) {
if (info == NULL ||
fabs(check->lin_cache_stddev) < fabs(info->lin_cache_stddev)
) {
info = check;
*best_off = info;
}
}
}
void
client_manage_polling_mode(struct server_info *info, int *didreconnect)
{
if (info->server_state == -2)
return;
if (info->poll_sleep)
return;
switch(info->poll_mode) {
case POLL_FIXED:
if (info->fd < 0) {
logdebuginfo(info, 2, "polling mode INIT, relookup & reconnect\n");
reconnect_server(info);
*didreconnect = 1;
if (info->fd < 0) {
if (info->poll_failed >= POLL_RECOVERY_RESTART * 5)
info->poll_sleep = max_sleep_opt;
else if (info->poll_failed >= POLL_RECOVERY_RESTART)
info->poll_sleep = nom_sleep_opt;
else
info->poll_sleep = min_sleep_opt;
break;
}
client_setserverstate(info, 0, "DNS lookup success");
if (info->poll_failed >= POLL_RECOVERY_RESTART * 5) {
logdebuginfo(info, 2, "polling mode INIT->STARTUP (very slow)\n");
info->poll_mode = POLL_STARTUP;
info->poll_sleep = max_sleep_opt;
info->poll_count = 0;
break;
} else if (info->poll_failed >= POLL_RECOVERY_RESTART) {
logdebuginfo(info, 2, "polling mode INIT->STARTUP (slow)\n");
info->poll_mode = POLL_STARTUP;
info->poll_count = 0;
break;
}
}
info->poll_mode = POLL_STARTUP;
logdebuginfo(info, 2, "polling mode INIT->STARTUP (normal)\n");
case POLL_STARTUP:
if (info->poll_failed >= POLL_FAIL_RESET) {
logdebuginfo(info, 2, "polling mode STARTUP->FAILED\n");
info->poll_mode = POLL_FAILED;
info->poll_count = 0;
break;
}
if (info->poll_count)
client_setserverstate(info, 1, "connected ok");
if (info->poll_count < POLL_STARTUP_MAX) {
info->poll_sleep = min_sleep_opt;
break;
}
info->poll_mode = POLL_ACQUIRE;
info->poll_count = 0;
logdebuginfo(info, 2, "polling mode STARTUP->ACQUIRE\n");
case POLL_ACQUIRE:
if (info->poll_failed >= POLL_FAIL_RESET) {
logdebuginfo(info, 2, "polling mode STARTUP->FAILED\n");
info->poll_mode = POLL_FAILED;
info->poll_count = 0;
break;
}
if (info->poll_count < POLL_ACQUIRE_MAX ||
info->lin_count < 8 ||
fabs(info->lin_cache_corr) < 0.85
) {
if (info->poll_count >= POLL_ACQUIRE_MAX &&
info->lin_count == LIN_RESTART - 2
) {
logdebuginfo(info, 2,
"WARNING: Unable to shift this source to "
"maintenance mode. Target correlation is awful\n");
}
break;
}
info->poll_mode = POLL_MAINTAIN;
info->poll_count = 0;
logdebuginfo(info, 2, "polling mode ACQUIRE->MAINTAIN\n");
case POLL_MAINTAIN:
if (info->poll_failed >= POLL_FAIL_RESET) {
logdebuginfo(info, 2, "polling mode STARTUP->FAILED\n");
info->poll_mode = POLL_FAILED;
info->poll_count = 0;
break;
}
if (info->lin_count >= LIN_RESTART / 2 &&
fabs(info->lin_cache_corr) < 0.70
) {
logdebuginfo(info, 2,
"polling mode MAINTAIN->ACQUIRE. Unable to maintain\n"
"the maintenance mode because the correlation went"
" bad!\n");
info->poll_mode = POLL_ACQUIRE;
info->poll_count = 0;
break;
}
info->poll_sleep = max_sleep_opt;
break;
case POLL_FAILED:
if (info->poll_count != 0) {
logdebuginfo(info, 2, "polling mode FAILED->ACQUIRE\n");
if (info->poll_failed >= POLL_FAIL_RESET)
info->poll_mode = POLL_STARTUP;
else
info->poll_mode = POLL_ACQUIRE;
break;
}
if (info->poll_failed >= POLL_RECOVERY_RESTART) {
logdebuginfo(info, 2, "polling mode FAILED->INIT\n");
client_setserverstate(info, 0, "FAILED");
disconnect_server(info);
info->poll_mode = POLL_FIXED;
break;
}
break;
}
if (info->poll_sleep == 0)
info->poll_sleep = nom_sleep_opt;
}
void
client_check_duplicate_ips(struct server_info **info_ary, int count)
{
server_info_t info1;
server_info_t info2;
int tries;
int i;
int j;
for (i = 0; i < count; ++i) {
info1 = info_ary[i];
if (info1->fd < 0 || info1->server_state != 0)
continue;
for (tries = 0; tries < 10; ++tries) {
for (j = 0; j < count; ++j) {
info2 = info_ary[j];
if (i == j || info2->fd < 0)
continue;
if (info1->fd < 0 ||
strcmp(info1->ipstr, info2->ipstr) == 0) {
reconnect_server(info1);
break;
}
}
if (j == count)
break;
}
if (tries == 10) {
disconnect_server(info1);
client_setserverstate(info1, -2,
"permanently disabling duplicate server");
}
}
}
static
int
client_insane(struct server_info **info_ary, int count, server_info_t best)
{
server_info_t info;
double best_offset;
double info_offset;
int good;
int bad;
int skip;
int quorum;
int i;
if (count < 2)
return(1);
best_offset = best->lin_sumoffset / best->lin_countoffset;
quorum = count;
for (i = 0; i < count; ++i) {
info = info_ary[i];
if (info->server_state == -2)
--quorum;
}
quorum = quorum / 2 + 1;
good = 0;
bad = 0;
skip = 0;
for (i = 0; i < count; ++i) {
info = info_ary[i];
if (info->lin_countoffset < 4 ||
info->lin_cache_stddev > insane_deviation
) {
++skip;
continue;
}
info_offset = info->lin_sumoffset / info->lin_countoffset;
info_offset -= best_offset;
if (info_offset < -insane_deviation || info_offset > insane_deviation)
++bad;
else
++good;
}
logdebuginfo(best, 5, "insanecheck good=%d bad=%d skip=%d "
"quorum=%d (allowed=%-+8.6f)\n",
good, bad, skip, quorum, insane_deviation);
if (good >= quorum)
return(1);
if (good + skip >= quorum)
return(0);
return(-1);
}
void
lin_regress(server_info_t info, struct timeval *ltv, struct timeval *lbtv,
double offset, int calc_offset_correction)
{
double time_axis;
double uncorrected_offset;
if (info->lin_count == 0) {
info->lin_tv = *ltv;
info->lin_btv = *lbtv;
time_axis = 0;
uncorrected_offset = offset;
} else {
time_axis = tv_delta_double(&info->lin_tv, ltv);
uncorrected_offset = offset - tv_delta_double(&info->lin_btv, lbtv);
}
++info->lin_count;
info->lin_sumx += time_axis;
info->lin_sumx2 += time_axis * time_axis;
info->lin_sumy += uncorrected_offset;
info->lin_sumy2 += uncorrected_offset * uncorrected_offset;
info->lin_sumxy += time_axis * uncorrected_offset;
if (calc_offset_correction) {
++info->lin_countoffset;
info->lin_sumoffset += offset;
info->lin_sumoffset2 += offset * offset;
}
if (info->lin_count > 1) {
info->lin_cache_slope =
(info->lin_count * info->lin_sumxy - info->lin_sumx * info->lin_sumy) /
(info->lin_count * info->lin_sumx2 - info->lin_sumx * info->lin_sumx);
info->lin_cache_yint =
(info->lin_sumy - info->lin_cache_slope * info->lin_sumx) /
(info->lin_count);
info->lin_cache_corr =
(info->lin_count * info->lin_sumxy - info->lin_sumx * info->lin_sumy) /
sqrt((info->lin_count * info->lin_sumx2 -
info->lin_sumx * info->lin_sumx) *
(info->lin_count * info->lin_sumy2 -
info->lin_sumy * info->lin_sumy)
);
}
if (info->lin_countoffset > 1) {
info->lin_cache_stddev =
sqrt((info->lin_sumoffset2 -
((info->lin_sumoffset * info->lin_sumoffset /
info->lin_countoffset))) /
(info->lin_countoffset - 1.0));
}
info->lin_cache_offset = offset;
info->lin_cache_freq = info->lin_cache_slope;
if (debug_level >= 4) {
logdebuginfo(info, 4, "iter=%2d time=%7.3f off=%+.6f uoff=%+.6f",
(int)info->lin_count,
time_axis, offset, uncorrected_offset);
if (info->lin_count > 1) {
logdebug(4, " slope %+7.6f"
" yint %+3.2f corr %+7.6f freq_ppm %+4.2f",
info->lin_cache_slope,
info->lin_cache_yint,
info->lin_cache_corr,
info->lin_cache_freq * 1000000.0);
}
if (info->lin_countoffset > 1) {
logdebug(4, " stddev %7.6f", info->lin_cache_stddev);
} else if (calc_offset_correction == 0) {
logdebug(4, " offset_ignored");
}
logdebug(4, "\n");
}
}
void
lin_reset(server_info_t info)
{
server_info_t scan;
info->lin_count = 0;
info->lin_sumx = 0;
info->lin_sumy = 0;
info->lin_sumxy = 0;
info->lin_sumx2 = 0;
info->lin_sumy2 = 0;
info->lin_countoffset = 0;
info->lin_sumoffset = 0;
info->lin_sumoffset2 = 0;
info->lin_cache_slope = 0;
info->lin_cache_yint = 0;
info->lin_cache_corr = 0;
info->lin_cache_offset = 0;
info->lin_cache_freq = 0;
while ((scan = info->altinfo) != NULL) {
info->altinfo = scan->altinfo;
free(scan);
}
}
void
lin_resetalloffsets(struct server_info **info_ary, int count)
{
server_info_t info;
int i;
for (i = 0; i < count; ++i) {
for (info = info_ary[i]; info; info = info->altinfo)
lin_resetoffsets(info);
}
}
void
lin_resetoffsets(server_info_t info)
{
info->lin_countoffset = 0;
info->lin_sumoffset = 0;
info->lin_sumoffset2 = 0;
}
void
client_setserverstate(server_info_t info, int state, const char *str)
{
if (info->server_state != state) {
info->server_state = state;
logdebuginfo(info, 1, "%s\n", str);
}
}