aviortui/essential/hidden/service/worker/clock/clock.go
2026-08-28 12:45:22 +08:00

515 lines
11 KiB
Go

/*
|--------------------------------------------------------------------------
| Description
|--------------------------------------------------------------------------
|
| Name:
| - Clock
|
| Purpose:
| - Provides the clock synchronization service for maintaining a
| consistent system time reference across automation tasks.
|
|--------------------------------------------------------------------------
*/
/*
|--------------------------------------------------------------------------
| Credit
|--------------------------------------------------------------------------
|
| Contributor:
| - Rajon Ahmed
| - Blockonite
|
|--------------------------------------------------------------------------
*/
package clock
import (
"encoding/binary"
"encoding/json"
"errors"
"net"
"os"
"path/filepath"
"runtime"
"sync"
"time"
)
const (
obrimClockSoftwareName = "obrim"
obrimClockNTPPort = 123
obrimClockNTPPacketSize = 48
obrimClockNTPVersion = 4
obrimClockNTPClientMode = 3
obrimClockNTPUnixOffset = 2208988800
obrimClockMaxRTT = 2 * time.Second
obrimClockSynchronizationInterval = 15 * time.Minute
obrimClockServerTimeout = 2 * time.Second
obrimClockSuccess = 0
obrimClockFailure = 1
)
var (
obrimClockMutex sync.RWMutex
obrimClockStop chan struct{}
obrimClockDone chan struct{}
obrimClockOffset int64
obrimClockLastSync int64
obrimClockStarted bool
)
// obrimClockConfig represents the persistent clock configuration.
type obrimClockConfig struct {
Clock *obrimClockState `json:"clock,omitempty"`
}
// obrimClockFile represents the persistent configuration document.
type obrimClockFile struct {
Config *obrimClockConfig `json:"config,omitempty"`
}
// obrimClockState represents the runtime and persistent synchronization state.
type obrimClockState struct {
ClockOffset int64 `json:"clockOffset"`
LastSync int64 `json:"lastSync"`
}
// obrimClockMeasurement represents one accepted or rejected NTP measurement.
type obrimClockMeasurement struct {
ServerTime time.Time
RequestTime time.Time
ResponseTime time.Time
RTT time.Duration
Offset time.Duration
Valid bool
}
// obrimClockResponse represents the timestamps extracted from an NTP response.
type obrimClockResponse struct {
ReceiveTimestamp uint64
TransmitTimestamp uint64
}
// ObrimClockStart starts the clock synchronization goroutine.
func ObrimClockStart() {
obrimClockMutex.Lock()
if obrimClockStarted {
obrimClockMutex.Unlock()
return
}
obrimClockStop = make(chan struct{})
obrimClockDone = make(chan struct{})
obrimClockStarted = true
obrimClockInitializeClock()
stop := obrimClockStop
done := obrimClockDone
obrimClockMutex.Unlock()
go func() {
defer close(done)
obrimClockSynchronizeClock()
ticker := time.NewTicker(obrimClockSynchronizationInterval)
defer ticker.Stop()
for {
select {
case <-ticker.C:
obrimClockSynchronizeClock()
case <-stop:
return
}
}
}()
}
// ObrimClockStop gracefully stops the clock synchronization goroutine.
func ObrimClockStop() {
obrimClockMutex.Lock()
if !obrimClockStarted {
obrimClockMutex.Unlock()
return
}
stop := obrimClockStop
done := obrimClockDone
obrimClockStop = nil
obrimClockDone = nil
obrimClockStarted = false
obrimClockMutex.Unlock()
close(stop)
<-done
}
// obrimClockInitializeClock initializes runtime clock state from persistent storage.
func obrimClockInitializeClock() {
state, result := obrimClockLoadClock()
obrimClockMutex.Lock()
defer obrimClockMutex.Unlock()
if result == obrimClockSuccess && state != nil {
obrimClockOffset = state.ClockOffset
obrimClockLastSync = state.LastSync
return
}
obrimClockOffset = 0
obrimClockLastSync = 0
_ = obrimClockPersistClock()
}
// obrimClockSynchronizeClock performs one NTP synchronization cycle.
func obrimClockSynchronizeClock() {
servers := []string{
"time.cloudflare.com",
"time.google.com",
"pool.ntp.org",
}
measurements := make([]obrimClockMeasurement, 0, len(servers))
for _, server := range servers {
measurement, result := obrimClockQueryServer(server)
if result != obrimClockSuccess {
continue
}
if obrimClockMeasureDelay(&measurement) != obrimClockSuccess {
continue
}
if obrimClockCalculateOffset(&measurement) != obrimClockSuccess {
continue
}
if obrimClockValidateMeasurement(&measurement) != obrimClockSuccess {
continue
}
measurements = append(measurements, measurement)
}
if len(measurements) == 0 {
return
}
offset, result := obrimClockCalculateAverageOffset(measurements)
if result != obrimClockSuccess {
return
}
clockOffset := offset.Nanoseconds()
lastSync := time.Now().UnixNano()
obrimClockUpdateClock(clockOffset, lastSync)
_ = obrimClockPersistClock()
}
// obrimClockQueryServer sends an NTP request and receives a server response.
func obrimClockQueryServer(server string) (obrimClockMeasurement, int) {
var measurement obrimClockMeasurement
address := net.JoinHostPort(server, "123")
connection, err := net.DialTimeout("udp", address, obrimClockServerTimeout)
if err != nil {
return measurement, obrimClockFailure
}
defer connection.Close()
request := make([]byte, obrimClockNTPPacketSize)
request[0] = (obrimClockNTPVersion << 3) | obrimClockNTPClientMode
requestTime := time.Now()
if _, err := connection.Write(request); err != nil {
return measurement, obrimClockFailure
}
if err := connection.SetReadDeadline(time.Now().Add(obrimClockServerTimeout)); err != nil {
return measurement, obrimClockFailure
}
response := make([]byte, obrimClockNTPPacketSize)
if _, err := connection.Read(response); err != nil {
return measurement, obrimClockFailure
}
responseTime := time.Now()
if len(response) < obrimClockNTPPacketSize {
return measurement, obrimClockFailure
}
receiveTimestamp := binary.BigEndian.Uint64(response[32:40])
transmitTimestamp := binary.BigEndian.Uint64(response[40:48])
if transmitTimestamp == 0 {
return measurement, obrimClockFailure
}
serverTime, err := obrimClockNTPToTime(transmitTimestamp)
if err != nil {
return measurement, obrimClockFailure
}
measurement = obrimClockMeasurement{
ServerTime: serverTime,
RequestTime: requestTime,
ResponseTime: responseTime,
}
_ = receiveTimestamp
return measurement, obrimClockSuccess
}
// obrimClockMeasureDelay calculates network round-trip time.
func obrimClockMeasureDelay(measurement *obrimClockMeasurement) int {
if measurement == nil {
return obrimClockFailure
}
if measurement.ResponseTime.Before(measurement.RequestTime) {
return obrimClockFailure
}
measurement.RTT = measurement.ResponseTime.Sub(measurement.RequestTime)
if measurement.RTT <= 0 || measurement.RTT > obrimClockMaxRTT {
return obrimClockFailure
}
return obrimClockSuccess
}
// obrimClockCalculateOffset calculates the local clock offset from an NTP response.
func obrimClockCalculateOffset(measurement *obrimClockMeasurement) int {
if measurement == nil {
return obrimClockFailure
}
midpoint := measurement.RequestTime.Add(measurement.RTT / 2)
measurement.Offset = measurement.ServerTime.Sub(midpoint)
return obrimClockSuccess
}
// obrimClockValidateMeasurement validates an NTP synchronization measurement.
func obrimClockValidateMeasurement(measurement *obrimClockMeasurement) int {
if measurement == nil {
return obrimClockFailure
}
if measurement.RTT <= 0 || measurement.RTT > obrimClockMaxRTT {
return obrimClockFailure
}
if measurement.ServerTime.IsZero() {
return obrimClockFailure
}
if measurement.RequestTime.IsZero() || measurement.ResponseTime.IsZero() {
return obrimClockFailure
}
if measurement.ResponseTime.Before(measurement.RequestTime) {
return obrimClockFailure
}
measurement.Valid = true
return obrimClockSuccess
}
// obrimClockCalculateAverageOffset calculates the trusted average clock offset.
func obrimClockCalculateAverageOffset(measurements []obrimClockMeasurement) (time.Duration, int) {
var total int64
var count int64
for _, measurement := range measurements {
if !measurement.Valid {
continue
}
total += measurement.Offset.Nanoseconds()
count++
}
if count == 0 {
return 0, obrimClockFailure
}
return time.Duration(total / count), obrimClockSuccess
}
// obrimClockUpdateClock updates the runtime clock state.
func obrimClockUpdateClock(clockOffset int64, lastSync int64) {
obrimClockMutex.Lock()
defer obrimClockMutex.Unlock()
obrimClockOffset = clockOffset
obrimClockLastSync = lastSync
}
// obrimClockPersistClock writes the runtime clock state to persistent storage.
func obrimClockPersistClock() int {
obrimClockMutex.RLock()
state := obrimClockState{
ClockOffset: obrimClockOffset,
LastSync: obrimClockLastSync,
}
obrimClockMutex.RUnlock()
configPath, err := obrimClockConfigPath()
if err != nil {
return obrimClockFailure
}
data, err := os.ReadFile(configPath)
if err != nil {
if errors.Is(err, os.ErrNotExist) {
return obrimClockFailure
}
return obrimClockFailure
}
var document obrimClockFile
if len(data) != 0 {
if err := json.Unmarshal(data, &document); err != nil {
return obrimClockFailure
}
}
if document.Config == nil {
document.Config = &obrimClockConfig{}
}
document.Config.Clock = &state
updatedData, err := json.MarshalIndent(document, "", " ")
if err != nil {
return obrimClockFailure
}
if err := os.WriteFile(configPath, updatedData, 0644); err != nil {
return obrimClockFailure
}
return obrimClockSuccess
}
// obrimClockCurrentClock returns the current synchronized UTC.
func obrimClockCurrentClock() time.Time {
obrimClockMutex.RLock()
offset := obrimClockOffset
obrimClockMutex.RUnlock()
return time.Now().UTC().Add(time.Duration(offset))
}
// obrimClockLoadClock loads the persisted clock state.
func obrimClockLoadClock() (*obrimClockState, int) {
configPath, err := obrimClockConfigPath()
if err != nil {
return nil, obrimClockFailure
}
data, err := os.ReadFile(configPath)
if err != nil {
return nil, obrimClockFailure
}
var document obrimClockFile
if err := json.Unmarshal(data, &document); err != nil {
return nil, obrimClockFailure
}
if document.Config == nil || document.Config.Clock == nil {
return nil, obrimClockFailure
}
return document.Config.Clock, obrimClockSuccess
}
// obrimClockConfigPath returns the platform-specific persistent configuration path.
func obrimClockConfigPath() (string, error) {
var basePath string
switch runtime.GOOS {
case "windows":
basePath = os.Getenv("APPDATA")
if basePath == "" {
return "", errors.New("APPDATA is not defined")
}
case "darwin":
homePath, err := os.UserHomeDir()
if err != nil {
return "", err
}
basePath = filepath.Join(homePath, "Library", "Application Support")
default:
homePath, err := os.UserHomeDir()
if err != nil {
return "", err
}
basePath = filepath.Join(homePath, ".config")
}
return filepath.Join(
basePath,
obrimClockSoftwareName,
"persistent",
"config",
"config.json",
), nil
}
// obrimClockNTPToTime converts an NTP timestamp into Unix UTC time.
func obrimClockNTPToTime(timestamp uint64) (time.Time, error) {
seconds := timestamp >> 32
fraction := timestamp & 0xffffffff
if seconds < obrimClockNTPUnixOffset {
return time.Time{}, errors.New("invalid NTP timestamp")
}
unixSeconds := int64(seconds) - obrimClockNTPUnixOffset
nanoseconds := int64((fraction * 1_000_000_000) >> 32)
return time.Unix(unixSeconds, nanoseconds).UTC(), nil
}