MCPcopy Create free account
hub / github.com/coder/wush / handleStreaming

Method handleStreaming

tsserver/server.go:492–570  ·  view source on GitHub ↗
(ctx context.Context, w http.ResponseWriter, req *tailcfg.MapRequest)

Source from the content-addressed store, hash-verified

490}
491
492func (ns *noiseServer) handleStreaming(ctx context.Context, w http.ResponseWriter, req *tailcfg.MapRequest) {
493 rc := http.NewResponseController(w)
494 // Longpolling will break if there is a write timeout, so it needs to be
495 // disabled.
496 rc.SetWriteDeadline(time.Time{})
497
498 node := ns.getSelfNode()
499
500 keepAlive := time.NewTicker(50 * time.Second)
501 defer keepAlive.Stop()
502
503 res := &tailcfg.MapResponse{
504 KeepAlive: false,
505 ControlTime: ptr.To(time.Now()),
506 Node: node,
507 DERPMap: ns.derpMap,
508 CollectServices: opt.NewBool(false),
509 Debug: &tailcfg.Debug{
510 DisableLogTail: true,
511 },
512 Peers: ns.peerMap(),
513 PacketFilter: tailcfg.FilterAllowAll,
514 }
515
516 err := writeMapResponse(w, req, res)
517 if err != nil {
518 ns.logger.Error("write map response", "err", err)
519 return
520 }
521 err = rc.Flush()
522 if err != nil {
523 ns.logger.Error("flush map response", "err", err)
524 return
525 }
526
527 for {
528 select {
529 case <-ctx.Done():
530 return
531 case upd := <-ns.peerUpdate:
532 res := &tailcfg.MapResponse{
533 KeepAlive: false,
534 ControlTime: ptr.To(time.Now()),
535 }
536 if upd.ty == updateTypeNewPeer {
537 ns.peers.Store(upd.node.ID, upd.node.Clone())
538 res.Peers = ns.peerMap()
539 } else if upd.ty == updateTypePeerUpdate {
540 res.PeersChangedPatch = []*tailcfg.PeerChange{upd.update}
541 }
542
543 err := writeMapResponse(w, req, res)
544 if err != nil {
545 ns.logger.Error("write map response", "err", err)
546 return
547 }
548 err = rc.Flush()
549 if err != nil {

Callers 1

Calls 3

getSelfNodeMethod · 0.95
peerMapMethod · 0.95
writeMapResponseFunction · 0.85

Tested by

no test coverage detected