Files
slide/decks/gen2D/pages/async-deep-dive.md
T
wonder 67ee76e58d
Deploy Slides / build-and-deploy (push) Successful in 1m8s
feat: expand gen2D deck with deep technical slides and interactive elements
- Add 5 new slide pages: eino-deep-dive, ratelimit-overview, ratelimit-lua, async-deep-dive, harness
- Add 4 new drawio SVGs: pipeline-detail, ratelimit, async-task, harness
- Update all existing slides with v-clicks progressive reveal
- Add presenter notes and interactive prompts
- Expand slides.md to register all new pages (16 → 25+ slides)
- Fix text overflow by trimming content per Item block
- Add Eino Graph deep dive: WithGenLocalState, Pre/Post Handler, branch routing
- Add rate limiting section: algorithm comparison, Lua script, Gin middleware
- Add async task deep dive: TaskQueue FIFO, context progress injection
- Add consistency section: style/pipeline/data/deployment 4-layer guarantees
2026-05-30 22:42:23 +08:00

1.9 KiB
Raw Blame History

异步任务 — 队列与调度

FIFO 串行执行,信号驱动调度

type TaskQueue struct {
    mu      sync.Mutex
    jobs    []*TaskJob
    ready   chan struct{}  // 新任务到达信号
    stop    chan struct{}  // 优雅关闭信号
}

func (q *TaskQueue) Enqueue(job *TaskJob) {
    q.mu.Lock()
    q.jobs = append(q.jobs, job)
    q.mu.Unlock()
    q.ready <- struct{}{}  // 唤醒 run() 协程
}
提交任务入队并发送信号。`run()` 协程阻塞等待信号,收到后调用 `processNext()` 取队首任务执行。FIFO 保证任务按提交顺序串行处理。 `TaskJob` 封装了 `context.Context` 和执行函数。闭包捕获任务参数,context 控制超时和取消,确保每个任务独立且可中断。

进度推送与 WebSocket

Context 注入 + 多通道推送

func WithProgressReporter(ctx context.Context, r ProgressReporter) context.Context {
    return context.WithValue(ctx, progressCtxKey, r)
}

// 管线节点中使用
reporter := GetProgressReporter(ctx)
reporter.Report(Progress{Stage: "prompt", Percent: 25})
5%(任务创建)→ 25%(提示词优化完成)→ 60%(素材生成完成)→ 80%(质检通过)→ 100%(格式适配 + 上传)。每个阶段由对应节点触发回调。 **WebSocket**:`ws://host/api/v1/tasks/:taskId/ws`,实时推送进度和结果。**HTTP 轮询**:`GET /api/v1/tasks/:taskId`,兼容降级方案。