(Almost) Implement misplaced recovery, integrating it into calc_rmw()
This commit is contained in:
+127
-66
@@ -1,4 +1,5 @@
|
||||
#include <malloc.h>
|
||||
#include <string.h>
|
||||
#include <assert.h>
|
||||
#include "xor.h"
|
||||
#include "osd_rmw.h"
|
||||
@@ -79,18 +80,21 @@ void reconstruct_stripe(osd_rmw_stripe_t *stripes, int pg_size, int role)
|
||||
}
|
||||
else if (prev >= 0)
|
||||
{
|
||||
assert(stripes[role].read_start >= stripes[prev].read_start &&
|
||||
stripes[role].read_start >= stripes[other].read_start);
|
||||
memxor(
|
||||
stripes[prev].read_buf + (stripes[prev].read_start - stripes[role].read_start),
|
||||
stripes[other].read_buf + (stripes[other].read_start - stripes[other].read_start),
|
||||
stripes[prev].read_buf + (stripes[role].read_start - stripes[prev].read_start),
|
||||
stripes[other].read_buf + (stripes[role].read_start - stripes[other].read_start),
|
||||
stripes[role].read_buf, stripes[role].read_end - stripes[role].read_start
|
||||
);
|
||||
prev = -1;
|
||||
}
|
||||
else
|
||||
{
|
||||
assert(stripes[role].read_start >= stripes[other].read_start);
|
||||
memxor(
|
||||
stripes[role].read_buf,
|
||||
stripes[other].read_buf + (stripes[other].read_start - stripes[role].read_start),
|
||||
stripes[other].read_buf + (stripes[role].read_start - stripes[other].read_start),
|
||||
stripes[role].read_buf, stripes[role].read_end - stripes[role].read_start
|
||||
);
|
||||
}
|
||||
@@ -156,10 +160,11 @@ void* alloc_read_buffer(osd_rmw_stripe_t *stripes, int read_pg_size, uint64_t ad
|
||||
return buf;
|
||||
}
|
||||
|
||||
void* calc_rmw_reads(void *write_buf, osd_rmw_stripe_t *stripes, uint64_t *osd_set, uint64_t pg_size, uint64_t pg_minsize, uint64_t pg_cursize)
|
||||
void* calc_rmw(void *request_buf, osd_rmw_stripe_t *stripes, uint64_t *read_osd_set,
|
||||
uint64_t pg_size, uint64_t pg_minsize, uint64_t pg_cursize, uint64_t *write_osd_set, uint64_t chunk_size)
|
||||
{
|
||||
// Generic parity modification (read-modify-write) algorithm
|
||||
// Reconstruct -> Read -> Calc parity -> Write
|
||||
// Read -> Reconstruct missing chunks -> Calc parity chunks -> Write
|
||||
// Now we always read continuous ranges. This means that an update of the beginning
|
||||
// of one data stripe and the end of another will lead to a read of full paired stripes.
|
||||
// FIXME: (Maybe) read small individual ranges in that case instead.
|
||||
@@ -174,64 +179,79 @@ void* calc_rmw_reads(void *write_buf, osd_rmw_stripe_t *stripes, uint64_t *osd_s
|
||||
stripes[role].write_end = stripes[role].req_end;
|
||||
}
|
||||
}
|
||||
for (int role = 0; role < pg_minsize; role++)
|
||||
{
|
||||
cover_read(start, end, stripes[role]);
|
||||
}
|
||||
int has_parity = 0;
|
||||
int write_parity = 0;
|
||||
for (int role = pg_minsize; role < pg_size; role++)
|
||||
{
|
||||
if (osd_set[role] != 0)
|
||||
if (write_osd_set[role] != 0)
|
||||
{
|
||||
has_parity++;
|
||||
write_parity = 1;
|
||||
stripes[role].write_start = start;
|
||||
stripes[role].write_end = end;
|
||||
}
|
||||
else
|
||||
stripes[role].missing = true;
|
||||
}
|
||||
if (write_parity)
|
||||
{
|
||||
for (int role = 0; role < pg_minsize; role++)
|
||||
{
|
||||
cover_read(start, end, stripes[role]);
|
||||
}
|
||||
}
|
||||
if (write_osd_set != read_osd_set)
|
||||
{
|
||||
// Object is degraded/misplaced and will be moved to <write_osd_set>
|
||||
for (int role = 0; role < pg_size; role++)
|
||||
{
|
||||
if (write_osd_set[role] != read_osd_set[role])
|
||||
{
|
||||
// We need to get data for any moved / recovered chunk
|
||||
// And we need a continuous write buffer so we'll only optimize
|
||||
// for the case when the whole chunk is ovewritten in the request
|
||||
if (stripes[role].req_start != 0 ||
|
||||
stripes[role].req_end != chunk_size)
|
||||
{
|
||||
stripes[role].read_start = 0;
|
||||
stripes[role].read_end = chunk_size;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
if (pg_cursize < pg_size)
|
||||
{
|
||||
if (has_parity == 0)
|
||||
// Some stripe(s) are missing, so we need to read parity
|
||||
for (int role = 0; role < pg_size; role++)
|
||||
{
|
||||
// Parity is missing, we don't need to read anything
|
||||
for (int role = 0; role < pg_minsize; role++)
|
||||
if (read_osd_set[role] == 0 && stripes[role].read_end != 0)
|
||||
{
|
||||
stripes[role].read_end = 0;
|
||||
}
|
||||
}
|
||||
else
|
||||
{
|
||||
// Other stripe(s) are missing
|
||||
for (int role = 0; role < pg_minsize; role++)
|
||||
{
|
||||
if (osd_set[role] == 0 && stripes[role].read_end != 0)
|
||||
stripes[role].missing = true;
|
||||
int found = 0;
|
||||
for (int r2 = 0; r2 < pg_size && found < pg_minsize; r2++)
|
||||
{
|
||||
stripes[role].missing = true;
|
||||
for (int r2 = 0; r2 < pg_size; r2++)
|
||||
// Read the non-covered range of <role> from at least <minsize> other stripes to reconstruct it
|
||||
if (read_osd_set[r2] != 0)
|
||||
{
|
||||
// Read the non-covered range of <role> from all other stripes to reconstruct it
|
||||
if (r2 != role && osd_set[r2] != 0)
|
||||
{
|
||||
extend_read(stripes[role].read_start, stripes[role].read_end, stripes[r2]);
|
||||
}
|
||||
extend_read(stripes[role].read_start, stripes[role].read_end, stripes[r2]);
|
||||
found++;
|
||||
}
|
||||
}
|
||||
if (found < pg_minsize)
|
||||
{
|
||||
// Incomplete object (FIXME)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
// Allocate read buffers
|
||||
void *rmw_buf = alloc_read_buffer(stripes, pg_size, has_parity * (end - start));
|
||||
// Position parity & write buffers
|
||||
void *rmw_buf = alloc_read_buffer(stripes, pg_size, (write_parity ? pg_size-pg_minsize : 0) * (end - start));
|
||||
// Position write buffers
|
||||
uint64_t buf_pos = 0, in_pos = 0;
|
||||
for (int role = 0; role < pg_size; role++)
|
||||
{
|
||||
if (stripes[role].req_end != 0)
|
||||
{
|
||||
stripes[role].write_buf = write_buf + in_pos;
|
||||
stripes[role].write_buf = request_buf + in_pos;
|
||||
in_pos += stripes[role].req_end - stripes[role].req_start;
|
||||
}
|
||||
else if (role >= pg_minsize && osd_set[role] != 0)
|
||||
else if (role >= pg_minsize && read_osd_set[role] != 0)
|
||||
{
|
||||
stripes[role].write_buf = rmw_buf + buf_pos;
|
||||
buf_pos += end - start;
|
||||
@@ -321,8 +341,9 @@ static void xor_multiple_buffers(buf_len_t *xor1, int n1, buf_len_t *xor2, int n
|
||||
}
|
||||
}
|
||||
|
||||
void reconstruct_stripes(osd_rmw_stripe_t *stripes, int pg_size)
|
||||
void calc_rmw_parity(osd_rmw_stripe_t *stripes, int pg_size, uint64_t *read_osd_set, uint64_t *write_osd_set, uint32_t chunk_size)
|
||||
{
|
||||
int pg_minsize = pg_size-1;
|
||||
for (int role = 0; role < pg_size; role++)
|
||||
{
|
||||
if (stripes[role].read_end != 0 && stripes[role].missing)
|
||||
@@ -332,41 +353,81 @@ void reconstruct_stripes(osd_rmw_stripe_t *stripes, int pg_size)
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
void calc_rmw_parity(osd_rmw_stripe_t *stripes, int pg_size)
|
||||
{
|
||||
if (stripes[pg_size-1].missing)
|
||||
uint32_t start = 0, end = 0;
|
||||
if (!stripes[pg_minsize].missing || write_osd_set != read_osd_set)
|
||||
{
|
||||
// Parity OSD is unavailable
|
||||
return;
|
||||
}
|
||||
reconstruct_stripes(stripes, pg_size);
|
||||
// Calculate new parity (EC k+1)
|
||||
int parity = pg_size-1, prev = -2;
|
||||
auto wr_end = stripes[parity].write_end;
|
||||
auto wr_start = stripes[parity].write_start;
|
||||
for (int other = 0; other < pg_size-1; other++)
|
||||
{
|
||||
if (prev == -2)
|
||||
for (int role = 0; role < pg_minsize; role++)
|
||||
{
|
||||
prev = other;
|
||||
}
|
||||
else
|
||||
{
|
||||
int n1 = 0, n2 = 0;
|
||||
buf_len_t xor1[3], xor2[3];
|
||||
if (prev == -1)
|
||||
if (stripes[role].req_end != 0)
|
||||
{
|
||||
xor1[n1++] = { .buf = stripes[parity].write_buf, .len = wr_end-wr_start };
|
||||
start = !end || stripes[role].req_start < start ? stripes[role].req_start : start;
|
||||
end = std::max(stripes[role].req_end, end);
|
||||
}
|
||||
}
|
||||
}
|
||||
if (write_osd_set != read_osd_set)
|
||||
{
|
||||
for (int role = 0; role < pg_minsize; role++)
|
||||
{
|
||||
if (write_osd_set[role] != read_osd_set[role] &&
|
||||
(stripes[role].req_start != 0 || stripes[role].req_end != chunk_size))
|
||||
{
|
||||
// Copy modified chunk into the read buffer to write it back
|
||||
memcpy(
|
||||
stripes[role].read_buf + stripes[role].req_start,
|
||||
stripes[role].write_buf,
|
||||
stripes[role].req_end - stripes[role].req_start
|
||||
);
|
||||
stripes[role].write_buf = stripes[role].read_buf;
|
||||
stripes[role].write_start = 0;
|
||||
stripes[role].write_end = chunk_size;
|
||||
}
|
||||
}
|
||||
}
|
||||
if (!stripes[pg_minsize].missing)
|
||||
{
|
||||
// Calculate new parity (EC k+1)
|
||||
int parity = pg_minsize, prev = -2;
|
||||
for (int other = 0; other < pg_minsize; other++)
|
||||
{
|
||||
if (prev == -2)
|
||||
{
|
||||
prev = other;
|
||||
}
|
||||
else
|
||||
{
|
||||
get_old_new_buffers(stripes[prev], wr_start, wr_end, xor1, n1);
|
||||
prev = -1;
|
||||
int n1 = 0, n2 = 0;
|
||||
buf_len_t xor1[3], xor2[3];
|
||||
if (prev == -1)
|
||||
{
|
||||
xor1[n1++] = { .buf = stripes[parity].write_buf, .len = end-start };
|
||||
}
|
||||
else
|
||||
{
|
||||
get_old_new_buffers(stripes[prev], start, end, xor1, n1);
|
||||
prev = -1;
|
||||
}
|
||||
get_old_new_buffers(stripes[other], start, end, xor2, n2);
|
||||
xor_multiple_buffers(xor1, n1, xor2, n2, stripes[parity].write_buf, end-start);
|
||||
}
|
||||
}
|
||||
}
|
||||
if (write_osd_set != read_osd_set)
|
||||
{
|
||||
for (int role = pg_minsize; role < pg_size; role++)
|
||||
{
|
||||
if (write_osd_set[role] != read_osd_set[role] && (start != 0 || end != chunk_size))
|
||||
{
|
||||
// Copy new parity into the read buffer to write it back
|
||||
memcpy(
|
||||
stripes[role].read_buf + start,
|
||||
stripes[role].write_buf,
|
||||
end - start
|
||||
);
|
||||
stripes[role].write_buf = stripes[role].read_buf;
|
||||
stripes[role].write_start = 0;
|
||||
stripes[role].write_end = chunk_size;
|
||||
}
|
||||
get_old_new_buffers(stripes[other], wr_start, wr_end, xor2, n2);
|
||||
xor_multiple_buffers(xor1, n1, xor2, n2, stripes[parity].write_buf, wr_end-wr_start);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user