QueryLatencySamples 查询时间范围内的采样,按 sampled_at 升序。
(nodeIDs []string, from, to time.Time)
| 34 | |
| 35 | // QueryLatencySamples 查询时间范围内的采样,按 sampled_at 升序。 |
| 36 | func (s *NodeStore) QueryLatencySamples(nodeIDs []string, from, to time.Time) ([]nodes.LatencySample, error) { |
| 37 | if len(nodeIDs) == 0 { |
| 38 | return nil, nil |
| 39 | } |
| 40 | ctx := context.Background() |
| 41 | rows, err := s.db.Query(ctx, |
| 42 | `SELECT node_id, isp, rtt_ms, sampled_at FROM node_latency_samples |
| 43 | WHERE sampled_at >= $1 AND sampled_at <= $2 |
| 44 | AND node_id = ANY($3) |
| 45 | ORDER BY sampled_at ASC`, |
| 46 | from, to, nodeIDs, |
| 47 | ) |
| 48 | if err != nil { |
| 49 | return nil, fmt.Errorf("query latency samples: %w", err) |
| 50 | } |
| 51 | defer rows.Close() |
| 52 | |
| 53 | var out []nodes.LatencySample |
| 54 | for rows.Next() { |
| 55 | var sm nodes.LatencySample |
| 56 | if err := rows.Scan(&sm.NodeID, &sm.ISP, &sm.RttMs, &sm.SampledAt); err != nil { |
| 57 | return nil, err |
| 58 | } |
| 59 | out = append(out, sm) |
| 60 | } |
| 61 | return out, rows.Err() |
| 62 | } |
| 63 | |
| 64 | // CleanupOldLatencySamples 删除 before 之前的采样。 |
| 65 | func (s *NodeStore) CleanupOldLatencySamples(before time.Time) error { |