2022-04-26 18:14:09 +00:00
|
|
|
#pragma once
|
|
|
|
|
|
|
|
#include <Compression/ICompressionCodec.h>
|
|
|
|
#include <qpl/qpl.h>
|
2022-07-07 14:04:17 +00:00
|
|
|
|
2022-04-26 18:14:09 +00:00
|
|
|
namespace Poco
|
|
|
|
{
|
|
|
|
class Logger;
|
|
|
|
}
|
|
|
|
|
|
|
|
namespace DB
|
|
|
|
{
|
|
|
|
|
|
|
|
class DeflateJobHWPool
|
|
|
|
{
|
|
|
|
public:
|
|
|
|
DeflateJobHWPool();
|
|
|
|
~DeflateJobHWPool();
|
|
|
|
static DeflateJobHWPool & instance();
|
2022-07-07 16:10:06 +00:00
|
|
|
static constexpr auto JOB_POOL_SIZE = 1024;
|
2022-04-26 18:14:09 +00:00
|
|
|
static constexpr qpl_path_t PATH = qpl_path_hardware;
|
2022-07-07 16:10:06 +00:00
|
|
|
static qpl_job * jobPool[JOB_POOL_SIZE];
|
|
|
|
static std::atomic_bool jobLocks[JOB_POOL_SIZE];
|
2022-06-08 15:47:44 +00:00
|
|
|
bool jobPoolEnabled;
|
2022-04-26 18:14:09 +00:00
|
|
|
|
2022-07-07 14:26:57 +00:00
|
|
|
bool jobPoolReady()
|
2022-06-08 15:47:44 +00:00
|
|
|
{
|
|
|
|
return jobPoolEnabled;
|
|
|
|
}
|
2022-07-07 14:26:57 +00:00
|
|
|
qpl_job * acquireJob(uint32_t * job_id)
|
2022-04-26 18:14:09 +00:00
|
|
|
{
|
2022-06-08 15:47:44 +00:00
|
|
|
if (jobPoolEnabled)
|
2022-04-26 18:14:09 +00:00
|
|
|
{
|
2022-06-08 13:28:35 +00:00
|
|
|
uint32_t retry = 0;
|
2022-07-07 16:10:06 +00:00
|
|
|
auto index = random(JOB_POOL_SIZE);
|
2022-06-08 13:28:35 +00:00
|
|
|
while (tryLockJob(index) == false)
|
2022-04-26 18:14:09 +00:00
|
|
|
{
|
2022-07-07 16:10:06 +00:00
|
|
|
index = random(JOB_POOL_SIZE);
|
2022-06-08 13:28:35 +00:00
|
|
|
retry++;
|
2022-07-07 16:10:06 +00:00
|
|
|
if (retry > JOB_POOL_SIZE)
|
2022-06-08 13:28:35 +00:00
|
|
|
{
|
|
|
|
return nullptr;
|
|
|
|
}
|
2022-04-26 18:14:09 +00:00
|
|
|
}
|
2022-07-07 16:10:06 +00:00
|
|
|
*job_id = JOB_POOL_SIZE - index;
|
2022-06-08 13:28:35 +00:00
|
|
|
return jobPool[index];
|
|
|
|
}
|
|
|
|
else
|
|
|
|
{
|
|
|
|
return nullptr;
|
2022-04-26 18:14:09 +00:00
|
|
|
}
|
|
|
|
}
|
2022-07-07 14:26:57 +00:00
|
|
|
qpl_job * releaseJob(uint32_t job_id)
|
2022-04-26 18:14:09 +00:00
|
|
|
{
|
2022-06-08 15:47:44 +00:00
|
|
|
if (jobPoolEnabled)
|
2022-06-08 13:28:35 +00:00
|
|
|
{
|
2022-07-07 16:10:06 +00:00
|
|
|
uint32_t index = JOB_POOL_SIZE - job_id;
|
2022-06-08 13:28:35 +00:00
|
|
|
ReleaseJobObjectGuard _(index);
|
|
|
|
return jobPool[index];
|
|
|
|
}
|
|
|
|
else
|
|
|
|
{
|
|
|
|
return nullptr;
|
|
|
|
}
|
2022-04-26 18:14:09 +00:00
|
|
|
}
|
2022-07-07 14:26:57 +00:00
|
|
|
qpl_job * getJobPtr(uint32_t job_id)
|
2022-04-26 18:14:09 +00:00
|
|
|
{
|
2022-06-08 15:47:44 +00:00
|
|
|
if (jobPoolEnabled)
|
2022-06-08 13:28:35 +00:00
|
|
|
{
|
2022-07-07 16:10:06 +00:00
|
|
|
uint32_t index = JOB_POOL_SIZE - job_id;
|
2022-06-08 13:28:35 +00:00
|
|
|
return jobPool[index];
|
|
|
|
}
|
|
|
|
else
|
|
|
|
{
|
|
|
|
return nullptr;
|
|
|
|
}
|
2022-04-26 18:14:09 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
private:
|
2022-07-07 14:26:57 +00:00
|
|
|
size_t random(uint32_t pool_size)
|
2022-04-26 18:14:09 +00:00
|
|
|
{
|
|
|
|
size_t tsc = 0;
|
|
|
|
unsigned lo, hi;
|
|
|
|
__asm__ volatile("rdtsc" : "=a"(lo), "=d"(hi) : :);
|
|
|
|
tsc = (((static_cast<uint64_t>(hi)) << 32) | (static_cast<uint64_t>(lo)));
|
|
|
|
return (static_cast<size_t>((tsc * 44485709377909ULL) >> 4)) % pool_size;
|
|
|
|
}
|
|
|
|
|
2022-07-07 14:26:57 +00:00
|
|
|
int32_t get_job_size_helper()
|
2022-04-26 18:14:09 +00:00
|
|
|
{
|
|
|
|
static uint32_t size = 0;
|
|
|
|
if (size == 0)
|
|
|
|
{
|
|
|
|
const auto status = qpl_get_job_size(PATH, &size);
|
|
|
|
if (status != QPL_STS_OK)
|
|
|
|
{
|
|
|
|
return -1;
|
|
|
|
}
|
|
|
|
}
|
|
|
|
return size;
|
|
|
|
}
|
|
|
|
|
2022-07-07 14:26:57 +00:00
|
|
|
int32_t init_job_helper(qpl_job * qpl_job_ptr)
|
2022-04-26 18:14:09 +00:00
|
|
|
{
|
|
|
|
if (qpl_job_ptr == nullptr)
|
|
|
|
{
|
|
|
|
return -1;
|
|
|
|
}
|
|
|
|
auto status = qpl_init_job(PATH, qpl_job_ptr);
|
|
|
|
if (status != QPL_STS_OK)
|
|
|
|
{
|
|
|
|
return -1;
|
|
|
|
}
|
|
|
|
return 0;
|
|
|
|
}
|
|
|
|
|
2022-07-07 14:26:57 +00:00
|
|
|
int32_t initJobPool()
|
2022-04-26 18:14:09 +00:00
|
|
|
{
|
|
|
|
static bool initialized = false;
|
|
|
|
|
|
|
|
if (initialized == false)
|
|
|
|
{
|
|
|
|
const int32_t size = get_job_size_helper();
|
|
|
|
if (size < 0)
|
|
|
|
return -1;
|
2022-07-07 16:10:06 +00:00
|
|
|
for (int i = 0; i < JOB_POOL_SIZE; ++i)
|
2022-04-26 18:14:09 +00:00
|
|
|
{
|
|
|
|
jobPool[i] = nullptr;
|
|
|
|
qpl_job * qpl_job_ptr = reinterpret_cast<qpl_job *>(new uint8_t[size]);
|
|
|
|
if (init_job_helper(qpl_job_ptr) < 0)
|
|
|
|
return -1;
|
|
|
|
jobPool[i] = qpl_job_ptr;
|
2022-07-07 16:10:06 +00:00
|
|
|
jobLocks[i].store(false);
|
2022-04-26 18:14:09 +00:00
|
|
|
}
|
|
|
|
initialized = true;
|
|
|
|
}
|
|
|
|
return 0;
|
|
|
|
}
|
|
|
|
|
2022-07-07 14:26:57 +00:00
|
|
|
bool tryLockJob(size_t index)
|
2022-04-26 18:14:09 +00:00
|
|
|
{
|
|
|
|
bool expected = false;
|
2022-07-07 16:10:06 +00:00
|
|
|
return jobLocks[index].compare_exchange_strong(expected, true);
|
2022-04-26 18:14:09 +00:00
|
|
|
}
|
|
|
|
|
2022-07-07 14:26:57 +00:00
|
|
|
void destroyJobPool()
|
2022-04-26 18:14:09 +00:00
|
|
|
{
|
|
|
|
const uint32_t size = get_job_size_helper();
|
2022-07-07 16:10:06 +00:00
|
|
|
for (uint32_t i = 0; i < JOB_POOL_SIZE && size > 0; ++i)
|
2022-04-26 18:14:09 +00:00
|
|
|
{
|
|
|
|
while (tryLockJob(i) == false)
|
|
|
|
{
|
|
|
|
}
|
|
|
|
if (jobPool[i])
|
|
|
|
{
|
|
|
|
qpl_fini_job(jobPool[i]);
|
|
|
|
delete[] jobPool[i];
|
|
|
|
}
|
|
|
|
jobPool[i] = nullptr;
|
2022-07-07 16:10:06 +00:00
|
|
|
jobLocks[i].store(false);
|
2022-04-26 18:14:09 +00:00
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
struct ReleaseJobObjectGuard
|
|
|
|
{
|
|
|
|
uint32_t index;
|
|
|
|
ReleaseJobObjectGuard() = delete;
|
|
|
|
|
|
|
|
public:
|
2022-07-07 14:26:57 +00:00
|
|
|
ReleaseJobObjectGuard(const uint32_t i) : index(i)
|
2022-04-26 18:14:09 +00:00
|
|
|
{
|
|
|
|
}
|
2022-07-07 14:26:57 +00:00
|
|
|
~ReleaseJobObjectGuard()
|
2022-04-26 18:14:09 +00:00
|
|
|
{
|
2022-07-07 16:10:06 +00:00
|
|
|
jobLocks[index].store(false);
|
2022-04-26 18:14:09 +00:00
|
|
|
}
|
|
|
|
};
|
2022-06-08 15:47:44 +00:00
|
|
|
Poco::Logger * log;
|
|
|
|
};
|
|
|
|
class SoftwareCodecDeflate
|
|
|
|
{
|
|
|
|
public:
|
|
|
|
SoftwareCodecDeflate();
|
|
|
|
~SoftwareCodecDeflate();
|
|
|
|
uint32_t doCompressData(const char * source, uint32_t source_size, char * dest, uint32_t dest_size);
|
|
|
|
void doDecompressData(const char * source, uint32_t source_size, char * dest, uint32_t uncompressed_size);
|
|
|
|
|
|
|
|
private:
|
|
|
|
qpl_job * jobSWPtr; //Software Job Codec Ptr
|
|
|
|
std::unique_ptr<uint8_t[]> jobSWbuffer;
|
|
|
|
qpl_job * getJobCodecPtr();
|
2022-04-26 18:14:09 +00:00
|
|
|
};
|
|
|
|
|
2022-06-08 15:47:44 +00:00
|
|
|
class HardwareCodecDeflate
|
|
|
|
{
|
|
|
|
public:
|
|
|
|
bool hwEnabled;
|
|
|
|
HardwareCodecDeflate();
|
|
|
|
~HardwareCodecDeflate();
|
|
|
|
uint32_t doCompressData(const char * source, uint32_t source_size, char * dest, uint32_t dest_size) const;
|
|
|
|
uint32_t doDecompressData(const char * source, uint32_t source_size, char * dest, uint32_t uncompressed_size) const;
|
|
|
|
uint32_t doDecompressDataReq(const char * source, uint32_t source_size, char * dest, uint32_t uncompressed_size);
|
2022-07-07 01:37:11 +00:00
|
|
|
void doFlushAsynchronousDecompressRequests();
|
2022-06-08 15:47:44 +00:00
|
|
|
|
|
|
|
private:
|
|
|
|
std::map<uint32_t, qpl_job *> jobDecompAsyncMap;
|
|
|
|
Poco::Logger * log;
|
|
|
|
};
|
2022-04-26 18:14:09 +00:00
|
|
|
class CompressionCodecDeflate : public ICompressionCodec
|
|
|
|
{
|
|
|
|
public:
|
|
|
|
CompressionCodecDeflate();
|
2022-06-08 15:47:44 +00:00
|
|
|
//~CompressionCodecDeflate() ;
|
2022-04-26 18:14:09 +00:00
|
|
|
uint8_t getMethodByte() const override;
|
|
|
|
void updateHash(SipHash & hash) const override;
|
|
|
|
|
|
|
|
protected:
|
|
|
|
bool isCompression() const override
|
|
|
|
{
|
|
|
|
return true;
|
|
|
|
}
|
|
|
|
bool isGenericCompression() const override
|
|
|
|
{
|
|
|
|
return true;
|
|
|
|
}
|
|
|
|
uint32_t doCompressData(const char * source, uint32_t source_size, char * dest) const override;
|
|
|
|
uint32_t doCompressDataSW(const char * source, uint32_t source_size, char * dest) const;
|
|
|
|
void doDecompressData(const char * source, uint32_t source_size, char * dest, uint32_t uncompressed_size) const override;
|
2022-07-07 01:37:11 +00:00
|
|
|
void doFlushAsynchronousDecompressRequests() override;
|
2022-04-26 18:14:09 +00:00
|
|
|
|
|
|
|
private:
|
|
|
|
uint32_t getMaxCompressedDataSize(uint32_t uncompressed_size) const override;
|
2022-06-08 15:47:44 +00:00
|
|
|
std::unique_ptr<HardwareCodecDeflate> hwCodec;
|
|
|
|
std::unique_ptr<SoftwareCodecDeflate> swCodec;
|
2022-04-26 18:14:09 +00:00
|
|
|
};
|
|
|
|
|
|
|
|
}
|