dagmodifier.go 10.5 KB
Newer Older
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22
package mod

import (
	"bytes"
	"errors"
	"io"
	"os"

	proto "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/goprotobuf/proto"
	mh "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multihash"
	context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/golang.org/x/net/context"

	chunk "github.com/jbenet/go-ipfs/importer/chunk"
	help "github.com/jbenet/go-ipfs/importer/helpers"
	trickle "github.com/jbenet/go-ipfs/importer/trickle"
	mdag "github.com/jbenet/go-ipfs/merkledag"
	pin "github.com/jbenet/go-ipfs/pin"
	ft "github.com/jbenet/go-ipfs/unixfs"
	uio "github.com/jbenet/go-ipfs/unixfs/io"
	u "github.com/jbenet/go-ipfs/util"
)

Jeromy's avatar
Jeromy committed
23 24 25 26
var ErrSeekFail = errors.New("failed to seek properly")
var ErrSeekEndNotImpl = errors.New("SEEK_END currently not implemented")
var ErrUnrecognizedWhence = errors.New("unrecognized whence")

27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63
// 2MB
var writebufferSize = 1 << 21

var log = u.Logger("dagio")

// DagModifier is the only struct licensed and able to correctly
// perform surgery on a DAG 'file'
// Dear god, please rename this to something more pleasant
type DagModifier struct {
	dagserv mdag.DAGService
	curNode *mdag.Node
	mp      pin.ManualPinner

	splitter   chunk.BlockSplitter
	ctx        context.Context
	readCancel func()

	writeStart uint64
	curWrOff   uint64
	wrBuf      *bytes.Buffer

	read *uio.DagReader
}

func NewDagModifier(ctx context.Context, from *mdag.Node, serv mdag.DAGService, mp pin.ManualPinner, spl chunk.BlockSplitter) (*DagModifier, error) {
	return &DagModifier{
		curNode:  from.Copy(),
		dagserv:  serv,
		splitter: spl,
		ctx:      ctx,
		mp:       mp,
	}, nil
}

// WriteAt will modify a dag file in place
func (dm *DagModifier) WriteAt(b []byte, offset int64) (int, error) {
	// TODO: this is currently VERY inneficient
Jeromy's avatar
Jeromy committed
64 65
	// each write that happens at an offset other than the current one causes a
	// flush to disk, and dag rewrite
66 67 68 69 70 71
	if offset == int64(dm.writeStart) && dm.wrBuf != nil {
		// If we would overwrite the previous write
		if len(b) >= dm.wrBuf.Len() {
			dm.wrBuf.Reset()
		}
	} else if uint64(offset) != dm.curWrOff {
72 73 74 75 76 77 78 79 80 81 82
		size, err := dm.Size()
		if err != nil {
			return 0, err
		}
		if offset > size {
			err := dm.expandSparse(offset - size)
			if err != nil {
				return 0, err
			}
		}

83
		err = dm.Sync()
84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102
		if err != nil {
			return 0, err
		}
		dm.writeStart = uint64(offset)
	}

	return dm.Write(b)
}

// A reader that just returns zeros
type zeroReader struct{}

func (zr zeroReader) Read(b []byte) (int, error) {
	for i, _ := range b {
		b[i] = 0
	}
	return len(b), nil
}

Jeromy's avatar
Jeromy committed
103 104
// expandSparse grows the file with zero blocks of 4096
// A small blocksize is chosen to aid in deduplication
105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120
func (dm *DagModifier) expandSparse(size int64) error {
	spl := chunk.SizeSplitter{4096}
	r := io.LimitReader(zeroReader{}, size)
	blks := spl.Split(r)
	nnode, err := dm.appendData(dm.curNode, blks)
	if err != nil {
		return err
	}
	_, err = dm.dagserv.Add(nnode)
	if err != nil {
		return err
	}
	dm.curNode = nnode
	return nil
}

Jeromy's avatar
Jeromy committed
121
// Write continues writing to the dag at the current offset
122 123 124 125 126 127 128
func (dm *DagModifier) Write(b []byte) (int, error) {
	if dm.read != nil {
		dm.read = nil
	}
	if dm.wrBuf == nil {
		dm.wrBuf = new(bytes.Buffer)
	}
Jeromy's avatar
Jeromy committed
129

130 131 132 133 134 135
	n, err := dm.wrBuf.Write(b)
	if err != nil {
		return n, err
	}
	dm.curWrOff += uint64(n)
	if dm.wrBuf.Len() > writebufferSize {
136
		err := dm.Sync()
137 138 139 140 141 142 143 144
		if err != nil {
			return n, err
		}
	}
	return n, nil
}

func (dm *DagModifier) Size() (int64, error) {
145
	pbn, err := ft.FromBytes(dm.curNode.Data)
146 147 148 149
	if err != nil {
		return 0, err
	}

150 151 152 153
	if dm.wrBuf != nil {
		if uint64(dm.wrBuf.Len())+dm.writeStart > pbn.GetFilesize() {
			return int64(dm.wrBuf.Len()) + int64(dm.writeStart), nil
		}
154 155 156 157 158
	}

	return int64(pbn.GetFilesize()), nil
}

159 160
// Sync writes changes to this dag to disk
func (dm *DagModifier) Sync() error {
Jeromy's avatar
Jeromy committed
161
	// No buffer? Nothing to do
162 163 164 165 166 167 168 169 170 171
	if dm.wrBuf == nil {
		return nil
	}

	// If we have an active reader, kill it
	if dm.read != nil {
		dm.read = nil
		dm.readCancel()
	}

Jeromy's avatar
Jeromy committed
172
	// Number of bytes we're going to write
173 174
	buflen := dm.wrBuf.Len()

175 176 177 178 179 180
	// Grab key for unpinning after mod operation
	curk, err := dm.curNode.Key()
	if err != nil {
		return err
	}

Jeromy's avatar
Jeromy committed
181
	// overwrite existing dag nodes
182
	thisk, done, err := dm.modifyDag(dm.curNode, dm.writeStart, dm.wrBuf)
183 184 185 186
	if err != nil {
		return err
	}

187
	nd, err := dm.dagserv.Get(thisk)
188 189 190 191 192 193
	if err != nil {
		return err
	}

	dm.curNode = nd

Jeromy's avatar
Jeromy committed
194
	// need to write past end of current dag
195 196 197 198 199 200 201
	if !done {
		blks := dm.splitter.Split(dm.wrBuf)
		nd, err = dm.appendData(dm.curNode, blks)
		if err != nil {
			return err
		}

202
		thisk, err = dm.dagserv.Add(nd)
203 204 205 206 207 208 209
		if err != nil {
			return err
		}

		dm.curNode = nd
	}

210 211 212 213 214 215 216 217
	// Finalize correct pinning, and flush pinner
	dm.mp.PinWithMode(thisk, pin.Recursive)
	dm.mp.RemovePinWithMode(curk, pin.Recursive)
	err = dm.mp.Flush()
	if err != nil {
		return err
	}

218 219 220 221 222 223
	dm.writeStart += uint64(buflen)

	dm.wrBuf = nil
	return nil
}

Jeromy's avatar
Jeromy committed
224
// modifyDag writes the data in 'data' over the data in 'node' starting at 'offset'
Jeromy's avatar
Jeromy committed
225 226
// returns the new key of the passed in node and whether or not all the data in the reader
// has been consumed.
Jeromy's avatar
Jeromy committed
227
func (dm *DagModifier) modifyDag(node *mdag.Node, offset uint64, data io.Reader) (u.Key, bool, error) {
228 229
	f, err := ft.FromBytes(node.Data)
	if err != nil {
Jeromy's avatar
Jeromy committed
230
		return "", false, err
231 232
	}

Jeromy's avatar
Jeromy committed
233 234
	// If we've reached a leaf node.
	if len(node.Links) == 0 {
235 236
		n, err := data.Read(f.Data[offset:])
		if err != nil && err != io.EOF {
Jeromy's avatar
Jeromy committed
237
			return "", false, err
238 239 240 241 242
		}

		// Update newly written node..
		b, err := proto.Marshal(f)
		if err != nil {
Jeromy's avatar
Jeromy committed
243
			return "", false, err
244 245 246 247 248
		}

		nd := &mdag.Node{Data: b}
		k, err := dm.dagserv.Add(nd)
		if err != nil {
Jeromy's avatar
Jeromy committed
249
			return "", false, err
250 251 252 253
		}

		// Hey look! we're done!
		var done bool
254
		if n < len(f.Data[offset:]) {
255 256 257
			done = true
		}

Jeromy's avatar
Jeromy committed
258
		return k, done, nil
259 260 261 262 263
	}

	var cur uint64
	var done bool
	for i, bs := range f.GetBlocksizes() {
264
		// We found the correct child to write into
265
		if cur+bs > offset {
266 267 268 269
			// Unpin block
			ckey := u.Key(node.Links[i].Hash)
			dm.mp.RemovePinWithMode(ckey, pin.Indirect)

270 271
			child, err := node.Links[i].GetNode(dm.dagserv)
			if err != nil {
Jeromy's avatar
Jeromy committed
272
				return "", false, err
273
			}
Jeromy's avatar
Jeromy committed
274
			k, sdone, err := dm.modifyDag(child, offset-cur, data)
275
			if err != nil {
Jeromy's avatar
Jeromy committed
276
				return "", false, err
277 278
			}

279 280 281
			// pin the new node
			dm.mp.PinWithMode(k, pin.Indirect)

282 283 284
			offset += bs
			node.Links[i].Hash = mh.Multihash(k)

285 286 287 288 289 290
			// Recache serialized node
			_, err = node.Encoded(true)
			if err != nil {
				return "", false, err
			}

291
			if sdone {
292
				// No more bytes to write!
293 294 295
				done = true
				break
			}
296
			offset = cur + bs
297 298 299 300 301
		}
		cur += bs
	}

	k, err := dm.dagserv.Add(node)
Jeromy's avatar
Jeromy committed
302
	return k, done, err
303 304
}

Jeromy's avatar
Jeromy committed
305
// appendData appends the blocks from the given chan to the end of this dag
306 307 308 309 310 311 312 313 314 315
func (dm *DagModifier) appendData(node *mdag.Node, blks <-chan []byte) (*mdag.Node, error) {
	dbp := &help.DagBuilderParams{
		Dagserv:  dm.dagserv,
		Maxlinks: help.DefaultLinksPerBlock,
		Pinner:   dm.mp,
	}

	return trickle.TrickleAppend(node, dbp.New(blks))
}

Jeromy's avatar
Jeromy committed
316
// Read data from this dag starting at the current offset
317
func (dm *DagModifier) Read(b []byte) (int, error) {
318
	err := dm.readPrep()
319 320 321 322
	if err != nil {
		return 0, err
	}

323 324 325 326 327 328 329 330 331 332 333
	n, err := dm.read.Read(b)
	dm.curWrOff += uint64(n)
	return n, err
}

func (dm *DagModifier) readPrep() error {
	err := dm.Sync()
	if err != nil {
		return err
	}

334
	if dm.read == nil {
335 336
		ctx, cancel := context.WithCancel(dm.ctx)
		dr, err := uio.NewDagReader(ctx, dm.curNode, dm.dagserv)
337
		if err != nil {
338
			return err
339 340 341 342
		}

		i, err := dr.Seek(int64(dm.curWrOff), os.SEEK_SET)
		if err != nil {
343
			return err
344 345 346
		}

		if i != int64(dm.curWrOff) {
347
			return ErrSeekFail
348 349
		}

350
		dm.readCancel = cancel
351 352 353
		dm.read = dr
	}

354 355 356 357 358 359 360 361 362 363 364
	return nil
}

// Read data from this dag starting at the current offset
func (dm *DagModifier) CtxReadFull(ctx context.Context, b []byte) (int, error) {
	err := dm.readPrep()
	if err != nil {
		return 0, err
	}

	n, err := dm.read.CtxReadFull(ctx, b)
365 366 367 368 369 370
	dm.curWrOff += uint64(n)
	return n, err
}

// GetNode gets the modified DAG Node
func (dm *DagModifier) GetNode() (*mdag.Node, error) {
371
	err := dm.Sync()
372 373 374 375 376 377
	if err != nil {
		return nil, err
	}
	return dm.curNode.Copy(), nil
}

Jeromy's avatar
Jeromy committed
378
// HasChanges returned whether or not there are unflushed changes to this dag
379 380 381 382 383
func (dm *DagModifier) HasChanges() bool {
	return dm.wrBuf != nil
}

func (dm *DagModifier) Seek(offset int64, whence int) (int64, error) {
384
	err := dm.Sync()
385 386 387 388 389 390 391 392 393 394 395 396
	if err != nil {
		return 0, err
	}

	switch whence {
	case os.SEEK_CUR:
		dm.curWrOff += uint64(offset)
		dm.writeStart = dm.curWrOff
	case os.SEEK_SET:
		dm.curWrOff = uint64(offset)
		dm.writeStart = uint64(offset)
	case os.SEEK_END:
Jeromy's avatar
Jeromy committed
397
		return 0, ErrSeekEndNotImpl
398
	default:
Jeromy's avatar
Jeromy committed
399
		return 0, ErrUnrecognizedWhence
400 401 402 403 404 405 406 407 408 409 410 411 412
	}

	if dm.read != nil {
		_, err = dm.read.Seek(offset, whence)
		if err != nil {
			return 0, err
		}
	}

	return int64(dm.curWrOff), nil
}

func (dm *DagModifier) Truncate(size int64) error {
413
	err := dm.Sync()
414 415 416 417 418 419 420 421 422
	if err != nil {
		return err
	}

	realSize, err := dm.Size()
	if err != nil {
		return err
	}

Jeromy's avatar
Jeromy committed
423
	// Truncate can also be used to expand the file
424
	if size > int64(realSize) {
Jeromy's avatar
Jeromy committed
425
		return dm.expandSparse(int64(size) - realSize)
426 427 428 429 430 431 432 433 434 435 436 437 438 439 440 441
	}

	nnode, err := dagTruncate(dm.curNode, uint64(size), dm.dagserv)
	if err != nil {
		return err
	}

	_, err = dm.dagserv.Add(nnode)
	if err != nil {
		return err
	}

	dm.curNode = nnode
	return nil
}

Jeromy's avatar
Jeromy committed
442
// dagTruncate truncates the given node to 'size' and returns the modified Node
443 444 445 446 447 448 449 450 451 452 453 454 455 456 457 458 459 460 461 462 463 464 465 466 467 468 469
func dagTruncate(nd *mdag.Node, size uint64, ds mdag.DAGService) (*mdag.Node, error) {
	if len(nd.Links) == 0 {
		// TODO: this can likely be done without marshaling and remarshaling
		pbn, err := ft.FromBytes(nd.Data)
		if err != nil {
			return nil, err
		}

		nd.Data = ft.WrapData(pbn.Data[:size])
		return nd, nil
	}

	var cur uint64
	end := 0
	var modified *mdag.Node
	ndata := new(ft.FSNode)
	for i, lnk := range nd.Links {
		child, err := lnk.GetNode(ds)
		if err != nil {
			return nil, err
		}

		childsize, err := ft.DataSize(child.Data)
		if err != nil {
			return nil, err
		}

Jeromy's avatar
Jeromy committed
470
		// found the child we want to cut
471 472 473 474 475 476 477 478 479 480 481 482 483 484 485 486 487 488 489 490 491 492 493 494 495 496 497 498 499 500 501 502 503 504
		if size < cur+childsize {
			nchild, err := dagTruncate(child, size-cur, ds)
			if err != nil {
				return nil, err
			}

			ndata.AddBlockSize(size - cur)

			modified = nchild
			end = i
			break
		}
		cur += childsize
		ndata.AddBlockSize(childsize)
	}

	_, err := ds.Add(modified)
	if err != nil {
		return nil, err
	}

	nd.Links = nd.Links[:end]
	err = nd.AddNodeLinkClean("", modified)
	if err != nil {
		return nil, err
	}

	d, err := ndata.GetBytes()
	if err != nil {
		return nil, err
	}

	nd.Data = d

505 506 507 508 509 510
	// invalidate cache and recompute serialized data
	_, err = nd.Encoded(true)
	if err != nil {
		return nil, err
	}

511 512
	return nd, nil
}