Line data Source code
1 : #ifndef HEADER_fd_src_discof_restore_utils_fd_ssctrl_h 2 : #define HEADER_fd_src_discof_restore_utils_fd_ssctrl_h 3 : 4 : #include "../../../util/net/fd_net_headers.h" 5 : #include "../../../waltz/fd_fqdn.h" 6 : #include "../../../flamenco/runtime/fd_runtime_const.h" 7 : 8 : /* The snapshot tiles have a somewhat involved state machine, which is 9 : controlled by snapct. Imagine first the following sequence: 10 : 11 : 1. snapct is reading a full snapshot from the network and sends some 12 : data to snapdc to be decompressed. 13 : 2. snapct hits a network error, and resets the connection to a new 14 : peer. 15 : 3. The decompressor fails on data from the old peer, and sends a 16 : malformed message to snapct. 17 : 4. snapct receives the malformed message, and abandons the new 18 : connection, even though it was not malformed. 19 : 20 : There are basically two ways to prevent this. Option A is the tiles 21 : can pass not just control messages to one another, but also tag them 22 : with some xid indicating which "attempt" the control message is for. 23 : 24 : This is pretty hard to reason about, and the state machine can grow 25 : quite complicated. 26 : 27 : There's an easier way: the tiles just are fully synchronized with 28 : snapct. Whatever "attempt" snapct is on, we ensure all other tiles 29 : are on it too. This means when any tile fails a snapshot, all tiles 30 : must fail it and fully flush all frags in the pipeline before snapct 31 : can proceed with a new attempt. 32 : 33 : The control flow then is basically, 34 : 35 : 1. All tiles start in the IDLE state. 36 : 2. snapct initializes the pipeline by sending an INIT message. 37 : Each tile enters the PROCESSING state and then forwards the INIT 38 : message down the pipeline. When snapct receives this INIT 39 : message, the entire pipeline is in PROCESSING state. 40 : 3. Tiles continue to process data / frags as applicable. If an 41 : error occurs, the tile enters the ERROR state and also sends an 42 : ERROR message downstream. All downstream tiles also enter the 43 : ERROR state and forward the message. Note that upstream tiles 44 : will not be in an ERROR state and will continue producing frags. 45 : When snapct receives the ERROR message, it will send a FAIL 46 : message. snapct then waits for this FAIL message to be 47 : propagated through the pipeline and received back. snapct 48 : also guarantees to flush any pending load data in the 49 : pipeline. It then knows that all tiles are synchronized back 50 : in an IDLE state and it can try again with a new INIT. 51 : 4. Once all snapshot data has entered the pipeline, snapct sends 52 : a FINI message through the pipeline and waits for it to be 53 : received back. This synchronizes tiles to finish any 54 : remaining in-flight work. It then sends either a NEXT or a 55 : DONE message through the pipeline, again waiting for it to be 56 : received back. NEXT means another snapshot follows, whereas 57 : DONE means a shutdown will follow. 58 : 59 : That keeps the tiles in lockstep, and simplifies the state machine 60 : to a manageable level. 61 : 62 : In more detail, all tiles in the feedback loop of the pipeline must 63 : comply with the following rules; snapct enforces them; tiles outside 64 : of the feedback loop may support a subset of these: 65 : - Allowed state transitions on given control message: 66 : IDLE to PROCESSING: on INIT_FULL / INIT_INCR ctrl msg. 67 : PROCESSING to FINISHING : on FINI ctrl msg (*). 68 : FINISHING to IDLE : on NEXT / DONE ctrl msg. 69 : IDLE to SHUTDOWN : on SHUTDOWN ctrl msg. 70 : - Control messages that apply to all states (except SHUTDOWN): 71 : On ERROR msg, always transition to ERROR. 72 : On FAIL msg, always transition to IDLE. 73 : - When in SHUTDOWN state, no data or control message is allowed. 74 : - Error handling: 75 : A tile that enters ERROR state on its own must forward an 76 : ERROR control message and discard any incoming message. When 77 : in ERROR state, data is discarded, and only a FAIL control 78 : message is processed and forwarded (all others are discarded). 79 : As a result, only one ERROR message can propagate through the 80 : pipeline at any given time. 81 : - Holding onto a control message: 82 : A tile may hold onto a control message and forward it later 83 : after performing some asynchronous routine. The tile's state 84 : transition may also be deferred until that control message is 85 : forwarded. 86 : - (*) A tile may self-transition from PROCESSING to FINISHING 87 : early, if it can detect end-of-stream. The FINI message is 88 : used to synchronize the transition. 89 : 90 : Each non-ERROR control message is only generated once in snapct 91 : and will not be re-sent. The pipeline will be locked on flushing 92 : that control message until all tiles forward it on, or an ERROR 93 : message is triggered by any of the tiles and forwarded. */ 94 : 95 : /* Parameters for the snapshot data stream links (snapld_dc, 96 : snapdc_in). Producers adapt to the link MTU, which just needs to 97 : exceed the control structs. */ 98 0 : #define FD_SNAPSHOT_DATA_DEPTH (256UL) 99 49725 : #define FD_SNAPSHOT_DATA_MTU (65408UL) 100 : 101 25017 : #define FD_SNAPSHOT_STATE_IDLE (0UL) /* Performing no work and should receive no data frags */ 102 24990 : #define FD_SNAPSHOT_STATE_PROCESSING (1UL) /* Performing usual work, no errors / EoF condition encountered */ 103 360 : #define FD_SNAPSHOT_STATE_FINISHING (2UL) /* Tile has observed EoF, expects no additional data frags */ 104 50268 : #define FD_SNAPSHOT_STATE_ERROR (3UL) /* Some error occurred, will wait for a FAIL command to reset */ 105 6 : #define FD_SNAPSHOT_STATE_SHUTDOWN (4UL) /* All work finished, tile can perform final cleanup and exit */ 106 : 107 29361 : #define FD_SNAPSHOT_MSG_DATA (0UL) /* Fragment represents some snapshot data */ 108 24 : #define FD_SNAPSHOT_MSG_META (1UL) /* Fragment represents a fd_ssctrl_meta_t message */ 109 : 110 74847 : #define FD_SNAPSHOT_MSG_CTRL_INIT_FULL (2UL) /* Pipeline should start processing a full snapshot */ 111 25035 : #define FD_SNAPSHOT_MSG_CTRL_INIT_INCR (3UL) /* Pipeline should start processing an incremental snapshot */ 112 441 : #define FD_SNAPSHOT_MSG_CTRL_FAIL (4UL) /* Current snapshot failed, undo work and reset to idle state */ 113 33 : #define FD_SNAPSHOT_MSG_CTRL_NEXT (5UL) /* Current snapshot succeeded, commit work, go idle, and expect another snapshot */ 114 30 : #define FD_SNAPSHOT_MSG_CTRL_DONE (6UL) /* Current snapshot succeeded, commit work, go idle, and expect shutdown */ 115 18 : #define FD_SNAPSHOT_MSG_CTRL_SHUTDOWN (7UL) /* Snapshot load successful, no work left to do, perform final cleanup and shut down*/ 116 108 : #define FD_SNAPSHOT_MSG_CTRL_ERROR (8UL) /* Some tile encountered an error with the current stream */ 117 171 : #define FD_SNAPSHOT_MSG_CTRL_FINI (9UL) /* Current snapshot has been fully loaded, finish processing */ 118 : 119 : /* snapld -> snapct (via snapld_dc) */ 120 21 : #define FD_SNAPSHOT_MSG_LOAD_COMPLETE (10UL) /* snapld finished reading/downloading all data */ 121 : 122 : /* Sent by snapct to tell snapld whether to load a local file or 123 : download from a particular external peer. */ 124 : typedef struct fd_ssctrl_init { 125 : int file; 126 : int zstd; 127 : ulong slot; /* slot advertised by the snapshot peer */ 128 : fd_ip4_port_t addr; 129 : uchar snapshot_hash[ FD_HASH_FOOTPRINT ]; /* advertised snapshot hash from snapshot file name */ 130 : char hostname[ FD_FQDN_BUF_MAX ]; 131 : char path[ PATH_MAX ]; 132 : ulong path_len; 133 : int is_https; 134 : int is_redirect; /* 1 if using well-known redirect path */ 135 : ulong file_sz; /* file size in bytes (file loads only, 0 for HTTP) */ 136 : } fd_ssctrl_init_t; 137 : 138 : /* Sent by snapld to tell snapct metadata about a downloaded snapshot. */ 139 : typedef struct fd_ssctrl_meta { 140 : ulong total_sz; 141 : ulong resolved_slot; 142 : uchar resolved_hash[ FD_HASH_FOOTPRINT ]; 143 : char resolved_name[ PATH_MAX ]; 144 : } fd_ssctrl_meta_t; 145 : 146 : struct fd_snapshot_account_hdr { 147 : uchar pubkey[ FD_PUBKEY_FOOTPRINT ]; 148 : uchar owner[ FD_PUBKEY_FOOTPRINT ]; 149 : ulong lamports; 150 : uchar executable; 151 : ulong data_len; 152 : }; 153 : typedef struct fd_snapshot_account_hdr fd_snapshot_account_hdr_t; 154 : 155 : /* fd_snapshot_account_hdr_init initializes a fd_snapshot_account_hdr_t struct 156 : with the appropriate account metadata fields. */ 157 : static inline void 158 : fd_snapshot_account_hdr_init( fd_snapshot_account_hdr_t * account, 159 : uchar const pubkey[ FD_PUBKEY_FOOTPRINT ], 160 : uchar const owner[ FD_PUBKEY_FOOTPRINT ], 161 : ulong lamports, 162 : uchar executable, 163 0 : ulong data_len ) { 164 0 : fd_memcpy( account->pubkey, pubkey, FD_PUBKEY_FOOTPRINT ); 165 0 : fd_memcpy( account->owner, owner, FD_PUBKEY_FOOTPRINT ); 166 0 : account->lamports = lamports; 167 0 : account->executable = executable; 168 0 : account->data_len = data_len; 169 0 : } 170 : 171 : static inline const char * 172 12 : fd_ssctrl_state_str( ulong state ) { 173 12 : switch( state ) { 174 0 : case FD_SNAPSHOT_STATE_IDLE: return "idle"; 175 6 : case FD_SNAPSHOT_STATE_PROCESSING: return "processing"; 176 6 : case FD_SNAPSHOT_STATE_FINISHING: return "finishing"; 177 0 : case FD_SNAPSHOT_STATE_ERROR: return "error"; 178 0 : case FD_SNAPSHOT_STATE_SHUTDOWN: return "shutdown"; 179 0 : default: return "unknown"; 180 12 : } 181 12 : } 182 : 183 : static inline const char * 184 0 : fd_ssctrl_msg_ctrl_str( ulong sig ) { 185 0 : switch( sig ) { 186 0 : case FD_SNAPSHOT_MSG_DATA: return "data"; 187 0 : case FD_SNAPSHOT_MSG_META: return "meta"; 188 0 : case FD_SNAPSHOT_MSG_CTRL_INIT_FULL: return "init_full"; 189 0 : case FD_SNAPSHOT_MSG_CTRL_INIT_INCR: return "init_incr"; 190 0 : case FD_SNAPSHOT_MSG_CTRL_FAIL: return "fail"; 191 0 : case FD_SNAPSHOT_MSG_CTRL_NEXT: return "next"; 192 0 : case FD_SNAPSHOT_MSG_CTRL_DONE: return "done"; 193 0 : case FD_SNAPSHOT_MSG_CTRL_SHUTDOWN: return "shutdown"; 194 0 : case FD_SNAPSHOT_MSG_CTRL_ERROR: return "error"; 195 0 : case FD_SNAPSHOT_MSG_CTRL_FINI: return "fini"; 196 0 : case FD_SNAPSHOT_MSG_LOAD_COMPLETE: return "load_complete"; 197 0 : default: return "unknown"; 198 0 : } 199 0 : } 200 : 201 : #endif /* HEADER_fd_src_discof_restore_utils_fd_ssctrl_h */