Merge pull request #5835 from Plagman/plagman/cache_thread_mr

DiskCache: offload Store to a WorkQueueThread
This commit is contained in:
LC authored and GitHub committed 2026-08-21 14:36:39 -04:00
commit b2b6b263bb
13 files changed
+518 -123

No files matched your search

+1
View File
@@ -70,6 +70,7 @@ set(SRCS
Utils/LongJump.cpp
Utils/Telemetry.cpp
Utils/Threads.cpp
Utils/WorkQueueThread.cpp
Utils/Profiler.cpp)
if (ARCHITECTURE_arm64)
+173 -110
View File
@@ -43,7 +43,7 @@ namespace DiskCache {
if (!ReadOnly) {
Modes = Modes | File::FileModes::WRITE | File::FileModes::CREATE;
}
FD = fextl::make_unique<File::File>(FileName.c_str(), Modes);
FD = fextl::make_unique<File::File>(FileName.c_str(), Modes, false);
if (!FD->IsValid()) {
FD.reset();
return false;
@@ -51,7 +51,7 @@ namespace DiskCache {
bool Valid = false;
bool TookLock = false;
ssize_t Size = FD->Seek(0, File::SeekOp::END);
ssize_t Size = FD->Size();
if (Size < FOZ_REF_MAGIC_SIZE && !ReadOnly) {
if (!FD->Lock(OPEN_LOCK_TIMEOUT_MS)) {
@@ -59,16 +59,15 @@ namespace DiskCache {
return false;
}
TookLock = true;
// seek in case someone else made it while we waited above
Size = FD->Seek(0, File::SeekOp::END);
// check size again in case someone else made it while we waited
Size = FD->Size();
}
if (Size == 0 && !ReadOnly) {
Valid = FD->Write(MesaFOZ::stream_reference_magic_and_version, FOZ_REF_MAGIC_SIZE) == FOZ_REF_MAGIC_SIZE;
Valid = FD->PWrite(MesaFOZ::stream_reference_magic_and_version, FOZ_REF_MAGIC_SIZE, 0) == FOZ_REF_MAGIC_SIZE;
} else {
FD->Seek(0, File::SeekOp::BEGIN);
uint8_t magic[FOZ_REF_MAGIC_SIZE];
if (FD->Read(magic, FOZ_REF_MAGIC_SIZE) == FOZ_REF_MAGIC_SIZE &&
if (FD->PRead(magic, FOZ_REF_MAGIC_SIZE, 0) == FOZ_REF_MAGIC_SIZE &&
memcmp(magic, MesaFOZ::stream_reference_magic_and_version, FOZ_REF_MAGIC_SIZE - 1) == 0) {
int version = magic[FOZ_REF_MAGIC_SIZE - 1];
Valid = version <= MesaFOZ::FOSSILIZE_FORMAT_VERSION && version >= MesaFOZ::FOSSILIZE_FORMAT_MIN_COMPAT_VERSION;
@@ -85,41 +84,34 @@ namespace DiskCache {
return Valid;
}
bool FOZFile::ReadNextBlob(MesaFOZ::foz_payload_key& OutKey, MesaFOZ::foz_payload_header& OutHeader, fextl::vector<uint8_t>& OutBlob) {
if (FD->Read(OutKey.bytes, sizeof(OutKey.bytes)) != sizeof(OutKey.bytes)) {
ssize_t FOZFile::Size() {
return FD ? FD->Size() : -1;
}
bool FOZFile::ReadAll(fextl::vector<uint8_t>& Out) {
ssize_t FileSize = Size();
if (FileSize < FOZ_REF_MAGIC_SIZE) {
return false;
}
if (FD->Read(&OutHeader, sizeof(OutHeader)) != sizeof(OutHeader)) {
return false;
}
OutBlob.resize(OutHeader.payload_size);
if (FD->Read(OutBlob.data(), OutBlob.size()) != (ssize_t)OutBlob.size()) {
return false;
}
return true;
Out.resize((size_t)FileSize - FOZ_REF_MAGIC_SIZE);
return FD->PRead(Out.data(), Out.size(), FOZ_REF_MAGIC_SIZE) == (ssize_t)Out.size();
}
bool FOZFile::ReadBlob(uint64_t Offset, std::span<uint8_t> OutBlob) {
ssize_t SeekRet = FD->Seek(Offset, File::SeekOp::BEGIN);
if (SeekRet < 0) {
if (FD->PRead(OutBlob.data(), OutBlob.size(), Offset) != (ssize_t)OutBlob.size()) {
return false;
}
if (FD->Read(OutBlob.data(), OutBlob.size()) != (ssize_t)OutBlob.size()) {
return false;
}
return true;
}
bool FOZFile::WriteBlob(const MesaFOZ::foz_payload_key& Key, std::span<const std::span<const uint8_t>> BlobChunks, uint64_t& OutBlobOffset) {
uint64_t WriteOffset = 0;
ssize_t SeekRet = FD->Seek(0, File::SeekOp::END);
if (SeekRet < 0) {
ssize_t FileSize = FD->Size();
if (FileSize < 0) {
return false;
}
WriteOffset = (uint64_t)SeekRet;
uint64_t WriteOffset = (uint64_t)FileSize;
if (FD->Write(Key.bytes, sizeof(Key.bytes)) != sizeof(Key.bytes)) {
if (FD->PWrite(Key.bytes, sizeof(Key.bytes), WriteOffset) != sizeof(Key.bytes)) {
return false;
}
WriteOffset += sizeof(Key.bytes);
@@ -134,7 +126,7 @@ namespace DiskCache {
.crc = 0, // todo? maybe
.uncompressed_size = (uint32_t)TotalBlobSize};
if (FD->Write(&ScratchHeader, sizeof(ScratchHeader)) != sizeof(ScratchHeader)) {
if (FD->PWrite(&ScratchHeader, sizeof(ScratchHeader), WriteOffset) != sizeof(ScratchHeader)) {
return false;
}
WriteOffset += sizeof(ScratchHeader);
@@ -145,9 +137,10 @@ namespace DiskCache {
if (Chunk.size() == 0) {
continue;
}
if (FD->Write(Chunk.data(), Chunk.size()) != (ssize_t)Chunk.size()) {
if (FD->PWrite(Chunk.data(), Chunk.size(), WriteOffset) != (ssize_t)Chunk.size()) {
return false;
}
WriteOffset += Chunk.size();
}
return true;
@@ -166,19 +159,40 @@ namespace DiskCache {
}
void IndexedDB::PopulateIndex(Index& CacheIndex) {
MesaFOZ::foz_payload_key Key;
MesaFOZ::foz_payload_header Header;
fextl::vector<uint8_t> Blob;
fextl::vector<uint8_t> Data;
if (!IndexFOZ.ReadAll(Data)) {
return;
}
while (IndexFOZ.ReadNextBlob(Key, Header, Blob)) {
if (Blob.size() != sizeof(MesaFOZ::mesa_index_db_file_entry)) {
ssize_t CacheFOZSize = CacheFOZ.Size();
if (CacheFOZSize < 0) {
return;
}
const uint8_t* IndexDataStart = Data.data();
const size_t IndexDataSize = Data.size();
size_t ReadOffset = 0;
while (ReadOffset + sizeof(MesaFOZ::foz_payload_key) + sizeof(MesaFOZ::foz_payload_header) <= IndexDataSize) {
const auto* FOZKey = reinterpret_cast<const MesaFOZ::foz_payload_key*>(IndexDataStart + ReadOffset);
ReadOffset += sizeof(MesaFOZ::foz_payload_key);
const auto* FOZHeader = reinterpret_cast<const MesaFOZ::foz_payload_header*>(IndexDataStart + ReadOffset);
ReadOffset += sizeof(MesaFOZ::foz_payload_header);
if (FOZHeader->payload_size != sizeof(MesaFOZ::mesa_index_db_file_entry) || ReadOffset + FOZHeader->payload_size > IndexDataSize) {
break;
}
MesaFOZ::mesa_index_db_file_entry* IndexEntry = (MesaFOZ::mesa_index_db_file_entry*)Blob.data();
if (IndexEntry->hash != XXH3_64bits(Key.bytes, FOSSILIZE_BLOB_HASH_LENGTH)) {
const auto* IndexBlobPayload = reinterpret_cast<const MesaFOZ::mesa_index_db_file_entry*>(IndexDataStart + ReadOffset);
ReadOffset += FOZHeader->payload_size;
if (IndexBlobPayload->hash != XXH3_64bits(FOZKey->bytes, FOSSILIZE_BLOB_HASH_LENGTH)) {
break;
}
CacheIndex.insert({IndexEntry->hash, {this, IndexEntry->cache_db_file_offset, IndexEntry->size}});
// skip corrupt (carefully) so we don't have to figure that out in the hot path later
if (IndexBlobPayload->cache_db_file_offset > (uint64_t)CacheFOZSize ||
IndexBlobPayload->size > (uint64_t)CacheFOZSize - IndexBlobPayload->cache_db_file_offset) {
continue;
}
CacheIndex.insert({IndexBlobPayload->hash, {this, IndexBlobPayload->cache_db_file_offset, IndexBlobPayload->size}});
}
// could truncate/delete index if we don't end up perfectly at end here
}
@@ -187,15 +201,18 @@ namespace DiskCache {
return CacheFOZ.ReadBlob(Offset, OutBlob);
}
bool IndexedDB::StoreCacheBlob(const MesaFOZ::foz_payload_key& Key, std::span<const std::span<const uint8_t>> BlobChunks, Index& Index) {
bool IndexedDB::StoreCacheBlob(const MesaFOZ::foz_payload_key& Key, std::span<const uint8_t> Blob, Index& Index, std::mutex& IndexMutex) {
if (ReadOnly) {
// shouldn't happen
return false;
}
uint64_t Hash = XXH3_64bits(Key.bytes, FOSSILIZE_BLOB_HASH_LENGTH);
if (Index.contains(Hash)) {
// shouldn't really happen.. assert or something?
return true;
{
std::lock_guard Guard(IndexMutex);
if (Index.contains(Hash)) {
// shouldn't really happen.. assert or something?
return true;
}
}
if (!CacheFOZ.Lock(STORE_LOCK_TIMEOUT_MS) || !IndexFOZ.Lock(STORE_LOCK_TIMEOUT_MS)) {
@@ -205,6 +222,7 @@ namespace DiskCache {
}
// write cache side first so we get offset for index
std::span<const uint8_t> BlobChunks[] = {Blob};
uint64_t BlobOffset = 0;
if (!CacheFOZ.WriteBlob(Key, BlobChunks, BlobOffset)) {
CacheFOZ.Unlock();
@@ -212,13 +230,8 @@ namespace DiskCache {
return false;
}
uint64_t TotalBlobSize = 0;
for (const std::span<const uint8_t>& Chunk : BlobChunks) {
TotalBlobSize += Chunk.size();
}
MesaFOZ::mesa_index_db_file_entry IndexEntry {.hash = Hash,
.size = (uint32_t)TotalBlobSize,
.size = (uint32_t)Blob.size(),
.last_access_time = 0, // todo..
.cache_db_file_offset = BlobOffset};
@@ -233,7 +246,8 @@ namespace DiskCache {
CacheFOZ.Unlock();
IndexFOZ.Unlock();
Index[Hash] = {this, BlobOffset, (uint32_t)TotalBlobSize};
std::lock_guard Guard(IndexMutex);
Index[Hash] = {this, BlobOffset, (uint32_t)Blob.size()};
return true;
}
@@ -300,13 +314,16 @@ namespace DiskCache {
// advance to next
RONames.remove_prefix(Delim + 1);
}
if (IsWritingDiskCache()) {
Writer = fextl::make_unique<WorkQueueThread>();
}
}
std::optional<CodeHitData> DiskCache::Lookup(Core::InternalThreadState* Thread, const ExecutableFileSectionInfo& Region, uint64_t GuestRIP) {
if (!IsReadingDiskCache()) {
return std::nullopt;
}
std::lock_guard Guard(Lock);
uint64_t ModuleOffset = GuestRIP - Region.FileStartVA;
// todo move key making to a helper once we have options and stuff (see Store)
@@ -314,12 +331,18 @@ namespace DiskCache {
memcpy(Key.bytes, &ModuleOffset, sizeof(ModuleOffset));
uint64_t Hash = XXH3_64bits(Key.bytes, FOSSILIZE_BLOB_HASH_LENGTH);
auto It = Index.find(Hash);
if (It == Index.end()) {
// definite miss
return std::nullopt;
IndexEntry Entry;
{
std::lock_guard Guard(IndexLock);
auto It = Index.find(Hash);
if (It == Index.end()) {
// definite miss
return std::nullopt;
}
// we can't hold onto the iterator, the map may shift while we don't hold the lock
Entry = It->second;
}
const IndexEntry& Entry = It->second;
// found a key hash match, could still be a miss, read the blob and verify more
CodeHitData HitData;
HitData.Blob.resize(Entry.Size);
@@ -405,6 +428,25 @@ namespace DiskCache {
return HitData;
}
static inline bool IsRelocationInBlock(const FEXCore::CPU::Relocation& Reloc, const CPU::CPUBackend::CompiledCode& CompiledCode) {
return Reloc.Header.Offset >= CompiledCode.HostCodeOffset && Reloc.Header.Offset < CompiledCode.HostCodeOffset + CompiledCode.Size;
}
struct DiskCache::CacheStoreWorkItem final : WorkQueueThread::WorkItem {
DiskCache* Self;
IndexedDB* DB;
MesaFOZ::foz_payload_key Key;
fextl::vector<uint8_t> Blob;
CacheStoreWorkItem(DiskCache* Self, IndexedDB* DB, const MesaFOZ::foz_payload_key& Key, fextl::vector<uint8_t>&& Blob)
: Self(Self)
, DB(DB)
, Key(Key)
, Blob(std::move(Blob)) {}
void Run() override {
DB->StoreCacheBlob(Key, Blob, Self->Index, Self->IndexLock);
}
};
bool DiskCache::Store(Core::InternalThreadState* Thread, const ExecutableFileSectionInfo& Region, uint64_t GuestRIP,
std::span<const uint8_t> GuestCode, const CPU::CPUBackend::CompiledCode& CompiledCode,
std::span<const FEXCore::CPU::Relocation> Relocations, const Frontend::Decoder::DecodedBlockInformation* DecodedBlockInfo) {
@@ -414,13 +456,12 @@ namespace DiskCache {
if (!DecodedBlockInfo) {
return false;
}
std::lock_guard Guard(Lock);
// check for any reloc targets outside of our jurisdiction
// todo what are they exactly? caching those blocks is great when it works, so need to figure this out and make finer-grained if we can
if (RelocationFilter) {
for (const auto& Reloc : Relocations) {
if (Reloc.Header.Offset < CompiledCode.HostCodeOffset || Reloc.Header.Offset >= CompiledCode.HostCodeOffset + CompiledCode.Size) {
if (!IsRelocationInBlock(Reloc, CompiledCode)) {
continue;
}
if (Reloc.Header.Type != CPU::RelocationTypes::RELOC_GUEST_RIP_LITERAL && Reloc.Header.Type != CPU::RelocationTypes::RELOC_GUEST_RIP_MOVE) {
@@ -438,23 +479,70 @@ namespace DiskCache {
}
}
// pack entrypoints to disk format
fextl::vector<BlobEntryPoint> CacheEntryPoints;
CacheEntryPoints.reserve(CompiledCode.EntryPoints.size());
for (auto [GuestAddr, HostAddr] : CompiledCode.EntryPoints) {
CacheEntryPoints.push_back({GuestAddr - Region.FileStartVA, uint32_t(HostAddr - CompiledCode.BlockBegin)});
}
// pack relocations to disk format
fextl::vector<BlobSmallRelocation> SmallRelocs;
fextl::vector<BlobThunkRelocation> ThunkRelocs;
// todo discover sizes first and reserve vecs?
uint32_t SmallRelocCount = 0;
uint32_t ThunkRelocCount = 0;
for (const auto& Reloc : Relocations) {
// relocs aren't cleared every time if IsGeneratingCache, so filter just in case
if (Reloc.Header.Offset < CompiledCode.HostCodeOffset || Reloc.Header.Offset >= CompiledCode.HostCodeOffset + CompiledCode.Size) {
if (!IsRelocationInBlock(Reloc, CompiledCode)) {
continue;
}
if (Reloc.Header.Type == CPU::RelocationTypes::RELOC_NAMED_THUNK_MOVE) {
ThunkRelocCount++;
} else {
SmallRelocCount++;
}
}
const uint32_t EntryPointCount = (uint32_t)CompiledCode.EntryPoints.size();
const uint32_t TouchedGuestPagesCount = DecodedBlockInfo ? (uint32_t)DecodedBlockInfo->CodePages.size() : 0;
const size_t HeaderOffset = 0;
const size_t HostCodeOffset = HeaderOffset + sizeof(BlobFixedHeader);
const size_t EntryPointsOffset = HostCodeOffset + CompiledCode.Size;
const size_t SmallRelocsOffset = EntryPointsOffset + EntryPointCount * sizeof(BlobEntryPoint);
const size_t ThunkRelocsOffset = SmallRelocsOffset + SmallRelocCount * sizeof(BlobSmallRelocation);
const size_t TouchedGuestPagesOffset = ThunkRelocsOffset + ThunkRelocCount * sizeof(BlobThunkRelocation);
const size_t GuestCodeOffset = TouchedGuestPagesOffset + TouchedGuestPagesCount * sizeof(int64_t);
const size_t TotalSize = GuestCodeOffset + GuestCode.size();
// we'll copy everything into here and pass it to the Writer, then return to caller quickly
fextl::vector<uint8_t> Blob;
Blob.resize(TotalSize);
uint8_t* BlobData = Blob.data();
uint64_t ModuleOffset = GuestRIP - Region.FileStartVA;
// todo also copy/hash options that affect codegen into the key
// todo should try to keep the key ascii i think?
MesaFOZ::foz_payload_key Key = {};
memcpy(Key.bytes, &ModuleOffset, sizeof(ModuleOffset));
BlobFixedHeader Header {
.GuestSize = (uint32_t)GuestCode.size(),
.HostSize = (uint32_t)CompiledCode.Size,
.EntryPointCount = EntryPointCount,
.SmallRelocCount = SmallRelocCount,
.ThunkRelocCount = ThunkRelocCount,
.TouchedGuestPagesCount = TouchedGuestPagesCount,
.GuestHash = XXH3_128bits(GuestCode.data(), GuestCode.size()),
};
memcpy(BlobData + HeaderOffset, &Header, sizeof(Header));
memcpy(BlobData + HostCodeOffset, CompiledCode.BlockBegin, CompiledCode.Size);
// pack and relocate entrypoints
auto* EntryPoints = reinterpret_cast<BlobEntryPoint*>(BlobData + EntryPointsOffset);
uint32_t EntryIdx = 0;
for (auto [GuestAddr, HostAddr] : CompiledCode.EntryPoints) {
EntryPoints[EntryIdx++] = {GuestAddr - Region.FileStartVA, uint32_t(HostAddr - CompiledCode.BlockBegin)};
}
// pack relocations
auto* SmallRelocs = reinterpret_cast<BlobSmallRelocation*>(BlobData + SmallRelocsOffset);
auto* ThunkRelocs = reinterpret_cast<BlobThunkRelocation*>(BlobData + ThunkRelocsOffset);
uint32_t SmallIdx = 0;
uint32_t ThunkIdx = 0;
for (const auto& Reloc : Relocations) {
if (!IsRelocationInBlock(Reloc, CompiledCode)) {
continue;
}
// re-relocate :harold:
@@ -468,7 +556,7 @@ namespace DiskCache {
SmallReloc.Offset = LocalOffset;
SmallReloc.Type = uint8_t(Reloc.Header.Type);
SmallReloc.Named.Symbol = uint32_t(Reloc.NamedSymbolLiteral.Symbol);
SmallRelocs.push_back(SmallReloc);
SmallRelocs[SmallIdx++] = SmallReloc;
break;
}
case CPU::RelocationTypes::RELOC_GUEST_RIP_LITERAL: {
@@ -476,7 +564,7 @@ namespace DiskCache {
SmallReloc.Offset = LocalOffset;
SmallReloc.Type = uint8_t(Reloc.Header.Type);
SmallReloc.RIPLiteral.GuestRIP = Reloc.GuestRIP.GuestRIP - GuestRIP;
SmallRelocs.push_back(SmallReloc);
SmallRelocs[SmallIdx++] = SmallReloc;
break;
}
case CPU::RelocationTypes::RELOC_GUEST_RIP_MOVE: {
@@ -485,7 +573,7 @@ namespace DiskCache {
SmallReloc.Type = uint8_t(Reloc.Header.Type);
SmallReloc.RIPMove.RegisterIndex = Reloc.GuestRIP.RegisterIndex;
SmallReloc.RIPMove.GuestRIP = Reloc.GuestRIP.GuestRIP - GuestRIP;
SmallRelocs.push_back(SmallReloc);
SmallRelocs[SmallIdx++] = SmallReloc;
break;
}
case CPU::RelocationTypes::RELOC_NAMED_THUNK_MOVE: {
@@ -493,50 +581,25 @@ namespace DiskCache {
BigReloc.Offset = LocalOffset;
BigReloc.RegisterIndex = Reloc.NamedThunkMove.RegisterIndex;
memcpy(BigReloc.SymbolHash, &Reloc.NamedThunkMove.Symbol, sizeof(BigReloc.SymbolHash));
ThunkRelocs.push_back(BigReloc);
ThunkRelocs[ThunkIdx++] = BigReloc;
break;
}
}
}
// pack touched pages, relative to GuestRIP
// relocate touched pages relative to GuestRIP
// in theory we could save some size here, unlikely we need all 64bits
fextl::vector<int64_t> GuestPageOffsets;
if (DecodedBlockInfo) {
GuestPageOffsets.reserve(DecodedBlockInfo->CodePages.size());
for (auto& GuestPage : DecodedBlockInfo->CodePages) {
GuestPageOffsets.push_back(GuestPage - GuestRIP);
}
auto* PageOffsets = reinterpret_cast<int64_t*>(BlobData + TouchedGuestPagesOffset);
uint32_t PageIdx = 0;
for (auto GuestPage : DecodedBlockInfo->CodePages) {
PageOffsets[PageIdx++] = GuestPage - GuestRIP;
}
uint64_t ModuleOffset = GuestRIP - Region.FileStartVA;
memcpy(BlobData + GuestCodeOffset, GuestCode.data(), GuestCode.size());
// todo also copy/hash options that affect codegen into the key
// todo should try to keep the key ascii i think?
MesaFOZ::foz_payload_key Key = {};
memcpy(Key.bytes, &ModuleOffset, sizeof(ModuleOffset));
BlobFixedHeader Header {
.GuestSize = (uint32_t)GuestCode.size(),
.HostSize = (uint32_t)CompiledCode.Size,
.EntryPointCount = (uint32_t)CacheEntryPoints.size(),
.SmallRelocCount = (uint32_t)SmallRelocs.size(),
.ThunkRelocCount = (uint32_t)ThunkRelocs.size(),
.TouchedGuestPagesCount = (uint32_t)GuestPageOffsets.size(),
.GuestHash = XXH3_128bits(GuestCode.data(), GuestCode.size()),
};
std::span<const uint8_t> BlobChunks[] = {
{(const uint8_t*)&Header, sizeof(Header)},
{(const uint8_t*)CompiledCode.BlockBegin, CompiledCode.Size},
{(const uint8_t*)CacheEntryPoints.data(), CacheEntryPoints.size() * sizeof(BlobEntryPoint)},
{(const uint8_t*)SmallRelocs.data(), SmallRelocs.size() * sizeof(BlobSmallRelocation)},
{(const uint8_t*)ThunkRelocs.data(), ThunkRelocs.size() * sizeof(BlobThunkRelocation)},
{(const uint8_t*)GuestPageOffsets.data(), GuestPageOffsets.size() * sizeof(int64_t)},
GuestCode,
};
return RWCacheDB->StoreCacheBlob(Key, BlobChunks, Index);
// hand the rest off to the writer thread
Writer->QueueWork(fextl::make_unique<CacheStoreWorkItem>(this, RWCacheDB.get(), Key, std::move(Blob)));
return true;
}
} // namespace DiskCache
+54
View File
@@ -0,0 +1,54 @@
// SPDX-License-Identifier: MIT
#include <FEXCore/Utils/WorkQueueThread.h>
#include <FEXCore/Utils/LogManager.h>
namespace FEXCore {
WorkQueueThread::WorkQueueThread() {
Thread = FEXCore::Threads::Thread::Create(ThreadEntry, this);
}
WorkQueueThread::~WorkQueueThread() {
{
std::unique_lock lk {Mutex};
Stop = true;
}
CV.notify_one();
if (Thread && Thread->joinable()) {
Thread->join(nullptr);
}
}
void WorkQueueThread::QueueWork(fextl::unique_ptr<WorkItem> Work) {
{
std::unique_lock lk {Mutex};
Queue.push_back(std::move(Work));
}
CV.notify_one();
}
void WorkQueueThread::ThreadProc() {
while (true) {
fextl::unique_ptr<WorkItem> Work;
{
std::unique_lock lk {Mutex};
while (!(Stop || !Queue.empty())) {
CV.wait(lk);
}
if (Queue.empty()) {
// nothing to do? must be stopping
LOGMAN_THROW_A_FMT(Stop, "WorkQueueThread wakes up empty but no Stop?");
return;
}
Work = std::move(Queue.front());
Queue.pop_front();
}
Work->Run();
// Work is destroyed here
}
}
} // namespace FEXCore
+11 -5
View File
@@ -7,6 +7,7 @@
#include "Interface/Core/CPUBackend.h"
#include "FEXCore/Config/Config.h"
#include "FEXCore/Utils/File.h"
#include "FEXCore/Utils/WorkQueueThread.h"
#include "FEXCore/fextl/memory.h"
#include <FEXCore/fextl/string.h>
#include <FEXCore/fextl/unordered_set.h>
@@ -123,7 +124,8 @@ namespace DiskCache {
}
return FD->Unlock();
}
bool ReadNextBlob(MesaFOZ::foz_payload_key& OutKey, MesaFOZ::foz_payload_header& OutHeader, fextl::vector<uint8_t>& OutBlob);
ssize_t Size();
bool ReadAll(fextl::vector<uint8_t>& Out); // from first blob
bool ReadBlob(uint64_t Offset, std::span<uint8_t> OutBlob);
bool WriteBlob(const MesaFOZ::foz_payload_key& Key, std::span<const std::span<const uint8_t>> BlobChunks, uint64_t& OutBlobOffset);
@@ -140,11 +142,11 @@ namespace DiskCache {
bool Open(const fextl::string& CacheDBName, bool ReadOnly);
void PopulateIndex(Index& CacheIndex);
bool ReadCacheBlob(uint64_t Offset, std::span<uint8_t> OutBlob);
bool StoreCacheBlob(const MesaFOZ::foz_payload_key& Key, std::span<const std::span<const uint8_t>> BlobChunks, Index& CacheIndex);
bool StoreCacheBlob(const MesaFOZ::foz_payload_key& Key, std::span<const uint8_t> Blob, Index& CacheIndex, std::mutex& IndexMutex);
private:
// give up after 2ms of trying to store - when we have an async thread we can increase this
static constexpr uint32_t STORE_LOCK_TIMEOUT_MS = 2;
// stores run on the Writer, so returning quick isn't as important
static constexpr uint32_t STORE_LOCK_TIMEOUT_MS = 1000;
FOZFile CacheFOZ;
FOZFile IndexFOZ;
@@ -174,7 +176,11 @@ namespace DiskCache {
fextl::vector<fextl::unique_ptr<IndexedDB>> ROCacheDBs;
fextl::unique_ptr<IndexedDB> RWCacheDB;
Index Index;
std::mutex Lock;
std::mutex IndexLock;
struct CacheStoreWorkItem;
// the Writer holds references to all this stuff above and needs to be last
fextl::unique_ptr<WorkQueueThread> Writer;
FEX_CONFIG_OPT(EnableDiskCache, DISKCACHE);
FEX_CONFIG_OPT(RelocationFilter, DISKCACHERELOCATIONFILTER);
+78 -3
View File
@@ -3,6 +3,7 @@
#include <FEXCore/fextl/allocator.h>
#include <FEXCore/fextl/string.h>
#include <FEXCore/Utils/EnumOperators.h>
#include "FEXCore/Utils/LogManager.h"
#include <chrono>
#include <thread>
@@ -11,6 +12,7 @@
#include <fcntl.h>
#include <unistd.h>
#include <sys/file.h>
#include <sys/stat.h>
#else
#define WIN32_LEAN_AND_MEAN
#include <windows.h>
@@ -43,7 +45,8 @@ public:
File() = default;
File(const char* Filepath, FileModes Modes) {
File(const char* Filepath, FileModes Modes, bool Seekable = true)
: Seekable {Seekable} {
#ifndef _WIN32
auto Disp = TranslateModes(Modes);
Handle = open(Filepath, Disp, DEFAULT_USER_PERMS);
@@ -75,6 +78,10 @@ public:
* @return The number of bytes actually written or -1 on error.
*/
ssize_t Write(const void* Buffer, size_t Bytes) {
if (!Seekable) {
LOGMAN_THROW_A_FMT(false, "Can't use non-positioned ops on a non-seekable file!");
return -1;
}
#ifndef _WIN32
return write(Handle, Buffer, Bytes);
#else
@@ -101,6 +108,10 @@ public:
* @return The number of bytes read or -1 on error.
*/
ssize_t Read(void* Buffer, size_t Bytes) {
if (!Seekable) {
LOGMAN_THROW_A_FMT(false, "Can't use non-positioned ops on a non-seekable file!");
return -1;
}
#ifndef _WIN32
return read(Handle, Buffer, Bytes);
#else
@@ -114,6 +125,48 @@ public:
#endif
}
ssize_t PRead(void* Buffer, size_t Bytes, uint64_t Offset) {
if (Seekable) {
LOGMAN_THROW_A_FMT(false, "Can't use positioned ops on a seekable file!");
return -1;
}
#ifndef _WIN32
return pread(Handle, Buffer, Bytes, Offset);
#else
DWORD BytesRead {};
OVERLAPPED Overlapped {};
Overlapped.Offset = static_cast<DWORD>(Offset);
Overlapped.OffsetHigh = static_cast<DWORD>(Offset >> 32);
auto Result = ReadFile(Handle, Buffer, Bytes, &BytesRead, &Overlapped);
if (Result) {
return BytesRead;
}
// Some error, match Linux side.
return -1;
#endif
}
ssize_t PWrite(const void* Buffer, size_t Bytes, uint64_t Offset) {
if (Seekable) {
LOGMAN_THROW_A_FMT(false, "Can't use positioned ops on a seekable file!");
return -1;
}
#ifndef _WIN32
return pwrite(Handle, Buffer, Bytes, Offset);
#else
DWORD BytesWritten {};
OVERLAPPED Overlapped {};
Overlapped.Offset = static_cast<DWORD>(Offset);
Overlapped.OffsetHigh = static_cast<DWORD>(Offset >> 32);
auto Result = WriteFile(Handle, Buffer, Bytes, &BytesWritten, &Overlapped);
if (Result) {
return BytesWritten;
}
// Some error, match Linux side.
return -1;
#endif
}
bool Lock(uint32_t TimeoutMS) {
for (uint32_t i = 0;; ++i) {
if (TryLock()) {
@@ -202,6 +255,22 @@ public:
#endif
}
ssize_t Size() {
#ifndef _WIN32
struct stat st;
if (fstat(Handle, &st) != 0) {
return -1;
}
return st.st_size;
#else
LARGE_INTEGER FileSize;
if (!GetFileSizeEx(Handle, &FileSize)) {
return -1;
}
return FileSize.QuadPart;
#endif
}
/**
* @brief Seek the file pointer location.
*
@@ -211,6 +280,10 @@ public:
* @return The current file pointer location or -1.
*/
ssize_t Seek(ssize_t Distance, SeekOp Op) {
if (!Seekable) {
LOGMAN_THROW_A_FMT(false, "Can't use non-positioned ops on a non-seekable file!");
return -1;
}
#ifndef _WIN32
return lseek(Handle, Distance, TranslateSeek(Op));
#else
@@ -227,10 +300,11 @@ public:
protected:
File(FileHandleType Handle, bool ShouldClose)
File(FileHandleType Handle, bool ShouldClose, bool Seekable = true)
: ShouldClose {ShouldClose}
, IsValidHandle {true}
, Handle {Handle} {}
, Handle {Handle}
, Seekable {Seekable} {}
private:
bool TryLock() {
if (Locked) {
@@ -257,6 +331,7 @@ private:
bool IsValidHandle {};
FileHandleType Handle {};
bool Seekable = true;
bool Locked = false;
#ifndef _WIN32
static constexpr int DEFAULT_USER_PERMS = S_IRWXU | S_IRWXG | S_IRWXO;
@@ -0,0 +1,41 @@
// SPDX-License-Identifier: MIT
#pragma once
#include <FEXCore/Utils/Threads.h>
#include <FEXCore/fextl/deque.h>
#include <FEXCore/fextl/memory.h>
#include <condition_variable>
#include <mutex>
namespace FEXCore {
class WorkQueueThread {
public:
// destroyed after Run()
struct WorkItem {
virtual ~WorkItem() = default;
virtual void Run() = 0;
};
WorkQueueThread();
~WorkQueueThread();
void QueueWork(fextl::unique_ptr<WorkItem> Work);
private:
// static function for the ::Thread to refer to
static void* ThreadEntry(void* Self) {
static_cast<WorkQueueThread*>(Self)->ThreadProc();
return nullptr;
}
void ThreadProc();
std::mutex Mutex;
std::condition_variable CV;
fextl::deque<fextl::unique_ptr<WorkItem>> Queue;
bool Stop = false;
fextl::unique_ptr<FEXCore::Threads::Thread> Thread;
};
} // namespace FEXCore
+2
View File
@@ -36,6 +36,7 @@ $end_info$
#include "Common/Exception.h"
#include "Common/ImageTracker.h"
#include "Common/InvalidationTracker.h"
#include "Common/Threads.h"
#include "Common/OvercommitTracker.h"
#include "Common/TSOHandlerConfig.h"
#include "Common/CPUFeatures.h"
@@ -581,6 +582,7 @@ NTSTATUS ProcessInit() {
InitSyscalls();
FEX::Windows::InitCRTProcess();
FEX::Windows::SetupThreadHandlers();
const auto ExecutablePath = FEX::Windows::GetExecutableFilePath();
const auto ExecutableName = FEX::Windows::BaseName(ExecutablePath);
FEX::Config::LoadConfig(fextl::string {ExecutableName}, _environ, FEX::ReadPortabilityInformation());
+2 -1
View File
@@ -14,7 +14,8 @@ add_library(CommonWindows STATIC
SHMStats.cpp
InvalidationTracker.cpp
ImageTracker.cpp
Logging.cpp)
Logging.cpp
Threads.cpp)
target_link_libraries(CommonWindows FEXCore_Base)
target_include_directories(CommonWindows PRIVATE "${CMAKE_SOURCE_DIR}/Source/Windows/include/")
+95
View File
@@ -0,0 +1,95 @@
// SPDX-License-Identifier: MIT
#include "Threads.h"
#include "CRT/CRT.h"
#include <FEXCore/Utils/LogManager.h>
#include <FEXCore/Utils/Threads.h>
#include <FEXCore/fextl/memory.h>
#include <winternl.h>
#include <windows.h>
namespace FEX::Windows {
namespace WinThreadImpl {
class Thread final : public FEXCore::Threads::Thread {
public:
Thread(FEXCore::Threads::ThreadFunc Func, void* Arg)
: UserFunc {Func}
, UserArg {Arg} {
// hide everything from guest, don't initialize anything, we'll do that manually in RunThread()
const ULONG CreateFlags = THREAD_CREATE_FLAGS_SKIP_THREAD_ATTACH | THREAD_CREATE_FLAGS_HIDE_FROM_DEBUGGER |
THREAD_CREATE_FLAGS_SKIP_LOADER_INIT | THREAD_CREATE_FLAGS_BYPASS_PROCESS_FREEZE;
NTSTATUS Status = NtCreateThreadEx(&Handle, THREAD_ALL_ACCESS, nullptr, GetCurrentProcess(),
reinterpret_cast<PRTL_THREAD_START_ROUTINE>(&Thread::RunThread), this, CreateFlags, 0, 0, 0, nullptr);
if (Status < 0) {
LogMan::Msg::EFmt("NtCreateThreadEx failed: 0x{:x}", static_cast<uint32_t>(Status));
Handle = nullptr;
}
}
bool joinable() override {
return Handle != nullptr;
}
bool join(void** ret) override {
if (!Handle) {
return false;
}
NtWaitForSingleObject(Handle, FALSE, nullptr);
if (ret) {
*ret = ReturnValue;
}
return true;
}
bool detach() override {
// almost no-op, it's already detached, just deref handle
if (Handle) {
NtClose(Handle);
Handle = nullptr;
}
return true;
}
bool IsSelf() override {
return TID == GetCurrentThreadId();
}
~Thread() override {
if (Handle) {
NtClose(Handle);
}
}
private:
static void RunThread(Thread* This) {
This->TID = GetCurrentThreadId();
// do initialization we skipped earlier here around the user entrypoint
FEX::Windows::InitCRTThread();
This->ReturnValue = This->UserFunc(This->UserArg);
FEX::Windows::DeinitCRTThread();
NtTerminateThread(GetCurrentThread(), 0);
}
FEXCore::Threads::ThreadFunc UserFunc;
void* UserArg;
HANDLE Handle {};
DWORD TID {};
void* ReturnValue {};
};
fextl::unique_ptr<FEXCore::Threads::Thread> CreateThread(FEXCore::Threads::ThreadFunc Func, void* Arg) {
return fextl::make_unique<Thread>(Func, Arg);
}
void CleanupAfterFork() {}
} // namespace WinThreadImpl
void SetupThreadHandlers() {
FEXCore::Threads::Pointers Ptrs = {
.CreateThread = WinThreadImpl::CreateThread,
.CleanupAfterFork = WinThreadImpl::CleanupAfterFork,
};
FEXCore::Threads::Thread::SetInternalPointers(Ptrs);
}
} // namespace FEX::Windows
+6
View File
@@ -0,0 +1,6 @@
// SPDX-License-Identifier: MIT
#pragma once
namespace FEX::Windows {
void SetupThreadHandlers();
} // namespace FEX::Windows
+29 -4
View File
@@ -111,10 +111,16 @@ DLLEXPORT_FUNC(HANDLE, CreateFileW,
DLLEXPORT_FUNC(WINBOOL, WriteFile,
(HANDLE hFile, const void* lpBuffer, DWORD nNumberOfBytesToWrite, LPDWORD lpNumberOfBytesWritten, LPOVERLAPPED lpOverlapped)) {
IO_STATUS_BLOCK IOSB;
LARGE_INTEGER ByteOffset;
PLARGE_INTEGER ByteOffsetPtr = nullptr;
if (lpOverlapped) {
UNIMPLEMENTED();
if (lpOverlapped->hEvent) {
UNIMPLEMENTED();
}
ByteOffset.QuadPart = (static_cast<LONGLONG>(lpOverlapped->OffsetHigh) << 32) | lpOverlapped->Offset;
ByteOffsetPtr = &ByteOffset;
}
NTSTATUS Status = NtWriteFile(hFile, nullptr, nullptr, nullptr, &IOSB, lpBuffer, nNumberOfBytesToWrite, nullptr, nullptr);
NTSTATUS Status = NtWriteFile(hFile, nullptr, nullptr, nullptr, &IOSB, lpBuffer, nNumberOfBytesToWrite, ByteOffsetPtr, nullptr);
if (lpNumberOfBytesWritten) {
*lpNumberOfBytesWritten = static_cast<DWORD>(IOSB.Information);
}
@@ -135,6 +141,18 @@ DLLEXPORT_FUNC(WINBOOL, WriteConsoleW,
UNIMPLEMENTED();
}
DLLEXPORT_FUNC(WINBOOL, GetFileSizeEx, (HANDLE hFile, PLARGE_INTEGER lpFileSize)) {
IO_STATUS_BLOCK IOSB;
FILE_STANDARD_INFORMATION StandardInfo;
if (NTSTATUS Status = NtQueryInformationFile(hFile, &IOSB, &StandardInfo, sizeof(StandardInfo), FileStandardInformation); Status) {
return WinAPIReturn(Status);
}
if (lpFileSize) {
*lpFileSize = StandardInfo.EndOfFile;
}
return true;
}
DLLEXPORT_FUNC(WINBOOL, SetFilePointerEx, (HANDLE hFile, LARGE_INTEGER liDistanceToMove, PLARGE_INTEGER lpNewFilePointer, DWORD dwMoveMethod)) {
IO_STATUS_BLOCK IOSB;
FILE_POSITION_INFORMATION PositionInfo;
@@ -173,10 +191,17 @@ DLLEXPORT_FUNC(WINBOOL, SetFilePointerEx, (HANDLE hFile, LARGE_INTEGER liDistanc
DLLEXPORT_FUNC(WINBOOL, ReadFile,
(HANDLE hFile, void* lpBuffer, DWORD nNumberOfBytesToRead, LPDWORD lpNumberOfBytesRead, LPOVERLAPPED lpOverlapped)) {
IO_STATUS_BLOCK IOSB;
LARGE_INTEGER ByteOffset;
PLARGE_INTEGER ByteOffsetPtr = nullptr;
if (lpOverlapped) {
UNIMPLEMENTED();
if (lpOverlapped->hEvent) {
// no actual async-overlapped, just enough for PRead
UNIMPLEMENTED();
}
ByteOffset.QuadPart = (static_cast<LONGLONG>(lpOverlapped->OffsetHigh) << 32) | lpOverlapped->Offset;
ByteOffsetPtr = &ByteOffset;
}
NTSTATUS Status = NtReadFile(hFile, nullptr, nullptr, nullptr, &IOSB, lpBuffer, nNumberOfBytesToRead, nullptr, nullptr);
NTSTATUS Status = NtReadFile(hFile, nullptr, nullptr, nullptr, &IOSB, lpBuffer, nNumberOfBytesToRead, ByteOffsetPtr, nullptr);
if (lpNumberOfBytesRead) {
*lpNumberOfBytesRead = static_cast<DWORD>(IOSB.Information);
}
+2
View File
@@ -38,6 +38,7 @@ $end_info$
#include "Common/TSOHandlerConfig.h"
#include "Common/ImageTracker.h"
#include "Common/InvalidationTracker.h"
#include "Common/Threads.h"
#include "Common/OvercommitTracker.h"
#include "Common/CPUFeatures.h"
#include "Common/Logging.h"
@@ -515,6 +516,7 @@ public:
void BTCpuProcessInit() {
FEX::Windows::InitCRTProcess();
FEX::Windows::SetupThreadHandlers();
const auto ExecutablePath = FEX::Windows::GetExecutableFilePath();
const auto ExecutableName = FEX::Windows::BaseName(ExecutablePath);
FEX::Config::LoadConfig(fextl::string {ExecutableName}, _environ, FEX::ReadPortabilityInformation());
+24
View File
@@ -407,6 +407,23 @@ typedef enum _SECTION_INHERIT {
ViewUnmap = 2,
} SECTION_INHERIT;
typedef NTSTATUS(WINAPI* PRTL_THREAD_START_ROUTINE)(void* Parameter);
typedef struct _PS_ATTRIBUTE {
ULONG_PTR Attribute;
SIZE_T Size;
union {
ULONG_PTR Value;
void* ValuePtr;
};
PSIZE_T ReturnLength;
} PS_ATTRIBUTE, *PPS_ATTRIBUTE;
typedef struct _PS_ATTRIBUTE_LIST {
SIZE_T TotalLength;
PS_ATTRIBUTE Attributes[1];
} PS_ATTRIBUTE_LIST, *PPS_ATTRIBUTE_LIST;
/* definitions of bits in the Feature set for the x86 processors */
#define CPU_FEATURE_VME 0x00000005 /* Virtual 86 Mode Extensions */
#define CPU_FEATURE_TSC 0x00000002 /* Time Stamp Counter available */
@@ -518,6 +535,12 @@ NTSTATUS WINAPI NtAllocateVirtualMemoryEx(HANDLE, PVOID*, SIZE_T*, ULONG, ULONG,
NTSTATUS WINAPI NtAllocateVirtualMemory(HANDLE, PVOID*, ULONG_PTR, SIZE_T*, ULONG, ULONG);
NTSTATUS WINAPI NtContinue(PCONTEXT, BOOLEAN);
NTSTATUS WINAPI NtCreateSection(HANDLE*, ACCESS_MASK, const OBJECT_ATTRIBUTES*, const LARGE_INTEGER*, ULONG, ULONG, HANDLE);
#define THREAD_CREATE_FLAGS_SKIP_THREAD_ATTACH 0x00000002
#define THREAD_CREATE_FLAGS_HIDE_FROM_DEBUGGER 0x00000004
#define THREAD_CREATE_FLAGS_SKIP_LOADER_INIT 0x00000020
#define THREAD_CREATE_FLAGS_BYPASS_PROCESS_FREEZE 0x00000040
NTSTATUS WINAPI NtCreateThreadEx(HANDLE*, ACCESS_MASK, OBJECT_ATTRIBUTES*, HANDLE, PRTL_THREAD_START_ROUTINE, void*, ULONG, ULONG_PTR,
SIZE_T, SIZE_T, PPS_ATTRIBUTE_LIST);
NTSTATUS WINAPI NtDelayExecution(BOOLEAN, const LARGE_INTEGER*);
NTSTATUS WINAPI NtDuplicateObject(HANDLE, HANDLE, HANDLE, PHANDLE, ACCESS_MASK, ULONG, ULONG);
NTSTATUS WINAPI NtFlushInstructionCache(HANDLE, LPCVOID, SIZE_T);
@@ -538,6 +561,7 @@ NTSTATUS WINAPI NtReadFile(HANDLE, HANDLE, PIO_APC_ROUTINE, PVOID, PIO_STATUS_BL
NTSTATUS WINAPI NtSetContextThread(HANDLE, const CONTEXT*);
NTSTATUS WINAPI NtSuspendThread(HANDLE, PULONG);
NTSTATUS WINAPI NtTerminateProcess(HANDLE, LONG);
NTSTATUS WINAPI NtTerminateThread(HANDLE, LONG);
NTSTATUS WINAPI NtUnlockFile(HANDLE, PIO_STATUS_BLOCK, PLARGE_INTEGER, PLARGE_INTEGER, ULONG);
NTSTATUS WINAPI NtWriteFile(HANDLE, HANDLE, PIO_APC_ROUTINE, PVOID, PIO_STATUS_BLOCK, const void*, ULONG, PLARGE_INTEGER, PULONG);
void WINAPI ProcessPendingCrossProcessEmulatorWork();