SaveStates/ZSTD: Improve termination and file handling
This commit is contained in:
@@ -814,7 +814,7 @@ void compressed_zstd_serialization_file_handler::initialize(utils::serial& ar)
|
|||||||
|
|
||||||
// Make sure at least one thread is free
|
// Make sure at least one thread is free
|
||||||
// Limit thread count in order to make sure memory limits are under control (TODO: scale with RAM size)
|
// Limit thread count in order to make sure memory limits are under control (TODO: scale with RAM size)
|
||||||
const usz thread_count = std::min<u32>(std::max<u32>(utils::get_thread_count(), 2) - 1, 16);
|
const usz thread_count = std::min<u32>(std::max<u32>(utils::get_thread_count(), 2) - 1, 32);
|
||||||
|
|
||||||
for (usz i = 0; i < thread_count; i++)
|
for (usz i = 0; i < thread_count; i++)
|
||||||
{
|
{
|
||||||
@@ -1142,42 +1142,38 @@ void compressed_zstd_serialization_file_handler::finalize(utils::serial& ar)
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
for (auto& context : m_compression_threads)
|
usz pending_threads = 0;
|
||||||
{
|
|
||||||
// Notify to abort
|
|
||||||
while (!context.notified)
|
|
||||||
{
|
|
||||||
const auto data = context.m_input.compare_and_swap(null_ptr, empty_data);
|
|
||||||
|
|
||||||
if (!data)
|
do
|
||||||
|
{
|
||||||
|
pending_threads = 0;
|
||||||
|
|
||||||
|
for (auto& context : m_compression_threads)
|
||||||
|
{
|
||||||
|
// Try to notify all in bulk
|
||||||
|
if (!context.notified && !context.m_input && context.m_input.compare_and_swap_test(null_ptr, empty_data))
|
||||||
{
|
{
|
||||||
context.notified = true;
|
context.notified = true;
|
||||||
context.m_input.notify_one();
|
context.m_input.notify_one();
|
||||||
break;
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// Wait until valid input is processed
|
|
||||||
thread_ctrl::wait_for(1000);
|
|
||||||
}
|
}
|
||||||
}
|
|
||||||
|
|
||||||
for (auto& context : m_compression_threads)
|
for (auto& context : m_compression_threads)
|
||||||
{
|
|
||||||
// Wait for notification to be consumed
|
|
||||||
while (context.m_input)
|
|
||||||
{
|
{
|
||||||
thread_ctrl::wait_for(1000);
|
// Wait for notification to be sent
|
||||||
|
// And wait for data to be written to be read by the thread
|
||||||
|
if (context.m_output || !context.notified)
|
||||||
|
{
|
||||||
|
pending_threads++;
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
|
||||||
|
|
||||||
for (auto& context : m_compression_threads)
|
if (pending_threads)
|
||||||
{
|
|
||||||
// Wait for data to be writen to be read by the thread
|
|
||||||
while (context.m_output)
|
|
||||||
{
|
{
|
||||||
thread_ctrl::wait_for(1000);
|
thread_ctrl::wait_for(500);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
while (pending_threads);
|
||||||
|
|
||||||
for (usz idx = m_output_buffer_index;;)
|
for (usz idx = m_output_buffer_index;;)
|
||||||
{
|
{
|
||||||
|
|||||||
Reference in New Issue
Block a user