MCPcopy Create free account
hub / github.com/antirez/ds4 / dist_worker_process_work_payload

Function dist_worker_process_work_payload

ds4_distributed.c:7148–7590  ·  view source on GitHub ↗

Source from the content-addressed store, hash-verified

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
7171static 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,

Callers 2

dist_worker_handle_workFunction · 0.85

Calls 15

dist_mem_readFunction · 0.85
dist_work_from_wireFunction · 0.85
dist_u64_from_halvesFunction · 0.85
dist_now_secFunction · 0.85
ds4_engine_layer_countFunction · 0.85
ds4_engine_vocab_sizeFunction · 0.85

Tested by

no test coverage detected