main.go 8.2 KB


  1. // Copyright (C) 2018 The Syncthing Authors.
  2. //
  3. // This Source Code Form is subject to the terms of the Mozilla Public
  4. // License, v. 2.0. If a copy of the MPL was not distributed with this file,
  5. // You can obtain one at https://mozilla.org/MPL/2.0/.
  6. package main
  7. import (
  8. "database/sql"
  9. "log"
  10. "os"
  11. "time"
  12. _ "github.com/lib/pq"
  13. )
  14. var dbConn = getEnvDefault("UR_DB_URL", "postgres://user:password@localhost/ur?sslmode=disable")
  15. func getEnvDefault(key, def string) string {
  16. if val := os.Getenv(key); val != "" {
  17. return val
  18. }
  19. return def
  20. }
  21. func main() {
  22. log.SetFlags(log.Ltime | log.Ldate)
  23. log.SetOutput(os.Stdout)
  24. db, err := sql.Open("postgres", dbConn)
  25. if err != nil {
  26. log.Fatalln("database:", err)
  27. }
  28. err = setupDB(db)
  29. if err != nil {
  30. log.Fatalln("database:", err)
  31. }
  32. for {
  33. runAggregation(db)
  34. // Sleep until one minute past next midnight
  35. sleepUntilNext(24*time.Hour, 1*time.Minute)
  36. }
  37. }
  38. func runAggregation(db *sql.DB) {
  39. since := maxIndexedDay(db, "VersionSummary")
  40. log.Println("Aggregating VersionSummary data since", since)
  41. rows, err := aggregateVersionSummary(db, since)
  42. if err != nil {
  43. log.Fatalln("aggregate:", err)
  44. }
  45. log.Println("Inserted", rows, "rows")
  46. log.Println("Aggregating UserMovement data")
  47. rows, err = aggregateUserMovement(db)
  48. if err != nil {
  49. log.Fatalln("aggregate:", err)
  50. }
  51. log.Println("Inserted", rows, "rows")
  52. log.Println("Aggregating Performance data")
  53. since = maxIndexedDay(db, "Performance")
  54. rows, err = aggregatePerformance(db, since)
  55. if err != nil {
  56. log.Fatalln("aggregate:", err)
  57. }
  58. log.Println("Inserted", rows, "rows")
  59. log.Println("Aggregating BlockStats data")
  60. since = maxIndexedDay(db, "BlockStats")
  61. rows, err = aggregateBlockStats(db, since)
  62. if err != nil {
  63. log.Fatalln("aggregate:", err)
  64. }
  65. log.Println("Inserted", rows, "rows")
  66. }
  67. func sleepUntilNext(intv, margin time.Duration) {
  68. now := time.Now().UTC()
  69. next := now.Truncate(intv).Add(intv).Add(margin)
  70. log.Println("Sleeping until", next)
  71. time.Sleep(next.Sub(now))
  72. }
  73. func setupDB(db *sql.DB) error {
  74. _, err := db.Exec(`CREATE TABLE IF NOT EXISTS VersionSummary (
  75. Day TIMESTAMP NOT NULL,
  76. Version VARCHAR(8) NOT NULL,
  77. Count INTEGER NOT NULL
  78. )`)
  79. if err != nil {
  80. return err
  81. }
  82. _, err = db.Exec(`CREATE TABLE IF NOT EXISTS UserMovement (
  83. Day TIMESTAMP NOT NULL,
  84. Added INTEGER NOT NULL,
  85. Bounced INTEGER NOT NULL,
  86. Removed INTEGER NOT NULL
  87. )`)
  88. if err != nil {
  89. return err
  90. }
  91. _, err = db.Exec(`CREATE TABLE IF NOT EXISTS Performance (
  92. Day TIMESTAMP NOT NULL,
  93. TotFiles INTEGER NOT NULL,
  94. TotMiB INTEGER NOT NULL,
  95. SHA256Perf DOUBLE PRECISION NOT NULL,
  96. MemorySize INTEGER NOT NULL,
  97. MemoryUsageMiB INTEGER NOT NULL
  98. )`)
  99. if err != nil {
  100. return err
  101. }
  102. _, err = db.Exec(`CREATE TABLE IF NOT EXISTS BlockStats (
  103. Day TIMESTAMP NOT NULL,
  104. Reports INTEGER NOT NULL,
  105. Total INTEGER NOT NULL,
  106. Renamed INTEGER NOT NULL,
  107. Reused INTEGER NOT NULL,
  108. Pulled INTEGER NOT NULL,
  109. CopyOrigin INTEGER NOT NULL,
  110. CopyOriginShifted INTEGER NOT NULL,
  111. CopyElsewhere INTEGER NOT NULL
  112. )`)
  113. if err != nil {
  114. return err
  115. }
  116. var t string
  117. row := db.QueryRow(`SELECT 'UniqueDayVersionIndex'::regclass`)
  118. if err := row.Scan(&t); err != nil {
  119. _, err = db.Exec(`CREATE UNIQUE INDEX UniqueDayVersionIndex ON VersionSummary (Day, Version)`)
  120. }
  121. row = db.QueryRow(`SELECT 'VersionDayIndex'::regclass`)
  122. if err := row.Scan(&t); err != nil {
  123. _, err = db.Exec(`CREATE INDEX VersionDayIndex ON VersionSummary (Day)`)
  124. }
  125. row = db.QueryRow(`SELECT 'MovementDayIndex'::regclass`)
  126. if err := row.Scan(&t); err != nil {
  127. _, err = db.Exec(`CREATE INDEX MovementDayIndex ON UserMovement (Day)`)
  128. }
  129. row = db.QueryRow(`SELECT 'PerformanceDayIndex'::regclass`)
  130. if err := row.Scan(&t); err != nil {
  131. _, err = db.Exec(`CREATE INDEX PerformanceDayIndex ON Performance (Day)`)
  132. }
  133. row = db.QueryRow(`SELECT 'BlockStatsDayIndex'::regclass`)
  134. if err := row.Scan(&t); err != nil {
  135. _, err = db.Exec(`CREATE INDEX BlockStatsDayIndex ON BlockStats (Day)`)
  136. }
  137. return err
  138. }
  139. func maxIndexedDay(db *sql.DB, table string) time.Time {
  140. var t time.Time
  141. row := db.QueryRow("SELECT MAX(Day) FROM " + table)
  142. err := row.Scan(&t)
  143. if err != nil {
  144. return time.Time{}
  145. }
  146. return t
  147. }
  148. func aggregateVersionSummary(db *sql.DB, since time.Time) (int64, error) {
  149. res, err := db.Exec(`INSERT INTO VersionSummary (
  150. SELECT
  151. DATE_TRUNC('day', Received) AS Day,
  152. SUBSTRING(Report->>'version' FROM '^v\d.\d+') AS Ver,
  153. COUNT(*) AS Count
  154. FROM ReportsJson
  155. WHERE
  156. DATE_TRUNC('day', Received) > $1
  157. AND DATE_TRUNC('day', Received) < DATE_TRUNC('day', NOW())
  158. AND Report->>'version' like 'v_.%'
  159. GROUP BY Day, Ver
  160. );
  161. `, since)
  162. if err != nil {
  163. return 0, err
  164. }
  165. return res.RowsAffected()
  166. }
  167. func aggregateUserMovement(db *sql.DB) (int64, error) {
  168. rows, err := db.Query(`SELECT
  169. DATE_TRUNC('day', Received) AS Day,
  170. Report->>'uniqueID'
  171. FROM ReportsJson
  172. WHERE
  173. DATE_TRUNC('day', Received) < DATE_TRUNC('day', NOW())
  174. AND Report->>'version' like 'v_.%'
  175. ORDER BY Day
  176. `)
  177. if err != nil {
  178. return 0, err
  179. }
  180. defer rows.Close()
  181. firstSeen := make(map[string]time.Time)
  182. lastSeen := make(map[string]time.Time)
  183. var minTs time.Time
  184. minTs = minTs.In(time.UTC)
  185. for rows.Next() {
  186. var ts time.Time
  187. var id string
  188. if err := rows.Scan(&ts, &id); err != nil {
  189. return 0, err
  190. }
  191. if minTs.IsZero() {
  192. minTs = ts
  193. }
  194. if _, ok := firstSeen[id]; !ok {
  195. firstSeen[id] = ts
  196. }
  197. lastSeen[id] = ts
  198. }
  199. type sumRow struct {
  200. day time.Time
  201. added int
  202. removed int
  203. bounced int
  204. }
  205. var sumRows []sumRow
  206. for t := minTs; t.Before(time.Now().Truncate(24 * time.Hour)); t = t.AddDate(0, 0, 1) {
  207. var added, removed, bounced int
  208. old := t.Before(time.Now().AddDate(0, 0, -30))
  209. for id, first := range firstSeen {
  210. last := lastSeen[id]
  211. if first.Equal(t) && last.Equal(t) && old {
  212. bounced++
  213. continue
  214. }
  215. if first.Equal(t) {
  216. added++
  217. }
  218. if last == t && old {
  219. removed++
  220. }
  221. }
  222. sumRows = append(sumRows, sumRow{t, added, removed, bounced})
  223. }
  224. tx, err := db.Begin()
  225. if err != nil {
  226. return 0, err
  227. }
  228. if _, err := tx.Exec("DELETE FROM UserMovement"); err != nil {
  229. tx.Rollback()
  230. return 0, err
  231. }
  232. for _, r := range sumRows {
  233. if _, err := tx.Exec("INSERT INTO UserMovement (Day, Added, Removed, Bounced) VALUES ($1, $2, $3, $4)", r.day, r.added, r.removed, r.bounced); err != nil {
  234. tx.Rollback()
  235. return 0, err
  236. }
  237. }
  238. return int64(len(sumRows)), tx.Commit()
  239. }
  240. func aggregatePerformance(db *sql.DB, since time.Time) (int64, error) {
  241. res, err := db.Exec(`INSERT INTO Performance (
  242. SELECT
  243. DATE_TRUNC('day', Received) AS Day,
  244. AVG((Report->>'totFiles')::numeric) As TotFiles,
  245. AVG((Report->>'totMiB')::numeric) As TotMiB,
  246. AVG((Report->>'sha256Perf')::numeric) As SHA256Perf,
  247. AVG((Report->>'memorySize')::numeric) As MemorySize,
  248. AVG((Report->>'memoryUsageMiB')::numeric) As MemoryUsageMiB
  249. FROM ReportsJson
  250. WHERE
  251. DATE_TRUNC('day', Received) > $1
  252. AND DATE_TRUNC('day', Received) < DATE_TRUNC('day', NOW())
  253. AND Report->>'version' like 'v_.%'
  254. GROUP BY Day
  255. );
  256. `, since)
  257. if err != nil {
  258. return 0, err
  259. }
  260. return res.RowsAffected()
  261. }
  262. func aggregateBlockStats(db *sql.DB, since time.Time) (int64, error) {
  263. // Filter out anything prior 0.14.41 as that has sum aggregations which
  264. // made no sense.
  265. res, err := db.Exec(`INSERT INTO BlockStats (
  266. SELECT
  267. DATE_TRUNC('day', Received) AS Day,
  268. COUNT(1) As Reports,
  269. SUM((Report->'blockStats'->>'total')::numeric) AS Total,
  270. SUM((Report->'blockStats'->>'renamed')::numeric) AS Renamed,
  271. SUM((Report->'blockStats'->>'reused')::numeric) AS Reused,
  272. SUM((Report->'blockStats'->>'pulled')::numeric) AS Pulled,
  273. SUM((Report->'blockStats'->>'copyOrigin')::numeric) AS CopyOrigin,
  274. SUM((Report->'blockStats'->>'copyOriginShifted')::numeric) AS CopyOriginShifted,
  275. SUM((Report->'blockStats'->>'copyElsewhere')::numeric) AS CopyElsewhere
  276. FROM ReportsJson
  277. WHERE
  278. DATE_TRUNC('day', Received) > $1
  279. AND DATE_TRUNC('day', Received) < DATE_TRUNC('day', NOW())
  280. AND (Report->>'urVersion')::numeric >= 3
  281. AND Report->>'version' like 'v_.%'
  282. AND Report->>'version' NOT LIKE 'v0.14.40%'
  283. AND Report->>'version' NOT LIKE 'v0.14.39%'
  284. AND Report->>'version' NOT LIKE 'v0.14.38%'
  285. GROUP BY Day
  286. );
  287. `, since)
  288. if err != nil {
  289. return 0, err
  290. }
  291. return res.RowsAffected()
  292. }