projects
/
fio.git
/ blobdiff
commit
grep
author
committer
pickaxe
?
search:
re
summary
|
shortlog
|
log
|
commit
|
commitdiff
|
tree
raw
|
inline
| side by side
Engines should not touch nr_open_files anymore
[fio.git]
/
fio.c
diff --git
a/fio.c
b/fio.c
index 0e1eadb35ce2d11555ee1a1ad3213083d136c5c5..43cc6af14e8ab632f731ae4106fce2b4e27ce5e7 100644
(file)
--- a/
fio.c
+++ b/
fio.c
@@
-60,17
+60,20
@@
static inline void td_set_runstate(struct thread_data *td, int runstate)
td->runstate = runstate;
}
td->runstate = runstate;
}
-static void terminate_threads(int group_id
, int forced_kill
)
+static void terminate_threads(int group_id)
{
struct thread_data *td;
int i;
for_each_td(td, i) {
if (group_id == TERMINATE_ALL || groupid == td->groupid) {
{
struct thread_data *td;
int i;
for_each_td(td, i) {
if (group_id == TERMINATE_ALL || groupid == td->groupid) {
+ /*
+ * if the thread is running, just let it exit
+ */
+ if (td->runstate < TD_RUNNING)
+ kill(td->pid, SIGQUIT);
td->terminate = 1;
td->start_delay = 0;
td->terminate = 1;
td->start_delay = 0;
- if (forced_kill)
- td_set_runstate(td, TD_EXITED);
}
}
}
}
}
}
@@
-86,7
+89,7
@@
static void sig_handler(int sig)
default:
printf("\nfio: terminating on signal %d\n", sig);
fflush(stdout);
default:
printf("\nfio: terminating on signal %d\n", sig);
fflush(stdout);
- terminate_threads(TERMINATE_ALL
, 0
);
+ terminate_threads(TERMINATE_ALL);
break;
}
}
break;
}
}
@@
-96,9
+99,15
@@
static void sig_handler(int sig)
*/
static int check_min_rate(struct thread_data *td, struct timeval *now)
{
*/
static int check_min_rate(struct thread_data *td, struct timeval *now)
{
+ unsigned long long bytes = 0;
unsigned long spent;
unsigned long rate;
unsigned long spent;
unsigned long rate;
- int ddir = td->ddir;
+
+ /*
+ * No minimum rate set, always ok
+ */
+ if (!td->ratemin)
+ return 0;
/*
* allow a 2 second settle period in the beginning
/*
* allow a 2 second settle period in the beginning
@@
-106,6
+115,11
@@
static int check_min_rate(struct thread_data *td, struct timeval *now)
if (mtime_since(&td->start, now) < 2000)
return 0;
if (mtime_since(&td->start, now) < 2000)
return 0;
+ if (td_read(td))
+ bytes += td->this_io_bytes[DDIR_READ];
+ if (td_write(td))
+ bytes += td->this_io_bytes[DDIR_WRITE];
+
/*
* if rate blocks is set, sample is running
*/
/*
* if rate blocks is set, sample is running
*/
@@
-114,14
+128,19
@@
static int check_min_rate(struct thread_data *td, struct timeval *now)
if (spent < td->ratecycle)
return 0;
if (spent < td->ratecycle)
return 0;
- rate = (td->this_io_bytes[ddir] - td->rate_bytes) / spent;
- if (rate < td->ratemin) {
- fprintf(f_out, "%s: min rate %u not met, got %luKiB/sec\n", td->name, td->ratemin, rate);
+ if (bytes < td->rate_bytes) {
+ fprintf(f_out, "%s: min rate %u not met\n", td->name, td->ratemin);
return 1;
return 1;
+ } else {
+ rate = (bytes - td->rate_bytes) / spent;
+ if (rate < td->ratemin || bytes < td->rate_bytes) {
+ fprintf(f_out, "%s: min rate %u not met, got %luKiB/sec\n", td->name, td->ratemin, rate);
+ return 1;
+ }
}
}
}
}
- td->rate_bytes =
td->this_io_bytes[ddir]
;
+ td->rate_bytes =
bytes
;
memcpy(&td->lastrate, now, sizeof(*now));
return 0;
}
memcpy(&td->lastrate, now, sizeof(*now));
return 0;
}
@@
-149,7
+168,7
@@
static void cleanup_pending_aio(struct thread_data *td)
/*
* get immediately available events, if any
*/
/*
* get immediately available events, if any
*/
- r = io_u_queued_complete(td, 0
, NULL
);
+ r = io_u_queued_complete(td, 0);
if (r < 0)
return;
if (r < 0)
return;
@@
-177,7
+196,7
@@
static void cleanup_pending_aio(struct thread_data *td)
}
if (td->cur_depth)
}
if (td->cur_depth)
- r = io_u_queued_complete(td, td->cur_depth
, NULL
);
+ r = io_u_queued_complete(td, td->cur_depth);
}
/*
}
/*
@@
-203,19
+222,19
@@
static int fio_io_sync(struct thread_data *td, struct fio_file *f)
requeue:
ret = td_io_queue(td, io_u);
if (ret < 0) {
requeue:
ret = td_io_queue(td, io_u);
if (ret < 0) {
- td_verror(td, io_u->error);
+ td_verror(td, io_u->error
, "td_io_queue"
);
put_io_u(td, io_u);
return 1;
} else if (ret == FIO_Q_QUEUED) {
put_io_u(td, io_u);
return 1;
} else if (ret == FIO_Q_QUEUED) {
- if (io_u_queued_complete(td, 1
, NULL
) < 0)
+ if (io_u_queued_complete(td, 1) < 0)
return 1;
} else if (ret == FIO_Q_COMPLETED) {
if (io_u->error) {
return 1;
} else if (ret == FIO_Q_COMPLETED) {
if (io_u->error) {
- td_verror(td, io_u->error);
+ td_verror(td, io_u->error
, "td_io_queue"
);
return 1;
}
return 1;
}
- if (io_u_sync_complete(td, io_u
, NULL
) < 0)
+ if (io_u_sync_complete(td, io_u) < 0)
return 1;
} else if (ret == FIO_Q_BUSY) {
if (td_io_commit(td))
return 1;
} else if (ret == FIO_Q_BUSY) {
if (td_io_commit(td))
@@
-227,7
+246,7
@@
requeue:
}
/*
}
/*
- * The main verify engine. Runs over the writes we previusly submitted,
+ * The main verify engine. Runs over the writes we previ
o
usly submitted,
* reads the blocks back in, and checks the crc/md5 of the data.
*/
static void do_verify(struct thread_data *td)
* reads the blocks back in, and checks the crc/md5 of the data.
*/
static void do_verify(struct thread_data *td)
@@
-254,6
+273,8
@@
static void do_verify(struct thread_data *td)
io_u = NULL;
while (!td->terminate) {
io_u = NULL;
while (!td->terminate) {
+ int ret2;
+
io_u = __get_io_u(td);
if (!io_u)
break;
io_u = __get_io_u(td);
if (!io_u)
break;
@@
-272,33
+293,37
@@
static void do_verify(struct thread_data *td)
put_io_u(td, io_u);
break;
}
put_io_u(td, io_u);
break;
}
-requeue:
- ret = td_io_queue(td, io_u);
+ io_u->end_io = verify_io_u;
+
+ ret = td_io_queue(td, io_u);
switch (ret) {
case FIO_Q_COMPLETED:
if (io_u->error)
ret = -io_u->error;
switch (ret) {
case FIO_Q_COMPLETED:
if (io_u->error)
ret = -io_u->error;
- if (io_u->xfer_buflen != io_u->resid && io_u->resid) {
+
else
if (io_u->xfer_buflen != io_u->resid && io_u->resid) {
int bytes = io_u->xfer_buflen - io_u->resid;
io_u->xfer_buflen = io_u->resid;
io_u->xfer_buf += bytes;
int bytes = io_u->xfer_buflen - io_u->resid;
io_u->xfer_buflen = io_u->resid;
io_u->xfer_buf += bytes;
- goto requeue;
+ requeue_io_u(td, &io_u);
+ } else {
+ ret = io_u_sync_complete(td, io_u);
+ if (ret < 0)
+ break;
}
}
- ret = io_u_sync_complete(td, io_u, verify_io_u);
- if (ret < 0)
- break;
continue;
case FIO_Q_QUEUED:
break;
case FIO_Q_BUSY:
requeue_io_u(td, &io_u);
continue;
case FIO_Q_QUEUED:
break;
case FIO_Q_BUSY:
requeue_io_u(td, &io_u);
- ret = td_io_commit(td);
+ ret2 = td_io_commit(td);
+ if (ret2 < 0)
+ ret = ret2;
break;
default:
assert(ret < 0);
break;
default:
assert(ret < 0);
- td_verror(td, -ret);
+ td_verror(td, -ret
, "td_io_queue"
);
break;
}
break;
}
@@
-321,11
+346,16
@@
requeue:
* Reap required number of io units, if any, and do the
* verification on them through the callback handler
*/
* Reap required number of io units, if any, and do the
* verification on them through the callback handler
*/
- if (io_u_queued_complete(td, min_events
, verify_io_u
) < 0)
+ if (io_u_queued_complete(td, min_events) < 0)
break;
}
break;
}
- if (td->cur_depth)
+ if (!td->error) {
+ min_events = td->cur_depth;
+
+ if (min_events)
+ ret = io_u_queued_complete(td, min_events);
+ } else
cleanup_pending_aio(td);
td_set_runstate(td, TD_RUNNING);
cleanup_pending_aio(td);
td_set_runstate(td, TD_RUNNING);
@@
-373,6
+403,7
@@
static void do_io(struct thread_data *td)
long bytes_done = 0;
int min_evts = 0;
struct io_u *io_u;
long bytes_done = 0;
int min_evts = 0;
struct io_u *io_u;
+ int ret2;
if (td->terminate)
break;
if (td->terminate)
break;
@@
-387,26
+418,24
@@
static void do_io(struct thread_data *td)
put_io_u(td, io_u);
break;
}
put_io_u(td, io_u);
break;
}
-requeue:
- ret = td_io_queue(td, io_u);
+ ret = td_io_queue(td, io_u);
switch (ret) {
case FIO_Q_COMPLETED:
switch (ret) {
case FIO_Q_COMPLETED:
- if (io_u->error) {
- ret = io_u->error;
- break;
- }
- if (io_u->xfer_buflen != io_u->resid && io_u->resid) {
+ if (io_u->error)
+ ret = -io_u->error;
+ else if (io_u->xfer_buflen != io_u->resid && io_u->resid) {
int bytes = io_u->xfer_buflen - io_u->resid;
io_u->xfer_buflen = io_u->resid;
io_u->xfer_buf += bytes;
int bytes = io_u->xfer_buflen - io_u->resid;
io_u->xfer_buflen = io_u->resid;
io_u->xfer_buf += bytes;
- goto requeue;
+ requeue_io_u(td, &io_u);
+ } else {
+ fio_gettime(&comp_time, NULL);
+ bytes_done = io_u_sync_complete(td, io_u);
+ if (bytes_done < 0)
+ ret = bytes_done;
}
}
- fio_gettime(&comp_time, NULL);
- bytes_done = io_u_sync_complete(td, io_u, NULL);
- if (bytes_done < 0)
- ret = bytes_done;
break;
case FIO_Q_QUEUED:
/*
break;
case FIO_Q_QUEUED:
/*
@@
-419,7
+448,9
@@
requeue:
break;
case FIO_Q_BUSY:
requeue_io_u(td, &io_u);
break;
case FIO_Q_BUSY:
requeue_io_u(td, &io_u);
- ret = td_io_commit(td);
+ ret2 = td_io_commit(td);
+ if (ret2 < 0)
+ ret = ret2;
break;
default:
assert(ret < 0);
break;
default:
assert(ret < 0);
@@
-443,7
+474,7
@@
requeue:
}
fio_gettime(&comp_time, NULL);
}
fio_gettime(&comp_time, NULL);
- bytes_done = io_u_queued_complete(td, min_evts
, NULL
);
+ bytes_done = io_u_queued_complete(td, min_evts);
if (bytes_done < 0)
break;
}
if (bytes_done < 0)
break;
}
@@
-458,12
+489,12
@@
requeue:
*/
usec = utime_since(&s, &comp_time);
*/
usec = utime_since(&s, &comp_time);
- rate_throttle(td, usec, bytes_done
, td->ddir
);
+ rate_throttle(td, usec, bytes_done);
if (check_min_rate(td, &comp_time)) {
if (exitall_on_terminate)
if (check_min_rate(td, &comp_time)) {
if (exitall_on_terminate)
- terminate_threads(td->groupid
, 0
);
- td_verror(td, ENODATA);
+ terminate_threads(td->groupid);
+ td_verror(td, ENODATA
, "check_min_rate"
);
break;
}
break;
}
@@
-487,15
+518,17
@@
requeue:
if (!td->error) {
struct fio_file *f;
if (!td->error) {
struct fio_file *f;
- if (td->cur_depth)
- cleanup_pending_aio(td);
+ i = td->cur_depth;
+ if (i)
+ ret = io_u_queued_complete(td, i);
if (should_fsync(td) && td->end_fsync) {
td_set_runstate(td, TD_FSYNCING);
for_each_file(td, f, i)
fio_io_sync(td, f);
}
if (should_fsync(td) && td->end_fsync) {
td_set_runstate(td, TD_FSYNCING);
for_each_file(td, f, i)
fio_io_sync(td, f);
}
- }
+ } else
+ cleanup_pending_aio(td);
}
static void cleanup_io_u(struct thread_data *td)
}
static void cleanup_io_u(struct thread_data *td)
@@
-585,7
+618,7
@@
static int switch_ioscheduler(struct thread_data *td)
f = fopen(tmp, "r+");
if (!f) {
f = fopen(tmp, "r+");
if (!f) {
- td_verror(td, errno);
+ td_verror(td, errno
, "fopen"
);
return 1;
}
return 1;
}
@@
-594,7
+627,7
@@
static int switch_ioscheduler(struct thread_data *td)
*/
ret = fwrite(td->ioscheduler, strlen(td->ioscheduler), 1, f);
if (ferror(f) || ret != 1) {
*/
ret = fwrite(td->ioscheduler, strlen(td->ioscheduler), 1, f);
if (ferror(f) || ret != 1) {
- td_verror(td, errno);
+ td_verror(td, errno
, "fwrite"
);
fclose(f);
return 1;
}
fclose(f);
return 1;
}
@@
-606,7
+639,7
@@
static int switch_ioscheduler(struct thread_data *td)
*/
ret = fread(tmp, 1, sizeof(tmp), f);
if (ferror(f) || ret < 0) {
*/
ret = fread(tmp, 1, sizeof(tmp), f);
if (ferror(f) || ret < 0) {
- td_verror(td, errno);
+ td_verror(td, errno
, "fread"
);
fclose(f);
return 1;
}
fclose(f);
return 1;
}
@@
-614,7
+647,7
@@
static int switch_ioscheduler(struct thread_data *td)
sprintf(tmp2, "[%s]", td->ioscheduler);
if (!strstr(tmp, tmp2)) {
log_err("fio: io scheduler %s not found\n", td->ioscheduler);
sprintf(tmp2, "[%s]", td->ioscheduler);
if (!strstr(tmp, tmp2)) {
log_err("fio: io scheduler %s not found\n", td->ioscheduler);
- td_verror(td, EINVAL);
+ td_verror(td, EINVAL
, "iosched_switch"
);
fclose(f);
return 1;
}
fclose(f);
return 1;
}
@@
-670,7
+703,7
@@
static void *thread_main(void *data)
goto err;
if (fio_setaffinity(td) == -1) {
goto err;
if (fio_setaffinity(td) == -1) {
- td_verror(td, errno);
+ td_verror(td, errno
, "cpu_set_affinity"
);
goto err;
}
goto err;
}
@@
-679,13
+712,13
@@
static void *thread_main(void *data)
if (td->ioprio) {
if (ioprio_set(IOPRIO_WHO_PROCESS, 0, td->ioprio) == -1) {
if (td->ioprio) {
if (ioprio_set(IOPRIO_WHO_PROCESS, 0, td->ioprio) == -1) {
- td_verror(td, errno);
+ td_verror(td, errno
, "ioprio_set"
);
goto err;
}
}
if (nice(td->nice) == -1) {
goto err;
}
}
if (nice(td->nice) == -1) {
- td_verror(td, errno);
+ td_verror(td, errno
, "nice"
);
goto err;
}
goto err;
}
@@
-701,16
+734,13
@@
static void *thread_main(void *data)
if (!td->create_serialize && setup_files(td))
goto err;
if (!td->create_serialize && setup_files(td))
goto err;
- if (open_files(td))
- goto err;
- /*
- * Do this late, as some IO engines would like to have the
- * files setup prior to initializing structures.
- */
if (td_io_init(td))
goto err;
if (td_io_init(td))
goto err;
+ if (open_files(td))
+ goto err;
+
if (td->exec_prerun) {
if (system(td->exec_prerun) < 0)
goto err;
if (td->exec_prerun) {
if (system(td->exec_prerun) < 0)
goto err;
@@
-736,10
+766,11
@@
static void *thread_main(void *data)
else
do_io(td);
else
do_io(td);
- runtime[td->ddir] += utime_since_now(&td->start);
- if (td_rw(td) && td->io_bytes[td->ddir ^ 1])
- runtime[td->ddir ^ 1] = runtime[td->ddir];
-
+ if (td_read(td) && td->io_bytes[DDIR_READ])
+ runtime[DDIR_READ] += utime_since_now(&td->start);
+ if (td_write(td) && td->io_bytes[DDIR_WRITE])
+ runtime[DDIR_WRITE] += utime_since_now(&td->start);
+
if (td->error || td->terminate)
break;
if (td->error || td->terminate)
break;
@@
-758,9
+789,11
@@
static void *thread_main(void *data)
}
update_rusage_stat(td);
}
update_rusage_stat(td);
- fio_gettime(&td->end_time, NULL);
- td->runtime[0] = runtime[0] / 1000;
- td->runtime[1] = runtime[1] / 1000;
+ td->ts.runtime[0] = runtime[0] / 1000;
+ td->ts.runtime[1] = runtime[1] / 1000;
+ td->ts.total_run_time = mtime_since_now(&td->epoch);
+ td->ts.io_bytes[0] = td->io_bytes[0];
+ td->ts.io_bytes[1] = td->io_bytes[1];
if (td->ts.bw_log)
finish_log(td, td->ts.bw_log, "bw");
if (td->ts.bw_log)
finish_log(td, td->ts.bw_log, "bw");
@@
-776,7
+809,7
@@
static void *thread_main(void *data)
}
if (exitall_on_terminate)
}
if (exitall_on_terminate)
- terminate_threads(td->groupid
, 0
);
+ terminate_threads(td->groupid);
err:
if (td->error)
err:
if (td->error)
@@
-862,7
+895,8
@@
static void reap_threads(int *nr_running, int *t_rate, int *m_rate)
if (WIFSIGNALED(status)) {
int sig = WTERMSIG(status);
if (WIFSIGNALED(status)) {
int sig = WTERMSIG(status);
- log_err("fio: pid=%d, got signal=%d\n", td->pid, sig);
+ if (sig != SIGQUIT)
+ log_err("fio: pid=%d, got signal=%d\n", td->pid, sig);
td_set_runstate(td, TD_REAPED);
goto reaped;
}
td_set_runstate(td, TD_REAPED);
goto reaped;
}
@@
-896,7
+930,7
@@
reaped:
}
if (*nr_running == cputhreads && !pending)
}
if (*nr_running == cputhreads && !pending)
- terminate_threads(TERMINATE_ALL
, 0
);
+ terminate_threads(TERMINATE_ALL);
}
/*
}
/*