Fsync & update metadata when block checksums are enabled
This commit is contained in:
@@ -18,7 +18,6 @@ journal_flusher_t::journal_flusher_t(blockstore_impl_t *bs)
|
|||||||
this->cur_flusher_count = bs->min_flusher_count;
|
this->cur_flusher_count = bs->min_flusher_count;
|
||||||
this->target_flusher_count = bs->min_flusher_count;
|
this->target_flusher_count = bs->min_flusher_count;
|
||||||
active_flushers = 0;
|
active_flushers = 0;
|
||||||
syncing_flushers = 0;
|
|
||||||
advance_lsn_counter = 0;
|
advance_lsn_counter = 0;
|
||||||
co = new journal_flusher_co[max_flusher_count];
|
co = new journal_flusher_co[max_flusher_count];
|
||||||
for (int i = 0; i < max_flusher_count; i++)
|
for (int i = 0; i < max_flusher_count; i++)
|
||||||
@@ -169,6 +168,8 @@ bool journal_flusher_co::loop()
|
|||||||
else if (wait_state == 21) goto resume_21;
|
else if (wait_state == 21) goto resume_21;
|
||||||
else if (wait_state == 22) goto resume_22;
|
else if (wait_state == 22) goto resume_22;
|
||||||
else if (wait_state == 23) goto resume_23;
|
else if (wait_state == 23) goto resume_23;
|
||||||
|
else if (wait_state == 24) goto resume_24;
|
||||||
|
else if (wait_state == 25) goto resume_25;
|
||||||
resume_0:
|
resume_0:
|
||||||
wait_state = 0;
|
wait_state = 0;
|
||||||
cur_oid = {};
|
cur_oid = {};
|
||||||
@@ -252,6 +253,10 @@ resume_1:
|
|||||||
{
|
{
|
||||||
// Read original checksum blocks to calculate padded checksums if required
|
// Read original checksum blocks to calculate padded checksums if required
|
||||||
fill_partial_checksum_blocks();
|
fill_partial_checksum_blocks();
|
||||||
|
if (read_to_fill_incomplete)
|
||||||
|
{
|
||||||
|
flusher->wanting_meta_fsync++;
|
||||||
|
}
|
||||||
}
|
}
|
||||||
// Read buffered data
|
// Read buffered data
|
||||||
cur_obj = NULL;
|
cur_obj = NULL;
|
||||||
@@ -259,25 +264,35 @@ resume_1:
|
|||||||
resume_2:
|
resume_2:
|
||||||
resume_3:
|
resume_3:
|
||||||
if (!read_buffered(2))
|
if (!read_buffered(2))
|
||||||
|
{
|
||||||
return false;
|
return false;
|
||||||
|
}
|
||||||
// Now, if csum_block_size is > bitmap_granularity and if we are doing partial checksum block updates,
|
// Now, if csum_block_size is > bitmap_granularity and if we are doing partial checksum block updates,
|
||||||
// perform a trick: clear bitmap bits in the metadata entry and recalculate block checksum with zeros
|
// perform a trick: clear bitmap bits in the metadata entry and recalculate block checksum with zeros
|
||||||
// in place of overwritten parts. Then, even if the actual partial update fully or partially fails,
|
// in place of overwritten parts. Then, even if the actual partial update fully or partially fails,
|
||||||
// we'll have a correct checksum because it won't include overwritten parts!
|
// we'll have a correct checksum because it won't include overwritten parts!
|
||||||
// The same thing actually happens even when csum_block_size == bitmap_granularity, but in that case
|
// The same thing actually happens even when csum_block_size == bitmap_granularity, but in that case
|
||||||
// we never need to read (and thus verify) overwritten parts from the data device.
|
// we never need to read (and thus verify) overwritten parts from the data device.
|
||||||
|
if (read_to_fill_incomplete)
|
||||||
|
{
|
||||||
|
flusher->wanting_meta_fsync--;
|
||||||
|
}
|
||||||
res = check_and_punch_checksums();
|
res = check_and_punch_checksums();
|
||||||
if (res == EBUSY)
|
if (res == EBUSY)
|
||||||
{
|
{
|
||||||
resume_4:
|
resume_4:
|
||||||
resume_5:
|
resume_5:
|
||||||
if (!write_meta_block(4))
|
if (!write_meta_block(4))
|
||||||
|
{
|
||||||
return false;
|
return false;
|
||||||
|
}
|
||||||
resume_6:
|
resume_6:
|
||||||
resume_7:
|
resume_7:
|
||||||
resume_8:
|
resume_8:
|
||||||
if (!fsync_batch(true, 6)) // FIXME: is it correct to batch here
|
if (!fsync_meta(6))
|
||||||
|
{
|
||||||
return false;
|
return false;
|
||||||
|
}
|
||||||
}
|
}
|
||||||
else if (res == ENOENT || res == EDOM)
|
else if (res == ENOENT || res == EDOM)
|
||||||
{
|
{
|
||||||
@@ -311,13 +326,22 @@ resume_10:
|
|||||||
// Mark the object compacted, but don't free and remove small_writes
|
// Mark the object compacted, but don't free and remove small_writes
|
||||||
// We'll free and remove them only when trimming
|
// We'll free and remove them only when trimming
|
||||||
// The only thing we modify here are big_write block checksums if >4k block is used
|
// The only thing we modify here are big_write block checksums if >4k block is used
|
||||||
cur_obj = bs->heap->read_entry(cur_oid, NULL);
|
cur_obj = bs->heap->read_entry(cur_oid, &modified_block);
|
||||||
if (!cur_obj)
|
if (!cur_obj)
|
||||||
{
|
{
|
||||||
// Abort compaction
|
// Abort compaction
|
||||||
goto release_oid;
|
goto release_oid;
|
||||||
}
|
}
|
||||||
calc_block_checksums();
|
calc_block_checksums();
|
||||||
|
if (read_to_fill_incomplete)
|
||||||
|
{
|
||||||
|
resume_24:
|
||||||
|
resume_25:
|
||||||
|
if (!write_meta_block(24))
|
||||||
|
{
|
||||||
|
return false;
|
||||||
|
}
|
||||||
|
}
|
||||||
bs->heap->mark_object_compacted(cur_obj, compact_lsn);
|
bs->heap->mark_object_compacted(cur_obj, compact_lsn);
|
||||||
// Done, free all buffers
|
// Done, free all buffers
|
||||||
free_buffers();
|
free_buffers();
|
||||||
@@ -635,66 +659,37 @@ resume_1:
|
|||||||
return true;
|
return true;
|
||||||
}
|
}
|
||||||
|
|
||||||
bool journal_flusher_co::fsync_batch(bool fsync_meta, int wait_base)
|
bool journal_flusher_co::fsync_meta(int wait_base)
|
||||||
{
|
{
|
||||||
if (wait_state == wait_base) goto resume_0;
|
if (wait_state == wait_base) goto resume_0;
|
||||||
else if (wait_state == wait_base+1) goto resume_1;
|
else if (wait_state == wait_base+1) goto resume_1;
|
||||||
else if (wait_state == wait_base+2) goto resume_2;
|
else if (wait_state == wait_base+2) goto resume_2;
|
||||||
if (!(fsync_meta ? bs->dsk.disable_meta_fsync : bs->dsk.disable_data_fsync))
|
resume_0:
|
||||||
|
if (bs->dsk.disable_meta_fsync)
|
||||||
{
|
{
|
||||||
cur_sync = flusher->syncs.end();
|
return true;
|
||||||
while (cur_sync != flusher->syncs.begin())
|
|
||||||
{
|
|
||||||
cur_sync--;
|
|
||||||
if (cur_sync->fsync_meta == fsync_meta && cur_sync->state == 0)
|
|
||||||
{
|
|
||||||
goto sync_found;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
cur_sync = flusher->syncs.emplace(flusher->syncs.end(), (flusher_sync_t){
|
|
||||||
.fsync_meta = fsync_meta,
|
|
||||||
.ready_count = 0,
|
|
||||||
.state = 0,
|
|
||||||
});
|
|
||||||
sync_found:
|
|
||||||
cur_sync->ready_count++;
|
|
||||||
flusher->syncing_flushers++;
|
|
||||||
resume_1:
|
|
||||||
if (!cur_sync->state)
|
|
||||||
{
|
|
||||||
if (flusher->syncing_flushers >= flusher->active_flushers || true /*FIXME*/)
|
|
||||||
{
|
|
||||||
// Sync batch is ready. Do it.
|
|
||||||
await_sqe(0);
|
|
||||||
data->iov = { 0 };
|
|
||||||
data->callback = simple_callback_w;
|
|
||||||
io_uring_prep_fsync(sqe, fsync_meta ? bs->dsk.meta_fd : bs->dsk.data_fd, IORING_FSYNC_DATASYNC);
|
|
||||||
cur_sync->state = 1;
|
|
||||||
wait_count++;
|
|
||||||
resume_2:
|
|
||||||
if (wait_count > 0)
|
|
||||||
{
|
|
||||||
wait_state = wait_base+2;
|
|
||||||
return false;
|
|
||||||
}
|
|
||||||
// Sync completed. All previous coroutines waiting for it must be resumed
|
|
||||||
cur_sync->state = 2;
|
|
||||||
bs->ringloop->wakeup();
|
|
||||||
}
|
|
||||||
else
|
|
||||||
{
|
|
||||||
// Wait until someone else sends and completes a sync.
|
|
||||||
wait_state = wait_base+1;
|
|
||||||
return false;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
flusher->syncing_flushers--;
|
|
||||||
cur_sync->ready_count--;
|
|
||||||
if (cur_sync->ready_count == 0)
|
|
||||||
{
|
|
||||||
flusher->syncs.erase(cur_sync);
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
if (flusher->wanting_meta_fsync || flusher->fsyncing_meta > 0)
|
||||||
|
{
|
||||||
|
wait_state = wait_base;
|
||||||
|
return false;
|
||||||
|
}
|
||||||
|
flusher->fsyncing_meta = true;
|
||||||
|
// Sync batch is ready. Do it.
|
||||||
|
await_sqe(1);
|
||||||
|
data->iov = { 0 };
|
||||||
|
data->callback = simple_callback_w;
|
||||||
|
io_uring_prep_fsync(sqe, bs->dsk.meta_fd, IORING_FSYNC_DATASYNC);
|
||||||
|
wait_count++;
|
||||||
|
resume_2:
|
||||||
|
if (wait_count > 0)
|
||||||
|
{
|
||||||
|
wait_state = wait_base+2;
|
||||||
|
return false;
|
||||||
|
}
|
||||||
|
// Sync completed. All previous coroutines waiting for it must be resumed
|
||||||
|
flusher->fsyncing_meta = false;
|
||||||
|
bs->ringloop->wakeup();
|
||||||
return true;
|
return true;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -17,13 +17,6 @@ struct meta_sector_t
|
|||||||
int usage_count;
|
int usage_count;
|
||||||
};
|
};
|
||||||
|
|
||||||
struct flusher_sync_t
|
|
||||||
{
|
|
||||||
bool fsync_meta;
|
|
||||||
int ready_count;
|
|
||||||
int state;
|
|
||||||
};
|
|
||||||
|
|
||||||
struct flusher_meta_write_t
|
struct flusher_meta_write_t
|
||||||
{
|
{
|
||||||
uint64_t sector, pos;
|
uint64_t sector, pos;
|
||||||
@@ -44,8 +37,6 @@ class journal_flusher_co
|
|||||||
struct io_uring_sqe *sqe;
|
struct io_uring_sqe *sqe;
|
||||||
struct ring_data_t *data;
|
struct ring_data_t *data;
|
||||||
|
|
||||||
std::list<flusher_sync_t>::iterator cur_sync;
|
|
||||||
std::map<object_id, uint64_t>::iterator repeat_it;
|
|
||||||
std::function<void(ring_data_t*)> simple_callback_r, simple_callback_w;
|
std::function<void(ring_data_t*)> simple_callback_r, simple_callback_w;
|
||||||
|
|
||||||
object_id cur_oid;
|
object_id cur_oid;
|
||||||
@@ -78,7 +69,7 @@ class journal_flusher_co
|
|||||||
void calc_block_checksums();
|
void calc_block_checksums();
|
||||||
bool write_meta_block(int wait_base);
|
bool write_meta_block(int wait_base);
|
||||||
bool read_buffered(int wait_base);
|
bool read_buffered(int wait_base);
|
||||||
bool fsync_batch(bool fsync_meta, int wait_base);
|
bool fsync_meta(int wait_base);
|
||||||
int fsync_buffer(int wait_base);
|
int fsync_buffer(int wait_base);
|
||||||
bool trim_lsn(int wait_base);
|
bool trim_lsn(int wait_base);
|
||||||
public:
|
public:
|
||||||
@@ -100,9 +91,9 @@ class journal_flusher_t
|
|||||||
uint64_t compact_counter = 0;
|
uint64_t compact_counter = 0;
|
||||||
|
|
||||||
int active_flushers = 0;
|
int active_flushers = 0;
|
||||||
int syncing_flushers = 0;
|
int wanting_meta_fsync = 0;
|
||||||
|
bool fsyncing_meta = false;
|
||||||
int syncing_buffer = 0;
|
int syncing_buffer = 0;
|
||||||
std::list<flusher_sync_t> syncs;
|
|
||||||
|
|
||||||
public:
|
public:
|
||||||
journal_flusher_t(blockstore_impl_t *bs);
|
journal_flusher_t(blockstore_impl_t *bs);
|
||||||
|
|||||||
@@ -390,7 +390,6 @@ skip_object:
|
|||||||
{
|
{
|
||||||
if (wr->is_compacted(this->compacted_lsn))
|
if (wr->is_compacted(this->compacted_lsn))
|
||||||
{
|
{
|
||||||
// FIXME block checksums should be modified in flusher in this case
|
|
||||||
to_compact = true;
|
to_compact = true;
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -269,7 +269,7 @@ void multilist_alloc_t::do_free(uint32_t pos)
|
|||||||
sizes[pos+size-1] = -size;
|
sizes[pos+size-1] = -size;
|
||||||
sizes[pos] = size;
|
sizes[pos] = size;
|
||||||
}
|
}
|
||||||
uint32_t ni = (size < maxn ? size : maxn)-1; // FIXME ni -> nb (next bucket)
|
uint32_t ni = (size < maxn ? size : maxn)-1;
|
||||||
nexts[pos] = heads[ni]+1;
|
nexts[pos] = heads[ni]+1;
|
||||||
prevs[pos] = 0;
|
prevs[pos] = 0;
|
||||||
if (heads[ni])
|
if (heads[ni])
|
||||||
|
|||||||
Reference in New Issue
Block a user