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
43 lines
1.0 KiB
Go
43 lines
1.0 KiB
Go
// Copyright (c) 2026 Lark Technologies Pte. Ltd.
|
|
// SPDX-License-Identifier: MIT
|
|
|
|
package consume
|
|
|
|
import (
|
|
"bufio"
|
|
"bytes"
|
|
"net"
|
|
"time"
|
|
|
|
"github.com/larksuite/cli/internal/event/protocol"
|
|
)
|
|
|
|
const preShutdownAckTimeout = 2 * time.Second
|
|
|
|
// checkLastForKey atomically reserves a cleanup lock; on any error defaults to true
|
|
// (cleanup-on-error is safer than leaking server state). Discards non-ack frames in flight.
|
|
func checkLastForKey(conn net.Conn, eventKey string, subscriptionID string) bool {
|
|
msg := protocol.NewPreShutdownCheck(eventKey, subscriptionID)
|
|
if err := protocol.EncodeWithDeadline(conn, msg, protocol.WriteTimeout); err != nil {
|
|
return true
|
|
}
|
|
|
|
if err := conn.SetReadDeadline(time.Now().Add(preShutdownAckTimeout)); err != nil {
|
|
return true
|
|
}
|
|
br := bufio.NewReader(conn)
|
|
for {
|
|
line, err := protocol.ReadFrame(br)
|
|
if err != nil {
|
|
return true
|
|
}
|
|
resp, err := protocol.Decode(bytes.TrimRight(line, "\n"))
|
|
if err != nil {
|
|
continue
|
|
}
|
|
if ack, ok := resp.(*protocol.PreShutdownAck); ok {
|
|
return ack.LastForKey
|
|
}
|
|
}
|
|
}
|