Repository navigation
Expand file tree
/
Copy pathprocess.go
More file actions
111 lines (95 loc) · 3.23 KB
/
Copy pathprocess.go
File metadata and controls
111 lines (95 loc) · 3.23 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
package importer
import (
"errors"
"fmt"
"log"
"time"
)
// Commit every 100 items to keep transactions small.
const commitInterval = 100
func (im *Importer) processItems() ([]ImportResponseData, error) {
canceled, err := im.s.isImportCanceled(im.tenantID)
if err != nil {
im.finalStatus = "FAILED"
return nil, err
}
if canceled {
im.finalStatus = "CANCELED"
log.Printf("import canceled detected: tenant=%s importID=%s", im.tenantID, im.importID)
return nil, errors.New("import canceled")
}
out := make([]ImportResponseData, 0, len(im.items))
lastLeaseRenewal := time.Now()
for i := range im.items {
if time.Since(lastLeaseRenewal) > leaseRenewalInterval {
// Check for user cancellation before renewing the lease.
canceled, err := im.s.isImportCanceled(im.tenantID)
if err != nil {
im.finalStatus = "FAILED"
return out, fmt.Errorf("check cancel status: %w", err)
} else if canceled {
im.finalStatus = "CANCELED"
log.Printf("import canceled mid-run: tenant=%s importID=%s item=%d", im.tenantID, im.importID, i)
return out, errors.New("import canceled")
}
if err := im.s.renewTenantImportLease(im.tenantID, im.importID); err != nil {
im.finalStatus = "FAILED"
return out, fmt.Errorf("renew import lease: %w", err)
}
lastLeaseRenewal = time.Now()
}
itemResult, err := im.processItem(i, &im.items[i])
if err != nil {
im.finalStatus = "FAILED"
return out, fmt.Errorf("process item %d: %w", i, err)
}
out = append(out, itemResult)
// Commit this batch and start a new transaction.
if (i+1)%commitInterval == 0 && i < len(im.items)-1 {
if err := im.flushStockAdjustments(); err != nil {
im.finalStatus = "FAILED"
return nil, fmt.Errorf("flush stock adjustments at item %d: %w", i, err)
}
if err := im.tx.Commit().Error; err != nil {
im.finalStatus = "FAILED"
return nil, fmt.Errorf("mini-batch commit failed at item %d: %w", i, err)
}
im.tx = im.s.Db.Begin()
if im.tx.Error != nil {
im.finalStatus = "FAILED"
return nil, fmt.Errorf("failed to begin transaction after mini-batch commit at item %d: %w", i, im.tx.Error)
}
}
}
if err := im.flushStockAdjustments(); err != nil {
im.finalStatus = "FAILED"
return nil, fmt.Errorf("flush stock adjustments at end: %w", err)
}
log.Printf("Processed %d items", len(im.items))
return out, nil
}
func (im *Importer) processItem(i int, item *ImportItem) (ImportResponseData, error) {
var errMsg string
// Keep one failed item from rolling back earlier items in the batch.
savepoint := fmt.Sprintf("sp_%d", i)
if err := im.tx.SavePoint(savepoint).Error; err != nil {
return ImportResponseData{}, fmt.Errorf("create savepoint %s: %w", savepoint, err)
}
itemID, err := im.handleItem(*item)
if err != nil {
if rbErr := im.tx.RollbackTo(savepoint).Error; rbErr != nil {
return ImportResponseData{}, fmt.Errorf("rollback to savepoint %s: %w", savepoint, rbErr)
}
log.Printf("import item failed (i=%d itemID=%s sku=%s barcode=%s): %v",
i, im.items[i].ItemID, im.items[i].SKU, im.items[i].Barcode, err,
)
errMsg = err.Error()
} else {
im.items[i].ItemID = itemID
}
return ImportResponseData{
ItemID: item.ItemID,
ItemName: item.ButtonName,
ErrorMessage: errMsg,
}, nil
}