| 1028 | } |
| 1029 | |
| 1030 | int cbm_parallel_extract_ex(cbm_pipeline_ctx_t *ctx, const cbm_file_info_t *files, int file_count, |
| 1031 | CBMFileResult **result_cache, _Atomic int64_t *shared_ids, |
| 1032 | int worker_count, const cbm_parallel_extract_opts_t *opts) { |
| 1033 | cbm_parallel_extract_opts_t resolved_opts = cbm_parallel_extract_resolve_opts(opts); |
| 1034 | |
| 1035 | if (file_count == 0) { |
| 1036 | return 0; |
| 1037 | } |
| 1038 | |
| 1039 | cbm_log_info("parallel.extract.start", "files", itoa_log(file_count), "workers", |
| 1040 | itoa_log(worker_count)); |
| 1041 | { |
| 1042 | size_t mb = (size_t)CBM_SZ_1K * CBM_SZ_1K; |
| 1043 | cbm_log_info("parallel.extract.retention", "retain_sources", |
| 1044 | resolved_opts.retain_sources ? "true" : "false", "total_mb", |
| 1045 | itoa_log((int)(resolved_opts.retain_total_budget_bytes / mb)), "per_file_mb", |
| 1046 | itoa_log((int)(resolved_opts.retain_per_file_max_bytes / mb))); |
| 1047 | } |
| 1048 | |
| 1049 | /* Log per-worker memory budget */ |
| 1050 | if (cbm_mem_budget() > 0) { |
| 1051 | size_t worker_budget = cbm_mem_worker_budget(worker_count); |
| 1052 | cbm_log_info("parallel.mem.budget", "total_mb", |
| 1053 | itoa_log((int)(cbm_mem_budget() / ((size_t)CBM_SZ_1K * CBM_SZ_1K))), |
| 1054 | "per_worker_mb", |
| 1055 | itoa_log((int)(worker_budget / ((size_t)CBM_SZ_1K * CBM_SZ_1K)))); |
| 1056 | } |
| 1057 | |
| 1058 | /* Sub-phase: Ensure extraction library is initialized */ |
| 1059 | CBM_PROF_START(t_init); |
| 1060 | cbm_init(); |
| 1061 | |
| 1062 | /* Slab allocator for tree-sitter (thread-safe via TLS). */ |
| 1063 | cbm_slab_install(); |
| 1064 | CBM_PROF_END("parallel_extract", "1_init_libs", t_init); |
| 1065 | |
| 1066 | /* Sub-phase: Sort files by descending size for tail-latency reduction */ |
| 1067 | CBM_PROF_START(t_sort); |
| 1068 | file_sort_entry_t *sorted = malloc((size_t)file_count * sizeof(file_sort_entry_t)); |
| 1069 | if (!sorted) { |
| 1070 | return CBM_NOT_FOUND; |
| 1071 | } |
| 1072 | for (int i = 0; i < file_count; i++) { |
| 1073 | sorted[i].idx = i; |
| 1074 | sorted[i].size = files[i].size; |
| 1075 | } |
| 1076 | qsort(sorted, file_count, sizeof(file_sort_entry_t), compare_by_size_desc); |
| 1077 | CBM_PROF_END_N("parallel_extract", "2_sort_files", t_sort, file_count); |
| 1078 | |
| 1079 | /* Allocate per-worker state (cache-line aligned via posix_memalign) */ |
| 1080 | extract_worker_state_t *workers = NULL; |
| 1081 | if (cbm_aligned_alloc((void **)&workers, CBM_CACHE_LINE, |
| 1082 | (size_t)worker_count * sizeof(extract_worker_state_t)) != 0) { |
| 1083 | free(sorted); |
| 1084 | return CBM_NOT_FOUND; |
| 1085 | } |
| 1086 | memset(workers, 0, (size_t)worker_count * sizeof(extract_worker_state_t)); |
| 1087 | |