| 7146 | session->token_hash = token_hash; |
| 7147 | session->token_hash_valid = true; |
| 7148 | } else { |
| 7149 | if (!err[0]) snprintf(err, sizeof(err), "failed to restore worker KV shard"); |
| 7150 | if (session) session->token_hash_valid = false; |
| 7151 | } |
| 7152 | pthread_mutex_unlock(&state->mu); |
| 7153 | } |
| 7154 | |
| 7155 | if (tmp) fclose(tmp); |
| 7156 | if (tmp) unlink(tmp_path); |
| 7157 | free(tokens); |
| 7158 | |
| 7159 | pthread_mutex_lock(&upstream->write_mu); |
| 7160 | rc = dist_send_snapshot_done(upstream->fd, request_id, err[0] ? 1u : 0u, |
| 7161 | err[0] ? err : NULL); |
| 7162 | pthread_mutex_unlock(&upstream->write_mu); |
| 7163 | if (err[0] && received < payload_bytes) return -1; |
| 7164 | return rc; |
| 7165 | } |
| 7166 | |
| 7167 | /* ========================================================================= |
| 7168 | * Worker Layer Execution |
| 7169 | * ========================================================================= */ |
| 7170 | |
| 7171 | static int dist_worker_process_work_payload( |
| 7172 | ds4_dist_worker_state *state, |
| 7173 | ds4_dist_worker_upstream *upstream, |
| 7174 | const void *payload, |
| 7175 | uint32_t bytes) { |
| 7176 | uint64_t request_id = 0; |
| 7177 | char err[256]; |
| 7178 | if (bytes < sizeof(ds4_dist_work_fixed)) { |
| 7179 | return dist_worker_upstream_send_work_error(upstream, request_id, "truncated distributed WORK frame"); |
| 7180 | } |
| 7181 | |
| 7182 | ds4_dist_mem_reader reader = { |
| 7183 | .p = payload, |
| 7184 | .remaining = bytes, |
| 7185 | }; |
| 7186 | ds4_dist_work_fixed work; |
| 7187 | int rc = dist_mem_read(&reader, &work, (uint32_t)sizeof(work)); |
| 7188 | if (rc <= 0) return -1; |
| 7189 | dist_work_from_wire(&work); |
| 7190 | const uint64_t session_id = dist_u64_from_halves(work.session_hi, work.session_lo); |
| 7191 | request_id = dist_u64_from_halves(work.request_hi, work.request_lo); |
| 7192 | const uint64_t work_prefix_hash = dist_u64_from_halves(work.prefix_hash_hi, |
| 7193 | work.prefix_hash_lo); |
| 7194 | const uint64_t work_result_hash = dist_u64_from_halves(work.result_hash_hi, |
| 7195 | work.result_hash_lo); |
| 7196 | const bool profile = dist_decode_profile_enabled() && work.n_tokens == 1; |
| 7197 | const double total_t0 = profile ? dist_now_sec() : 0.0; |
| 7198 | DIST_DEBUG("worker work request=%llu layers=%u:%u tokens=%u pos=%u flags=0x%x token_bytes=%u input_hc=%u/%ub route_count=%u route_index=%u route_bytes=%u", |
| 7199 | (unsigned long long)request_id, |
| 7200 | work.layer_start, |
| 7201 | work.layer_end, |
| 7202 | work.n_tokens, |
| 7203 | work.pos0, |
| 7204 | work.flags, |
| 7205 | work.token_bytes, |
no test coverage detected