#ifndef ETHREAD_H
#define ETHREAD_H

#ifndef _WIN32
 #include <pthread.h>
#else
 #include <winsock2.h>
 #include <Windows.h>
 #define pthread_t HANDLE
// #define pthread_mutex_t HANDLE   // can't use event mutex because of condition signal requires a SRW type lock
 #define pthread_mutex_t SRWLOCK
 #define pthread_cond_t CONDITION_VARIABLE
#endif

#include "efunc.h"
#include <queue>

class econdsig;

enum EMUTEX_TYPE { EMUTEX_RECURSIVE };

class emutex
{
 private:
  pthread_mutex_t _mutex;
 public:
  emutex(int type);
  emutex();
  ~emutex();

  void lock() const;
  void unlock() const;
  bool trylock() const;

  friend class econdsig;
};

class econdsig
{
 private:
  pthread_cond_t _cond;
 public:
  econdsig();
  ~econdsig();

  void wait(emutex& mutex);
  void broadcast();
  void signal();
};

//void* ethread_run(void *ptinfo);
//pthread_t ethread_create(const efunc& func,const evararray& args);

class ethread
{
 protected:
#ifdef _WIN32
  static DWORD WINAPI entrypoint(LPVOID);
#else
  static void *entrypoint(void*);
#endif
  emutex mutex;
  pthread_t _pthread;

  bool _stopThread;
  bool _pausedThread;
  econdsig condReady;
  econdsig condCanRun;

  int _runThread();
  virtual void _runTask()=0;
 public:
  ethread();
  ~ethread();

  void wake();
  bool trywake();

  void stop();

  void wait();
  bool isBusy();
};

class ethreadFunc : public ethread
{
 private:
  efunc _func;
  evararray _args;
  virtual void _runTask();
 public:
  bool tryrun(const efunc& func,const evararray& args=evararray());
  void run(const efunc& func,const evararray& args=evararray());
};

template <class T>
class emsgQueue
{
 public:
  emutex mutex;
  econdsig sigCanAdd;
  econdsig sigCanGet;
  econdsig sigWait;

  int maxSize;
  emsgQueue();

  queue<T> _queue;
  void add(const T& msg);
  T get();
};

template <class T>
emsgQueue<T>::emsgQueue(): maxSize(20) {}

template <class T>
void emsgQueue<T>::add(const T& msg)
{
  mutex.lock();
  while (_queue.size()>=maxSize)
    sigCanAdd.wait(mutex);
  _queue.push(msg);
  sigCanGet.signal();
  mutex.unlock();
}

template <class T>
T emsgQueue<T>::get()
{
  mutex.lock();
  while (_queue.size()==0)
    sigCanGet.wait(mutex);
  T res=_queue.pop();
  sigCanAdd.broadcast();
  mutex.unlock();
}


class ethreads
{
 public:
  ebasicarray<ethreadFunc*> threads;

  ~ethreads();

  void run(const efunc& func,const evararray& args=evararray(),int nthreads=0);
  void setThreads(int nthreads=1);
  void wait();
  void stop();
};


class etaskQueue;
class etaskBase;

class eworker
{
 protected:
  etaskQueue *tqueue;
 public:
  int i;

  virtual ~eworker();
  void setQueue(etaskQueue& tqueue);
  virtual void dispatch()=0;
  virtual void execute(etaskBase& task,const efunc& func,const evararray& args)=0;
};

class ethreadWorker : public ethread,public eworker
{
 private:
  virtual void _runTask();
 public:
  virtual void dispatch();
  virtual void execute(etaskBase& task,const efunc& func,const evararray& args);
};

class etaskBase
{
 public:
  emutex mutex;
  econdsig condWait;

  etaskQueue *_tqueue;

  etaskBase();

  void setQueue(etaskQueue& tqueue);
  bool hasQueue() const;
  evar waitResult();

  virtual int queuedCount() const=0;
  virtual int runningCount() const=0;
  virtual bool isDone() const=0;
  virtual void run(eworker& worker)=0;
  virtual void result(eworker& worker,const evararray& args,const evar& result)=0;
  virtual void error(eworker& worker,const evararray& args)=0;
  virtual void wait()=0;
  virtual evar getResult() const=0;
};

class etask : public etaskBase
{
 public:
  int _complete;
  int _running;
  int _total;

  evar _result;
  efunc _func;
  evararray _args;

 
  etask(const efunc& func,const evararray& args);

  virtual int queuedCount() const;
  virtual int runningCount() const;
  virtual bool isDone() const;
  virtual void run(eworker& worker);
  virtual void result(eworker& worker,const evararray& args,const evar& result);
  virtual void error(eworker& worker,const evararray& args);
  virtual void wait();
  virtual evar getResult() const;
};

class etaskArray : public etask
{
 public:
  evararray _resultarr;

  etaskArray(const efunc& func,const evararray& args,int count);

  virtual void run(eworker& worker);
  virtual void result(eworker& worker,const evararray& args,const evar& result);
  virtual void error(eworker& worker,const evararray& args);
};

class etaskApply : public etask
{
 public:
  evararray& _resultarr;

  etaskApply(const efunc& func,const evararray& args,evararray& arr);

  virtual void run(eworker& worker);
  virtual void result(eworker& worker,const evararray& args,const evar& result);
  virtual void error(eworker& worker,const evararray& args);
};


ostream& operator<<(ostream& stream,const etaskBase& task);


class etaskQueue
{
 public:
  list<etaskBase*> queue;
  ebasicarray<eworker*> workers;

  emutex mutex;
  econdsig condWait;

  void setThreads(int count);

  void add(etaskBase* task);
  bool get(eworker& worker);
  void wake();
  void wait();

//  void checkCompleted();
  void taskCompleted(etaskBase& task);
};

/*
class etask
{
 public:
  int status;
  efunc func;
  evararray args;
  evar result;

  inline void setRunning(){ status=1; }
  inline void setDone(){ status=2; }

  inline bool isPending(){ return(status==0); }
  inline bool isRunning(){ return(status==1); }
  inline bool isDone(){ return(status==2); }
  
  etask(const efunc& func,const evararray& args);
};

class etaskman;

class etaskthread
{
 protected:
  static void *entrypoint(void*);
 private:
  etaskman& taskman;

  pthread_t _pthread;

  evar _runJob(etask& task);
  int _runThread();
 public:
  etaskthread(etaskman& taskman);
  ~etaskthread();

  friend class etaskman;
};

class etaskman
{
 private:
  emutex runThreadsMutex;
  econdsig finishedThreadsCond;
  econdsig runThreadsCond;

  int runningThreads;
  int firstPendingTask;
 public:
  earray<etaskthread> threads;
  earray<etask>   tasks;

  etaskman();
  ~etaskman();

  efunc onTaskDone;
  efunc onAllDone;

  void doAllTasksDone();

  void createThread(int n=1);
  etask& addTask(const efunc& func,const evararray& args);
  etask* getTask(etaskthread& thread);
  void wait();
};
*/

#endif

