Threading: do not keep a slots post_pq locked while processing the packets.

remotes/origin/master-1.2.x
Victor Julien 15 years ago
parent e81f94cd83
commit cd987ae7a5

@ -142,7 +142,11 @@ void *TmThreadsSlot1NoIn(void *td)
/* handle error */ /* handle error */
if (r == TM_ECODE_FAILED) { if (r == TM_ECODE_FAILED) {
TmqhReleasePacketsToPacketPool(&s->slot_pre_pq); TmqhReleasePacketsToPacketPool(&s->slot_pre_pq);
SCMutexLock(&s->slot_post_pq.mutex_q);
TmqhReleasePacketsToPacketPool(&s->slot_post_pq); TmqhReleasePacketsToPacketPool(&s->slot_post_pq);
SCMutexUnlock(&s->slot_post_pq.mutex_q);
if (p != NULL) if (p != NULL)
TmqhOutputPacketpool(tv, p); TmqhOutputPacketpool(tv, p);
TmThreadsSetFlag(tv, THV_FAILED); TmThreadsSetFlag(tv, THV_FAILED);
@ -383,7 +387,11 @@ void *TmThreadsSlot1(void *td)
/* handle error */ /* handle error */
if (r == TM_ECODE_FAILED) { if (r == TM_ECODE_FAILED) {
TmqhReleasePacketsToPacketPool(&s->slot_pre_pq); TmqhReleasePacketsToPacketPool(&s->slot_pre_pq);
SCMutexLock(&s->slot_post_pq.mutex_q);
TmqhReleasePacketsToPacketPool(&s->slot_post_pq); TmqhReleasePacketsToPacketPool(&s->slot_post_pq);
SCMutexUnlock(&s->slot_post_pq.mutex_q);
TmqhOutputPacketpool(tv, p); TmqhOutputPacketpool(tv, p);
TmThreadsSetFlag(tv, THV_FAILED); TmThreadsSetFlag(tv, THV_FAILED);
break; break;
@ -464,7 +472,11 @@ TmEcode TmThreadsSlotVarRun(ThreadVars *tv, Packet *p,
if (unlikely(r == TM_ECODE_FAILED)) { if (unlikely(r == TM_ECODE_FAILED)) {
/* Encountered error. Return packets to packetpool and return */ /* Encountered error. Return packets to packetpool and return */
TmqhReleasePacketsToPacketPool(&s->slot_pre_pq); TmqhReleasePacketsToPacketPool(&s->slot_pre_pq);
SCMutexLock(&s->slot_post_pq.mutex_q);
TmqhReleasePacketsToPacketPool(&s->slot_post_pq); TmqhReleasePacketsToPacketPool(&s->slot_post_pq);
SCMutexUnlock(&s->slot_post_pq.mutex_q);
TmThreadsSetFlag(tv, THV_FAILED); TmThreadsSetFlag(tv, THV_FAILED);
return TM_ECODE_FAILED; return TM_ECODE_FAILED;
} }
@ -480,7 +492,11 @@ TmEcode TmThreadsSlotVarRun(ThreadVars *tv, Packet *p,
r = TmThreadsSlotVarRun(tv, extra_p, s->slot_next); r = TmThreadsSlotVarRun(tv, extra_p, s->slot_next);
if (unlikely(r == TM_ECODE_FAILED)) { if (unlikely(r == TM_ECODE_FAILED)) {
TmqhReleasePacketsToPacketPool(&s->slot_pre_pq); TmqhReleasePacketsToPacketPool(&s->slot_pre_pq);
SCMutexLock(&s->slot_post_pq.mutex_q);
TmqhReleasePacketsToPacketPool(&s->slot_post_pq); TmqhReleasePacketsToPacketPool(&s->slot_post_pq);
SCMutexUnlock(&s->slot_post_pq.mutex_q);
TmqhOutputPacketpool(tv, extra_p); TmqhOutputPacketpool(tv, extra_p);
TmThreadsSetFlag(tv, THV_FAILED); TmThreadsSetFlag(tv, THV_FAILED);
return TM_ECODE_FAILED; return TM_ECODE_FAILED;
@ -667,20 +683,26 @@ void *TmThreadsSlotVar(void *td)
tv->tmqh_out(tv, p); tv->tmqh_out(tv, p);
} /* if (p != NULL) */ } /* if (p != NULL) */
/* now handle the post_pq packets */ /* now handle the post_pq packets */
TmSlot *slot; TmSlot *slot;
for (slot = s; slot != NULL; slot = slot->slot_next) { for (slot = s; slot != NULL; slot = slot->slot_next) {
if (slot->slot_post_pq.top != NULL) { if (slot->slot_post_pq.top != NULL) {
SCMutexLock(&slot->slot_post_pq.mutex_q); while (1) {
while (slot->slot_post_pq.top != NULL) { SCMutexLock(&slot->slot_post_pq.mutex_q);
Packet *extra_p = PacketDequeue(&slot->slot_post_pq); Packet *extra_p = PacketDequeue(&slot->slot_post_pq);
SCMutexUnlock(&slot->slot_post_pq.mutex_q);
if (extra_p == NULL) if (extra_p == NULL)
break; break;
if (slot->slot_next != NULL) { if (slot->slot_next != NULL) {
r = TmThreadsSlotVarRun(tv, extra_p, slot->slot_next); r = TmThreadsSlotVarRun(tv, extra_p, slot->slot_next);
if (r == TM_ECODE_FAILED) { if (r == TM_ECODE_FAILED) {
SCMutexLock(&slot->slot_post_pq.mutex_q);
TmqhReleasePacketsToPacketPool(&slot->slot_post_pq); TmqhReleasePacketsToPacketPool(&slot->slot_post_pq);
SCMutexUnlock(&slot->slot_post_pq.mutex_q);
TmqhOutputPacketpool(tv, extra_p); TmqhOutputPacketpool(tv, extra_p);
TmThreadsSetFlag(tv, THV_FAILED); TmThreadsSetFlag(tv, THV_FAILED);
break; break;
@ -689,7 +711,6 @@ void *TmThreadsSlotVar(void *td)
/* output the packet */ /* output the packet */
tv->tmqh_out(tv, extra_p); tv->tmqh_out(tv, extra_p);
} /* while */ } /* while */
SCMutexUnlock(&slot->slot_post_pq.mutex_q);
} /* if */ } /* if */
} /* for */ } /* for */

@ -131,7 +131,10 @@ static inline TmEcode TmThreadsSlotProcessPkt(ThreadVars *tv, TmSlot *s, Packet
TmqhOutputPacketpool(tv, p); TmqhOutputPacketpool(tv, p);
TmSlot *slot = s; TmSlot *slot = s;
while (slot != NULL) { while (slot != NULL) {
SCMutexLock(&slot->slot_post_pq.mutex_q);
TmqhReleasePacketsToPacketPool(&slot->slot_post_pq); TmqhReleasePacketsToPacketPool(&slot->slot_post_pq);
SCMutexUnlock(&slot->slot_post_pq.mutex_q);
slot = slot->slot_next; slot = slot->slot_next;
} }
TmThreadsSetFlag(tv, THV_FAILED); TmThreadsSetFlag(tv, THV_FAILED);
@ -139,27 +142,33 @@ static inline TmEcode TmThreadsSlotProcessPkt(ThreadVars *tv, TmSlot *s, Packet
} else { } else {
tv->tmqh_out(tv, p); tv->tmqh_out(tv, p);
/* post process pq */ /* post process pq */
TmSlot *slot = s; TmSlot *slot = s;
while (slot != NULL) { while (slot != NULL) {
if (slot->slot_post_pq.top != NULL) { if (slot->slot_post_pq.top != NULL) {
SCMutexLock(&slot->slot_post_pq.mutex_q); while (1) {
while (slot->slot_post_pq.top != NULL) { SCMutexLock(&slot->slot_post_pq.mutex_q);
Packet *extra_p = PacketDequeue(&slot->slot_post_pq); Packet *extra_p = PacketDequeue(&slot->slot_post_pq);
if (extra_p != NULL) { SCMutexUnlock(&slot->slot_post_pq.mutex_q);
if (slot->slot_next != NULL) {
r = TmThreadsSlotVarRun(tv, extra_p, slot->slot_next); if (extra_p == NULL)
if (r == TM_ECODE_FAILED) { break;
TmqhReleasePacketsToPacketPool(&slot->slot_post_pq);
TmqhOutputPacketpool(tv, extra_p); if (slot->slot_next != NULL) {
TmThreadsSetFlag(tv, THV_FAILED); r = TmThreadsSlotVarRun(tv, extra_p, slot->slot_next);
break; if (r == TM_ECODE_FAILED) {
} SCMutexLock(&slot->slot_post_pq.mutex_q);
TmqhReleasePacketsToPacketPool(&slot->slot_post_pq);
SCMutexUnlock(&slot->slot_post_pq.mutex_q);
TmqhOutputPacketpool(tv, extra_p);
TmThreadsSetFlag(tv, THV_FAILED);
break;
} }
tv->tmqh_out(tv, extra_p);
} }
} /* while (slot->slot_post_pq.top != NULL) */ tv->tmqh_out(tv, extra_p);
SCMutexUnlock(&slot->slot_post_pq.mutex_q); }
} /* if (slot->slot_post_pq.top != NULL) */ } /* if (slot->slot_post_pq.top != NULL) */
slot = slot->slot_next; slot = slot->slot_next;
} /* while (slot != NULL) */ } /* while (slot != NULL) */

Loading…
Cancel
Save