working on timeout filesystem layer ...

This commit is contained in:
Glenn Maynard
2005-01-27 00:14:02 +00:00
parent c19144bb0c
commit 582a2ec27a
+271 -38
View File
@@ -50,6 +50,10 @@ enum ThreadRequest
REQ_OPEN, REQ_OPEN,
REQ_CLOSE, REQ_CLOSE,
REQ_GET_FILE_SIZE, REQ_GET_FILE_SIZE,
REQ_READ,
REQ_WRITE,
REQ_FLUSH,
REQ_COPY,
REQ_POPULATE_FILE_SET, REQ_POPULATE_FILE_SET,
REQ_FLUSH_DIR_CACHE, REQ_FLUSH_DIR_CACHE,
@@ -62,7 +66,7 @@ enum ThreadRequest
class ThreadedFileWorker class ThreadedFileWorker
{ {
public: public:
ThreadedFileWorker( CString sDriverName, RageFileDriver *pChildDriver ); ThreadedFileWorker( CString sPath );
~ThreadedFileWorker(); ~ThreadedFileWorker();
void SetTimeout( float fSeconds ); void SetTimeout( float fSeconds );
@@ -73,6 +77,10 @@ public:
RageFileBasic *Open( const CString &sPath, int iMode, int &iErr ); RageFileBasic *Open( const CString &sPath, int iMode, int &iErr );
void Close( RageFileBasic *pFile ); void Close( RageFileBasic *pFile );
int GetFileSize( RageFileBasic *&pFile ); int GetFileSize( RageFileBasic *&pFile );
int Read( RageFileBasic *&pFile, void *pBuf, int iSize, CString &sError );
int Write( RageFileBasic *&pFile, const void *pBuf, int iSize, CString &sError );
int Flush( RageFileBasic *&pFile, CString &sError );
RageFileBasic *Copy( RageFileBasic *&pFile, CString &sError );
bool FlushDirCache( const CString &sPath ); bool FlushDirCache( const CString &sPath );
bool PopulateFileSet( FileSet &fs, const CString &sPath ); bool PopulateFileSet( FileSet &fs, const CString &sPath );
@@ -104,22 +112,30 @@ private:
/* REQ_OPEN, REQ_POPULATE_FILE_SET, REQ_FLUSH_DIR_CACHE: */ /* REQ_OPEN, REQ_POPULATE_FILE_SET, REQ_FLUSH_DIR_CACHE: */
CString m_sRequestPath; /* in */ CString m_sRequestPath; /* in */
/* REQ_POPULATE_FILE_SET: */ /* REQ_OPEN, REQ_COPY: */
RageFileBasic *m_pResultFile; /* out */ RageFileBasic *m_pResultFile; /* out */
/* REQ_POPULATE_FILE_SET: */ /* REQ_POPULATE_FILE_SET: */
FileSet m_ResultFileSet; /* out */ FileSet m_ResultFileSet; /* out */
/* REQ_OPEN: */ /* REQ_OPEN: */
int m_iResultErr; /* out */
int m_iRequestMode; /* in */ int m_iRequestMode; /* in */
/* REQ_GET_FILE_SIZE: */ /* REQ_GET_FILE_SIZE, REQ_COPY: */
RageFileBasic *m_pRequestFile; /* in */ RageFileBasic *m_pRequestFile; /* in */
// CString m_sResultError; /* out */ /* REQ_OPEN, REQ_GET_FILE_SIZE, REQ_READ */
/* REQ_GET_FILE_SIZE */ int m_iResultRequest; /* out */
int m_iResultReturn; /* out */
/* REQ_READ, REQ_WRITE */
int m_iRequestSize; /* in */
CString m_sResultError; /* out */
/* REQ_READ */
char *m_pResultBuffer; /* out */
/* REQ_WRITE */
char *m_pRequestBuffer; /* in */
}; };
static vector<ThreadedFileWorker *> g_apWorkers; static vector<ThreadedFileWorker *> g_apWorkers;
@@ -135,16 +151,21 @@ void RageFileDriverTimeout::SetTimeout( float fSeconds )
} }
ThreadedFileWorker::ThreadedFileWorker( CString sDriverName, RageFileDriver *pChildDriver ): ThreadedFileWorker::ThreadedFileWorker( CString sPath ):
m_WorkerEvent( sDriverName + "worker event" ) m_WorkerEvent( sPath + "worker event" )
{ {
/* Grab a reference to the child driver. We'll operate on it directly. */
m_pChildDriver = FILEMAN->GetFileDriver( sPath );
if( m_pChildDriver == NULL )
LOG->Warn( "ThreadedFileWorker: Mountpoint \"%s\" not found", sPath.c_str() );
m_Request = REQ_INVALID; m_Request = REQ_INVALID;
m_bTimedOut = false; m_bTimedOut = false;
m_pChildDriver = pChildDriver;
m_pResultFile = NULL; m_pResultFile = NULL;
m_pRequestFile = NULL; m_pRequestFile = NULL;
m_pResultBuffer = NULL;
m_WorkerThread.SetName( "ThreadedFileWorker fileset worker \"" + sDriverName + "\"" ); m_WorkerThread.SetName( "ThreadedFileWorker fileset worker \"" + sPath + "\"" );
m_WorkerThread.Create( StartWorkerMain, this ); m_WorkerThread.Create( StartWorkerMain, this );
g_apWorkersMutex.Lock(); g_apWorkersMutex.Lock();
@@ -172,6 +193,9 @@ ThreadedFileWorker::~ThreadedFileWorker()
DoRequest( REQ_SHUTDOWN ); DoRequest( REQ_SHUTDOWN );
m_WorkerThread.Wait(); m_WorkerThread.Wait();
if( m_pChildDriver != NULL )
FILEMAN->ReleaseFileDriver( m_pChildDriver );
/* Unregister ourself. */ /* Unregister ourself. */
g_apWorkersMutex.Lock(); g_apWorkersMutex.Lock();
for( unsigned i = 0; i < g_apWorkers.size(); ++i ) for( unsigned i = 0; i < g_apWorkers.size(); ++i )
@@ -223,6 +247,11 @@ bool ThreadedFileWorker::DoRequest( ThreadRequest r )
* completes. */ * completes. */
m_bTimedOut = true; m_bTimedOut = true;
} }
else
{
/* The operation completed; make sure that the worker thread doesn't erase the file. */
m_pRequestFile = NULL;
}
m_WorkerEvent.Unlock(); m_WorkerEvent.Unlock();
return bRequestFinished; return bRequestFinished;
@@ -256,7 +285,7 @@ void ThreadedFileWorker::WorkerMain()
case REQ_OPEN: case REQ_OPEN:
ASSERT( m_pResultFile == NULL ); ASSERT( m_pResultFile == NULL );
ASSERT( !m_sRequestPath.empty() ); ASSERT( !m_sRequestPath.empty() );
m_pResultFile = m_pChildDriver->Open( m_sRequestPath, m_iRequestMode, m_iResultErr ); m_pResultFile = m_pChildDriver->Open( m_sRequestPath, m_iRequestMode, m_iResultRequest );
break; break;
case REQ_CLOSE: case REQ_CLOSE:
@@ -266,7 +295,32 @@ void ThreadedFileWorker::WorkerMain()
case REQ_GET_FILE_SIZE: case REQ_GET_FILE_SIZE:
ASSERT( m_pRequestFile != NULL ); ASSERT( m_pRequestFile != NULL );
m_iResultReturn = m_pRequestFile->GetFileSize(); m_iResultRequest = m_pRequestFile->GetFileSize();
break;
case REQ_READ:
ASSERT( m_pRequestFile != NULL );
ASSERT( m_pResultBuffer != NULL );
m_iResultRequest = m_pRequestFile->Read( m_pResultBuffer, m_iRequestSize );
m_sResultError = m_pRequestFile->GetError();
break;
case REQ_WRITE:
ASSERT( m_pRequestFile != NULL );
ASSERT( m_pRequestBuffer != NULL );
m_iResultRequest = m_pRequestFile->Write( m_pRequestBuffer, m_iRequestSize );
m_sResultError = m_pRequestFile->GetError();
break;
case REQ_FLUSH:
ASSERT( m_pRequestFile != NULL );
m_iResultRequest = m_pRequestFile->Flush();
m_sResultError = m_pRequestFile->GetError();
break;
case REQ_COPY:
ASSERT( m_pRequestFile != NULL );
m_pResultFile = m_pRequestFile->Copy();
break; break;
case REQ_POPULATE_FILE_SET: case REQ_POPULATE_FILE_SET:
@@ -290,6 +344,13 @@ void ThreadedFileWorker::WorkerMain()
if( m_bTimedOut ) if( m_bTimedOut )
{ {
/* The event timed out. Clean up any residue from the last action. */ /* The event timed out. Clean up any residue from the last action. */
if( m_pRequestFile != NULL )
{
m_apDeletedFiles.push_back( m_pResultFile );
m_pResultFile = NULL;
break;
}
if( m_pResultFile != NULL ) if( m_pResultFile != NULL )
{ {
m_apDeletedFiles.push_back( m_pResultFile ); m_apDeletedFiles.push_back( m_pResultFile );
@@ -297,6 +358,11 @@ void ThreadedFileWorker::WorkerMain()
break; break;
} }
delete m_pRequestBuffer;
m_pRequestBuffer = NULL;
delete m_pResultBuffer;
m_pResultBuffer = NULL;
/* Clear the time-out flag, indicating that we can work again. */ /* Clear the time-out flag, indicating that we can work again. */
m_bTimedOut = false; m_bTimedOut = false;
} }
@@ -314,6 +380,12 @@ void ThreadedFileWorker::WorkerMain()
RageFileBasic *ThreadedFileWorker::Open( const CString &sPath, int iMode, int &iErr ) RageFileBasic *ThreadedFileWorker::Open( const CString &sPath, int iMode, int &iErr )
{ {
if( m_pChildDriver == NULL )
{
iErr = ENODEV;
return NULL;
}
/* If we're currently in a timed-out state, fail. */ /* If we're currently in a timed-out state, fail. */
if( m_bTimedOut ) if( m_bTimedOut )
{ {
@@ -330,7 +402,7 @@ RageFileBasic *ThreadedFileWorker::Open( const CString &sPath, int iMode, int &i
return false; return false;
} }
iErr = m_iResultErr; iErr = m_iResultRequest;
RageFileBasic *pRet = m_pResultFile; RageFileBasic *pRet = m_pResultFile;
m_pResultFile = NULL; m_pResultFile = NULL;
@@ -339,6 +411,8 @@ RageFileBasic *ThreadedFileWorker::Open( const CString &sPath, int iMode, int &i
void ThreadedFileWorker::Close( RageFileBasic *pFile ) void ThreadedFileWorker::Close( RageFileBasic *pFile )
{ {
ASSERT( m_pChildDriver != NULL ); /* how did you get a file to begin with? */
/* We need to delete files passed to us by Close, even if we're in a timed-out state. */ /* We need to delete files passed to us by Close, even if we're in a timed-out state. */
m_WorkerEvent.Lock(); m_WorkerEvent.Lock();
m_apDeletedFiles.push_back( pFile ); m_apDeletedFiles.push_back( pFile );
@@ -348,14 +422,18 @@ void ThreadedFileWorker::Close( RageFileBasic *pFile )
int ThreadedFileWorker::GetFileSize( RageFileBasic *&pFile ) int ThreadedFileWorker::GetFileSize( RageFileBasic *&pFile )
{ {
ASSERT( m_pChildDriver != NULL ); /* how did you get a file to begin with? */
/* If we're currently in a timed-out state, fail. */ /* If we're currently in a timed-out state, fail. */
if( m_bTimedOut ) if( m_bTimedOut )
{ {
this->Close( pFile ); this->Close( pFile );
pFile = NULL; pFile = NULL;
return -1;
} }
if( pFile == NULL )
return -1;
m_pRequestFile = pFile; m_pRequestFile = pFile;
if( !DoRequest(REQ_GET_FILE_SIZE) ) if( !DoRequest(REQ_GET_FILE_SIZE) )
@@ -365,11 +443,149 @@ int ThreadedFileWorker::GetFileSize( RageFileBasic *&pFile )
return -1; return -1;
} }
return m_iResultReturn; return m_iResultRequest;
} }
int ThreadedFileWorker::Read( RageFileBasic *&pFile, void *pBuf, int iSize, CString &sError )
{
ASSERT( m_pChildDriver != NULL ); /* how did you get a file to begin with? */
/* If we're currently in a timed-out state, fail. */
if( m_bTimedOut )
{
this->Close( pFile );
pFile = NULL;
}
if( pFile == NULL )
{
sError = "Operation timed out";
return -1;
}
m_pRequestFile = pFile;
m_iRequestSize = iSize;
m_pResultBuffer = new char[iSize];
if( !DoRequest(REQ_READ) )
{
/* If we time out, we can no longer access pFile. */
sError = "Operation timed out";
pFile = NULL;
return -1;
}
int iGot = m_iResultRequest;
memcpy( pBuf, m_pResultBuffer, iGot );
delete [] m_pResultBuffer;
return iGot;
}
int ThreadedFileWorker::Write( RageFileBasic *&pFile, const void *pBuf, int iSize, CString &sError )
{
ASSERT( m_pChildDriver != NULL ); /* how did you get a file to begin with? */
/* If we're currently in a timed-out state, fail. */
if( m_bTimedOut )
{
this->Close( pFile );
pFile = NULL;
}
if( pFile == NULL )
{
sError = "Operation timed out";
return -1;
}
m_pRequestFile = pFile;
m_iRequestSize = iSize;
m_pRequestBuffer = new char[iSize];
memcpy( m_pRequestBuffer, pBuf, iSize );
if( !DoRequest(REQ_WRITE) )
{
/* If we time out, we can no longer access pFile. */
sError = "Operation timed out";
pFile = NULL;
return -1;
}
int iGot = m_iResultRequest;
delete [] m_pResultBuffer;
return iGot;
}
int ThreadedFileWorker::Flush( RageFileBasic *&pFile, CString &sError )
{
ASSERT( m_pChildDriver != NULL ); /* how did you get a file to begin with? */
/* If we're currently in a timed-out state, fail. */
if( m_bTimedOut )
{
this->Close( pFile );
pFile = NULL;
}
if( pFile == NULL )
{
sError = "Operation timed out";
return -1;
}
m_pRequestFile = pFile;
if( !DoRequest(REQ_FLUSH) )
{
/* If we time out, we can no longer access pFile. */
sError = "Operation timed out";
pFile = NULL;
return -1;
}
return m_iResultRequest;
}
RageFileBasic *ThreadedFileWorker::Copy( RageFileBasic *&pFile, CString &sError )
{
ASSERT( m_pChildDriver != NULL ); /* how did you get a file to begin with? */
/* If we're currently in a timed-out state, fail. */
if( m_bTimedOut )
{
this->Close( pFile );
pFile = NULL;
}
if( pFile == NULL )
{
sError = "Operation timed out";
return NULL;
}
m_pRequestFile = pFile;
if( !DoRequest(REQ_COPY) )
{
/* If we time out, we can no longer access pFile. */
sError = "Operation timed out";
pFile = NULL;
return NULL;
}
RageFileBasic *pRet = m_pResultFile;
m_pResultFile = NULL;
return pRet;
}
bool ThreadedFileWorker::PopulateFileSet( FileSet &fs, const CString &sPath ) bool ThreadedFileWorker::PopulateFileSet( FileSet &fs, const CString &sPath )
{ {
if( m_pChildDriver == NULL )
return false;
/* If we're currently in a timed-out state, fail. */ /* If we're currently in a timed-out state, fail. */
if( m_bTimedOut ) if( m_bTimedOut )
return false; return false;
@@ -391,6 +607,9 @@ bool ThreadedFileWorker::PopulateFileSet( FileSet &fs, const CString &sPath )
bool ThreadedFileWorker::FlushDirCache( const CString &sPath ) bool ThreadedFileWorker::FlushDirCache( const CString &sPath )
{ {
if( m_pChildDriver == NULL )
return false;
/* If we're currently in a timed-out state, fail. */ /* If we're currently in a timed-out state, fail. */
if( m_bTimedOut ) if( m_bTimedOut )
return false; return false;
@@ -432,57 +651,80 @@ public:
RageFileBasic *Copy() const RageFileBasic *Copy() const
{ {
// XXX CString sError;
return NULL; RageFileBasic *pRet = m_pWorker->Copy( m_pFile, sError );
if( m_pFile == NULL )
{
// SetError( "Operation timed out" );
return NULL;
}
if( pRet == NULL )
{
// SetError( sError );
return NULL;
}
return pRet;
} }
protected: protected:
int ReadInternal( void *pBuffer, size_t iBytes ) int ReadInternal( void *pBuffer, size_t iBytes )
{ {
// XXX CString sError;
int iRet = m_pWorker->Read( m_pFile, pBuffer, iBytes, sError );
if( m_pFile == NULL ) if( m_pFile == NULL )
{ {
SetError( "Operation timed out" ); SetError( "Operation timed out" );
return -1; return -1;
} }
int iRet = m_pFile->Read( pBuffer, iBytes );
if( iRet == -1 ) if( iRet == -1 )
m_pFile->GetError(); SetError( sError );
return iRet; return iRet;
} }
int WriteInternal( const void *pBuffer, size_t iBytes ) int WriteInternal( const void *pBuffer, size_t iBytes )
{ {
// XXX CString sError;
int iRet = m_pWorker->Write( m_pFile, pBuffer, iBytes, sError );
if( m_pFile == NULL ) if( m_pFile == NULL )
{ {
SetError( "Operation timed out" ); SetError( "Operation timed out" );
return -1; return -1;
} }
int iRet = m_pFile->Write( pBuffer, iBytes );
if( iRet == -1 ) if( iRet == -1 )
m_pFile->GetError(); SetError( sError );
return iRet; return iRet;
} }
int FlushInternal() int FlushInternal()
{ {
// XXX CString sError;
int iRet = m_pWorker->Flush( m_pFile, sError );
if( m_pFile == NULL ) if( m_pFile == NULL )
{ {
SetError( "Operation timed out" ); SetError( "Operation timed out" );
return -1; return -1;
} }
int iRet = m_pFile->Flush();
if( iRet == -1 ) if( iRet == -1 )
m_pFile->GetError(); SetError( sError );
return iRet; return iRet;
} }
RageFileBasic *m_pFile; /* Mutable because the const operation Copy() timing out results in the original
* object no longer being available. */
mutable RageFileBasic *m_pFile;
ThreadedFileWorker *m_pWorker; ThreadedFileWorker *m_pWorker;
/* GetFileSize isn't allowed to fail, so cache the file size on load. */ /* GetFileSize isn't allowed to fail, so cache the file size on load. */
@@ -518,10 +760,7 @@ private:
RageFileDriverTimeout::RageFileDriverTimeout( CString sPath ): RageFileDriverTimeout::RageFileDriverTimeout( CString sPath ):
RageFileDriver( new TimedFilenameDB() ) RageFileDriver( new TimedFilenameDB() )
{ {
/* Grab a reference to the child driver. We'll operate on it directly. */ m_pWorker = new ThreadedFileWorker( sPath );
// XXX let ThreadedFileWorker own m_pChild
m_pChild = FILEMAN->GetFileDriver( sPath );
m_pWorker = new ThreadedFileWorker( sPath, m_pChild );
((TimedFilenameDB *) FDB)->SetWorker( m_pWorker ); ((TimedFilenameDB *) FDB)->SetWorker( m_pWorker );
} }
@@ -552,18 +791,12 @@ RageFileBasic *RageFileDriverTimeout::Open( const CString &sPath, int iMode, int
void RageFileDriverTimeout::FlushDirCache( const CString &sPath ) void RageFileDriverTimeout::FlushDirCache( const CString &sPath )
{ {
if( m_pChild == NULL )
return;
m_pWorker->FlushDirCache( sPath ); m_pWorker->FlushDirCache( sPath );
} }
RageFileDriverTimeout::~RageFileDriverTimeout() RageFileDriverTimeout::~RageFileDriverTimeout()
{ {
delete m_pWorker; delete m_pWorker;
if( m_pChild != NULL )
FILEMAN->ReleaseFileDriver( m_pChild );
} }
static struct FileDriverEntry_Timeout: public FileDriverEntry static struct FileDriverEntry_Timeout: public FileDriverEntry