This repository has no description
1package mill
2
3import (
4 "fmt"
5 "tangled.org/core/spindle/db"
6 millproto "tangled.org/core/spindle/mill/proto"
7 millv1 "tangled.org/core/spindle/mill/proto/gen"
8 "tangled.org/core/spindle/models"
9 "time"
10)
11
12const (
13 leaseRowReserved = "reserved"
14 leaseRowRunning = "running"
15)
16
17func (m *Mill) persistLease(lease *RemoteLease, state string) error {
18 if m.db == nil {
19 return nil
20 }
21 return m.db.SaveMillLease(db.MillLease{
22 LeaseID: lease.id,
23 NodeID: lease.nodeID,
24 Epoch: lease.epoch,
25 Engine: lease.engine,
26 Knot: lease.wid.Knot,
27 Rkey: lease.wid.Rkey,
28 Workflow: lease.wid.Name,
29 State: state,
30 })
31}
32
33func (m *Mill) RestoreState() error {
34 if m.db == nil {
35 return nil
36 }
37 cursors, err := m.db.ListExecutorCursors()
38 if err != nil {
39 return err
40 }
41 rows, err := m.db.ListMillLeases()
42 if err != nil {
43 return err
44 }
45
46 m.mu.Lock()
47 for _, c := range cursors {
48 m.nodeSeqno[c.NodeID+"/"+c.Epoch] = c.AckedSeqno
49 }
50 for _, r := range rows {
51 lease := newLease(r.LeaseID, r.NodeID, r.Epoch, r.Engine)
52 lease.wid = models.WorkflowId{
53 PipelineId: models.PipelineId{Knot: r.Knot, Rkey: r.Rkey},
54 Name: r.Workflow,
55 }
56 // restored leases start as orphans, an executor must reclaim it via
57 // its first snapshot, or the sweep will fail it
58 lease.orphaned = true
59 lease.claimed = false
60 if r.State == leaseRowRunning {
61 lease.state = leaseRunning
62 }
63 m.leases[r.LeaseID] = lease
64 }
65 restored := len(rows)
66 m.mu.Unlock()
67
68 if restored > 0 {
69 m.l.Info("restored mill leases from previous run", "leases", restored, "cursors", len(cursors))
70 // executors get a grace window to reconnect and claim their leases
71 time.AfterFunc(m.cfg.ReconnectGrace, m.sweepUnclaimedOrphans)
72 }
73 return nil
74}
75func (m *Mill) sweepUnclaimedOrphans() {
76 m.mu.Lock()
77 var unclaimed []*RemoteLease
78 for _, lease := range m.leases {
79 if !lease.orphaned {
80 continue
81 }
82 if lease.claimed {
83 continue
84 }
85 if lease.getState() == leaseDone {
86 continue
87 }
88 if sess := m.sessions[lease.nodeID]; sess == nil || sess.disconnected {
89 unclaimed = append(unclaimed, lease)
90 }
91 }
92 m.mu.Unlock()
93
94 retry := false
95 reason := "executor did not reconnect after mill restart"
96 for _, lease := range unclaimed {
97 m.l.Warn("failing unclaimed restored lease", "lease", lease.id, "node", lease.nodeID)
98 if err := m.finishOrphan(lease, string(models.StatusKindFailed), &reason, nil); err != nil {
99 m.l.Error("finish unclaimed restored lease", "lease", lease.id, "err", err)
100 retry = true
101 }
102 }
103 if retry {
104 time.AfterFunc(5*time.Second, m.sweepUnclaimedOrphans)
105 }
106}
107
108func (m *Mill) reconcileLeases(sess *millSession, activeLeaseIDs []string) error {
109 active := make(map[string]struct{}, len(activeLeaseIDs))
110 for _, id := range activeLeaseIDs {
111 active[id] = struct{}{}
112 }
113
114 known := make(map[string]struct{})
115 m.mu.Lock()
116 var gone []*RemoteLease
117 for _, lease := range m.leases {
118 if lease.nodeID == sess.nodeID {
119 known[lease.id] = struct{}{}
120 if lease.epoch != sess.epoch {
121 gone = append(gone, lease)
122 } else if _, ok := active[lease.id]; !ok {
123 gone = append(gone, lease)
124 }
125 }
126 }
127 for _, lease := range m.reservations {
128 if lease.nodeID == sess.nodeID && lease.epoch == sess.epoch {
129 known[lease.id] = struct{}{}
130 }
131 }
132 m.mu.Unlock()
133 var unknown []string
134 for id := range active {
135 if _, ok := known[id]; !ok {
136 unknown = append(unknown, id)
137 }
138 }
139
140 for _, lease := range gone {
141 status := string(models.StatusKindFailed)
142 reason := "executor no longer holds lease"
143 if cancelled, cancelReason := lease.cancelRequested(); cancelled {
144 status = string(models.StatusKindCancelled)
145 reason = cancelReason
146 }
147 m.l.Warn("finishing reconciled lease", "lease", lease.id, "node", sess.nodeID, "leaseInc", lease.epoch, "sessInc", sess.epoch)
148 if lease.orphaned {
149 if err := m.finishOrphan(lease, status, &reason, nil); err != nil {
150 return err
151 }
152 } else if err := m.finishLiveLease(lease, status, reason); err != nil {
153 return err
154 }
155 }
156 for _, id := range unknown {
157 if err := sess.send(&millproto.Message{CancelAttempt: &millv1.CancelAttempt{LeaseId: id, Reason: "lease is not owned by this mill"}}); err != nil {
158 return fmt.Errorf("cancel unknown executor lease %q: %w", id, err)
159 }
160 }
161 return nil
162}
163
164func (m *Mill) completeLeaseRow(lease *RemoteLease, status string, errMsg *string, exitCode *int64) error {
165 if m.db == nil {
166 return nil
167 }
168 return m.db.CompleteMillLease(
169 lease.id,
170 string(lease.wid.PipelineId.AtUri()),
171 lease.wid.Name,
172 status,
173 errMsg,
174 exitCode,
175 m.n,
176 )
177}
178
179func (m *Mill) finishLiveLease(lease *RemoteLease, status, reason string) error {
180 lease.finishMu.Lock()
181 defer lease.finishMu.Unlock()
182 if lease.getState() == leaseDone {
183 return nil
184 }
185 if err := m.completeLeaseRow(lease, status, &reason, nil); err != nil {
186 return err
187 }
188 lease.deliverTerminal(&millv1.AttemptResult{
189 Status: mapTerminalStatusString(status),
190 Error: reason,
191 })
192 return m.cleanupLeaseLocked(lease)
193}
194
195func (m *Mill) finishOrphan(lease *RemoteLease, status string, errMsg *string, exitCode *int64) error {
196 lease.finishMu.Lock()
197 defer lease.finishMu.Unlock()
198 if lease.getState() == leaseDone {
199 return nil
200 }
201 if err := m.completeLeaseRow(lease, status, errMsg, exitCode); err != nil {
202 return err
203 }
204 lease.markDone()
205 return m.cleanupLeaseLocked(lease)
206}
207
208func mapTerminalStatusString(s string) millv1.TerminalStatus {
209 switch s {
210 case "success":
211 return millv1.TerminalStatus_SUCCESS
212 case "failed":
213 return millv1.TerminalStatus_FAILED
214 case "timeout":
215 return millv1.TerminalStatus_TIMEOUT
216 case "cancelled":
217 return millv1.TerminalStatus_CANCELLED
218 default:
219 return millv1.TerminalStatus_TERMINAL_STATUS_UNSPECIFIED
220 }
221}