summaryrefslogtreecommitdiffstats
diff options
context:
space:
mode:
-rw-r--r--obj.c4
-rw-r--r--queue.c129
-rw-r--r--queue.h20
-rwxr-xr-xstringbuf.c1
4 files changed, 129 insertions, 25 deletions
diff --git a/obj.c b/obj.c
index eb4a7b32..420850be 100644
--- a/obj.c
+++ b/obj.c
@@ -131,8 +131,8 @@ finalize_it:
rsRetVal objSerializeProp(rsCStrObj *pCStr, uchar *pszPropName, propertyType_t propType, void *pUsr)
{
DEFiRet;
- uchar *pszBuf;
- size_t lenBuf;
+ uchar *pszBuf = NULL;
+ size_t lenBuf = 0;
uchar szBuf[64];
assert(pCStr != NULL);
diff --git a/queue.c b/queue.c
index 17e33675..8b06c882 100644
--- a/queue.c
+++ b/queue.c
@@ -190,10 +190,75 @@ rsRetVal qDelLinkedList(queue_t *pThis, void **ppUsr)
/* -------------------- disk -------------------- */
+/* first, disk-queue internal utility functions */
+
+
+/* open a queue file */
+rsRetVal qDiskOpenFile(queue_t *pThis, queueFileDescription_t *pFile, int flags, mode_t mode)
+{
+ DEFiRet;
+ uchar *pszFile = NULL;
+
+ assert(pThis != NULL);
+ assert(pFile != NULL);
+ /* now open the write file */
+ CHKiRet(genFileName(&pszFile, pThis->tVars.disk.pszSpoolDir, pThis->tVars.disk.lenSpoolDir,
+ (uchar*) "mainq", 5, pFile->iCurrFileNum, (uchar*) "qf", 2));
+
+ pFile->fd = open((char*)pszFile, flags, mode); // TODO: open modes!
+ pFile->iCurrOffs = 0;
+
+ dbgprintf("Queue 0x%lx: opened file '%s' for %d as %d\n", (unsigned long) pThis, pszFile, flags, pFile->fd);
+
+finalize_it:
+ if(pszFile != NULL)
+ free(pszFile); /* no longer needed in any case (just for open) */
+
+ return iRet;
+}
+
+
+/* close a queue file */
+rsRetVal qDiskCloseFile(queue_t *pThis, queueFileDescription_t *pFile)
+{
+ DEFiRet;
+
+ assert(pThis != NULL);
+ assert(pFile != NULL);
+ dbgprintf("Queue 0x%lx: closing file %d\n", (unsigned long) pThis, pFile->fd);
+
+ close(pFile->fd); // TODO: error check
+ pFile->fd = -1;
+
+ return iRet;
+}
+
+
+/* switch to next queue file */
+rsRetVal qDiskNextFile(queue_t *pThis, queueFileDescription_t *pFile)
+{
+ DEFiRet;
+
+ assert(pThis != NULL);
+ assert(pFile != NULL);
+ CHKiRet(qDiskCloseFile(pThis, pFile));
+
+ /* we do modulo 1,000,000 so that the file number is always at most 6 digits. If we have a million
+ * or more queue files, something is awfully wrong and it is OK if we run into problems in that
+ * situation ;) -- rgerhards, 2008-01-09
+ */
+ pFile->iCurrFileNum = (pFile->iCurrFileNum + 1) % 1000000;
+
+finalize_it:
+ return iRet;
+}
+
+
+/* now come the disk mode queue driver functions */
+
rsRetVal qConstructDisk(queue_t *pThis)
{
DEFiRet;
- uchar *pszFile;
assert(pThis != NULL);
@@ -201,20 +266,17 @@ rsRetVal qConstructDisk(queue_t *pThis)
ABORT_FINALIZE(RS_RET_OUT_OF_MEMORY);
pThis->tVars.disk.lenSpoolDir = strlen((char*)pThis->tVars.disk.pszSpoolDir);
+ pThis->tVars.disk.iMaxFileSize = 1024; // TODO: configurable!
- /* now open the file */
- CHKiRet(genFileName(&pszFile, pThis->tVars.disk.pszSpoolDir, pThis->tVars.disk.lenSpoolDir,
- (uchar*) "mainq", 5, 1, (uchar*) "qf", 2));
-
- dbgprintf("Queue 0x%lx: opening file '%s'\n", (unsigned long) pThis, pszFile);
+ pThis->tVars.disk.fWrite.iCurrFileNum = 1;
+ pThis->tVars.disk.fWrite.iCurrOffs = 0;
+ pThis->tVars.disk.fWrite.fd = -1;
- pThis->tVars.disk.fd = open((char*)pszFile, O_RDWR|O_CREAT, 0600);
- dbgprintf("opened file %d\n", pThis->tVars.disk.fd);
+ pThis->tVars.disk.fRead.iCurrFileNum = 1;
+ pThis->tVars.disk.fRead.fd = -1;
+ pThis->tVars.disk.fWrite.iCurrOffs = 0;
finalize_it:
- if(pThis->tVars.disk.pszSpoolDir != NULL)
- free(pThis->tVars.disk.pszSpoolDir);
-
return iRet;
}
@@ -225,7 +287,13 @@ rsRetVal qDestructDisk(queue_t *pThis)
assert(pThis != NULL);
- close(pThis->tVars.disk.fd);
+ if(pThis->tVars.disk.fWrite.fd != -1)
+ qDiskCloseFile(pThis, &pThis->tVars.disk.fWrite);
+ if(pThis->tVars.disk.fRead.fd != -1)
+ qDiskCloseFile(pThis, &pThis->tVars.disk.fRead);
+
+ if(pThis->tVars.disk.pszSpoolDir != NULL)
+ free(pThis->tVars.disk.pszSpoolDir);
return iRet;
}
@@ -233,14 +301,23 @@ rsRetVal qDestructDisk(queue_t *pThis)
rsRetVal qAddDisk(queue_t *pThis, void* pUsr)
{
DEFiRet;
- int i;
+ int iWritten;
rsCStrObj *pCStr;
assert(pThis != NULL);
- dbgprintf("writing to file %d\n", pThis->tVars.disk.fd);
- CHKiRet((objSerialize(pUsr))(pUsr, &pCStr)); // TODO: hier weiter machen!
- i = write(pThis->tVars.disk.fd, rsCStrGetBufBeg(pCStr), rsCStrLen(pCStr));
- dbgprintf("write wrote %d bytes, errno: %d, err %s\n", i, errno, strerror(errno));
+
+ if(pThis->tVars.disk.fWrite.fd == -1)
+ CHKiRet(qDiskOpenFile(pThis, &pThis->tVars.disk.fWrite, O_RDWR|O_CREAT, 0600)); // TODO: open modes!
+
+ CHKiRet((objSerialize(pUsr))(pUsr, &pCStr));
+ iWritten = write(pThis->tVars.disk.fWrite.fd, rsCStrGetBufBeg(pCStr), rsCStrLen(pCStr));
+ dbgprintf("Queue 0x%lx: write wrote %d bytes, errno: %d, err %s\n", (unsigned long) pThis,
+ iWritten, errno, strerror(errno));
+ /* TODO: handle error case -- rgerhards, 2008-01-07 */
+
+ pThis->tVars.disk.fWrite.iCurrOffs += iWritten;
+ if(pThis->tVars.disk.fWrite.iCurrOffs >= pThis->tVars.disk.iMaxFileSize)
+ CHKiRet(qDiskNextFile(pThis, &pThis->tVars.disk.fWrite));
finalize_it:
return iRet;
@@ -250,8 +327,22 @@ rsRetVal qDelDisk(queue_t __attribute__((unused)) *pThis, void __attribute__((un
{
DEFiRet;
+ assert(pThis != NULL);
+
+ if(pThis->tVars.disk.fRead.fd == -1)
+ CHKiRet(qDiskOpenFile(pThis, &pThis->tVars.disk.fRead, O_RDONLY, 0600)); // TODO: open modes!
+
+ /* read here */
+ /* de-serialize here */
+
+ /* switch to next file when EOF is reached. We may also delete the last file in that case.
+ pThis->tVars.disk.fWrite.iCurrOffs += iWritten;
+ if(pThis->tVars.disk.fWrite.iCurrOffs >= pThis->tVars.disk.iMaxFileSize)
+ CHKiRet(qDiskNextFile(pThis, &pThis->tVars.disk.fWrite));
+ */
iRet = RS_RET_ERR;
+finalize_it:
return iRet;
}
@@ -434,7 +525,7 @@ rsRetVal queueConstruct(queue_t **ppThis, queueType_t qType, int iWorkerThreads,
pThis->iNumWorkerThreads = iWorkerThreads;
pThis->qType = qType;
- /* set type-specific handlers */
+ /* set type-specific handlers and other very type-specific things (we can not totally hide it...) */
switch(qType) {
case QUEUETYPE_FIXED_ARRAY:
pThis->qConstruct = qConstructFixedArray;
@@ -453,6 +544,8 @@ rsRetVal queueConstruct(queue_t **ppThis, queueType_t qType, int iWorkerThreads,
pThis->qDestruct = qDestructDisk;
pThis->qAdd = qAddDisk;
pThis->qDel = qDelDisk;
+ /* special handling */
+ pThis->iNumWorkerThreads = 1; /* we need exactly one worker */
break;
case QUEUETYPE_DIRECT:
pThis->qConstruct = qConstructDirect;
diff --git a/queue.h b/queue.h
index 8c8782d3..f0740c44 100644
--- a/queue.h
+++ b/queue.h
@@ -26,6 +26,18 @@
#include <pthread.h>
#include "obj.h"
+/* some information about disk files used by the queue. In the long term, we may
+ * export this settings to a separate file module - or not (if they are too
+ * queue-specific. I just thought I mention it here so that everyone is aware
+ * of this possibility. -- rgerhards, 2008-01-07
+ */
+typedef struct {
+ int fd; /* the file descriptor, -1 if closed */
+ int iCurrFileNum;/* current file number (NOT descriptor, but the number in the file name!) */
+ int iCurrOffs; /* current offset */
+} queueFileDescription_t;
+
+
/* queue types */
typedef enum {
QUEUETYPE_FIXED_ARRAY = 0,/* a simple queue made out of a fixed (initially malloced) array fast but memoryhog */
@@ -73,10 +85,10 @@ typedef struct queue_s {
size_t lenSpoolDir;
uchar *pszFilePrefix;
size_t lenFilePrefix;
- int iCurrFileNum; /* number of file currently processed */
- int fd; /* current file descriptor */
- long iWritePos; /* next write position offset */
- long iReadPos; /* next read position offset */
+ int iNumberFiles; /* how many files make up the queue? */
+ int iMaxFileSize; /* max size for a single queue file */
+ queueFileDescription_t fWrite; /* current file to be written */
+ queueFileDescription_t fRead; /* current file to be read */
} disk;
} tVars;
} queue_t;
diff --git a/stringbuf.c b/stringbuf.c
index 040ede77..b777c40d 100755
--- a/stringbuf.c
+++ b/stringbuf.c
@@ -168,7 +168,6 @@ static rsRetVal rsCStrExtendBuf(rsCStrObj *pThis, size_t iMinNeeded)
}
iNewSize += pThis->iBufSize; /* add current size */
-dbgprintf("extending string buffer, old %d, new %d\n", pThis->iBufSize, iNewSize);
/* and then allocate and copy over */
/* DEV debugging only: dbgprintf("extending string buffer, old %d, new %d\n", pThis->iBufSize, iNewSize); */
if((pNewBuf = (uchar*) malloc(iNewSize * sizeof(uchar))) == NULL)