xref: /freebsd/sys/contrib/openzfs/cmd/zstream/zstream_redup.c (revision 2f10ffc003be396f3fc23cd2888023896560252b)
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