tuic_journal.go 3.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135
  1. package job
  2. import (
  3. "encoding/json"
  4. "fmt"
  5. "os"
  6. "path/filepath"
  7. "runtime"
  8. "sync"
  9. "time"
  10. "github.com/google/uuid"
  11. "github.com/mhsanaei/3x-ui/v3/internal/config"
  12. "github.com/mhsanaei/3x-ui/v3/internal/tuic"
  13. )
  14. var tuicJournalMu sync.Mutex
  15. type tuicTrafficBatch struct {
  16. ID string `json:"id"`
  17. Deltas []tuic.ClientTrafficDelta `json:"deltas"`
  18. }
  19. func tuicJournalDir() string { return filepath.Join(config.GetDBFolderPath(), "tuic-traffic-journal") }
  20. func syncJournalDir(dir string) error {
  21. if runtime.GOOS == "windows" {
  22. return nil
  23. }
  24. file, err := os.Open(dir)
  25. if err != nil {
  26. return err
  27. }
  28. defer file.Close()
  29. return file.Sync()
  30. }
  31. func storeTuicBatch(batch tuicTrafficBatch) (string, error) {
  32. dir := tuicJournalDir()
  33. if err := os.MkdirAll(dir, 0o700); err != nil {
  34. return "", err
  35. }
  36. if err := syncJournalDir(filepath.Dir(dir)); err != nil {
  37. return "", err
  38. }
  39. data, err := json.Marshal(batch)
  40. if err != nil {
  41. return "", err
  42. }
  43. file, err := os.CreateTemp(dir, ".pending-")
  44. if err != nil {
  45. return "", err
  46. }
  47. tmp := file.Name()
  48. defer os.Remove(tmp)
  49. if _, err := file.Write(data); err != nil {
  50. file.Close()
  51. return "", err
  52. }
  53. if err := file.Sync(); err != nil {
  54. file.Close()
  55. return "", err
  56. }
  57. if err := file.Close(); err != nil {
  58. return "", err
  59. }
  60. path := filepath.Join(dir, batch.ID+".json")
  61. if err := os.Rename(tmp, path); err != nil {
  62. return "", err
  63. }
  64. return path, syncJournalDir(dir)
  65. }
  66. func (j *TuicJob) replayTuicJournal() error {
  67. entries, err := os.ReadDir(tuicJournalDir())
  68. if os.IsNotExist(err) {
  69. return nil
  70. }
  71. if err != nil {
  72. return err
  73. }
  74. for _, entry := range entries {
  75. if entry.IsDir() || filepath.Ext(entry.Name()) != ".json" {
  76. continue
  77. }
  78. path := filepath.Join(tuicJournalDir(), entry.Name())
  79. data, err := os.ReadFile(path)
  80. if err != nil {
  81. return err
  82. }
  83. var batch tuicTrafficBatch
  84. if err := json.Unmarshal(data, &batch); err != nil {
  85. return fmt.Errorf("TUIC journal %s: %w", entry.Name(), err)
  86. }
  87. if batch.ID+".json" != entry.Name() {
  88. return fmt.Errorf("TUIC journal batch ID mismatch")
  89. }
  90. if err := j.inboundService.AddTuicTrafficBatch(batch.ID, aggregateTuicClientTraffic(batch.Deltas, nil)); err != nil {
  91. return err
  92. }
  93. if err := os.Remove(path); err != nil {
  94. return err
  95. }
  96. if err := syncJournalDir(tuicJournalDir()); err != nil {
  97. return err
  98. }
  99. }
  100. return nil
  101. }
  102. func (j *TuicJob) flushTuicJournal() error {
  103. tuicJournalMu.Lock()
  104. defer tuicJournalMu.Unlock()
  105. manager := tuic.GetManager()
  106. _, deltas := manager.CollectAllTraffic()
  107. if len(deltas) > 0 {
  108. if path, err := storeTuicBatch(tuicTrafficBatch{ID: uuid.NewString(), Deltas: deltas}); err != nil {
  109. if path == "" {
  110. manager.RequeueClientTraffic(deltas)
  111. }
  112. return fmt.Errorf("persist TUIC shutdown journal: %w", err)
  113. }
  114. }
  115. var err error
  116. for attempt := 0; attempt < 3; attempt++ {
  117. if err = j.replayTuicJournal(); err == nil {
  118. return nil
  119. }
  120. if attempt < 2 {
  121. time.Sleep(100 * time.Millisecond)
  122. }
  123. }
  124. return fmt.Errorf("TUIC traffic retained in durable journal: %w", err)
  125. }