Compare commits

...
8 changed files with 128 additions and 60 deletions
+1
View File
@@ -313,6 +313,7 @@ void osd_t::parse_config(bool init)
pg_reshard_chunk_pause_ms = config["pg_reshard_chunk_pause_ms"].uint64_value();
if (!pg_reshard_chunk_pause_ms)
pg_reshard_chunk_pause_ms = 100;
use_degraded_write = !json_is_false(config["use_degraded_write"]);
if (!old_auto_scrub && auto_scrub)
{
// Schedule scrubbing
+3 -1
View File
@@ -111,6 +111,7 @@ class osd_t
bool no_rebalance = false;
bool no_recovery = false;
bool no_scrub = false;
bool use_degraded_write = true;
bool allow_net_split = false;
std::vector<std::string> cfg_bind_addresses;
int bind_port, listen_backlog = 128;
@@ -343,6 +344,7 @@ class osd_t
void continue_primary_sync(osd_op_t *cur_op);
void continue_primary_del(osd_op_t *cur_op);
bool check_write_queue(osd_op_t *cur_op, pg_t & pg);
bool wr_degraded(osd_op_t *cur_op);
pg_osd_set_state_t* add_object_to_set(pg_t & pg, const object_id oid, const pg_osd_set_t & osd_set,
uint64_t old_pg_state, int log_at_level);
bool remove_object_from_state(object_id & oid, pg_osd_set_state_t **object_state, pg_t &pg, bool report = true);
@@ -352,7 +354,7 @@ class osd_t
osd_rmw_stripe_t *stripes, bool ref);
pg_osd_set_state_t *mark_partial_write(pg_t & pg, osd_op_t *cur_op);
void deref_object_state(pg_t & pg, pg_osd_set_state_t **object_state, bool deref);
bool remember_unstable_write(osd_op_t *cur_op, pg_t & pg, pg_osd_set_t & loc_set, int base_state);
bool remember_unstable_write(osd_op_t *cur_op, pg_t & pg, uint64_t *cur_set, int base_state);
void handle_primary_subop(osd_op_t *subop, osd_op_t *cur_op);
void handle_primary_bs_subop(osd_op_t *subop);
void add_bs_subop_stats(osd_op_t *subop, bool recovery_related = false);
+1
View File
@@ -294,6 +294,7 @@ void osd_t::submit_recovery_op(osd_recovery_op_t *op)
.inode = op->oid.inode,
.offset = op->oid.stripe,
.len = 0,
.flags = OSD_OP_RECOVERY_RELATED,
},
};
if (log_level > 2)
+9
View File
@@ -335,8 +335,10 @@ pg_osd_set_state_t* pg_t::add_object_to_state(const object_id oid, const uint64_
{
std::vector<osd_num_t> read_target;
bool found = false;
bool has_extra;
uint32_t bad_mask = (LOC_OUTDATED | LOC_CORRUPTED);
retry:
has_extra = false;
if (scheme == POOL_SCHEME_REPLICATED)
{
for (auto & o: osd_set)
@@ -352,6 +354,10 @@ retry:
// FIXME: This is because we then use .data() and assume it's at least <pg_size> long
read_target.resize(pg_size);
}
else if (read_target.size() > pg_size)
{
has_extra = true;
}
}
else
{
@@ -360,6 +366,8 @@ retry:
{
if (!(o.loc_bad & bad_mask))
{
if (read_target[o.role])
has_extra = true;
read_target[o.role] = o.osd_num;
found = true;
}
@@ -375,6 +383,7 @@ retry:
state_dict[osd_set] = {
.read_target = read_target,
.osd_set = osd_set,
.has_extra = has_extra,
.state = state,
.object_count = 1,
};
+1
View File
@@ -29,6 +29,7 @@ struct pg_osd_set_state_t
std::vector<osd_num_t> read_target;
// full OSD set including additional OSDs where the object is misplaced
pg_osd_set_t osd_set;
bool has_extra = false;
uint64_t state = 0;
uint64_t object_count = 0;
uint64_t ref_count = 0;
+1 -1
View File
@@ -490,7 +490,7 @@ void osd_t::handle_primary_subop(osd_op_t *subop, osd_op_t *cur_op)
}
if ((op_data->errors + op_data->done) >= op_data->n_subops)
{
if (!op_data->errors || !op_data->done || opcode != OSD_OP_SEC_WRITE && opcode != OSD_OP_SEC_WRITE_STABLE)
if (opcode != OSD_OP_SEC_WRITE && opcode != OSD_OP_SEC_WRITE_STABLE)
{
delete[] op_data->subops;
op_data->subops = NULL;
+103 -58
View File
@@ -37,6 +37,11 @@ bool osd_t::check_write_queue(osd_op_t *cur_op, pg_t & pg)
return true;
}
bool osd_t::wr_degraded(osd_op_t *cur_op)
{
return use_degraded_write && !(cur_op->req.rw.flags & OSD_OP_RECOVERY_RELATED);
}
void osd_t::continue_primary_write(osd_op_t *cur_op)
{
if (!cur_op->op_data && !prepare_primary_rw(cur_op))
@@ -62,10 +67,13 @@ void osd_t::continue_primary_write(osd_op_t *cur_op)
{
return;
}
if ((cur_op->req.rw.flags & OSD_OP_RECOVERY_RELATED) && cur_op->client_id != 0)
{
cur_op->req.rw.flags = 0;
}
resume_1:
// Determine blocks to read and write
// Missing chunks are allowed to be overwritten even in incomplete objects
// FIXME: Allow to do small writes to the old (degraded/misplaced) OSD set for lower performance impact
op_data->prev_set = get_object_osd_set(pg, op_data->oid, &op_data->object_state);
if (op_data->object_state)
{
@@ -88,18 +96,23 @@ retry_1:
cur_op->reply.hdr.retval = -EIO;
goto continue_others;
}
// Object is degraded/misplaced and will be moved to <write_osd_set>
op_data->stripes[0].read_start = 0;
op_data->stripes[0].read_end = bs_block_size;
assert(!cur_op->rmw_buf);
cur_op->rmw_buf = op_data->stripes[0].read_buf = memalign_or_die(MEM_ALIGNMENT, bs_block_size);
if (!wr_degraded(cur_op))
{
// Object is degraded/misplaced and will be moved to <write_osd_set>
op_data->stripes[0].read_start = 0;
op_data->stripes[0].read_end = bs_block_size;
assert(!cur_op->rmw_buf);
cur_op->rmw_buf = op_data->stripes[0].read_buf = memalign_or_die(MEM_ALIGNMENT, bs_block_size);
}
}
}
else
{
assert(!cur_op->rmw_buf);
cur_op->rmw_buf = calc_rmw(cur_op->buf, op_data->stripes, op_data->prev_set,
pg.pg_size, pg.pg_data_size, pg.pg_cursize, pg.cur_set.data(), bs_block_size, clean_entry_bitmap_size);
pg.pg_size, pg.pg_data_size, pg.pg_cursize,
wr_degraded(cur_op) ? op_data->prev_set : pg.cur_set.data(),
bs_block_size, clean_entry_bitmap_size);
if (!cur_op->rmw_buf)
{
// Refuse partial overwrite of an incomplete object
@@ -152,8 +165,7 @@ resume_3:
bitmap_set(op_data->stripes[0].bmp_buf, op_data->stripes[0].write_start,
op_data->stripes[0].write_end-op_data->stripes[0].write_start, bs_bitmap_granularity);
// Possibly copy new data from the request into the recovery buffer
if (pg.cur_set.data() != op_data->prev_set && (op_data->stripes[0].write_start != 0 ||
op_data->stripes[0].write_end != bs_block_size))
if (cur_op->rmw_buf)
{
memcpy(
(uint8_t*)op_data->stripes[0].read_buf + op_data->stripes[0].req_start,
@@ -173,11 +185,15 @@ resume_3:
// Recover missing stripes, calculate parity
if (pg.scheme == POOL_SCHEME_XOR)
{
calc_rmw_parity_xor(op_data->stripes, pg.pg_size, op_data->prev_set, pg.cur_set.data(), bs_block_size, clean_entry_bitmap_size);
calc_rmw_parity_xor(op_data->stripes, pg.pg_size, op_data->prev_set,
wr_degraded(cur_op) ? op_data->prev_set : pg.cur_set.data(),
bs_block_size, clean_entry_bitmap_size);
}
else if (pg.scheme == POOL_SCHEME_EC)
{
calc_rmw_parity_ec(op_data->stripes, pg.pg_size, pg.pg_data_size, op_data->prev_set, pg.cur_set.data(), bs_block_size, clean_entry_bitmap_size);
calc_rmw_parity_ec(op_data->stripes, pg.pg_size, pg.pg_data_size, op_data->prev_set,
wr_degraded(cur_op) ? op_data->prev_set : pg.cur_set.data(),
bs_block_size, clean_entry_bitmap_size);
}
}
// Send writes
@@ -204,9 +220,10 @@ resume_3:
{
// Check that current OSD set is in history and/or add it there
std::vector<osd_num_t> history_set;
for (auto peer_osd: pg.cur_set)
if (peer_osd != 0)
history_set.push_back(peer_osd);
auto new_set = wr_degraded(cur_op) ? op_data->prev_set : pg.cur_set.data();
for (int i = 0; i < pg.pg_size; i++)
if (new_set[i] != 0)
history_set.push_back(new_set[i]);
std::sort(history_set.begin(), history_set.end());
auto it = std::lower_bound(pg.target_history.begin(), pg.target_history.end(), history_set);
if (it == pg.target_history.end() || *it != history_set)
@@ -229,7 +246,8 @@ resume_10:
pg_cancel_write_queue(pg, cur_op, op_data->oid, -EPIPE);
return;
}
submit_primary_subops(SUBMIT_WRITE, op_data->target_ver, pg.cur_set.data(), cur_op);
submit_primary_subops(SUBMIT_WRITE, op_data->target_ver,
wr_degraded(cur_op) ? op_data->prev_set : pg.cur_set.data(), cur_op);
resume_4:
if (op_data->n_subops > 0)
{
@@ -248,7 +266,8 @@ resume_5:
{
if (pg.scheme != POOL_SCHEME_REPLICATED)
{
submit_primary_rollback_subops(cur_op, pg.cur_set.data());
submit_primary_rollback_subops(cur_op,
wr_degraded(cur_op) ? op_data->prev_set : pg.cur_set.data());
resume_11:
if (op_data->n_subops > 0)
{
@@ -264,6 +283,11 @@ resume_12:
{
pg.ver_override.erase(op_data->oid);
mark_partial_write(pg, cur_op);
if (op_data->subops)
{
delete[] op_data->subops;
op_data->subops = NULL;
}
pg_cancel_write_queue(pg, cur_op, op_data->oid, op_data->errcode);
return;
}
@@ -278,12 +302,24 @@ resume_12:
pg_cancel_write_queue(pg, cur_op, op_data->oid, op_data->errcode);
return;
}
if (op_data->object_state &&
op_data->object_state->has_extra &&
wr_degraded(cur_op))
{
// Mark possible extra copies which were not overwritten as outdated
mark_partial_write(pg, cur_op);
}
if (op_data->subops)
{
delete[] op_data->subops;
op_data->subops = NULL;
}
if (pg.scheme != POOL_SCHEME_REPLICATED)
{
// Remove version override just after the write, but before stabilizing
pg.ver_override.erase(op_data->oid);
}
if (op_data->object_state)
if (op_data->object_state && !wr_degraded(cur_op))
{
// Any kind of a non-clean object can have extra chunks, because we don't record objects
// as degraded & misplaced or incomplete & misplaced at the same time. So try to remove extra chunks
@@ -334,8 +370,7 @@ resume_12:
}
resume_6:
resume_7:
op_data->n_subops = 0;
if (!remember_unstable_write(cur_op, pg, pg.cur_loc_set, 6))
if (!remember_unstable_write(cur_op, pg, wr_degraded(cur_op) ? op_data->prev_set : pg.cur_set.data(), 6))
{
return;
}
@@ -345,7 +380,7 @@ resume_7:
pg.clean_count++;
pg.total_count++;
}
if (op_data->object_state)
if (op_data->object_state && !wr_degraded(cur_op))
{
{
int recovery_type = op_data->object_state->state & (OBJ_DEGRADED|OBJ_INCOMPLETE) ? 0 : 1;
@@ -370,16 +405,21 @@ resume_7:
);
recovery_stat[recovery_type].usec += usec;
}
if (immediate_commit == IMMEDIATE_ALL)
if (immediate_commit != IMMEDIATE_ALL)
{
// Do not delete anything, it will be deleted as part of copies_to_delete_after_sync
deref_object_state(pg, &op_data->object_state, true);
}
else
{
submit_primary_del_subops(cur_op, pg.cur_set.data(), pg.pg_size, op_data->object_state->osd_set);
}
deref_object_state(pg, &op_data->object_state, true);
if (op_data->n_subops > 0)
{
deref_object_state(pg, &op_data->object_state, true);
if (op_data->n_subops > 0)
{
resume_8:
op_data->st = 8;
return;
op_data->st = 8;
return;
}
resume_9:
if (op_data->errors > 0)
{
@@ -458,17 +498,13 @@ void osd_t::on_change_pg_history_hook(pool_id_t pool_id, pg_num_t pg_num)
}
}
bool osd_t::remember_unstable_write(osd_op_t *cur_op, pg_t & pg, pg_osd_set_t & loc_set, int base_state)
bool osd_t::remember_unstable_write(osd_op_t *cur_op, pg_t & pg, uint64_t *cur_set, int base_state)
{
osd_primary_op_data_t *op_data = cur_op->op_data;
if (op_data->st == base_state)
{
goto resume_6;
}
else if (op_data->st == base_state+1)
{
goto resume_7;
}
if (immediate_commit == IMMEDIATE_ALL)
{
immediate:
@@ -476,24 +512,27 @@ immediate:
{
// Send STABILIZE ops immediately
op_data->unstable_write_osds = new std::vector<unstable_osd_num_t>();
op_data->unstable_writes = new obj_ver_id[loc_set.size()];
op_data->unstable_writes = new obj_ver_id[pg.pg_size];
{
int last_start = 0;
for (auto & chunk: loc_set)
for (int role = 0; role < pg.pg_size; role++)
{
op_data->unstable_writes[last_start] = (obj_ver_id){
.oid = {
.inode = op_data->oid.inode,
.stripe = op_data->oid.stripe | chunk.role,
},
.version = op_data->fact_ver,
};
op_data->unstable_write_osds->push_back((unstable_osd_num_t){
.osd_num = chunk.osd_num,
.start = last_start,
.len = 1,
});
last_start++;
if (cur_set[role] != 0)
{
op_data->unstable_writes[last_start] = (obj_ver_id){
.oid = {
.inode = op_data->oid.inode,
.stripe = op_data->oid.stripe | role,
},
.version = op_data->fact_ver,
};
op_data->unstable_write_osds->push_back((unstable_osd_num_t){
.osd_num = cur_set[role],
.start = last_start,
.len = 1,
});
last_start++;
}
}
}
submit_primary_stab_subops(cur_op);
@@ -537,24 +576,30 @@ lazy:
if (pg.scheme != POOL_SCHEME_REPLICATED)
{
// Remember version as unstable for EC/XOR
for (auto & chunk: loc_set)
for (int role = 0; role < pg.pg_size; role++)
{
this->dirty_osds.insert(chunk.osd_num);
this->unstable_writes[(osd_object_id_t){
.osd_num = chunk.osd_num,
.oid = {
.inode = op_data->oid.inode,
.stripe = op_data->oid.stripe | chunk.role,
},
}] = op_data->fact_ver;
if (cur_set[role] != 0)
{
this->dirty_osds.insert(cur_set[role]);
this->unstable_writes[(osd_object_id_t){
.osd_num = cur_set[role],
.oid = {
.inode = op_data->oid.inode,
.stripe = op_data->oid.stripe | role,
},
}] = op_data->fact_ver;
}
}
}
else
{
// Only remember to sync OSDs for replicated pools
for (auto & chunk: loc_set)
for (int role = 0; role < pg.pg_size; role++)
{
this->dirty_osds.insert(chunk.osd_num);
if (cur_set[role] != 0)
{
this->dirty_osds.insert(cur_set[role]);
}
}
}
// Remember PG as dirty to drop the connection when PG goes offline
+9
View File
@@ -1,10 +1,19 @@
#!/bin/bash -ex
OSD_SIZE=200
if [[ $SCHEME = replicated ]]; then
OSD_COUNT=${OSD_COUNT:-2}
fi
GLOBAL_CONFIG=',"client_retry_enospc":false'
. `dirname $0`/run_3osds.sh
kill -9 $OSD1_PID
truncate -s 0 ./testdata/bin/test_osd1.bin
dd if=/dev/zero of=./testdata/bin/test_osd1.bin bs=1024 count=1 seek=$((300*1024-1))
$ETCDCTL del /vitastor/osd/state/1
start_osd 1
# Should fail with ENOSPC
if $VITASTOR_FIO -bs=1M -direct=1 -iodepth=4 \
-rw=write -pool=1 -inode=1 -size=500M -cluster_log_level=10; then