This repository has no description
0

Configure Feed

Select the types of activity you want to include in your feed.

core / spindle / mill / restore.go
5.7 kB 221 lines
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}