[PR #29408] fix: fix race condition #32411

Open
opened 2026-02-21 20:51:21 -05:00 by yindo · 0 comments
Owner

Original Pull Request: https://github.com/langgenius/dify/pull/29408

State: open
Merged: No


Important

  1. Make sure you have read our contribution guidelines
  2. Ensure there is an associated issue and you have been assigned to it
  3. Use the correct syntax to link this PR: Fixes #<issue number>.

Summary

fix #29334

Goal

  • Eliminate a race condition between streaming completion and output moderation so the final “replace” event isn’t lost before the pipeline shuts down.

Core changes

  • Introduced ModerationCoordinator

    • New class providing thread-safe coordination: marks when the stream end is seen and exposes an async_done Event to signal moderation completion.
  • Queue manager integration

    • AppQueueManager: added a moderation_coordinator slot and a set_moderation_coordinator abstract method.

    • MessageBasedAppQueueManager:

      • Added set_moderation_coordinator and a stop_when_ready(poll_ms) helper that waits until ModerationCoordinator.ready_to_close() before stopping.
  • Adjusted publish path logic:

    • On QueueMessageEndEvent/QueueAdvancedChatMessageEndEvent: wait (stop_when_ready) then stop.
    • On QueueStopEvent/QueueErrorEvent: stop immediately.
  • Pipeline wiring

    • BasedGenerateTaskPipeline:
      • _init_output_moderation now accepts and passes the queue manager’s moderation_coordinator to OutputModeration.
    • EasyUIBasedGenerateTaskPipeline:
    • Instantiates a ModerationCoordinator, injects it into the queue manager if supported.
    • After the end event, marks “stream end seen” and waits up to 2s for async_done before yielding the final response.
  • Output moderation robustness

    • OutputModeration:
    • Added coordinator field.
    • New flush_and_stop with bounded thread join.
    • stop_thread now joins with a timeout and sets async_done when no thread ran.
    • worker wrapped with try/finally to always set async_done on exit; logs errors without breaking shutdown.
    • Behavior for publishing replace events unchanged; coordination guarantees they’re not dropped.

Screenshots

Before After
... ...

Checklist

  • This change requires a documentation update, included: Dify Document
  • I understand that this PR may be closed in case there was no previous discussion or issues. (This doesn't apply to typos!)
  • I've added a test for each change that was introduced, and I tried as much as possible to make a single atomic change.
  • I've updated the documentation accordingly.
  • I ran dev/reformat(backend) and cd web && npx lint-staged(frontend) to appease the lint gods
**Original Pull Request:** https://github.com/langgenius/dify/pull/29408 **State:** open **Merged:** No --- > [!IMPORTANT] > > 1. Make sure you have read our [contribution guidelines](https://github.com/langgenius/dify/blob/main/CONTRIBUTING.md) > 1. Ensure there is an associated issue and you have been assigned to it > 1. Use the correct syntax to link this PR: `Fixes #<issue number>`. ## Summary fix #29334 Goal - Eliminate a race condition between streaming completion and output moderation so the final “replace” event isn’t lost before the pipeline shuts down. Core changes - Introduced ModerationCoordinator - New class providing thread-safe coordination: marks when the stream end is seen and exposes an async_done Event to signal moderation completion. - Queue manager integration - AppQueueManager: added a moderation_coordinator slot and a set_moderation_coordinator abstract method. - MessageBasedAppQueueManager: - Added set_moderation_coordinator and a stop_when_ready(poll_ms) helper that waits until ModerationCoordinator.ready_to_close() before stopping. - Adjusted publish path logic: - On QueueMessageEndEvent/QueueAdvancedChatMessageEndEvent: wait (stop_when_ready) then stop. - On QueueStopEvent/QueueErrorEvent: stop immediately. - Pipeline wiring - BasedGenerateTaskPipeline: - _init_output_moderation now accepts and passes the queue manager’s moderation_coordinator to OutputModeration. - EasyUIBasedGenerateTaskPipeline: - Instantiates a ModerationCoordinator, injects it into the queue manager if supported. - After the end event, marks “stream end seen” and waits up to 2s for async_done before yielding the final response. - Output moderation robustness - OutputModeration: - Added coordinator field. - New flush_and_stop with bounded thread join. - stop_thread now joins with a timeout and sets async_done when no thread ran. - worker wrapped with try/finally to always set async_done on exit; logs errors without breaking shutdown. - Behavior for publishing replace events unchanged; coordination guarantees they’re not dropped. ## Screenshots | Before | After | |--------|-------| | ... | ... | ## Checklist - [ ] This change requires a documentation update, included: [Dify Document](https://github.com/langgenius/dify-docs) - [x] I understand that this PR may be closed in case there was no previous discussion or issues. (This doesn't apply to typos!) - [x] I've added a test for each change that was introduced, and I tried as much as possible to make a single atomic change. - [x] I've updated the documentation accordingly. - [x] I ran `dev/reformat`(backend) and `cd web && npx lint-staged`(frontend) to appease the lint gods
yindo added the pull-request label 2026-02-21 20:51:21 -05:00
Sign in to join this conversation.
1 Participants
Notifications
Due Date
No due date set.
Dependencies

No dependencies set.

Reference: langgenius/dify#32411