#include <sys/socket.h>
#include <sys/stat.h>
#include <sys/un.h>
+#include <sys/uio.h>
#include <netinet/in.h>
#include <arpa/inet.h>
#include <netdb.h>
#include <syslog.h>
#include <signal.h>
+#ifdef CONFIG_ZLIB
#include <zlib.h>
+#endif
#include "fio.h"
#include "server.h"
#include "crc/crc16.h"
#include "lib/ieee754.h"
-#include "fio_version.h"
-
int fio_net_port = FIO_NET_PORT;
int exit_backend = 0;
static struct sockaddr_in saddr_in;
static struct sockaddr_in6 saddr_in6;
static int use_ipv6;
+#ifdef CONFIG_ZLIB
+static unsigned int has_zlib = 1;
+#else
+static unsigned int has_zlib = 0;
+#endif
+static unsigned int use_zlib;
struct fio_fork_item {
struct flist_head list;
return cmdret;
}
+static void add_reply(uint64_t tag, struct flist_head *list)
+{
+ struct fio_net_cmd_reply *reply;
+
+ reply = (struct fio_net_cmd_reply *) (uintptr_t) tag;
+ flist_add_tail(&reply->list, list);
+}
+
+static uint64_t alloc_reply(uint64_t tag, uint16_t opcode)
+{
+ struct fio_net_cmd_reply *reply;
+
+ reply = calloc(1, sizeof(*reply));
+ INIT_FLIST_HEAD(&reply->list);
+ gettimeofday(&reply->tv, NULL);
+ reply->saved_tag = tag;
+ reply->opcode = opcode;
+
+ return (uintptr_t) reply;
+}
+
+static void free_reply(uint64_t tag)
+{
+ struct fio_net_cmd_reply *reply;
+
+ reply = (struct fio_net_cmd_reply *) (uintptr_t) tag;
+ free(reply);
+}
+
void fio_net_cmd_crc_pdu(struct fio_net_cmd *cmd, const void *pdu)
{
uint32_t pdu_len;
}
int fio_net_send_cmd(int fd, uint16_t opcode, const void *buf, off_t size,
- uint64_t tag)
+ uint64_t *tagptr, struct flist_head *list)
{
struct fio_net_cmd *cmd = NULL;
size_t this_len, cur_len = 0;
+ uint64_t tag;
int ret;
+ if (list) {
+ assert(tagptr);
+ tag = *tagptr = alloc_reply(*tagptr, opcode);
+ } else
+ tag = tagptr ? *tagptr : 0;
+
do {
this_len = size;
if (this_len > FIO_SERVER_MAX_FRAGMENT_PDU)
buf += this_len;
} while (!ret && size);
+ if (list) {
+ if (ret)
+ free_reply(tag);
+ else
+ add_reply(tag, list);
+ }
+
if (cmd)
free(cmd);
int fio_net_send_simple_cmd(int sk, uint16_t opcode, uint64_t tag,
struct flist_head *list)
{
- struct fio_net_int_cmd *cmd;
int ret;
- if (!list)
- return fio_net_send_simple_stack_cmd(sk, opcode, tag);
-
- cmd = malloc(sizeof(*cmd));
+ if (list)
+ tag = alloc_reply(tag, opcode);
- fio_init_net_cmd(&cmd->cmd, opcode, NULL, 0, (uintptr_t) cmd);
- fio_net_cmd_crc(&cmd->cmd);
-
- INIT_FLIST_HEAD(&cmd->list);
- gettimeofday(&cmd->tv, NULL);
- cmd->saved_tag = tag;
-
- ret = fio_send_data(sk, &cmd->cmd, sizeof(cmd->cmd));
+ ret = fio_net_send_simple_stack_cmd(sk, opcode, tag);
if (ret) {
- free(cmd);
+ if (list)
+ free_reply(tag);
+
return ret;
}
- flist_add_tail(&cmd->list, list);
+ if (list)
+ add_reply(tag, list);
+
return 0;
}
return fio_net_send_simple_cmd(sk, FIO_NET_CMD_QUIT, 0, NULL);
}
-int fio_net_send_stop(int sk, int error, int signal)
+static int fio_net_send_ack(int sk, struct fio_net_cmd *cmd, int error,
+ int signal)
{
struct cmd_end_pdu epdu;
+ uint64_t tag = 0;
- dprint(FD_NET, "server: sending stop (%d, %d)\n", error, signal);
+ if (cmd)
+ tag = cmd->tag;
epdu.error = __cpu_to_le32(error);
epdu.signal = __cpu_to_le32(signal);
- return fio_net_send_cmd(sk, FIO_NET_CMD_STOP, &epdu, sizeof(epdu), 0);
+ return fio_net_send_cmd(sk, FIO_NET_CMD_STOP, &epdu, sizeof(epdu), &tag, NULL);
+}
+
+int fio_net_send_stop(int sk, int error, int signal)
+{
+ dprint(FD_NET, "server: sending stop (%d, %d)\n", error, signal);
+ return fio_net_send_ack(sk, NULL, error, signal);
}
static void fio_server_add_fork_item(pid_t pid, struct flist_head *list)
}
ret = fio_backend();
+ free_threads_shm();
_exit(ret);
}
}
spdu.jobs = cpu_to_le32(thread_number);
- fio_net_send_cmd(server_fd, FIO_NET_CMD_START, &spdu, sizeof(spdu), 0);
+ spdu.stat_outputs = cpu_to_le32(stat_number);
+ fio_net_send_cmd(server_fd, FIO_NET_CMD_START, &spdu, sizeof(spdu), NULL, NULL);
return 0;
}
free(argv);
spdu.jobs = cpu_to_le32(thread_number);
- fio_net_send_cmd(server_fd, FIO_NET_CMD_START, &spdu, sizeof(spdu), 0);
+ spdu.stat_outputs = cpu_to_le32(stat_number);
+ fio_net_send_cmd(server_fd, FIO_NET_CMD_START, &spdu, sizeof(spdu), NULL, NULL);
return 0;
}
static int handle_probe_cmd(struct fio_net_cmd *cmd)
{
- struct cmd_probe_pdu probe;
+ struct cmd_client_probe_pdu *pdu = (struct cmd_client_probe_pdu *) cmd->payload;
+ struct cmd_probe_reply_pdu probe;
+ uint64_t tag = cmd->tag;
dprint(FD_NET, "server: sending probe reply\n");
memset(&probe, 0, sizeof(probe));
gethostname((char *) probe.hostname, sizeof(probe.hostname));
-#ifdef FIO_BIG_ENDIAN
+#ifdef CONFIG_BIG_ENDIAN
probe.bigendian = 1;
#endif
- probe.fio_major = FIO_MAJOR;
- probe.fio_minor = FIO_MINOR;
- probe.fio_patch = FIO_PATCH;
+ strncpy((char *) probe.fio_version, fio_version_string, sizeof(probe.fio_version));
probe.os = FIO_OS;
probe.arch = FIO_ARCH;
-
probe.bpp = sizeof(void *);
+ probe.cpus = __cpu_to_le32(cpus_online());
- return fio_net_send_cmd(server_fd, FIO_NET_CMD_PROBE, &probe, sizeof(probe), cmd->tag);
+ /*
+ * If the client supports compression and we do too, then enable it
+ */
+ if (has_zlib && le64_to_cpu(pdu->flags) & FIO_PROBE_FLAG_ZLIB) {
+ probe.flags = __cpu_to_le64(FIO_PROBE_FLAG_ZLIB);
+ use_zlib = 1;
+ } else {
+ probe.flags = 0;
+ use_zlib = 0;
+ }
+
+ return fio_net_send_cmd(server_fd, FIO_NET_CMD_PROBE, &probe, sizeof(probe), &tag, NULL);
}
static int handle_send_eta_cmd(struct fio_net_cmd *cmd)
{
struct jobs_eta *je;
size_t size;
+ uint64_t tag = cmd->tag;
int i;
if (!thread_number)
je->nr_running = cpu_to_le32(je->nr_running);
je->nr_ramp = cpu_to_le32(je->nr_ramp);
je->nr_pending = cpu_to_le32(je->nr_pending);
+ je->nr_setting_up = cpu_to_le32(je->nr_setting_up);
je->files_open = cpu_to_le32(je->files_open);
- for (i = 0; i < 2; i++) {
+ for (i = 0; i < DDIR_RWDIR_CNT; i++) {
je->m_rate[i] = cpu_to_le32(je->m_rate[i]);
je->t_rate[i] = cpu_to_le32(je->t_rate[i]);
je->m_iops[i] = cpu_to_le32(je->m_iops[i]);
je->elapsed_sec = cpu_to_le64(je->elapsed_sec);
je->eta_sec = cpu_to_le64(je->eta_sec);
je->nr_threads = cpu_to_le32(je->nr_threads);
+ je->is_pow2 = cpu_to_le32(je->is_pow2);
+ je->unit_base = cpu_to_le32(je->unit_base);
- fio_net_send_cmd(server_fd, FIO_NET_CMD_ETA, je, size, cmd->tag);
+ fio_net_send_cmd(server_fd, FIO_NET_CMD_ETA, je, size, &tag, NULL);
free(je);
return 0;
}
+static int send_update_job_reply(int fd, uint64_t __tag, int error)
+{
+ uint64_t tag = __tag;
+ uint32_t pdu_error;
+
+ pdu_error = __cpu_to_le32(error);
+ return fio_net_send_cmd(fd, FIO_NET_CMD_UPDATE_JOB, &pdu_error, sizeof(pdu_error), &tag, NULL);
+}
+
+static int handle_update_job_cmd(struct fio_net_cmd *cmd)
+{
+ struct cmd_add_job_pdu *pdu = (struct cmd_add_job_pdu *) cmd->payload;
+ struct thread_data *td;
+ uint32_t tnumber;
+
+ tnumber = le32_to_cpu(pdu->thread_number);
+
+ dprint(FD_NET, "server: updating options for job %u\n", tnumber);
+
+ if (!tnumber || tnumber > thread_number) {
+ send_update_job_reply(server_fd, cmd->tag, ENODEV);
+ return 0;
+ }
+
+ td = &threads[tnumber - 1];
+ convert_thread_options_to_cpu(&td->o, &pdu->top);
+ send_update_job_reply(server_fd, cmd->tag, 0);
+ return 0;
+}
+
static int handle_command(struct fio_net_cmd *cmd)
{
int ret;
- dprint(FD_NET, "server: got op [%s], pdu=%u, tag=%lx\n",
- fio_server_op(cmd->opcode), cmd->pdu_len, cmd->tag);
+ dprint(FD_NET, "server: got op [%s], pdu=%u, tag=%llx\n",
+ fio_server_op(cmd->opcode), cmd->pdu_len,
+ (unsigned long long) cmd->tag);
switch (cmd->opcode) {
case FIO_NET_CMD_QUIT:
case FIO_NET_CMD_RUN:
ret = handle_run_cmd(cmd);
break;
+ case FIO_NET_CMD_UPDATE_JOB:
+ ret = handle_update_job_cmd(cmd);
+ break;
default:
log_err("fio: unknown opcode: %s\n", fio_server_op(cmd->opcode));
ret = 1;
static int accept_loop(int listen_sk)
{
struct sockaddr_in addr;
- fio_socklen_t len = sizeof(addr);
+ socklen_t len = sizeof(addr);
struct pollfd pfd;
int ret = 0, sk, flags, exitval = 0;
memcpy(pdu->buf, buf, len);
- fio_net_send_cmd(server_fd, FIO_NET_CMD_TEXT, pdu, tlen, 0);
+ fio_net_send_cmd(server_fd, FIO_NET_CMD_TEXT, pdu, tlen, NULL, NULL);
free(pdu);
return len;
}
{
int i;
- for (i = 0; i < 2; i++) {
+ for (i = 0; i < DDIR_RWDIR_CNT; i++) {
dst->max_run[i] = cpu_to_le64(src->max_run[i]);
dst->min_run[i] = cpu_to_le64(src->min_run[i]);
dst->max_bw[i] = cpu_to_le64(src->max_bw[i]);
}
dst->kb_base = cpu_to_le32(src->kb_base);
+ dst->unit_base = cpu_to_le32(src->unit_base);
dst->groupid = cpu_to_le32(src->groupid);
+ dst->unified_rw_rep = cpu_to_le32(src->unified_rw_rep);
}
/*
p.ts.groupid = cpu_to_le32(ts->groupid);
p.ts.pid = cpu_to_le32(ts->pid);
p.ts.members = cpu_to_le32(ts->members);
+ p.ts.unified_rw_rep = cpu_to_le32(ts->unified_rw_rep);
- for (i = 0; i < 2; i++) {
+ for (i = 0; i < DDIR_RWDIR_CNT; i++) {
convert_io_stat(&p.ts.clat_stat[i], &ts->clat_stat[i]);
convert_io_stat(&p.ts.slat_stat[i], &ts->slat_stat[i]);
convert_io_stat(&p.ts.lat_stat[i], &ts->lat_stat[i]);
p.ts.io_u_lat_m[i] = cpu_to_le32(ts->io_u_lat_m[i]);
}
- for (i = 0; i < 2; i++)
+ for (i = 0; i < DDIR_RWDIR_CNT; i++)
for (j = 0; j < FIO_IO_U_PLAT_NR; j++)
p.ts.io_u_plat[i][j] = cpu_to_le32(ts->io_u_plat[i][j]);
- for (i = 0; i < 3; i++) {
+ for (i = 0; i < DDIR_RWDIR_CNT; i++) {
p.ts.total_io_u[i] = cpu_to_le64(ts->total_io_u[i]);
p.ts.short_io_u[i] = cpu_to_le64(ts->short_io_u[i]);
}
p.ts.total_submit = cpu_to_le64(ts->total_submit);
p.ts.total_complete = cpu_to_le64(ts->total_complete);
- for (i = 0; i < 2; i++) {
+ for (i = 0; i < DDIR_RWDIR_CNT; i++) {
p.ts.io_bytes[i] = cpu_to_le64(ts->io_bytes[i]);
p.ts.runtime[i] = cpu_to_le64(ts->runtime[i]);
}
p.ts.total_err_count = cpu_to_le64(ts->total_err_count);
p.ts.first_error = cpu_to_le32(ts->first_error);
p.ts.kb_base = cpu_to_le32(ts->kb_base);
+ p.ts.unit_base = cpu_to_le32(ts->unit_base);
convert_gs(&p.rs, rs);
- fio_net_send_cmd(server_fd, FIO_NET_CMD_TS, &p, sizeof(p), 0);
+ fio_net_send_cmd(server_fd, FIO_NET_CMD_TS, &p, sizeof(p), NULL, NULL);
}
void fio_server_send_gs(struct group_run_stats *rs)
dprint(FD_NET, "server sending group run stats\n");
convert_gs(&gs, rs);
- fio_net_send_cmd(server_fd, FIO_NET_CMD_GS, &gs, sizeof(gs), 0);
+ fio_net_send_cmd(server_fd, FIO_NET_CMD_GS, &gs, sizeof(gs), NULL, NULL);
}
static void convert_agg(struct disk_util_agg *dst, struct disk_util_agg *src)
convert_dus(&pdu.dus, &du->dus);
convert_agg(&pdu.agg, &du->agg);
- fio_net_send_cmd(server_fd, FIO_NET_CMD_DU, &pdu, sizeof(pdu), 0);
+ fio_net_send_cmd(server_fd, FIO_NET_CMD_DU, &pdu, sizeof(pdu), NULL, NULL);
}
}
return fio_sendv_data(sk, iov, 2);
}
-int fio_send_iolog(struct thread_data *td, struct io_log *log, const char *name)
+static int fio_send_iolog_gz(struct cmd_iolog_pdu *pdu, struct io_log *log)
{
- struct cmd_iolog_pdu pdu;
+ int ret = 0;
+#ifdef CONFIG_ZLIB
z_stream stream;
void *out_pdu;
- int i, ret = 0;
-
- pdu.thread_number = cpu_to_le32(td->thread_number);
- pdu.nr_samples = __cpu_to_le32(log->nr_samples);
- pdu.log_type = cpu_to_le32(log->log_type);
- strcpy((char *) pdu.name, name);
-
- for (i = 0; i < log->nr_samples; i++) {
- struct io_sample *s = &log->log[i];
-
- s->time = cpu_to_le64(s->time);
- s->val = cpu_to_le64(s->val);
- s->ddir = cpu_to_le32(s->ddir);
- s->bs = cpu_to_le32(s->bs);
- }
/*
* Dirty - since the log is potentially huge, compress it into
goto err;
}
- /*
- * Send header first, it's not compressed.
- */
- ret = fio_send_cmd_ext_pdu(server_fd, FIO_NET_CMD_IOLOG, &pdu,
- sizeof(pdu), 0, FIO_NET_CMD_F_MORE);
- if (ret)
- goto err_zlib;
-
stream.next_in = (void *) log->log;
stream.avail_in = log->nr_samples * sizeof(struct io_sample);
deflateEnd(&stream);
err:
free(out_pdu);
+#endif
return ret;
}
+int fio_send_iolog(struct thread_data *td, struct io_log *log, const char *name)
+{
+ struct cmd_iolog_pdu pdu;
+ int i, ret = 0;
+
+ pdu.thread_number = cpu_to_le32(td->thread_number);
+ pdu.nr_samples = __cpu_to_le32(log->nr_samples);
+ pdu.log_type = cpu_to_le32(log->log_type);
+ pdu.compressed = cpu_to_le32(use_zlib);
+ strcpy((char *) pdu.name, name);
+
+ for (i = 0; i < log->nr_samples; i++) {
+ struct io_sample *s = &log->log[i];
+
+ s->time = cpu_to_le64(s->time);
+ s->val = cpu_to_le64(s->val);
+ s->ddir = cpu_to_le32(s->ddir);
+ s->bs = cpu_to_le32(s->bs);
+ }
+
+ /*
+ * Send header first, it's not compressed.
+ */
+ ret = fio_send_cmd_ext_pdu(server_fd, FIO_NET_CMD_IOLOG, &pdu,
+ sizeof(pdu), 0, FIO_NET_CMD_F_MORE);
+ if (ret)
+ return ret;
+
+ /*
+ * Now send actual log, compress if we can, otherwise just plain
+ */
+ if (use_zlib)
+ return fio_send_iolog_gz(&pdu, log);
+
+ return fio_send_cmd_ext_pdu(server_fd, FIO_NET_CMD_IOLOG, log->log,
+ log->nr_samples * sizeof(struct io_sample), 0, 0);
+}
+
void fio_server_send_add_job(struct thread_data *td)
{
struct cmd_add_job_pdu pdu;
pdu.groupid = cpu_to_le32(td->groupid);
convert_thread_options_to_net(&pdu.top, &td->o);
- fio_net_send_cmd(server_fd, FIO_NET_CMD_ADD_JOB, &pdu, sizeof(pdu), 0);
+ fio_net_send_cmd(server_fd, FIO_NET_CMD_ADD_JOB, &pdu, sizeof(pdu), NULL, NULL);
}
void fio_server_send_start(struct thread_data *td)
static int fio_init_server_ip(void)
{
struct sockaddr *addr;
- fio_socklen_t socklen;
+ socklen_t socklen;
int sk, opt;
if (use_ipv6)
}
opt = 1;
- if (setsockopt(sk, SOL_SOCKET, SO_REUSEADDR, &opt, sizeof(opt)) < 0) {
+ if (setsockopt(sk, SOL_SOCKET, SO_REUSEADDR, (void *)&opt, sizeof(opt)) < 0) {
log_err("fio: setsockopt: %s\n", strerror(errno));
close(sk);
return -1;
static int fio_init_server_sock(void)
{
struct sockaddr_un addr;
- fio_socklen_t len;
+ socklen_t len;
mode_t mode;
int sk;
host++;
lport = atoi(host);
if (!lport || lport > 65535) {
- log_err("fio: bad server port %u\n", port);
+ log_err("fio: bad server port %u\n", lport);
return 1;
}
/* no hostname given, we are done */
portp++;
lport = atoi(portp);
if (!lport || lport > 65535) {
- log_err("fio: bad server port %u\n", port);
+ log_err("fio: bad server port %u\n", lport);
return 1;
}
}