#include "StdAfx.hpp"
#include "BackgroundJob.hpp"
#include "Framework.hpp" // Dla CUSTOM_EXCEPTIONS_TRY, CUSTOM_EXCEPTIONS_CATCH
#include <algorithm>

//HHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHH
// BkgJob

BkgJob::BkgJob(TYPE Type, MODE Mode) :
m_Type(Type),
m_Mode(Mode),
m_Priority(PRIORITY_NORMAL),
m_State(STATE_NONE)
{
}

BkgJob::~BkgJob()
{
	ASSERT_INT3_DEBUG(GetState() == STATE_NONE || IsFinished());
}

void BkgJob::SetPriority(PRIORITY Priority)
{
	PRIORITY OldPriority = m_Priority;
	if (Priority == OldPriority) return;

	MUTEX_LOCK(&g_BkgJobManager->m_JobQueueMutex);

	TYPE Type = GetType();
	bool Found = false;
	for (
		std::deque<BkgJob*>::iterator it = g_BkgJobManager->m_JobQueue[Type][OldPriority].begin();
		it != g_BkgJobManager->m_JobQueue[Type][OldPriority].end();
		++it)
	{
		if (*it == this)
		{
			g_BkgJobManager->m_JobQueue[Type][OldPriority].erase(it);
			Found = true;
			break;
		}
	}

	if (Found)
	{
		m_Priority = Priority;
		g_BkgJobManager->m_JobQueue[Type][Priority].push_back(this);
	}
}

void BkgJob::BoostPriority()
{
	PRIORITY Priority = PRIORITY_VERY_HIGH;
	PRIORITY OldPriority = m_Priority;

	MUTEX_LOCK(&g_BkgJobManager->m_JobQueueMutex);

	TYPE Type = GetType();
	bool Found = false;
	for (
		std::deque<BkgJob*>::iterator it = g_BkgJobManager->m_JobQueue[Type][OldPriority].begin();
		it != g_BkgJobManager->m_JobQueue[Type][OldPriority].end();
		++it)
	{
		if (*it == this)
		{
			g_BkgJobManager->m_JobQueue[Type][OldPriority].erase(it);
			Found = true;
			break;
		}
	}

	if (Found)
	{
		m_Priority = Priority;
		g_BkgJobManager->m_JobQueue[Type][Priority].push_front(this);
	}
}

void BkgJob::Join()
{
	BoostPriority();

	ASSERT_INT3_DEBUG(GetMode() == MODE_JOINABLE);
	ASSERT_INT3_DEBUG(GetState() != STATE_NONE);
	while (!IsFinished())
	{
		Sleep(0);
		g_BkgJobManager->Frame();
	}
}


//HHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHH
// BkgJobManagerThread

class BkgJobManagerThread : public common::Thread
{
public:
	BkgJobManagerThread(BkgJobManager *Manager, BkgJob::TYPE Type, uint Index);
	virtual ~BkgJobManagerThread();

protected:
	virtual void Run();

private:
	BkgJob::TYPE m_Type;
	uint m_Index;
};


BkgJobManagerThread::BkgJobManagerThread(BkgJobManager *Manager, BkgJob::TYPE Type, uint Index) :
m_Type(Type),
m_Index(Index)
{
}

BkgJobManagerThread::~BkgJobManagerThread()
{
}

void BkgJobManagerThread::Run()
{
	for (;;)
	{
		g_BkgJobManager->m_JobQueueSemaphores[m_Type]->P();
		BkgJob *Job = NULL;
		float Found = false;
		{
			MUTEX_LOCK(&g_BkgJobManager->m_JobQueueMutex);
			for (uint i = 0; i < BkgJob::PRIORITY_COUNT; i++)
			{
				if (!g_BkgJobManager->m_JobQueue[m_Type][i].empty())
				{
					Job = g_BkgJobManager->m_JobQueue[m_Type][i].front();
					g_BkgJobManager->m_JobQueue[m_Type][i].pop_front();
					Found = true;
					break;
				}
			}
		}
		ASSERT_INT3_DEBUG(Found);
		
		// Koniec wątku

		if (Job == NULL)
			break;

		// Wykonaj job
		ASSERT_INT3_DEBUG(Job->m_State == BkgJob::STATE_QUEUED);
		Job->m_State = BkgJob::STATE_WORKING;

		try
		{
			CUSTOM_EXCEPTIONS_TRY
			{
				Job->OnWork();
			}
			CUSTOM_EXCEPTIONS_CATCH
		}
		catch (common::Error &e)
		{
			e.Push("Exception caught in BkgJobManagerThread::Run", __TFILE__, __LINE__);
			Job->m_SavedError.reset(new common::Error(e));
		}

		// Zadanie skończone
		{
			MUTEX_LOCK(&g_BkgJobManager->m_JobQueueMutex);
			g_BkgJobManager->m_JobsDoneQueue.push_back(Job);
		}
	}
}


//HHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHHH
// BkgJobManager

BkgJobManager *g_BkgJobManager = NULL;

BkgJobManager::BkgJobManager() :
m_JobQueueMutex(0)
{
	ZeroMemory(&m_JobQueueSemaphores, sizeof(m_JobQueueSemaphores));

	for (uint i = 0; i < _countof(m_JobQueueSemaphores); i++)
		m_JobQueueSemaphores[i] = new common::Semaphore(0);
}

void BkgJobManager::Init()
{
	// Odpalenie wątków
	uint ComputationThreadCount = 2; // TODO - tyle ile rdzeni
	BkgJobManagerThread *NewThread;
	for (uint i = 0; i < ComputationThreadCount; i++)
	{
		NewThread = new BkgJobManagerThread(this, BkgJob::TYPE_COMPUTATION, i);
		m_Threads[BkgJob::TYPE_COMPUTATION].push_back(NewThread);
		NewThread->Start();
	}
	NewThread = new BkgJobManagerThread(this, BkgJob::TYPE_IO, 0);
	m_Threads[BkgJob::TYPE_IO].push_back(NewThread);
	NewThread->Start();
}

BkgJobManager::~BkgJobManager()
{
	// Wyślij wątkom polecenie zakończenia
	{
		MUTEX_LOCK(&m_JobQueueMutex);
		for (uint Type_i = 0; Type_i < BkgJob::TYPE_COUNT; Type_i++)
		{
			for (uint Thread_i = 0; Thread_i < m_Threads[Type_i].size(); Thread_i++)
			{
				m_JobQueue[Type_i][BkgJob::PRIORITY_VERY_LOW].push_back(NULL);
				m_JobQueueSemaphores[Type_i]->V();
			}
		}
	}

	// Poczekaj na zakończenie watków, zniszcz wątki
	for (uint Type_i = BkgJob::TYPE_COUNT; Type_i--; )
	{
		for (uint Thread_i = m_Threads[Type_i].size(); Thread_i--; )
		{
			m_Threads[Type_i][Thread_i]->Join();
			delete m_Threads[Type_i][Thread_i];
		}
	}

	// Dokończ pozostałe przetworzone joby
	Frame();

	for (uint i = _countof(m_JobQueueSemaphores); i--; )
		delete m_JobQueueSemaphores[i];
}

void BkgJobManager::Frame()
{
	std::vector<BkgJob*> DoneJobs;
	{
		MUTEX_LOCK(&m_JobQueueMutex);
		DoneJobs.resize(m_JobsDoneQueue.size());
		std::copy(m_JobsDoneQueue.begin(), m_JobsDoneQueue.end(), DoneJobs.begin());
		m_JobsDoneQueue.clear();
	}

	for (uint i = 0; i < DoneJobs.size(); i++)
	{
		BkgJob *Job = DoneJobs[i];
		try
		{
			CUSTOM_EXCEPTIONS_TRY
			{
				Job->OnWorkDone();
			}
			CUSTOM_EXCEPTIONS_CATCH
		}
		catch (common::Error &e)
		{
			e.Push("Exception caught in BkgJobManager::Frame", __TFILE__, __LINE__);
			ASSERT_INT3_DEBUG(Job->m_SavedError.get() == NULL);
			Job->m_SavedError.reset(new common::Error(e));
		}

		Job->m_State = (Job->m_SavedError.get() == NULL ? BkgJob::STATE_DONE : BkgJob::STATE_DONE_ERROR);

		if (Job->GetMode() == BkgJob::MODE_DETACHED)
			delete Job;
	}
}

void BkgJobManager::AddJob(BkgJob *Job)
{
	ASSERT_INT3_DEBUG(Job->GetState() == BkgJob::STATE_NONE);

	BkgJob::TYPE Type = Job->GetType();
	BkgJob::PRIORITY Priority = Job->GetPriority();

	{
		MUTEX_LOCK(&m_JobQueueMutex);
		m_JobQueue[Type][Priority].push_back(Job);
	}

	Job->m_State = BkgJob::STATE_QUEUED;

	m_JobQueueSemaphores[Type]->V();
}

void BkgJobManager::WaitForJobs()
{
	for (;;)
	{
		bool Empty = true;
		{
			MUTEX_LOCK(&m_JobQueueMutex);
			for (uint Type_i = 0; Type_i < BkgJob::TYPE_COUNT && Empty; Type_i++)
			{
				for (uint Priority_i = 0; Priority_i < BkgJob::PRIORITY_COUNT && Empty; Priority_i++)
				{
					if (!m_JobQueue[Type_i][Priority_i].empty())
						Empty = false;
				}
			}
		}

		if (Empty) return;

		Sleep(0);
		Frame();
	}
}

void BkgJobManager::WaitForJobs(BkgJob::TYPE Type)
{
	for (;;)
	{
		bool Empty = true;
		{
			MUTEX_LOCK(&m_JobQueueMutex);
			for (uint Priority_i = 0; Priority_i < BkgJob::PRIORITY_COUNT && Empty; Priority_i++)
			{
				if (!m_JobQueue[Type][Priority_i].empty())
					Empty = false;
			}
		}

		if (Empty) return;

		Sleep(0);
		Frame();
	}
}
