Skip to content

Commit de1e7b0

Browse files
committed
improved error handling and reliability for v2 migration
1 parent 2d695c0 commit de1e7b0

1 file changed

Lines changed: 42 additions & 16 deletions

File tree

‎Hub/internal/store/migrations.go‎

Lines changed: 42 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -4,6 +4,7 @@ import (
44
"database/sql"
55
"fmt"
66
"log"
7+
"strings"
78
"time"
89
)
910

@@ -131,10 +132,14 @@ var migrations = []Migration{
131132
// 2. Migrate existing data
132133
rows, err := tx.Query(`SELECT node_id, sensor_id, time_bucket FROM sensor_heartbeats ORDER BY node_id, sensor_id, time_bucket ASC`)
133134
if err != nil {
134-
// Table might not exist or be empty
135+
// Table might not exist (fresh install) - only swallow that specific case,
136+
// let any other query failure propagate and abort the migration.
137+
if !strings.Contains(err.Error(), "no such table") {
138+
return fmt.Errorf("querying sensor_heartbeats: %w", err)
139+
}
135140
return nil
136141
}
137-
142+
138143
type hb struct {
139144
NodeID string
140145
SensorID string
@@ -143,12 +148,20 @@ var migrations = []Migration{
143148
var heartbeats []hb
144149
for rows.Next() {
145150
var n, s, b string
146-
if err := rows.Scan(&n, &s, &b); err == nil {
147-
t, err := time.Parse(time.RFC3339, b)
148-
if err == nil {
149-
heartbeats = append(heartbeats, hb{NodeID: n, SensorID: s, Bucket: t})
150-
}
151+
if err := rows.Scan(&n, &s, &b); err != nil {
152+
rows.Close()
153+
return fmt.Errorf("scanning heartbeat row: %w", err)
151154
}
155+
t, err := time.Parse(time.RFC3339, b)
156+
if err != nil {
157+
rows.Close()
158+
return fmt.Errorf("parsing heartbeat timestamp %q: %w", b, err)
159+
}
160+
heartbeats = append(heartbeats, hb{NodeID: n, SensorID: s, Bucket: t})
161+
}
162+
if err := rows.Err(); err != nil {
163+
rows.Close()
164+
return fmt.Errorf("iterating heartbeat rows: %w", err)
152165
}
153166
rows.Close()
154167

@@ -163,38 +176,51 @@ var migrations = []Migration{
163176
var lastTime time.Time
164177
nowTime := time.Now().UTC()
165178

166-
closeSensorSequence := func() {
179+
closeSensorSequence := func() error {
167180
if lastNode != "" && lastSensor != "" {
168181
if nowTime.Sub(lastTime) > time.Minute {
169182
offlineTime := lastTime.Add(time.Minute)
170-
stmt.Exec(lastNode, lastSensor, "offline", offlineTime.Format(time.RFC3339))
183+
if _, err := stmt.Exec(lastNode, lastSensor, "offline", offlineTime.Format(time.RFC3339)); err != nil {
184+
return fmt.Errorf("inserting offline status for %s/%s: %w", lastNode, lastSensor, err)
185+
}
171186
}
172187
}
188+
return nil
173189
}
174190

175191
for _, h := range heartbeats {
176192
if h.NodeID != lastNode || h.SensorID != lastSensor {
177193
// Close previous sensor sequence if it died before now
178-
closeSensorSequence()
194+
if err := closeSensorSequence(); err != nil {
195+
return err
196+
}
179197

180198
// New sensor sequence: insert 'online'
181-
stmt.Exec(h.NodeID, h.SensorID, "online", h.Bucket.Format(time.RFC3339))
199+
if _, err := stmt.Exec(h.NodeID, h.SensorID, "online", h.Bucket.Format(time.RFC3339)); err != nil {
200+
return fmt.Errorf("inserting online status for %s/%s: %w", h.NodeID, h.SensorID, err)
201+
}
182202
} else {
183203
// Existing sequence: check gap
184204
gap := h.Bucket.Sub(lastTime)
185205
if gap > time.Minute { // Missed at least one minute bucket (>60s gap)
186206
offlineTime := lastTime.Add(time.Minute)
187-
stmt.Exec(h.NodeID, h.SensorID, "offline", offlineTime.Format(time.RFC3339))
188-
stmt.Exec(h.NodeID, h.SensorID, "online", h.Bucket.Format(time.RFC3339))
207+
if _, err := stmt.Exec(h.NodeID, h.SensorID, "offline", offlineTime.Format(time.RFC3339)); err != nil {
208+
return fmt.Errorf("inserting offline status for %s/%s: %w", h.NodeID, h.SensorID, err)
209+
}
210+
if _, err := stmt.Exec(h.NodeID, h.SensorID, "online", h.Bucket.Format(time.RFC3339)); err != nil {
211+
return fmt.Errorf("inserting online status for %s/%s: %w", h.NodeID, h.SensorID, err)
212+
}
189213
}
190214
}
191215
lastNode = h.NodeID
192216
lastSensor = h.SensorID
193217
lastTime = h.Bucket
194218
}
195-
219+
196220
// Close the final sensor sequence
197-
closeSensorSequence()
221+
if err := closeSensorSequence(); err != nil {
222+
return err
223+
}
198224
}
199225

200226
// 3. Drop old table
@@ -298,4 +324,4 @@ func RunMigrations(db *sql.DB) error {
298324

299325
log.Println("[DB] Database up to date.")
300326
return nil
301-
}
327+
}

0 commit comments

Comments
 (0)