Files
wehub-resource-sync bf9395e022
CI / license-header (push) Has been skipped
CI / e2e-dry-run (push) Has been skipped
CI / fast-gate (push) Failing after 0s
Test PR Label Logic / test-pr-labels (push) Failing after 1s
Skill Format Check / check-format (push) Failing after 2s
CI / security (push) Failing after 5s
CI / unit-test (push) Has been skipped
CI / lint (push) Has been skipped
CI / script-test (push) Has been skipped
CI / deterministic-gate (push) Has been skipped
CI / coverage (push) Has been skipped
CI / results (push) Has been cancelled
CI / deadcode (push) Has been cancelled
CI / e2e-live (push) Has been cancelled
chore: import upstream snapshot with attribution
2026-07-13 12:22:54 +08:00

47 lines
1.4 KiB
Go

// Copyright (c) 2026 Lark Technologies Pte. Ltd.
// SPDX-License-Identifier: MIT
package consume
import (
"bufio"
"bytes"
"fmt"
"net"
"os"
"time"
"github.com/larksuite/cli/internal/event/protocol"
)
const helloAckTimeout = 5 * time.Second // symmetric with bus-side hello read deadline
// doHello returns a bufio.Reader holding any bytes already pulled off conn so events
// buffered with the ack in one TCP segment aren't dropped.
func doHello(conn net.Conn, eventKey string, eventTypes []string, subscriptionID string) (*protocol.HelloAck, *bufio.Reader, error) {
hello := protocol.NewHello(os.Getpid(), eventKey, eventTypes, "v1", subscriptionID)
if err := protocol.EncodeWithDeadline(conn, hello, protocol.WriteTimeout); err != nil {
return nil, nil, err
}
if err := conn.SetReadDeadline(time.Now().Add(helloAckTimeout)); err != nil {
return nil, nil, fmt.Errorf("set hello_ack deadline: %w", err)
}
br := bufio.NewReader(conn)
line, err := protocol.ReadFrame(br)
if err != nil {
return nil, nil, fmt.Errorf("no hello_ack received: %w", err)
}
// best-effort clear; if the conn is already broken, the loop's first read will surface it
_ = conn.SetReadDeadline(time.Time{})
msg, err := protocol.Decode(bytes.TrimRight(line, "\n"))
if err != nil {
return nil, nil, err
}
ack, ok := msg.(*protocol.HelloAck)
if !ok {
return nil, nil, fmt.Errorf("expected hello_ack, got %T", msg)
}
return ack, br, nil
}