Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 7 additions & 3 deletions main.go
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
package main

Check warning on line 1 in main.go

View workflow job for this annotation

GitHub Actions / fmt_vet_lint

should have a package comment

import (
"context"
Expand Down Expand Up @@ -32,7 +32,7 @@
Path string `arg:""`
Bucket string `help:"Remote bucket"`
Metadata bool `help:"Print only the JSON metadata"`
HeaderJson bool `help:"Print a JSON representation of part of the header information"`

Check warning on line 35 in main.go

View workflow job for this annotation

GitHub Actions / fmt_vet_lint

struct field HeaderJson should be HeaderJSON
Tilejson bool `help:"Print the TileJSON"`
PublicURL string `help:"Public base URL of tile endpoint for TileJSON e.g. https://example.com/tiles"`
} `cmd:"" help:"Inspect a local or remote archive"`
Expand All @@ -52,7 +52,7 @@

Edit struct {
Input string `arg:"" help:"Input archive" type:"existingfile"`
HeaderJson string `help:"Input header JSON file (written by show --header-json)" type:"existingfile"`

Check warning on line 55 in main.go

View workflow job for this annotation

GitHub Actions / fmt_vet_lint

struct field HeaderJson should be HeaderJSON
Metadata string `help:"Input metadata JSON file (written by show --metadata)" type:"existingfile"`
} `cmd:"" help:"Edit JSON metadata or parts of the header"`

Expand All @@ -70,9 +70,8 @@
} `cmd:"" help:"Create an archive from a larger archive for a subset of zoom levels or geographic region"`

Merge struct {
Output string `arg:"" help:"Output archive" type:"path"`
Input []string `arg:"" help:"Input archives"`
} `cmd:"" help:"Merge multiple archives into a single archive" hidden:""`
Archives []string `arg:"" name:"inputs_then_output" help:"One or more disjoint input archives, followed by the output filename."`
} `cmd:"" help:"Merge multiple disjoint archives into one: INPUT1.pmtiles INPUT2.pmtiles OUTPUT.pmtiles"`

Convert struct {
Input string `arg:"" help:"Input archive" type:"existingfile"`
Expand Down Expand Up @@ -217,6 +216,11 @@
if err != nil {
logger.Fatalf("Failed to convert %s, %v", path, err)
}
case "merge <inputs_then_output>":
err := pmtiles.Merge(logger, cli.Merge.Archives)
if err != nil {
logger.Fatalf("Failed to merge, %v", err)
}
case "upload <input-pmtiles> <remote-pmtiles>":
err := pmtiles.Upload(logger, cli.Upload.InputPmtiles, cli.Upload.Bucket, cli.Upload.RemotePmtiles, cli.Upload.MaxConcurrency)

Expand Down
51 changes: 51 additions & 0 deletions pmtiles/directory_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,9 +4,60 @@ import (
"bytes"
"github.com/stretchr/testify/assert"
"math/rand"
"sort"
"testing"
)

func fakeArchive(header HeaderV3, metadata map[string]interface{}, tiles map[Zxy][]byte, leaves bool, internalCompression Compression) []byte {
byTileID := make(map[uint64][]byte)
keys := make([]uint64, 0, len(tiles))
for zxy, bytes := range tiles {
header.MaxZoom = max(header.MaxZoom, zxy.Z)
id := ZxyToID(zxy.Z, zxy.X, zxy.Y)
byTileID[id] = bytes
keys = append(keys, id)
}
sort.Slice(keys, func(i, j int) bool { return keys[i] < keys[j] })
resolver := newResolver(false, false)
tileDataBytes := make([]byte, 0)
for _, id := range keys {
tileBytes := byTileID[id]
resolver.AddTileIsNew(id, tileBytes, 1)
tileDataBytes = append(tileDataBytes, tileBytes...)
}

metadataBytes, _ := SerializeMetadata(metadata, internalCompression)
var rootBytes []byte
var leavesBytes []byte
if leaves {
rootBytes, leavesBytes, _ = buildRootsLeaves(resolver.Entries, 1, internalCompression)
} else {
rootBytes = SerializeEntries(resolver.Entries, internalCompression)
leavesBytes = make([]byte, 0)
}

header.InternalCompression = internalCompression
if header.TileType == Mvt {
header.TileCompression = Gzip
}

header.RootOffset = HeaderV3LenBytes
header.RootLength = uint64(len(rootBytes))
header.MetadataOffset = header.RootOffset + header.RootLength
header.MetadataLength = uint64(len(metadataBytes))
header.LeafDirectoryOffset = header.MetadataOffset + header.MetadataLength
header.LeafDirectoryLength = uint64(len(leavesBytes))
header.TileDataOffset = header.LeafDirectoryOffset + header.LeafDirectoryLength
header.TileDataLength = resolver.Offset

archiveBytes := SerializeHeader(header)
archiveBytes = append(archiveBytes, rootBytes...)
archiveBytes = append(archiveBytes, metadataBytes...)
archiveBytes = append(archiveBytes, leavesBytes...)
archiveBytes = append(archiveBytes, tileDataBytes...)
return archiveBytes
}

func TestDirectoryRoundtrip(t *testing.T) {
entries := make([]EntryV3, 0)
entries = append(entries, EntryV3{0, 0, 0, 0})
Expand Down
307 changes: 307 additions & 0 deletions pmtiles/merge.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,307 @@
package pmtiles

import (
"fmt"
"github.com/RoaringBitmap/roaring/roaring64"
"github.com/schollz/progressbar/v3"
"io"
"log"
"math"
"os"
"slices"
"sort"
)

type MergeEntry struct {
Entry EntryV3
InputIdx int // the index of the input archive 0...N
InputOffset uint64 // the original offset of the entry in the archive's tile section
}

type MergeOp struct {
InputIdx int
Length uint64
}

type Remapping struct {
SrcOffset uint64
DstOffset uint64
}

// load N archives, validating that they are mergeable and disjoint.
// returns a formatted error and index of the mismatched archive if not.
// if valid, returns a sorted list of MergeEntry records, each containing a directory Entry
// but with offset values referring to the original input archive.
func prepareInputs(inputs []io.ReadSeeker) ([]HeaderV3, []MergeEntry, error, int) {
var headers []HeaderV3
var mergedEntries []MergeEntry
union := roaring64.New()

for inputIdx, input := range inputs {
buf := make([]byte, HeaderV3LenBytes)
_, err := input.Read(buf)
if err != nil {
return nil, nil, err, inputIdx
}
h, err := DeserializeHeader(buf)
if err != nil {
return nil, nil, err, inputIdx
}
headers = append(headers, h)

if !h.Clustered {
return nil, nil, fmt.Errorf("must be clustered"), inputIdx
}

if inputIdx > 0 {
if h.TileType != headers[0].TileType {
return nil, nil, fmt.Errorf("tile type %s does not match %s", tileTypeToString(h.TileType), tileTypeToString(headers[0].TileType)), inputIdx
}
if h.TileCompression != headers[0].TileCompression {
c1, _ := compressionToString(h.TileCompression)
c2, _ := compressionToString(headers[0].TileCompression)
return nil, nil, fmt.Errorf("tile compression %s does not match %s", c1, c2), inputIdx
}
if h.InternalCompression != headers[0].InternalCompression {
c1, _ := compressionToString(h.InternalCompression)
c2, _ := compressionToString(headers[0].InternalCompression)
return nil, nil, fmt.Errorf("internal compression %s does not match %s", c1, c2), inputIdx
}
}

tileset := roaring64.New()
err = IterateEntries(h,
func(offset uint64, length uint64) ([]byte, error) {
input.Seek(int64(offset), io.SeekStart)
return io.ReadAll(io.LimitReader(input, int64(length)))
},
func(e EntryV3) {
tileset.AddRange(e.TileID, e.TileID+uint64(e.RunLength))
mergedEntries = append(mergedEntries, MergeEntry{Entry: e, InputOffset: e.Offset, InputIdx: inputIdx})
})

if err != nil {
return nil, nil, err, inputIdx
}

if union.Intersects(tileset) {
tmp := union.Clone()
tmp.And(tileset)
iz, ix, iy := IDToZxy(tmp.Minimum())
return nil, nil, fmt.Errorf("%d overlapping tiles, starting with %d %d %d. Inputs must be disjoint", tmp.GetCardinality(), iz, ix, iy), inputIdx
}
union.Or(tileset)
}

sort.Slice(mergedEntries, func(i, j int) bool {
return mergedEntries[i].Entry.TileID < mergedEntries[j].Entry.TileID
})

return headers, mergedEntries, nil, 0
}

// remaps a sorted slice of MergeEntry
// changes each Entry to be contiguous in the new archive.
// also handles deduplicated backreferences
func remapMergeEntries(entries []MergeEntry, numInputs int) ([]MergeEntry, uint64, uint64, uint64, error) {

acc := uint64(0)
addressedTiles := uint64(0)
tileContents := 0
remappings := make([][]Remapping, numInputs)

for idx, me := range entries {
remapping := remappings[me.InputIdx]
if len(remapping) > 0 && me.InputOffset <= remapping[len(remapping)-1].SrcOffset {
// find the original offset in the remapping slice
i, ok := slices.BinarySearchFunc(remapping, me.InputOffset, func(r Remapping, k uint64) int {
switch {
case r.SrcOffset < k:
return -1
case r.SrcOffset > k:
return 1
default:
return 0
}
})
if ok {
entries[idx].Entry.Offset = remapping[i].DstOffset
} else {
return nil, 0, 0, 0, fmt.Errorf("clustered archive has out-of-order entries")
}
} else {
entries[idx].Entry.Offset = acc
remappings[me.InputIdx] = append(remappings[me.InputIdx], Remapping{SrcOffset: me.InputOffset, DstOffset: acc})
acc += uint64(me.Entry.Length)
tileContents += 1
}

addressedTiles += uint64(entries[idx].Entry.RunLength)
}

return entries, addressedTiles, uint64(tileContents), acc, nil
}

// combines contiguous I/O operations and eliminate backreferences from the copy operation.
func batchMergeEntries(entries []MergeEntry, numInputs int) []MergeOp {
lastOffset := make([]int64, numInputs)
for i := range lastOffset {
lastOffset[i] = -1
}
lastLength := make([]uint32, numInputs)
var mergeOps []MergeOp
for _, me := range entries {
if int64(me.InputOffset) > lastOffset[me.InputIdx] {
last := len(mergeOps) - 1
entryLength := uint64(me.Entry.Length)
if last >= 0 && (mergeOps[last].InputIdx == me.InputIdx) && (int64(me.InputOffset) == lastOffset[me.InputIdx]+int64(lastLength[me.InputIdx])) {
mergeOps[last].Length += entryLength
} else {
mergeOps = append(mergeOps, MergeOp{InputIdx: me.InputIdx, Length: entryLength})
}
lastOffset[me.InputIdx] = int64(me.InputOffset)
lastLength[me.InputIdx] = me.Entry.Length
}
}
return mergeOps
}

func zoomBounds(entries []MergeEntry) (uint8, uint8) {
firstZ, _, _ := IDToZxy(entries[0].Entry.TileID)
lastEntry := entries[len(entries)-1].Entry
lastZ, _, _ := IDToZxy(lastEntry.TileID + uint64(lastEntry.RunLength) - 1)
return uint8(firstZ), uint8(lastZ)
}

func bounds(headers []HeaderV3) (int32, int32, int32, int32) {
minLonE7 := int32(math.MaxInt32)
minLatE7 := int32(math.MaxInt32)
maxLonE7 := int32(math.MinInt32)
maxLatE7 := int32(math.MinInt32)

for _, h := range headers {
if h.MinLonE7 < minLonE7 {
minLonE7 = h.MinLonE7
}
if h.MinLatE7 < minLatE7 {
minLatE7 = h.MinLatE7
}
if h.MaxLonE7 > maxLonE7 {
maxLonE7 = h.MaxLonE7
}
if h.MaxLatE7 > maxLatE7 {
maxLatE7 = h.MaxLatE7
}
}

return minLonE7, minLatE7, maxLonE7, maxLatE7
}

func Merge(logger *log.Logger, inputs []string) error {
var handles []io.ReadSeeker

if len(inputs) < 2 {
return fmt.Errorf("Too few inputs")
}

for _, name := range inputs[:len(inputs)-1] {
f, err := os.OpenFile(name, os.O_RDONLY, 0666)
if err != nil {
return err
}
handles = append(handles, f)
defer f.Close()
}

headers, mergedEntries, err, errIdx := prepareInputs(handles)
if err != nil {
return fmt.Errorf("%s: %w", inputs[errIdx], err)
}

renumberedEntries, addressedTiles, numTileContents, tileDataLength, err := remapMergeEntries(mergedEntries, len(headers))
if err != nil {
return err
}

tmp := make([]EntryV3, len(renumberedEntries))
for i := range renumberedEntries {
tmp[i] = renumberedEntries[i].Entry
}
rootBytes, leavesBytes, _ := optimizeDirectories(tmp, 16384-HeaderV3LenBytes, Gzip)

logger.Printf("Copying center and JSON metadata from first input %s", inputs[0])

var header HeaderV3
header.Clustered = true
header.RootOffset = HeaderV3LenBytes
header.RootLength = uint64(len(rootBytes))
header.MetadataOffset = header.RootOffset + header.RootLength
header.MetadataLength = headers[0].MetadataLength
header.TileType = headers[0].TileType
header.InternalCompression = headers[0].InternalCompression
header.TileCompression = headers[0].TileCompression
header.LeafDirectoryOffset = header.MetadataOffset + header.MetadataLength
header.LeafDirectoryLength = uint64(len(leavesBytes))
header.TileDataOffset = header.LeafDirectoryOffset + header.LeafDirectoryLength
header.TileDataLength = tileDataLength
header.AddressedTilesCount = addressedTiles
header.TileEntriesCount = uint64(len(renumberedEntries))
header.TileContentsCount = numTileContents

minZoom, maxZoom := zoomBounds(renumberedEntries)
header.MinZoom = minZoom
header.MaxZoom = maxZoom
minLonE7, minLatE7, maxLonE7, maxLatE7 := bounds(headers)
header.MinLonE7 = minLonE7
header.MinLatE7 = minLatE7
header.MaxLonE7 = maxLonE7
header.MaxLatE7 = maxLatE7
header.CenterLonE7 = headers[0].CenterLonE7
header.CenterLatE7 = headers[0].CenterLatE7
header.CenterZoom = headers[0].CenterZoom

mergeOps := batchMergeEntries(renumberedEntries, len(headers))

output, err := os.Create(inputs[len(inputs)-1])
if err != nil {
return err
}
defer output.Close()

headerBytes := SerializeHeader(header)
_, err = output.Write(headerBytes)
if err != nil {
return err
}
_, err = output.Write(rootBytes)
if err != nil {
return err
}
firstHandle := handles[0]
firstHandle.Seek(int64(headers[0].MetadataOffset), io.SeekStart)
io.CopyN(output, firstHandle, int64(headers[0].MetadataLength))

_, err = output.Write(leavesBytes)
if err != nil {
return err
}

for idx, handle := range handles {
handle.Seek(int64(headers[idx].TileDataOffset), io.SeekStart)
}

bar := progressbar.DefaultBytes(
int64(tileDataLength),
"merging tile data",
)

for _, op := range mergeOps {
handle := handles[op.InputIdx]
_, err := io.CopyN(io.MultiWriter(output, bar), handle, int64(op.Length))
if err != nil {
return err
}
}

return nil
}
Loading
Loading