handleLogsStream 同步阻塞,把 nodeapi.LogsChannel 的每行包成 {"line":"..."} 通过 sender 推回 server,event="log"。 ctx 取消(caller Close 或 server cancel_id)或日志通道关闭时返回 nil; session 会在 Handle 返回后回一条普通 Ok 响应,作为流终态。
(ctx context.Context)
| 213 | // ctx 取消(caller Close 或 server cancel_id)或日志通道关闭时返回 nil; |
| 214 | // session 会在 Handle 返回后回一条普通 Ok 响应,作为流终态。 |
| 215 | func (d *APIDispatcher) handleLogsStream(ctx context.Context) error { |
| 216 | if d.sender == nil { |
| 217 | return errors.New("nodeagent: sender not initialized for streaming") |
| 218 | } |
| 219 | reqID := ReqIDFromContext(ctx) |
| 220 | if reqID == "" { |
| 221 | return errors.New("nodeagent: LogsStream missing reqID in ctx") |
| 222 | } |
| 223 | ch := d.api.LogsChannel(ctx) |
| 224 | for line := range ch { |
| 225 | body, err := json.Marshal(map[string]string{"line": line}) |
| 226 | if err != nil { |
| 227 | return fmt.Errorf("marshal log frame: %w", err) |
| 228 | } |
| 229 | if err := d.sender.PushEvent(reqID, "log", body, 0); err != nil { |
| 230 | return err |
| 231 | } |
| 232 | if ctx.Err() != nil { |
| 233 | return nil |
| 234 | } |
| 235 | } |
| 236 | return nil |
| 237 | } |
| 238 | |
| 239 | // handleTracerouteStream 同步阻塞,把 nodeapi.TracerouteHops 的每跳推回 server, |
| 240 | // event="traceroute_hop"。错误事件会作为终态错误返回(session 会回一条 Ok=false 响应)。 |
no test coverage detected