add combine

This commit is contained in:
zkx 2023-12-18 10:48:21 +08:00
parent 1596c796b3
commit cba67e10a1
5 changed files with 112 additions and 194 deletions

View File

@ -1,190 +0,0 @@
package ec
import (
"errors"
"io"
"io/ioutil"
"gitlink.org.cn/cloudream/common/pkgs/ipfs"
"gitlink.org.cn/cloudream/common/pkgs/logger"
stgglb "gitlink.org.cn/cloudream/storage/common/globals"
)
type BlockReader struct {
ipfsCli *ipfs.PoolClient
/*将文件分块相关的属性*/
//fileHash
fileHash string
//fileSize
fileSize int64
//ecK将文件的分块数
ecK int
//chunkSize
chunkSize int64
/*可选项*/
//fastRead,true的时候直接通过hash读block
jumpReadOpt bool
}
func NewBlockReader() (*BlockReader, error) {
ipfsClient, err := stgglb.IPFSPool.Acquire()
if err != nil {
return nil, err
}
//default:fast模式通过hash直接获取
return &BlockReader{ipfsCli: ipfsClient, chunkSize: 256 * 1024, jumpReadOpt: false}, nil
}
func (r *BlockReader) Close() {
r.ipfsCli.Close()
}
func (r *BlockReader) SetJumpRead(fileHash string, fileSize int64, ecK int) {
r.fileHash = fileHash
r.fileSize = fileSize
r.ecK = ecK
r.jumpReadOpt = true
}
func (r *BlockReader) SetchunkSize(size int64) {
r.chunkSize = size
}
func (r *BlockReader) FetchBLock(blockHash string) (io.ReadCloser, error) {
return r.ipfsCli.OpenRead(blockHash)
}
func (r *BlockReader) FetchBLocks(blockHashs []string) ([]io.ReadCloser, error) {
readers := make([]io.ReadCloser, len(blockHashs))
for i, hash := range blockHashs {
var err error
readers[i], err = r.ipfsCli.OpenRead(hash)
if err != nil {
return nil, err
}
}
return readers, nil
}
func (r *BlockReader) JumpFetchBlock(innerID int) (io.ReadCloser, error) {
if !r.jumpReadOpt {
return nil, nil
}
pipeReader, pipeWriter := io.Pipe()
go func() {
for i := int64(r.chunkSize * int64(innerID)); i < r.fileSize; i += int64(r.ecK) * r.chunkSize {
reader, err := r.ipfsCli.OpenRead(r.fileHash, ipfs.ReadOption{Offset: i, Length: r.chunkSize})
if err != nil {
pipeWriter.CloseWithError(err)
return
}
data, err := ioutil.ReadAll(reader)
if err != nil {
pipeWriter.CloseWithError(err)
return
}
reader.Close()
_, err = pipeWriter.Write(data)
if err != nil {
pipeWriter.CloseWithError(err)
return
}
}
//如果文件大小不是分块的整数倍,可能需要补0
if r.fileSize%(r.chunkSize*int64(r.ecK)) != 0 {
//pktNum_1:chunkNum-1
pktNum_1 := r.fileSize / (r.chunkSize * int64(r.ecK))
offset := (r.fileSize - int64(pktNum_1)*int64(r.ecK)*r.chunkSize)
count0 := int64(innerID)*int64(r.ecK)*r.chunkSize - offset
if count0 > 0 {
add0 := make([]byte, count0)
pipeWriter.Write(add0)
}
}
pipeWriter.Close()
}()
return pipeReader, nil
}
// FetchBlock1这个函数废弃了
func (r *BlockReader) FetchBlock1(input interface{}, errMsg chan error) (io.ReadCloser, error) {
/*两种模式下传入第一个参数但是input的类型不同
jumpReadOpt-true传入blcokHash, string型通过哈希直接读
jumpReadOpt->false: 传入innerIDint型选择需要获取的数据块的id
*/
var innerID int
var blockHash string
switch input.(type) {
case int:
// 执行针对整数的逻辑分支
if r.jumpReadOpt {
return nil, errors.New("conflict, wrong input type and jumpReadOpt:true")
} else {
innerID = input.(int)
}
case string:
if !r.jumpReadOpt {
return nil, errors.New("conflict, wrong input type and jumpReadOpt:false")
} else {
blockHash = input.(string)
}
default:
return nil, errors.New("wrong input type")
}
//开始执行
if r.jumpReadOpt { //快速读
ipfsCli, err := stgglb.IPFSPool.Acquire()
if err != nil {
logger.Warnf("new ipfs client: %s", err.Error())
return nil, err
}
defer ipfsCli.Close()
return ipfsCli.OpenRead(blockHash)
} else { //跳跃读
ipfsCli, err := stgglb.IPFSPool.Acquire()
if err != nil {
logger.Warnf("new ipfs client: %s", err.Error())
return nil, err
}
defer ipfsCli.Close()
pipeReader, pipeWriter := io.Pipe()
go func() {
for i := int64(r.chunkSize * int64(innerID)); i < r.fileSize; i += int64(r.ecK) * r.chunkSize {
reader, err := ipfsCli.OpenRead(r.fileHash, ipfs.ReadOption{i, r.chunkSize})
if err != nil {
pipeWriter.Close()
errMsg <- err
return
}
data, err := ioutil.ReadAll(reader)
if err != nil {
pipeWriter.Close()
errMsg <- err
return
}
reader.Close()
_, err = pipeWriter.Write(data)
if err != nil {
pipeWriter.Close()
errMsg <- err
return
}
}
//如果文件大小不是分块的整数倍,可能需要补0
if r.fileSize%(r.chunkSize*int64(r.ecK)) != 0 {
//pktNum_1:chunkNum-1
pktNum_1 := r.fileSize / (r.chunkSize * int64(r.ecK))
offset := (r.fileSize - int64(pktNum_1)*int64(r.ecK)*r.chunkSize)
count0 := int64(innerID)*int64(r.ecK)*r.chunkSize - offset
if count0 > 0 {
add0 := make([]byte, count0)
pipeWriter.Write(add0)
}
}
pipeWriter.Close()
errMsg <- nil
}()
return pipeReader, nil
}
}

View File

@ -0,0 +1,36 @@
package ec
import (
"io"
"os"
"testing"
"gitlink.org.cn/cloudream/common/pkgs/ipfs"
stgglb "gitlink.org.cn/cloudream/storage/common/globals"
//"gitlink.org.cn/cloudream/common/pkgs/ipfs"
//"gitlink.org.cn/cloudream/storage/agent/internal/config"
//stgglb "gitlink.org.cn/cloudream/storage/common/globals"
)
// QmW8wiHes4qHCav5jvWZ4356BpKtbLjsZEXiJj3F2SPhju :74475
// QmZfn2XW4TA6a59abAYsp6RjCjzuJh1wLVSheYRknLdzyR :66455
func test_pdf(t *testing.T) {
chunkSize := int64(1024 * 1024 * 13)
blkReader, _ := NewBlockReader()
defer blkReader.Close()
blkReader.SetJumpRead("QmepgjSRu3ELERvYuEG9iAhM4ZgUVeHD2tM2LzYis8VRh4", 131433640, 3)
blkReader.SetchunkSize(chunkSize)
dataBlocks := make([]io.ReadCloser, 3)
dataBlocks[0], _ = blkReader.JumpFetchBlock(0)
dataBlocks[1], _ = blkReader.JumpFetchBlock(1)
dataBlocks[2], _ = blkReader.JumpFetchBlock(2)
enc, _ := NewRs(3, 5, chunkSize)
fw_1, _ := os.Create("my1.pptx")
fptr, _ := enc.Combine(dataBlocks, 131433640)
io.Copy(fw_1, fptr)
defer fw_1.Close()
}
func Test_file(t *testing.T) {
stgglb.InitIPFSPool(&ipfs.Config{Port: 5001})
test_pdf(t)
}

View File

@ -14,6 +14,21 @@ import (
stgglb "gitlink.org.cn/cloudream/storage/common/globals"
)
func test_Combile(t *testing.T) {
chunkSize := int64(6)
blkReader, _ := NewBlockReader()
defer blkReader.Close()
blkReader.SetJumpRead("QmcN1EJm2w9XT62Q9YqA5Ym7YDzjmnqJYc565bzRs5VosW", 46, 3)
blkReader.SetchunkSize(chunkSize)
dataBlocks := make([]io.ReadCloser, 3)
dataBlocks[0], _ = blkReader.JumpFetchBlock(0)
dataBlocks[1], _ = blkReader.JumpFetchBlock(1)
dataBlocks[2], _ = blkReader.JumpFetchBlock(2)
enc, _ := NewRs(3, 5, chunkSize)
fptr, _ := enc.Combine(dataBlocks, 46)
print_ioreaders(t, []io.ReadCloser{fptr}, chunkSize)
}
func test_Encode(t *testing.T) {
enc, _ := NewRs(3, 5, 10)
rc := make([]io.ReadCloser, 3)
@ -204,7 +219,8 @@ func Test_main(t *testing.T) {
//test_Fetch_and_Encode(t)
//test_Fetch_and_Encode_and_Degraded(t)
//test_pin_data_blocks(t)
test_reconstructData(t)
//test_reconstructData(t)
test_Combile(t)
}
/*

View File

@ -213,8 +213,64 @@ func (r *Rs) ReconstructSome(input []io.ReadCloser, inBlockIdx []int, outBlockId
outWriter[i].Close()
}
}()
if err != nil {
return nil, err
}
return outReader, nil
}
func (r *Rs) Combine(input []io.ReadCloser, fileSize int64) (io.ReadCloser, error) {
outReader, outWriter := io.Pipe()
buf := make([]byte, r.chunkSize)
go func() {
//直接写入flagtrue时直接写入不必考虑写入后超过文件大小
DirectWrite := true
for {
if fileSize == 0 {
outWriter.Close()
break
}
if fileSize < int64(r.ecK)*r.chunkSize {
DirectWrite = false
}
for i := 0; i < r.ecK; i++ {
_, err := input[i].Read(buf)
if err == io.EOF {
break
} else if err != nil {
outReader.CloseWithError(err)
return
}
//直接写,不用检查
if DirectWrite {
_, err = outWriter.Write(buf)
if err != nil {
outWriter.CloseWithError(err)
return
}
fileSize -= r.chunkSize
//检查是否超过fileSize
} else {
if fileSize == 0 {
break
}
if fileSize >= r.chunkSize {
_, err = outWriter.Write(buf)
if err != nil {
outWriter.CloseWithError(err)
return
}
fileSize -= r.chunkSize
} else {
_, err = outWriter.Write(buf[:fileSize])
if err != nil {
outWriter.CloseWithError(err)
return
}
fileSize = 0
}
}
}
}
}()
return outReader, nil
}

0
common/pkgs/ec/test.txt Normal file
View File