#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