1 // SPDX-License-Identifier: CDDL-1.0 2 /* 3 * This file and its contents are supplied under the terms of the 4 * Common Development and Distribution License ("CDDL"), version 1.0. 5 * You may only use this file in accordance with the terms of version 6 * 1.0 of the CDDL. 7 * 8 * A full copy of the text of the CDDL should have accompanied this 9 * source. A copy of the CDDL is also available via the Internet at 10 * https://opensource.org/license/CDDL-1.0. 11 */ 12 13 /* 14 * Copyright (c) 2020 by Delphix. All rights reserved. 15 */ 16 17 #include <assert.h> 18 #include <cityhash.h> 19 #include <err.h> 20 #include <errno.h> 21 #include <libzutil.h> 22 #include <stdint.h> 23 #include <stdio.h> 24 #include <stdlib.h> 25 #include <string.h> 26 #include <sys/bitops.h> 27 #include <sys/param.h> 28 #include <sys/stdtypes.h> 29 #include <sys/sysmacros.h> 30 #include <sys/zfs_ioctl.h> 31 #include <umem.h> 32 #include <unistd.h> 33 34 #include "zstream.h" 35 #include "zstream_modules.h" 36 #include "zstream_util.h" 37 38 #define MAX_RDT_PHYSMEM_PERCENT 20 39 #define SMALLEST_POSSIBLE_MAX_RDT_MB 128 40 41 typedef struct redup_entry { 42 struct redup_entry *rde_next; 43 uint64_t rde_guid; 44 uint64_t rde_object; 45 uint64_t rde_offset; 46 uint64_t rde_stream_offset; 47 } redup_entry_t; 48 49 typedef struct redup_table { 50 redup_entry_t **redup_hash_array; 51 umem_cache_t *ddecache; 52 uint64_t ddt_count; 53 int numhashbits; 54 } redup_table_t; 55 56 typedef struct { 57 redup_table_t rc_rdt; 58 FILE *rc_fp; 59 } redup_context_t; 60 61 static void 62 rdt_insert(redup_table_t *rdt, 63 uint64_t guid, uint64_t object, uint64_t offset, uint64_t stream_offset) 64 { 65 uint64_t ch = cityhash3(guid, object, offset); 66 uint64_t hashcode = BF64_GET(ch, 0, rdt->numhashbits); 67 redup_entry_t **rdepp; 68 69 rdepp = &(rdt->redup_hash_array[hashcode]); 70 redup_entry_t *rde = umem_cache_alloc(rdt->ddecache, UMEM_NOFAIL); 71 rde->rde_next = *rdepp; 72 rde->rde_guid = guid; 73 rde->rde_object = object; 74 rde->rde_offset = offset; 75 rde->rde_stream_offset = stream_offset; 76 *rdepp = rde; 77 rdt->ddt_count++; 78 } 79 80 static void 81 rdt_lookup(redup_table_t *rdt, 82 uint64_t guid, uint64_t object, uint64_t offset, 83 uint64_t *stream_offsetp) 84 { 85 uint64_t ch = cityhash3(guid, object, offset); 86 uint64_t hashcode = BF64_GET(ch, 0, rdt->numhashbits); 87 88 for (redup_entry_t *rde = rdt->redup_hash_array[hashcode]; 89 rde != NULL; rde = rde->rde_next) { 90 if (rde->rde_guid == guid && 91 rde->rde_object == object && 92 rde->rde_offset == offset) { 93 *stream_offsetp = rde->rde_stream_offset; 94 return; 95 } 96 } 97 assert(!"could not find expected redup table entry"); 98 } 99 100 static disposition_t 101 chain_redup_writes(void *item_in, void *context_in) 102 { 103 drr_packet_t *item = (drr_packet_t *)item_in; 104 redup_context_t *context = (redup_context_t *)context_in; 105 106 if (item == NULL) { 107 return (D_OK); 108 } 109 110 dmu_replay_record_t *drr = &item->dp_drr; 111 struct drr_write *drrw = &drr->drr_u.drr_write; 112 struct drr_begin *drrb = &drr->drr_u.drr_begin; 113 114 switch (drr->drr_type) { 115 116 case DRR_BEGIN: 117 { 118 uint64_t flags = DMU_GET_FEATUREFLAGS(drrb->drr_versioninfo); 119 flags &= ~(DMU_BACKUP_FEATURE_DEDUP | 120 DMU_BACKUP_FEATURE_DEDUPPROPS); 121 DMU_SET_FEATUREFLAGS(drrb->drr_versioninfo, flags); 122 break; 123 } 124 125 case DRR_WRITE_BYREF: 126 { 127 struct drr_write_byref drrwb = drr->drr_u.drr_write_byref; 128 129 /* 130 * Look up in hash table by drrwb->drr_refguid, 131 * drr_refobject, drr_refoffset. Replace this 132 * record with the found WRITE record, but with 133 * drr_object,drr_offset,drr_toguid replaced with ours. 134 */ 135 uint64_t stream_offset = 0; 136 rdt_lookup(&context->rc_rdt, drrwb.drr_refguid, 137 drrwb.drr_refobject, drrwb.drr_refoffset, 138 &stream_offset); 139 140 if (fseeko(context->rc_fp, stream_offset, SEEK_SET) != 0) { 141 err(1, "seek into source file failed, offset %llu", 142 (u_longlong_t)stream_offset); 143 } 144 if (fread(drr, sizeof (*drr), 1, context->rc_fp) != 1) { 145 err(1, "read of prior write failed"); 146 } 147 if (ATTR_IS_SET(CA_BYTESWAPPED)) { 148 byteswap_record(drr, BSWAP_32(drr->drr_type)); 149 } 150 151 VERIFY3U(drr->drr_type, ==, DRR_WRITE); 152 VERIFY3U(drrw->drr_toguid, ==, drrwb.drr_refguid); 153 VERIFY3U(drrw->drr_object, ==, drrwb.drr_refobject); 154 VERIFY3U(drrw->drr_offset, ==, drrwb.drr_refoffset); 155 156 uint64_t size = DRR_WRITE_PAYLOAD_SIZE(drrw); 157 uint8_t *buff = safe_malloc(size); 158 size_t n_read = fread(buff, size, 1, context->rc_fp); 159 if (n_read != 1) 160 err(1, "read of prior payload failed"); 161 set_payload(item, buff, size); 162 163 drrw->drr_toguid = drrwb.drr_toguid; 164 drrw->drr_object = drrwb.drr_object; 165 drrw->drr_offset = drrwb.drr_offset; 166 break; 167 } 168 169 case DRR_WRITE: 170 rdt_insert(&context->rc_rdt, drrw->drr_toguid, drrw->drr_object, 171 drrw->drr_offset, item->dp_stream_offset); 172 break; 173 174 default: 175 break; 176 } 177 return (D_OK); 178 } 179 180 static chain_step_t 181 serial_redup_writes(redup_context_t *context) 182 { 183 chain_step_t step = { 184 .cs_type = CS_SERIAL, 185 .cs_in_size = sizeof (drr_packet_t), 186 .cs_out_size = sizeof (drr_packet_t), 187 .cs_context = context, 188 .cs_serial = { 189 .process = chain_redup_writes 190 } 191 }; 192 return (step); 193 } 194 195 int 196 zstream_do_redup(int argc, char *argv[]) 197 { 198 int c; 199 chain_attrs_t attrs = {0}; 200 redup_context_t context = {0}; 201 uint64_t numbuckets; 202 203 while ((c = getopt(argc, argv, "v")) != -1) { 204 switch (c) { 205 case 'v': 206 ENABLE_OPTION(&attrs, CA_VERBOSE); 207 break; 208 case '?': 209 warnx("invalid option '%c'", optopt); 210 zstream_usage(); 211 } 212 } 213 214 argc -= optind; 215 argv += optind; 216 217 if (argc != 1) 218 zstream_usage(); 219 220 context.rc_fp = fopen(argv[0], "rb"); 221 if (context.rc_fp == NULL) { 222 err(1, "unable to open %s", argv[0]); 223 } 224 225 #ifdef _ILP32 226 uint64_t max_rde_size = SMALLEST_POSSIBLE_MAX_RDT_MB << 20; 227 #else 228 uint64_t physbytes = sysconf(_SC_PHYS_PAGES) * sysconf(_SC_PAGESIZE); 229 uint64_t max_rde_size = MAX((physbytes * MAX_RDT_PHYSMEM_PERCENT) / 100, 230 SMALLEST_POSSIBLE_MAX_RDT_MB << 20); 231 #endif 232 233 numbuckets = max_rde_size / (sizeof (redup_entry_t)); 234 if (!ISP2(numbuckets)) 235 numbuckets = 1ULL << highbit64(numbuckets); 236 237 context.rc_rdt.redup_hash_array = 238 safe_calloc(numbuckets * sizeof (redup_entry_t *)); 239 context.rc_rdt.ddecache = umem_cache_create("rde", 240 sizeof (redup_entry_t), 0, NULL, NULL, NULL, NULL, NULL, 0); 241 context.rc_rdt.numhashbits = highbit64(numbuckets) - 1; 242 context.rc_rdt.ddt_count = 0; 243 244 zstream_chain_t redup_chain = { 245 STANDARD_INPUT_STACK(argv[0]), 246 serial_redup_writes(&context), 247 STANDARD_OUTPUT_STACK(NULL) 248 }; 249 zstream_chain_exec(redup_chain, &attrs); 250 251 if (attrs.ca_command_opts & CA_VERBOSE) { 252 char mem_str[16]; 253 record_stats_t *acsi = attrs.ca_stats_in; 254 zfs_nicenum(context.rc_rdt.ddt_count * sizeof (redup_entry_t), 255 mem_str, sizeof (mem_str)); 256 fprintf(stderr, "Converted stream with %llu total records, " 257 "including %llu dedup records, using %sB memory.\n", 258 (u_longlong_t)attrs.ca_totals_in.rs_num_records, 259 (u_longlong_t)acsi[DRR_WRITE_BYREF].rs_num_records, 260 mem_str); 261 } 262 263 fclose(context.rc_fp); 264 umem_cache_destroy(context.rc_rdt.ddecache); 265 free(context.rc_rdt.redup_hash_array); 266 return (0); 267 } 268