Sync branch with HEAD of trunk

This commit is contained in:
alexhudson
2007-07-26 22:01:34 +00:00
parent 422036f5a7
commit 3a873c389c
33 changed files with 982 additions and 886 deletions
+193 -155
View File
@@ -65,6 +65,49 @@ static int HandleDSN(FILE *data, FILE *control);
#define MAX_CHARS_IN_PDBSEARCH 512
#define FOPEN_CHECK(handle, path, mode) fopen_check(&(handle), (path), (mode), __LINE__)
FILE *
fopen_check(FILE **handle, char *path, char *mode, int line)
{
LogAssertF(*handle == NULL, "File handle already open on line %d", line);
*handle = fopen(path, mode);
return *handle;
}
#define FCLOSE_CHECK(f) fclose_check(&(f), __LINE__)
int
fclose_check(FILE **fh, int line)
{
int ret;
ret = fclose(*fh);
if (ret == 0) {
*fh = NULL;
} else {
LogFailureF("File close failed on line %d: %d", line, errno);
}
return ret;
}
#define UNLINK_CHECK(path) unlink_check((path), __LINE__)
int
unlink_check(char *path, int line)
{
int ret;
ret = unlink(path);
LogAssertF(ret == 0, "Unable to delete file %s on line %d: %d", path, line, errno);
return ret;
}
#define RENAME_CHECK(oldpath, newpath) rename_check((oldpath), (newpath), __LINE__)
int
rename_check(const char *oldpath, const char *newpath, int line)
{
int ret;
ret = rename(oldpath, newpath);
LogAssertF(ret == 0, "Unable to rename %s to %s on line %d: %d", oldpath, newpath, line, errno);
return ret;
}
static int
PDBSearch(char *doc, char *searchString)
{
@@ -236,9 +279,7 @@ SpoolEntryIDLock(unsigned long id)
XplSignalLocalSemaphore(Queue.spoolLocks.semaphores[SPOOL_LOCK_ARRAY_MASK & id]);
LoggerEvent(Agent.agent.loggingHandle, LOGGER_SUBSYSTEM_GENERAL, LOGGER_EVENT_PQE_LOCK_FULL, LOG_INFO, 0, NULL, NULL, id, 0, NULL, 0);
XplConsolePrintf("bongoqueue: Unable to lock spool entry %x; table full.\r\n", (unsigned int)id);
Log(LOG_INFO, "Unable to lock spool entry %x, table full", (unsigned int) id);
}
return(NULL);
@@ -300,7 +341,7 @@ AddPushAgent(QueueClient *client,
int count;
QueuePushClient *temp;
LoggerEvent(Agent.agent.loggingHandle, LOGGER_SUBSYSTEM_GENERAL, LOGGER_EVENT_ADD_QUEUE_AGENT, LOG_INFO, 0, NULL, NULL, XplHostToLittle(client->conn->socketAddress.sin_addr.s_addr), queue, &port, sizeof(port));
Log(LOG_INFO, "Adding client on host %s:%d to queue %d", LOGIP(client->conn->socketAddress), port, queue);
XplMutexLock(Queue.pushClients.lock);
@@ -317,10 +358,7 @@ AddPushAgent(QueueClient *client,
Queue.pushClients.array = temp;
Queue.pushClients.allocated += PUSHCLIENTALLOC;
} else {
LoggerEvent(Agent.agent.loggingHandle, LOGGER_SUBSYSTEM_GENERAL, LOGGER_EVENT_NMAP_OUT_OF_MEMORY, LOG_CRITICAL, 0, "Queue", identifier, (Queue.pushClients.allocated + PUSHCLIENTALLOC) * sizeof(QueuePushClient), 0, NULL, 0);
XplConsolePrintf("bongoqueue: Out of memory processing user %s, mailbox %s; request size: %d bytes\r\n", "bongoqueue:AddPushAgent",
identifier, (int)((Queue.pushClients.allocated + PUSHCLIENTALLOC) * sizeof(QueuePushClient)));
LogFailureF("Out of memory processing mailbox %s", identifier);
return(-1);
}
}
@@ -360,9 +398,9 @@ RemovePushAgentIndex(int index, BOOL force)
if ((Queue.pushClients.array[index].errorCount > MAX_PUSHCLIENTS_ERRORS) || force) {
if (force) {
LoggerEvent(Agent.agent.loggingHandle, LOGGER_SUBSYSTEM_GENERAL, LOGGER_EVENT_REREGISTER_QUEUE_AGENT, LOG_INFO, 0, NULL, NULL, XplHostToLittle(Queue.pushClients.array[index].address), Queue.pushClients.array[index].queue, &(Queue.pushClients.array[index].port), sizeof(Queue.pushClients.array[index].port));
Log(LOG_INFO, "Reregistered queue agent");
} else {
LoggerEvent(Agent.agent.loggingHandle, LOGGER_SUBSYSTEM_GENERAL, LOGGER_EVENT_REMOVED_QUEUE_AGENT, LOG_INFO, 0, NULL, NULL, XplHostToLittle(Queue.pushClients.array[index].address), Queue.pushClients.array[index].queue, &(Queue.pushClients.array[index].port), sizeof(Queue.pushClients.array[index].port));
Log(LOG_INFO, "Removed queue agent");
}
if (index < (Queue.pushClients.count - 1)) {
@@ -520,7 +558,7 @@ DeliverToStore(NMAPConnections *list,
(ccode = ConnFlush(nmap->conn)) == -1) {
nmap->error = TRUE;
fclose(fh);
FCLOSE_CHECK(fh);
return DELIVER_TRY_LATER;
}
@@ -620,9 +658,9 @@ ProcessQueueEntry(unsigned char *entryIn)
time_t date;
struct sockaddr_in saddr;
struct stat sb;
FILE *fh;
FILE *data;
FILE *newFH;
FILE *fh = NULL;
FILE *data = NULL;
FILE *newFH = NULL;
MIMEReportStruct *report = NULL;
MDBValueStruct *vs;
QueueClient *client;
@@ -647,12 +685,12 @@ StartOver:
entryID = strtol(entry, NULL, 16);
LoggerEvent(Agent.agent.loggingHandle, LOGGER_SUBSYSTEM_QUEUE, LOGGER_EVENT_PQE_START, LOG_DEBUG, entryID, NULL, NULL, queue, entryID, NULL, 0);
Log(LOG_DEBUG, "Processing entry %ld on queue %d", entryID, queue);
idLock = SpoolEntryIDLock(entryID);
if (idLock) {
sprintf(path, "%s/c%s.%03d", Conf.spoolPath, entry, queue);
fh = fopen(path, "r+b");
FOPEN_CHECK(fh, path, "r+b");
} else {
ProcessQueueEntryCleanUp(NULL, report);
return(FALSE);
@@ -661,7 +699,7 @@ StartOver:
if (fh) {
fgets(line, CONN_BUFSIZE, fh);
date = atoi(line + 1);
fclose(fh);
FCLOSE_CHECK(fh);
} else {
ProcessQueueEntryCleanUp(idLock, report);
return(FALSE);
@@ -670,6 +708,7 @@ StartOver:
/* We've got pre and post processing off queue entries - this is pre */
switch(queue) {
case Q_INCOMING: {
FILE *temp = NULL;
sb.st_size = -1;
qDate = NULL;
qFlags = NULL;
@@ -678,15 +717,15 @@ StartOver:
qFrom = NULL;
sprintf(path, "%s/w%s.%03d", Conf.spoolPath, entry, queue);
newFH = fopen(path, "wb");
FOPEN_CHECK(newFH, path, "wb");
sprintf(path, "%s/c%s.%03d",Conf.spoolPath, entry, queue);
if (newFH
&& (stat(path, &sb) == 0)
&& (sb.st_size > 8)
&& ((qEnvelope = (unsigned char *)MemMalloc(sb.st_size + 1)) != NULL)
&& ((fh = fopen(path, "rb")) != NULL)
&& (fread(qEnvelope, sizeof(unsigned char), sb.st_size, fh) == (size_t)sb.st_size)) {
&& ((temp = fopen(path, "rb")) != NULL)
&& (fread(qEnvelope, sizeof(unsigned char), sb.st_size, temp) == (size_t)sb.st_size)) {
/* Sort the control file as follows:
QUEUE_DATE
QUEUE_FLAGS
@@ -702,7 +741,7 @@ StartOver:
QUEUE_RECIP_REMOTE
QUEUE_THIRD_PARTY
*/
fclose(fh);
fclose(temp);
qEnvelope[sb.st_size] = '\0';
@@ -799,15 +838,15 @@ StartOver:
/* fixme - if a new queue entry has at least QUEUE_FROM
but no recipients we should bounce the message rather
than consuming it! */
fclose(newFH);
FCLOSE_CHECK(newFH);
unlink(path);
UNLINK_CHECK(path);
sprintf(path, "%s/w%s.%03d",Conf.spoolPath, entry, queue);
unlink(path);
UNLINK_CHECK(path);
sprintf(path, "%s/d%s.msg",Conf.spoolPath, entry);
unlink(path);
UNLINK_CHECK(path);
MemFree(qEnvelope);
@@ -868,7 +907,7 @@ StartOver:
MDBFreeValues(vs);
} else {
LoggerEvent(Agent.agent.loggingHandle, LOGGER_SUBSYSTEM_QUEUE, LOGGER_EVENT_PQE_FIND_OBJECT_FAILED, LOG_INFO, entryID, cur + 1, "", queue, 0, NULL, 0);
Log(LOG_INFO, "Entry %ld queue %d, can't find %s", entryID, queue, cur + 1);
if (ptr2) {
*ptr2 = ' ';
@@ -925,25 +964,25 @@ StartOver:
MemFree(qEnvelope);
fclose(newFH);
FCLOSE_CHECK(newFH);
unlink(path);
UNLINK_CHECK(path);
sprintf(path2, "%s/w%s.%03d", Conf.spoolPath, entry, queue);
rename(path2, path);
RENAME_CHECK(path2, path);
break;
}
if (!newFH) {
sprintf(path, "%s/w%s.%03d",Conf.spoolPath, entry, queue);
LoggerEvent(Agent.agent.loggingHandle, LOGGER_SUBSYSTEM_GENERAL, LOGGER_EVENT_FILE_OPEN_FAILURE, LOG_CRITICAL, entryID, path, __FILE__, __LINE__, 0, NULL, 0);
LogFailureF("File open failure. Entry %ld, path %s", entryID, path);
} else if (!qEnvelope) {
LoggerEvent(Agent.agent.loggingHandle, LOGGER_SUBSYSTEM_GENERAL, LOGGER_EVENT_NMAP_OUT_OF_MEMORY, LOG_CRITICAL, entryID, __FILE__, NULL, sb.st_size, 0, NULL, 0);
LogFailureF("Out of memory. Entry %ld, size %ld", entryID, sb.st_size);
} else if (!fh) {
LoggerEvent(Agent.agent.loggingHandle, LOGGER_SUBSYSTEM_GENERAL, LOGGER_EVENT_FILE_OPEN_FAILURE, LOG_CRITICAL, entryID, path, __FILE__, __LINE__, 0, NULL, 0);
LogFailureF("File open failure. Entry %ld, path %s", entryID, path);
} else {
LoggerEvent(Agent.agent.loggingHandle, LOGGER_SUBSYSTEM_GENERAL, LOGGER_EVENT_FILE_OPEN_FAILURE, LOG_CRITICAL, entryID, path, __FILE__, __LINE__, 0, NULL, 0);
LogFailureF("Event file open failure. Entry %ld, path %s", entryID, path);
}
if (qEnvelope) {
@@ -951,17 +990,17 @@ StartOver:
}
if (fh) {
fclose(fh);
FCLOSE_CHECK(fh);
}
if (newFH) {
fclose(newFH);
FCLOSE_CHECK(newFH);
}
LoggerEvent(Agent.agent.loggingHandle, LOGGER_SUBSYSTEM_GENERAL, LOGGER_EVENT_WRITE_ERROR, LOG_WARNING, 0, "Spool", "Queue", __LINE__, 0, NULL, 0);
Log(LOG_WARNING, "Write error in queue");
sprintf(path, "%s/w%s.%03d",Conf.spoolPath, entry, queue);
unlink(path);
UNLINK_CHECK(path);
ProcessQueueEntryCleanUp(idLock, report);
return(TRUE);
@@ -974,7 +1013,7 @@ StartOver:
/* We move it to the Q_RTS queue and the linger code there will bounce it for us */
sprintf(path, "%s/c%s.%03d", Conf.spoolPath, entry, queue);
sprintf(path2, "%s/c%s.%03d", Conf.spoolPath, entry, Q_RTS);
rename(path, path2);
RENAME_CHECK(path, path2);
sprintf(path, "%03d%s", Q_RTS, entry);
SpoolEntryIDUnlock(idLock);
@@ -1022,21 +1061,22 @@ StartOver:
bounce = FALSE;
sprintf(path, "%s/c%s.%03d", Conf.spoolPath, entry, queue);
fh = fopen(path, "rb");
FOPEN_CHECK(fh, path, "rb");
sprintf(path, "%s/w%s.%03d", Conf.spoolPath, entry, queue);
newFH = fopen(path, "wb");
FOPEN_CHECK(newFH, path, "wb");
if (!fh || !newFH) {
if (fh) {
fclose(fh);
FCLOSE_CHECK(fh);
}
if (newFH) {
fclose(newFH);
FCLOSE_CHECK(newFH);
}
sprintf(path, "%s/w%s.%03d",Conf.spoolPath, entry, queue);
unlink(path);
UNLINK_CHECK(path);
ProcessQueueEntryCleanUp(idLock, report);
return(TRUE);
@@ -1075,8 +1115,6 @@ StartOver:
authenticatedSender[1]='\0';
}
/* LoggerEvent(NMAP.handle.logging, LOGGER_SUBSYSTEM_QUEUE, LOGGER_EVENT_PQE_FROM, LOG_DEBUG, entryID, entry, sender, queue, 0, MIME_TEXT_PLAIN, authenticatedSender? strlen(authenticatedSender) + 1: 1, authenticatedSender? (char *)authenticatedSender: "", NULL, 0); */
fprintf(newFH, "%s\r\n", line);
break;
}
@@ -1085,12 +1123,12 @@ StartOver:
struct sockaddr_in siaddr;
if (!data) {
data = fopen(path, "rb");
FOPEN_CHECK(data, path, "rb");
if (!data) {
fclose(fh);
FCLOSE_CHECK(fh);
sprintf(path, "%s/w%s.%03d",Conf.spoolPath, entry, queue);
fclose(newFH);
unlink(path);
FCLOSE_CHECK(newFH);
UNLINK_CHECK(path);
ProcessQueueEntryCleanUp(idLock, report);
return(TRUE);
}
@@ -1117,11 +1155,10 @@ StartOver:
vs = MDBCreateValueStruct(Agent.agent.directoryHandle, NULL);
if (!MsgFindObject(recipient, NULL, NULL, &siaddr, vs)) {
LoggerEvent(Agent.agent.loggingHandle, LOGGER_SUBSYSTEM_QUEUE, LOGGER_EVENT_PQE_FIND_OBJECT_FAILED, LOG_WARNING, entryID, recipient, "", queue, 0, NULL, 0);
Log(LOG_WARNING, "User %s unknown, entry %ld", recipient, entryID);
status = DELIVER_USER_UNKNOWN;
} else {
LoggerEvent(Agent.agent.loggingHandle, LOGGER_SUBSYSTEM_QUEUE, LOGGER_EVENT_PQE_REMOTE_NMAP, LOG_DEBUG, entryID, entry, line + 1, queue, XplHostToLittle(siaddr.sin_addr.s_addr), NULL, 0);
Log(LOG_DEBUG, "Delivering %s on queue %d to %s", entry, queue, line+1);
status = DeliverToStore(&list, &siaddr, NMAP_DOCTYPE_CAL, sender, authenticatedSender, dataFilename, data, dSize, recipient, mailbox, flags);
if (Agent.agent.state == BONGO_AGENT_STATE_STOPPING) {
status = DELIVER_TRY_LATER;
@@ -1135,7 +1172,7 @@ StartOver:
}
if (status<DELIVER_SUCCESS) {
LoggerEvent(Agent.agent.loggingHandle, LOGGER_SUBSYSTEM_QUEUE, LOGGER_EVENT_PQE_DELIVERY_FAILED, LOG_NOTICE, entryID, entry, NULL, queue, status, NULL, 0);
Log(LOG_NOTICE, "Delivery failed: entry %ld, queue %d, status %d", entryID, queue, status);
if (status==DELIVER_USER_UNKNOWN || status==DELIVER_INTERNAL_ERROR || status==DELIVER_QUOTA_EXCEEDED) {
XplSafeIncrement(Queue.localDeliveryFailed);
if (flags & DSN_FAILURE) {
@@ -1191,13 +1228,13 @@ StartOver:
struct sockaddr_in siaddr;
if (!data) {
sprintf(path, "%s/d%s.msg", Conf.spoolPath, entry);
data = fopen(path, "rb");
FOPEN_CHECK(data, path, "rb");
if (!data) {
fclose(fh);
fclose(newFH);
FCLOSE_CHECK(fh);
FCLOSE_CHECK(newFH);
sprintf(path, "%s/w%s.%03d", Conf.spoolPath, entry, queue);
unlink(path);
UNLINK_CHECK(path);
ProcessQueueEntryCleanUp(idLock, report);
return(TRUE);
@@ -1224,7 +1261,7 @@ StartOver:
/* Attempt delivery, check if local or remote */
vs = MDBCreateValueStruct(Agent.agent.directoryHandle, NULL);
if (MsgFindObject(recipient, NULL, NULL, &siaddr, vs)) {
LoggerEvent(Agent.agent.loggingHandle, LOGGER_SUBSYSTEM_QUEUE, LOGGER_EVENT_PQE_REMOTE_NMAP, LOG_DEBUG, entryID, entry, line + 1, queue, XplHostToLittle(siaddr.sin_addr.s_addr), NULL, 0);
Log(LOG_DEBUG, "Deliver to store entry %s in queue %d for host %s", entry, queue, LOGIP(siaddr));
status = DeliverToStore(&list, &siaddr, NMAP_DOCTYPE_MAIL, sender, authenticatedSender, dataFilename, data, dSize, recipient, mailbox, messageFlags);
if (Agent.agent.state == BONGO_AGENT_STATE_STOPPING) {
@@ -1256,7 +1293,7 @@ StartOver:
flags = 0;
keep = TRUE; /* recipient prevent it from getting removed */
} else {
LoggerEvent(Agent.agent.loggingHandle, LOGGER_SUBSYSTEM_QUEUE, LOGGER_EVENT_PQE_FIND_OBJECT_FAILED, LOG_INFO, entryID, recipient, (authenticatedSender && (authenticatedSender[0] != '-') && (authenticatedSender[1] != ' '))? authenticatedSender: sender, queue, saddr.sin_addr.s_addr, NULL, 0);
Log(LOG_INFO, "Can't forward undeliverable entry %ld in queue %d for user %s on host %s", entry, queue, sender, LOGIP(saddr));
XplRWReadLockRelease(&Conf.lock);
}
@@ -1265,7 +1302,7 @@ StartOver:
MDBDestroyValueStruct(vs);
if (status < DELIVER_SUCCESS) {
LoggerEvent(Agent.agent.loggingHandle, LOGGER_SUBSYSTEM_QUEUE, LOGGER_EVENT_PQE_DELIVERY_FAILED, LOG_WARNING, entryID, entry, NULL, queue, status, NULL, 0);
Log(LOG_WARNING, "Couldn't deliver entry %s on queue %d, status %d", entry, queue, status);
if ((status == DELIVER_USER_UNKNOWN) || (status == DELIVER_INTERNAL_ERROR) || (status == DELIVER_QUOTA_EXCEEDED)) {
XplSafeIncrement(Queue.localDeliveryFailed);
if (flags & DSN_FAILURE) {
@@ -1360,7 +1397,7 @@ StartOver:
}
default:
LoggerEvent(Agent.agent.loggingHandle, LOGGER_SUBSYSTEM_QUEUE, LOGGER_EVENT_PQE_FIND_OBJECT_FAILED, LOG_INFO, entryID, line, entry, queue, 0, NULL, 0);
Log(LOG_INFO, "Unknown command: %s (entry: %ld)", line, entryID);
fprintf(newFH, "%s\r\n", line);
break;
}
@@ -1370,33 +1407,33 @@ StartOver:
EndStoreDelivery(&list);
memset(&list, 0, sizeof(NMAPConnections));
}
fclose(newFH);
fclose(fh);
FCLOSE_CHECK(newFH);
FCLOSE_CHECK(fh);
if (bounce) {
unsigned char Path2[XPL_MAX_PATH+1];
/* First, rename the work file into a control file */
sprintf(path, "%s/c%s.%03d",Conf.spoolPath, entry, queue);
unlink(path);
UNLINK_CHECK(path);
sprintf(Path2, "%s/w%s.%03d",Conf.spoolPath, entry, queue);
rename(Path2, path);
RENAME_CHECK(Path2, path);
if (!data) {
sprintf(path, "%s/d%s.msg",Conf.spoolPath, entry);
data=fopen(path,"rb");
FOPEN_CHECK(data, path, "rb");
} else {
fseek(data, 0, SEEK_SET);
}
sprintf(path, "%s/c%s.%03d",Conf.spoolPath, entry, queue);
fh=fopen(path,"rb");
FOPEN_CHECK(fh, path, "rb");
if (fh && data && 0 == HandleDSN(data, fh)) {
/* Now bounce the thing */
fseek(fh, 0, SEEK_SET);
sprintf(path, "%s/w%s.%03d",Conf.spoolPath, entry, queue);
newFH=fopen(path, "wb");
FOPEN_CHECK(newFH, path, "wb");
if (newFH) {
/* Now remove the bounced entries */
while (!feof(fh) && !ferror(fh)) {
@@ -1411,42 +1448,42 @@ StartOver:
}
}
}
fclose(newFH);
FCLOSE_CHECK(newFH);
}
fclose(fh);
FCLOSE_CHECK(fh);
} else {
if (fh) {
fclose(fh);
FCLOSE_CHECK(fh);
}
if (data) {
fclose(data);
FCLOSE_CHECK(data);
}
LoggerEvent(Agent.agent.loggingHandle, LOGGER_SUBSYSTEM_GENERAL, LOGGER_EVENT_FILE_OPEN_FAILURE, LOG_CRITICAL, entryID, path, __FILE__, __LINE__, 0, NULL, 0);
Log(LOG_CRITICAL, "File open error for entry %ld, path: %s", entryID, path);
ProcessQueueEntryCleanUp(idLock, report);
return(FALSE);
}
}
if (data) {
fclose(data);
FCLOSE_CHECK(data);
}
if (keep) {
unsigned char Path2[XPL_MAX_PATH+1];
sprintf(path, "%s/c%s.%03d", Conf.spoolPath, entry, queue);
unlink(path);
UNLINK_CHECK(path);
sprintf(Path2, "%s/w%s.%03d", Conf.spoolPath, entry, queue);
rename(Path2, path);
RENAME_CHECK(Path2, path);
} else {
sprintf(path, "%s/w%s.%03d", Conf.spoolPath, entry, queue);
unlink(path);
UNLINK_CHECK(path);
sprintf(path, "%s/c%s.%03d", Conf.spoolPath, entry, queue);
unlink(path);
UNLINK_CHECK(path);
sprintf(path, "%s/d%s.msg", Conf.spoolPath, entry);
unlink(path);
UNLINK_CHECK(path);
XplSafeDecrement(Queue.queuedLocal);
ProcessQueueEntryCleanUp(idLock, report);
@@ -1476,14 +1513,17 @@ StartOver:
for (used = 0; (used < (unsigned long)Queue.pushClients.count) && (Agent.agent.state < BONGO_AGENT_STATE_STOPPING); used++) {
if ((Queue.pushClients.array[used].queue == queue) && (Queue.pushClients.array[used].errorCount <= MAX_PUSHCLIENTS_ERRORS)) {
sprintf(path, "%s/c%s.%03d", Conf.spoolPath, entry, queue);
if ((stat(path, &sb) == 0) && ((fh = fopen(path, "rb")) != NULL)) {
/* Count the number of lines */
do {
if (fgets(line, CONN_BUFSIZE, fh)
&& ((line[0] == QUEUE_RECIP_REMOTE) || (line[0] == QUEUE_RECIP_LOCAL) || (line[0] == QUEUE_RECIP_MBOX_LOCAL))) {
lines++;
}
} while (!feof(fh) && !ferror(fh));
if (stat(path, &sb) == 0) {
FOPEN_CHECK(fh, path, "rb");
if (fh) {
/* Count the number of lines */
do {
if (fgets(line, CONN_BUFSIZE, fh)
&& ((line[0] == QUEUE_RECIP_REMOTE) || (line[0] == QUEUE_RECIP_LOCAL) || (line[0] == QUEUE_RECIP_MBOX_LOCAL))) {
lines++;
}
} while (!feof(fh) && !ferror(fh));
}
} else {
ProcessQueueEntryCleanUp(idLock, report);
return(TRUE);
@@ -1501,9 +1541,9 @@ StartOver:
QueueClientFree(client);
}
LoggerEvent(Agent.agent.loggingHandle, LOGGER_SUBSYSTEM_GENERAL, LOGGER_EVENT_OUT_OF_MEMORY, LOG_ERROR, entryID, __FILE__, NULL, sizeof(QueueClient), __LINE__, NULL, 0);
LogFailureF("Cannot allocate %d bytes memory (entry %ld)", sizeof(QueueClient), entryID);
if (fh) {
fclose(fh);
FCLOSE_CHECK(fh);
}
fh = NULL;
@@ -1536,7 +1576,7 @@ StartOver:
QueueClientFree(client);
if (fh) {
fclose(fh);
FCLOSE_CHECK(fh);
}
fh = NULL;
@@ -1560,7 +1600,7 @@ StartOver:
Queue.pushClients.array[used].errorCount = 0;
} else {
LoggerEvent(Agent.agent.loggingHandle, LOGGER_SUBSYSTEM_QUEUE, LOGGER_EVENT_PQE_FAILED_CONNECT, LOG_WARNING, entryID, entry, NULL, queue, XplHostToLittle(saddr.sin_addr.s_addr), &(Queue.pushClients.array[used].port), sizeof(Queue.pushClients.array[used].port));
LogFailureF("Couldn't connect client %s", LOGIP(saddr));
ConnClose(client->conn, 0);
ConnFree(client->conn);
@@ -1573,7 +1613,7 @@ StartOver:
QueueClientFree(client);
if (fh) {
fclose(fh);
FCLOSE_CHECK(fh);
}
fh = NULL;
@@ -1586,7 +1626,7 @@ StartOver:
ConnWriteF(client->conn, "6020 %03d-%s %ld %ld %ld\r\n", queue, entry, (unsigned long)sb.st_size, dSize, lines);
ConnWriteFile(client->conn, fh);
fclose(fh);
FCLOSE_CHECK(fh);
fh = NULL;
sprintf(client->entry.workQueue, "%03d-%s", queue, entry);
@@ -1598,7 +1638,7 @@ StartOver:
client->entry.report = report;
if (!HandleCommand(client)) {
LoggerEvent(Agent.agent.loggingHandle, LOGGER_SUBSYSTEM_QUEUE, LOGGER_EVENT_PQE_FAILED, LOG_WARNING, entryID, entry, NULL, queue, XplHostToLittle(saddr.sin_addr.s_addr), &(Queue.pushClients.array[used].port), sizeof(Queue.pushClients.array[used].port));
LogFailureF("Couldn't handle command on entry %s for host %s", entry, LOGIP(saddr));
}
/* fixme - evaluate this section.
@@ -1673,7 +1713,7 @@ StartOver:
sprintf(path, "%s/c%s.%03d", Conf.spoolPath, entry, queue);
sprintf(path2, "%s/c%s.%03d", Conf.spoolPath, entry, i);
rename(path, path2);
RENAME_CHECK(path, path2);
sprintf(path, "%03d%s", i, entry);
SpoolEntryIDUnlock(idLock);
@@ -1693,7 +1733,7 @@ StartOver:
QDBHandleRelease(handle);
}
fh = fopen(path,"rb");
FOPEN_CHECK(fh, path, "rb");
keep = FALSE;
if (fh) {
do {
@@ -1723,9 +1763,7 @@ StartOver:
if (ptr2) {
if ((handle = QDBHandleAlloc()) != NULL) {
ccode = QDBAdd(handle, ptr2 + 1, entryID, queue);
if (ccode == -1) {
LoggerEvent(Agent.agent.loggingHandle, LOGGER_SUBSYSTEM_DATABASE, LOGGER_EVENT_DATABASE_INSERT_ERROR, LOG_ERROR, 0, ptr2 + 1, NULL, ccode, 0, NULL, 0);
}
LogAssertF(ccode != -1, "Database insert error: %s", ptr2 + 1);
QDBHandleRelease(handle);
}
@@ -1741,7 +1779,7 @@ StartOver:
}
}
} while (!feof(fh) && !ferror(fh));
fclose(fh);
FCLOSE_CHECK(fh);
} else {
keep = TRUE;
}
@@ -1750,16 +1788,16 @@ StartOver:
if (bounce) {
/* Call the bouncing code */
sprintf(path2, "%s/d%s.msg", Conf.spoolPath, entry);
data = fopen(path2, "rb");
fh = fopen(path, "rb");
FOPEN_CHECK(data, path2, "rb");
FOPEN_CHECK(fh, path, "rb");
if (fh && data && 0 == HandleDSN(data, fh)) {
fclose(data);
FCLOSE_CHECK(data);
/* If we're not keeping the file, we can ignore its contents */
if (keep) {
sprintf(path2, "%s/w%s.%03d", Conf.spoolPath, entry, queue);
newFH = fopen(path2, "wb");
FOPEN_CHECK(newFH, path2, "wb");
if (newFH) {
/* Now remove the bounced entries */
while (!feof(fh) && !ferror(fh)) {
@@ -1777,27 +1815,27 @@ StartOver:
}
}
fclose(fh);
fclose(newFH);
FCLOSE_CHECK(fh);
FCLOSE_CHECK(newFH);
unlink(path);
rename(path2, path);
UNLINK_CHECK(path);
RENAME_CHECK(path2, path);
} else {
fclose(fh);
FCLOSE_CHECK(fh);
}
} else {
fclose(fh);
FCLOSE_CHECK(fh);
}
} else {
if (fh) {
fclose(fh);
FCLOSE_CHECK(fh);
}
if (data) {
fclose(data);
FCLOSE_CHECK(data);
}
LoggerEvent(Agent.agent.loggingHandle, LOGGER_SUBSYSTEM_GENERAL, LOGGER_EVENT_FILE_OPEN_FAILURE, LOG_CRITICAL, entryID, path, __FILE__, __LINE__, 0, NULL, 0);
LogFailureF("File open failure: entry %ld, path %s", entryID, path);
ProcessQueueEntryCleanUp(idLock, report);
return(TRUE);
@@ -1807,10 +1845,10 @@ StartOver:
if (!keep) {
XplSafeDecrement(Queue.queuedLocal);
unlink(path);
UNLINK_CHECK(path);
sprintf(path, "%s/d%s.msg",Conf.spoolPath, entry);
unlink(path);
UNLINK_CHECK(path);
break;
}
@@ -1827,13 +1865,13 @@ StartOver:
}
sprintf(path, "%s/w%s.%03d", Conf.spoolPath, entry, queue);
newFH = fopen(path, "wb");
FOPEN_CHECK(newFH, path, "wb");
sprintf(path, "%s/d%s.msg", Conf.spoolPath, entry);
data = fopen(path, "rb");
FOPEN_CHECK(data, path, "rb");
sprintf(path, "%s/c%s.%03d", Conf.spoolPath, entry, queue);
fh = fopen(path, "rb");
FOPEN_CHECK(fh, path, "rb");
if (fh && newFH && data) {
XplSafeIncrement(Queue.remoteDeliveryFailed);
@@ -1878,42 +1916,42 @@ StartOver:
}
}
fclose(newFH);
fclose(fh);
FCLOSE_CHECK(newFH);
FCLOSE_CHECK(fh);
/* Still got path from above */
unlink(path);
UNLINK_CHECK(path);
sprintf(path2, "%s/w%s.%03d", Conf.spoolPath, entry, queue);
rename(path2, path);
RENAME_CHECK(path2, path);
fh = fopen(path, "rb");
FOPEN_CHECK(fh, path, "rb");
if (fh) {
HandleDSN(data, fh);
fclose(fh);
FCLOSE_CHECK(fh);
} else {
LoggerEvent(Agent.agent.loggingHandle, LOGGER_SUBSYSTEM_GENERAL, LOGGER_EVENT_FILE_OPEN_FAILURE, LOG_CRITICAL, entryID, path, __FILE__, __LINE__, 0, NULL, 0);
LogFailureF("Couldn't open path %s (%ld)", path, entryID);
}
fclose(data);
FCLOSE_CHECK(data);
sprintf(path, "%s/c%s.%03d", Conf.spoolPath, entry, queue);
unlink(path);
UNLINK_CHECK(path);
sprintf(path, "%s/d%s.msg", Conf.spoolPath, entry);
unlink(path);
UNLINK_CHECK(path);
XplSafeDecrement(Queue.queuedLocal);
} else {
if (data) {
fclose(data);
FCLOSE_CHECK(data);
}
if (fh) {
fclose(fh);
FCLOSE_CHECK(fh);
}
if (newFH) {
fclose(newFH);
FCLOSE_CHECK(newFH);
}
}
@@ -1923,7 +1961,7 @@ StartOver:
count = 0;
sprintf(path, "%s/c%s.%03d", Conf.spoolPath, entry, Q_RTS);
sprintf(path2, "%s/c%s.%03d", Conf.spoolPath, entry, Q_INCOMING);
rename(path, path2);
RENAME_CHECK(path, path2);
sprintf(path, "%03d%s", Q_INCOMING, entry);
}
@@ -2443,7 +2481,7 @@ HandleDSN(FILE *data, FILE *control)
fprintf (stderr, "could not open rtsData\n");
fclose(rtsControl);
sprintf(path, "%s/d%07lx.msg", Conf.spoolPath, id);
unlink(path);
UNLINK_CHECK(path);
return -1;
}
@@ -2452,10 +2490,10 @@ HandleDSN(FILE *data, FILE *control)
fclose(rtsData);
sprintf(path, "%s/d%07lx.msg", Conf.spoolPath, id);
unlink(path);
UNLINK_CHECK(path);
sprintf(path, "%s/c%07lx.%03d", Conf.spoolPath, id, Q_INCOMING);
unlink(path);
UNLINK_CHECK(path);
SpoolEntryIDUnlock(idLock);
return 0;
@@ -2797,7 +2835,7 @@ CreateQueueThreads(BOOL failed)
current++;
CHOP_NEWLINE(path);
unlink(path);
UNLINK_CHECK(path);
}
}
@@ -2805,7 +2843,7 @@ CreateQueueThreads(BOOL failed)
}
sprintf(path, "%s/killfile", MsgGetDBFDir(NULL));
unlink(path);
UNLINK_CHECK(path);
XplCloseDir(dirP);
@@ -3032,7 +3070,7 @@ CommandQabrt(void *param)
client->entry.control = NULL;
sprintf(client->path,"%s/c%07lx.in", Conf.spoolPath, client->entry.id);
unlink(client->path);
UNLINK_CHECK(client->path);
}
if (client->entry.data) {
@@ -3040,7 +3078,7 @@ CommandQabrt(void *param)
client->entry.data = NULL;
sprintf(client->path,"%s/d%07lx.msg",Conf.spoolPath, client->entry.id);
unlink(client->path);
UNLINK_CHECK(client->path);
}
if (client->entry.work) {
fclose(client->entry.work);
@@ -3048,7 +3086,7 @@ CommandQabrt(void *param)
if (client->entry.workQueue[0] != '\0') {
sprintf(client->path,"%s/w%s.%s",Conf.spoolPath, client->entry.workQueue + 4, client->entry.workQueue);
unlink(client->path);
UNLINK_CHECK(client->path);
}
}
@@ -3333,7 +3371,7 @@ CommandQcrea(void *param)
} else {
fclose(client->entry.control);
sprintf(client->path, "%s/c%07lx.in", Conf.spoolPath, id);
unlink(client->path);
UNLINK_CHECK(client->path);
return(ConnWrite(client->conn, MSG5221SPACELOW, sizeof(MSG5221SPACELOW) - 1));
}
@@ -3368,10 +3406,10 @@ CommandQdele(void *param)
LockQueueEntry(ptr+4, atoi(ptr)); */
sprintf(client->path, "%s/c%s.%s", Conf.spoolPath, ptr + 4, ptr);
unlink(client->path);
UNLINK_CHECK(client->path);
sprintf(client->path, "%s/d%s.msg", Conf.spoolPath, ptr + 4);
unlink(client->path);
UNLINK_CHECK(client->path);
/* FIXME: close out the file handles? */
@@ -3400,15 +3438,15 @@ CommandQdone(void *param)
sprintf(client->path, "%s/w%s.%s", Conf.spoolPath, client->entry.workQueue + 4, client->entry.workQueue);
sprintf(path, "%s/c%s.%s", Conf.spoolPath, client->entry.workQueue + 4, client->entry.workQueue);
if ((stat(client->path, &sb1) == 0) && (stat(path, &sb2) == 0)) {
unlink(path);
rename(client->path, path);
UNLINK_CHECK(path);
RENAME_CHECK(client->path, path);
} else {
if (stat(client->path, &sb) == 0) {
/* Our new queue file exists, everything still ok */
rename(client->path, path);
RENAME_CHECK(client->path, path);
} else {
/* got to keep the old one */
unlink(client->path);
UNLINK_CHECK(client->path);
}
}
@@ -3883,12 +3921,12 @@ CommandQmove(void *param)
if (stat(client->path, &sb) == 0) {
sprintf(path, "%s/c%s.%03d", Conf.spoolPath, ptr + 4, atoi(ptr));
rename(client->path, path);
RENAME_CHECK(client->path, path);
ccode = ConnWrite(client->conn, MSG1000OK, sizeof(MSG1000OK) - 1);
} else {
sprintf(client->path, "%s/d%s.msg", Conf.spoolPath, ptr + 4);
unlink(client->path);
UNLINK_CHECK(client->path);
ccode = ConnWrite(client->conn, MSG4224NOENTRY, sizeof(MSG4224NOENTRY) - 1);
}
@@ -3953,7 +3991,7 @@ CommandQrcp(void *param)
sprintf(client->path,"%s/c%07lx.in",Conf.spoolPath, client->entry.id);
sprintf(path, "%s/c%07lx.%03ld", Conf.spoolPath, client->entry.id, client->entry.target);
rename(client->path, path);
RENAME_CHECK(client->path, path);
sprintf(client->path, "%03ld%lx", client->entry.target, client->entry.id);
@@ -3982,7 +4020,7 @@ CommandQrcp(void *param)
client->entry.control = NULL;
sprintf(client->path, "%s/c%07lx.in", Conf.spoolPath, id);
unlink(client->path);
UNLINK_CHECK(client->path);
}
sprintf(client->path, "%s/c%07lx.in", Conf.spoolPath, client->entry.id);
@@ -4143,7 +4181,7 @@ CommandQrun(void *param)
sprintf(client->path,"%s/c%07lx.in",Conf.spoolPath, client->entry.id);
sprintf(path, "%s/c%07lx.%03ld", Conf.spoolPath, client->entry.id, client->entry.target);
rename(client->path, path);
RENAME_CHECK(client->path, path);
XplSafeIncrement(Queue.queuedLocal);
@@ -4213,11 +4251,11 @@ CommandQsrchDomain(void *param)
ccode = ConnWrite(client->conn, MSG1000OK, sizeof(MSG1000OK) - 1);
}
} else {
LoggerEvent(Agent.agent.loggingHandle, LOGGER_SUBSYSTEM_DATABASE, LOGGER_EVENT_DATABASE_FIND_ERROR, LOG_ERROR, 0, ptr, NULL, ccode, 0, NULL, 0);
LogFailureF("Couldn't find %s in QDB", ptr);
ccode = ConnWrite(client->conn, MSG4261NODOMAIN, sizeof(MSG4261NODOMAIN) - 1);
}
} else {
LoggerEvent(Agent.agent.loggingHandle, LOGGER_SUBSYSTEM_GENERAL, LOGGER_EVENT_NMAP_OUT_OF_MEMORY, LOG_CRITICAL, 0, client->buffer, NULL, sizeof(MDBValueStruct), 0, NULL, 0);
LogFailure("Out of memory");
ccode = ConnWrite(client->conn, MSG5230NOMEMORYERR, sizeof(MSG5230NOMEMORYERR) - 1);
}
@@ -4776,6 +4814,6 @@ CommandQflush(void *param)
QueueClient *client = (QueueClient *)param;
Queue.flushNeeded = TRUE;
return (ConnWrite(client->conn, MSG1000OK, sizeof(MSG1000OK) - 1));
}