From 5f262852b2bc4e7d16e5d5f5ede2c3bf1b451379 Mon Sep 17 00:00:00 2001 From: Jinshan Xiong Date: Fri, 23 Mar 2018 23:00:06 -0700 Subject: [PATCH 4/4] async IO --- lustre/include/cl_object.h | 1 + lustre/llite/llite_internal.h | 8 +-- lustre/llite/lloop.c | 101 +++++++++++++++++++++++-------------- lustre/llite/rw26.c | 114 +++++++++++++++++++----------------------- lustre/obdclass/cl_io.c | 4 +- 5 files changed, 123 insertions(+), 105 deletions(-) diff --git a/lustre/include/cl_object.h b/lustre/include/cl_object.h index e224c95..9cf4045 100644 --- a/lustre/include/cl_object.h +++ b/lustre/include/cl_object.h @@ -2584,6 +2584,7 @@ void cl_req_completion(const struct lu_env *env, struct cl_req *req, int ioret); * anchor and wakes up waiting thread when transfer is complete. */ struct cl_sync_io { + struct cl_page_list csi_page_list; /** number of pages yet to be transferred. */ atomic_t csi_sync_nr; /** error code. */ diff --git a/lustre/llite/llite_internal.h b/lustre/llite/llite_internal.h index 7ccc1aa..e32576f 100644 --- a/lustre/llite/llite_internal.h +++ b/lustre/llite/llite_internal.h @@ -1045,6 +1045,7 @@ struct vvp_thread_info { struct ra_io_arg vti_ria; struct kiocb vti_kiocb; struct ll_cl_context vti_io_ctx; + struct cl_sync_io vti_anchor; }; extern struct lu_context_key vvp_key; @@ -1410,9 +1411,10 @@ static inline void cl_stats_tally(struct cl_device *dev, enum cl_req_type crt, ll_stats_ops_tally(ll_s2sbi(cl2vvp_dev(dev)->vdv_sb), opc, rc); } -extern ssize_t ll_direct_rw_pages(const struct lu_env *env, struct cl_io *io, - int rw, struct inode *inode, - struct ll_dio_pages *pv); +extern int +ll_direct_rw_pages(const struct lu_env *env, struct cl_io *io, + int rw, struct inode *inode, struct ll_dio_pages *pv, + struct cl_sync_io *anchor); static inline int ll_file_nolock(const struct file *file) { diff --git a/lustre/llite/lloop.c b/lustre/llite/lloop.c index 83efcaf..2c7d058 100644 --- a/lustre/llite/lloop.c +++ b/lustre/llite/lloop.c @@ -165,6 +165,12 @@ enum { LO_FLAGS_READ_ONLY = 1, }; +struct lloop_bio_data { + struct lloop_device *lbd_dev; + struct bio *lbd_bio; + struct cl_sync_io lbd_anchor; +}; + static int lloop_major; #define MAX_LOOP_DEFAULT 16 static int max_loop = MAX_LOOP_DEFAULT; @@ -191,22 +197,29 @@ static inline sector_t bio_sector(struct bio *bio) #endif } -static loff_t get_loop_size(struct lloop_device *lo, struct file *file) +static void lloop_end_bio(const struct lu_env *env, struct cl_sync_io *anchor) { - loff_t size, offset, loopsize; + struct lloop_bio_data *data = container_of(anchor, typeof(*data), lbd_anchor); + struct bio *bio = data->lbd_bio; + int ret = data->lbd_anchor.csi_sync_rc; - /* Compute loopsize in bytes */ - size = i_size_read(file->f_mapping->host); - offset = lo->lo_offset; - loopsize = size - offset; - if (lo->lo_sizelimit > 0 && lo->lo_sizelimit < loopsize) - loopsize = lo->lo_sizelimit; + cl_page_list_discard(env, NULL, &anchor->csi_page_list); + cl_page_list_fini(env, &anchor->csi_page_list); - /* - * Unfortunately, if we want to do I/O on the device, - * the number of 512-byte sectors has to fit into a sector_t. - */ - return loopsize >> 9; + while (bio) { + struct bio *tmp = bio->bi_next; + + bio->bi_next = NULL; +#ifdef HAVE_BIO_ENDIO_USES_ONE_ARG + bio->bi_error = ret; + bio_endio(bio); +#else + bio_endio(bio, ret); +#endif + bio = tmp; + } + + OBD_FREE_PTR(data); } static int loop_bio_rw(struct bio *bio) @@ -227,19 +240,28 @@ static int do_bio_lustrebacked(struct lloop_device *lo, struct bio *head) int rw; size_t page_count = 0; struct bio *bio; - ssize_t bytes; struct ll_dio_pages *pvec = &lo->lo_pvec; + struct lloop_bio_data *data; vfs_fsync(lo->lo_backing_file, 0); truncate_inode_pages(inode->i_mapping, 0); + OBD_ALLOC_PTR(data); + if (!data) + return -ENOMEM; + + data->lbd_dev = lo; + data->lbd_bio = head; + cl_sync_io_init(&data->lbd_anchor, 1, lloop_end_bio); + /* initialize the IO */ memset(io, 0, sizeof(*io)); io->ci_obj = obj; ret = cl_io_init(env, io, CIT_MISC, obj); if (ret) - return io->ci_result; + GOTO(out, ret = io->ci_result); + io->ci_lockreq = CILR_NEVER; LASSERT(head != NULL); @@ -269,11 +291,15 @@ static int do_bio_lustrebacked(struct lloop_device *lo, struct bio *head) CDEBUG(D_INFO, "%s pvec pages: %zd\n", rw ? "Writing" : "Reading", page_count); mutex_lock(&inode->i_mutex); - bytes = ll_direct_rw_pages(env, io, rw, inode, pvec); + ret = ll_direct_rw_pages(env, io, rw, inode, pvec, &data->lbd_anchor); mutex_unlock(&inode->i_mutex); cl_io_fini(env, io); - CDEBUG(D_INFO, "I/O to pvec done, bytes: %zd\n", bytes); - return (bytes == pvec->ldp_size) ? 0 : (int)bytes; + + cl_sync_io_note(env, &data->lbd_anchor, ret); + +out: + CDEBUG(D_INFO, "I/O to pvec done, rc: %d\n", ret); + return ret; } /* @@ -381,25 +407,6 @@ static void loop_unplug(struct request_queue *q) } #endif -static inline void loop_handle_bio(struct lloop_device *lo, struct bio *bio) -{ - int ret; - - ret = do_bio_lustrebacked(lo, bio); - while (bio) { - struct bio *tmp = bio->bi_next; - - bio->bi_next = NULL; -#ifdef HAVE_BIO_ENDIO_USES_ONE_ARG - bio->bi_error = ret; - bio_endio(bio); -#else - bio_endio(bio, ret); -#endif - bio = tmp; - } -} - static inline int loop_active(struct lloop_device *lo) { return atomic_read(&lo->lo_pending) || lo->lo_state == LLOOP_RUNDOWN; @@ -468,7 +475,7 @@ static int loop_thread(void *data) LASSERT(bio != NULL); LASSERT(count <= atomic_read(&lo->lo_pending)); - loop_handle_bio(lo, bio); + (void) do_bio_lustrebacked(lo, bio); atomic_sub(count, &lo->lo_pending); if (need_resched()) @@ -481,6 +488,24 @@ out: return ret; } +static loff_t get_loop_size(struct lloop_device *lo, struct file *file) +{ + loff_t size, offset, loopsize; + + /* Compute loopsize in bytes */ + size = i_size_read(file->f_mapping->host); + offset = lo->lo_offset; + loopsize = size - offset; + if (lo->lo_sizelimit > 0 && lo->lo_sizelimit < loopsize) + loopsize = lo->lo_sizelimit; + + /* + * Unfortunately, if we want to do I/O on the device, + * the number of 512-byte sectors has to fit into a sector_t. + */ + return loopsize >> 9; +} + static int loop_set_fd(struct lloop_device *lo, struct file *unused, struct block_device *bdev, struct file *file) { diff --git a/lustre/llite/rw26.c b/lustre/llite/rw26.c index 428d72e..8aa1dda 100644 --- a/lustre/llite/rw26.c +++ b/lustre/llite/rw26.c @@ -200,9 +200,12 @@ static void ll_free_user_pages(struct page **pages, int npages, int do_dirty) OBD_FREE_LARGE(pages, npages * sizeof(*pages)); } -ssize_t ll_direct_rw_pages(const struct lu_env *env, struct cl_io *io, - int rw, struct inode *inode, - struct ll_dio_pages *pv) +/* + * if anchor is provided, it should have been initialized. + */ +int ll_direct_rw_pages(const struct lu_env *env, struct cl_io *io, int rw, + struct inode *inode, struct ll_dio_pages *pv, + struct cl_sync_io *anchor) { struct cl_page *clp; struct cl_2queue *queue; @@ -211,13 +214,19 @@ ssize_t ll_direct_rw_pages(const struct lu_env *env, struct cl_io *io, ssize_t rc = 0; loff_t file_offset = pv->ldp_start_offset; size_t size = pv->ldp_size; - int page_count = pv->ldp_nr; - struct page **pages = pv->ldp_pages; + size_t page_count = pv->ldp_nr; size_t page_size = cl_page_size(obj); - bool do_io; - int io_pages = 0; + bool sync_io = false; + unsigned io_pages = 0; ENTRY; + rw = rw == READ ? CRT_READ : CRT_WRITE; + if (!anchor) { + anchor = &vvp_env_info(env)->vti_anchor; + cl_sync_io_init(anchor, 0, &cl_sync_io_end); + sync_io = true; + } + queue = &io->ci_queue; cl_2queue_init(queue); for (i = 0; i < page_count; i++) { @@ -227,63 +236,25 @@ ssize_t ll_direct_rw_pages(const struct lu_env *env, struct cl_io *io, LASSERT(!(file_offset & (page_size - 1))); clp = cl_page_find(env, obj, cl_index(obj, file_offset), pv->ldp_pages[i], CPT_TRANSIENT); - if (IS_ERR(clp)) { - rc = PTR_ERR(clp); - break; - } + if (IS_ERR(clp)) + GOTO(out, rc = PTR_ERR(clp)); rc = cl_page_own(env, io, clp); if (rc) { LASSERT(clp->cp_state == CPS_FREEING); cl_page_put(env, clp); - break; - } - - do_io = true; - - /* check the page type: if the page is a host page, then do - * write directly */ - if (clp->cp_type == CPT_CACHEABLE) { - struct page *vmpage = cl_page_vmpage(clp); - struct page *src_page; - struct page *dst_page; - void *src; - void *dst; - - src_page = (rw == WRITE) ? pages[i] : vmpage; - dst_page = (rw == WRITE) ? vmpage : pages[i]; - - src = ll_kmap_atomic(src_page, KM_USER0); - dst = ll_kmap_atomic(dst_page, KM_USER1); - memcpy(dst, src, min(page_size, size)); - ll_kunmap_atomic(dst, KM_USER1); - ll_kunmap_atomic(src, KM_USER0); - - /* make sure page will be added to the transfer by - * cl_io_submit()->...->vvp_page_prep_write(). */ - if (rw == WRITE) - set_page_dirty(vmpage); - - if (rw == READ) { - /* do not issue the page for read, since it - * may reread a ra page which has NOT uptodate - * bit set. */ - cl_page_disown(env, io, clp); - do_io = false; - } + GOTO(out, rc); } - if (likely(do_io)) { - cl_2queue_add(queue, clp); + clp->cp_sync_io = anchor; + cl_2queue_add(queue, clp); - /* - * Set page clip to tell transfer formation engine - * that page has to be sent even if it is beyond KMS. - */ - cl_page_clip(env, clp, 0, min(size, page_size)); - - ++io_pages; - } + /* + * Set page clip to tell transfer formation engine + * that page has to be sent even if it is beyond KMS. + */ + cl_page_clip(env, clp, 0, min(size, page_size)); + ++io_pages; /* drop the reference count for cl_page_find */ cl_page_put(env, clp); @@ -291,14 +262,27 @@ ssize_t ll_direct_rw_pages(const struct lu_env *env, struct cl_io *io, file_offset += page_size; } - if (rc == 0 && io_pages) { - rc = cl_io_submit_sync(env, io, - rw == READ ? CRT_READ : CRT_WRITE, - queue, 0); + if (io_pages > 0) { + atomic_add(io_pages, &anchor->csi_sync_nr); + rc = cl_io_submit_rw(env, io, rw, queue); + if (!rc) { + cl_page_list_for_each(clp, &queue->c2_qin) { + clp->cp_sync_io = NULL; + cl_sync_io_note(env, anchor, 1); + } + + /* wait for the IO to be finished. */ + if (sync_io) { + rc = cl_sync_io_wait(env, anchor, 0); + cl_page_list_assume(env, io, &queue->c2_qout); + } else { + cl_page_list_splice(&queue->c2_qout, + &anchor->csi_page_list); + } + } } - if (rc == 0) - rc = pv->ldp_size; +out: cl_2queue_discard(env, io, queue); cl_2queue_disown(env, io, queue); cl_2queue_fini(env, queue); @@ -317,8 +301,12 @@ ll_direct_IO_seg(const struct lu_env *env, struct cl_io *io, int rw, .ldp_offsets = NULL, .ldp_start_offset = file_offset }; + ssize_t rc; - return ll_direct_rw_pages(env, io, rw, inode, &pvec); + rc = ll_direct_rw_pages(env, io, rw, inode, &pvec, NULL); + if (!rc) + rc = size; + return rc; } #ifdef KMALLOC_MAX_SIZE diff --git a/lustre/obdclass/cl_io.c b/lustre/obdclass/cl_io.c index 3cc7973..ba46096 100644 --- a/lustre/obdclass/cl_io.c +++ b/lustre/obdclass/cl_io.c @@ -887,7 +887,6 @@ void cl_page_list_del(const struct lu_env *env, struct cl_page_list *plist, struct cl_page *page) { LASSERT(plist->pl_nr > 0); - LASSERT(cl_page_is_vmlocked(env, page)); LINVRNT(plist->pl_owner == current); ENTRY; @@ -1045,6 +1044,7 @@ void cl_page_list_assume(const struct lu_env *env, cl_page_list_for_each(page, plist) cl_page_assume(env, io, page); } +EXPORT_SYMBOL(cl_page_list_assume); /** * Discards all pages in a queue. @@ -1060,6 +1060,7 @@ void cl_page_list_discard(const struct lu_env *env, struct cl_io *io, cl_page_discard(env, io, page); EXIT; } +EXPORT_SYMBOL(cl_page_list_discard); /** * Initialize dual page queue. @@ -1434,6 +1435,7 @@ void cl_sync_io_init(struct cl_sync_io *anchor, int nr, { ENTRY; memset(anchor, 0, sizeof(*anchor)); + cl_page_list_init(&anchor->csi_page_list); init_waitqueue_head(&anchor->csi_waitq); atomic_set(&anchor->csi_sync_nr, nr); atomic_set(&anchor->csi_barrier, nr > 0); -- 1.9.1