#ifndef DPMULTIPLEX_HEADER #define DPMULTIPLEX_HEADER 1 //////////////////////////////////////////////////////// //! rcsid="$Id: Multiplex.hh,v 1.12 2000/12/13 12:17:15 ees1cg Exp $" //! example=exDPMultiplex.cc //! docentry="Data Processing" //! file="amma/StdType/DataProc/Multiplex.hh" //! lib=DataProc //! author="Charles Galambos" //! date="04/07/98" #include "amma/DP/Port.hh" #include "amma/SArray1d.hh" #define DPMDEBUG 0 #if DPMDEBUG #define ONDEBUG(x) x #else #define ONDEBUG(x) #endif /////////////////////////////// //! userlevel=Develop //: Multiplex an operation body. // This is indended to be used with DPThreadC to // do multiprocessing.

// This will preserve the order of the data in the // stream. template class DPMultiplexBodyC : public DPIOPortBodyC { public: DPMultiplexBodyC(IntT num,const DPIOPortC & nproc); //: Constructor. // NB. This only uses nproc as a template, it is not actually // used for processing. This avoids possible problems of copying // data out of a running thread.

// ** nproc MUST not have a running thread **. DPMultiplexBodyC(const SArray1dC > &nprocs); //: Constructor. DPMultiplexBodyC(istream &in); //: Stream constructor. DPMultiplexBodyC(const DPMultiplexBodyC &oth); //: Copy constructor. // Makes a deep copy of procs array. virtual BooleanT IsPutReady() const; //: Is some data ready ? // TRUE = yes. virtual BooleanT IsGetReady() const; //: Is some data ready ? // TRUE = yes. virtual void PutEOS(); //: Put End Of Stream marker. virtual BooleanT Put(const OutT &dat); //: Put data. virtual IntT PutArray(const SArray1dC &data); //: Put an array of data to stream. // returns the number of elements processed. virtual InT Get(); //: Get next piece of data. // May block if not ready, or it will return a constructed // with the default constructor. inline IntT GetArray(SArray1dC &data); //: Get an array of data from stream. // returns the number of elements processed. virtual BooleanT Get(InT &buff); //: Get next piece of data. // May block if not ready, or it will return a constructed // with the default constructor. virtual BooleanT IsGetEOS() const; //: Is get EOS ? virtual BodyRefCounterVC &Copy() const; //: Make a deep copy of object. virtual BooleanT Save(ostream &out) const; //: Save to ostream. private: IntT in,out; SArray1dC > procs; }; ////////////////////////////////// //! userlevel=Normal //: Multiplex an operation. // This is indended to be used with DPThreadC to // do multiprocessing. template class DPMultiplexC : public DPIOPortC { public: DPMultiplexC(); //: Default constructor. DPMultiplexC(IntT num,const DPIOPortC & nproc); //: Constructor. DPMultiplexC(const SArray1dC > &procs); //: Constructor. DPMultiplexC(const DPMultiplexC &oth); //: Copy constructor. }; /////////////////////////////// template DPIOPortC DPMultiplex(IntT num,const DPIOPortC & nproc) { return DPMultiplexC(num,nproc); } template DPIOPortC DPMultiplex(const SArray1dC > &procs) { return DPMultiplexC(procs); } //////////////////////////////////// template DPMultiplexBodyC::DPMultiplexBodyC(IntT num,const DPIOPortC & nproc) : in(0), out(0), procs(num) { assert(nproc.IsValid()); for(IntT i = 0;i < num;i++) procs[i] = nproc.Copy(); } template DPMultiplexBodyC::DPMultiplexBodyC(const SArray1dC > &nprocs) : in(0), out(0), procs(nprocs) {} template DPMultiplexBodyC::DPMultiplexBodyC(istream &strmin) : DPIOPortBodyC(strmin) { strmin >> in >> out >> procs; } template DPMultiplexBodyC::DPMultiplexBodyC(const DPMultiplexBodyC &oth) : in(oth.in), out(oth.out), procs(oth.procs.Size()) { for(IntT i = 0;i < (IntT) oth.procs.Size();i++) procs[i] = oth.procs[i].Copy(); } template BooleanT DPMultiplexBodyC::IsPutReady() const { return procs[in].IsPutReady(); } template BooleanT DPMultiplexBodyC::IsGetReady() const { return procs[out].IsGetReady(); } template void DPMultiplexBodyC::PutEOS() { IntT i; // Put EOS's in proper sequence. for(i = in;i < (IntT) procs.Size();i++) procs[i].PutEOS(); for(i = 0;i < in;i++) procs[i].PutEOS(); } template BooleanT DPMultiplexBodyC::Put(const OutT &dat) { ONDEBUG(cerr << "DPMultiplexBodyC::Put(), in:" << in << " Started\n"); if(!procs[in].Put(dat)) { cerr << "WARNING: DPMultiplexBodyC::Put(), Failed. \n"; return FALSE; } ONDEBUG(cerr << "DPMultiplexBodyC::Put(), in:" << in << " Done\n"); in++; if(in >= (IntT) procs.Size()) in = 0; return TRUE; } template IntT DPMultiplexBodyC::PutArray(const SArray1dC &data) { for(SArray1dIterC it(data);it;it++) { if(!procs[in].Put(*it)) { cerr << "WARNING: DPMultiplexBodyC::PutArray(), Failed. \n"; return it.Index().V(); } in++; if(in >= (IntT) procs.Size()) in = 0; } return data.Size(); } template InT DPMultiplexBodyC::Get() { ONDEBUG(cerr << "DPMultiplexBodyC::Get(), out:" << out << " Started\n"); InT ret = procs[out].Get(); ONDEBUG(cerr << "DPMultiplexBodyC::Get(), out:" << out << " Done.\n"); out++; if(out >= (IntT) procs.Size()) out = 0; return ret; } template BooleanT DPMultiplexBodyC::Get(InT &buff) { ONDEBUG(cerr << "DPMultiplexBodyC::Get(), out:" << out << " Started\n"); if(!procs[out].Get(buff)) { cerr << "WARNING: DPMultiplexBodyC::Get(), Failed. \n"; return FALSE; } ONDEBUG(cerr << "DPMultiplexBodyC::Get(), out:" << out << " Done.\n"); out++; if(out >= (IntT) procs.Size()) out = 0; return TRUE; } template IntT DPMultiplexBodyC::GetArray(SArray1dC &data) { for(SArray1dIterC it(data);it;it++) { if(!procs[out].Get(*it)) { cerr << "WARNING: DPMultiplexBodyC::PutArray(), Failed. \n"; return it.Index().V(); } out++; if(out >= (IntT) procs.Size()) out = 0; } return data.Size(); } template BooleanT DPMultiplexBodyC::IsGetEOS() const { // No point in advancing the stream ptr. return procs[out].IsGetEOS(); } template BodyRefCounterVC & DPMultiplexBodyC::Copy() const { return *new DPMultiplexBodyC(*this); } template BooleanT DPMultiplexBodyC::Save(ostream &out) const { DPIOPortBodyC::Save(out); out << in << " " << out << " " << procs; return TRUE; } ////////////////////////////////////////////// template DPMultiplexC::DPMultiplexC() {} template DPMultiplexC::DPMultiplexC(IntT num,const DPIOPortC & nproc) : DPEntityC(*new DPMultiplexBodyC(num,nproc)) {} template DPMultiplexC::DPMultiplexC(const SArray1dC > &procs) : DPEntityC(*new DPMultiplexBodyC(procs)) {} template DPMultiplexC::DPMultiplexC(const DPMultiplexC &oth) : DPEntityC(oth) {} #ifdef ONDEBUG #undef ONDEBUG #endif #endif