From 0dd093e537cc9e70fda25a0bb9cb8b26164d8c26 Mon Sep 17 00:00:00 2001 From: Quadrubo <71718414+Quadrubo@users.noreply.github.com> Date: Sun, 7 Jun 2026 17:27:56 +0200 Subject: [PATCH] feat: split tracks at POI markers --- .gitignore | 4 + README.md | 44 +++++ justfile | 8 +- nixos/tracksync.nix | 5 +- server/.env.example | 4 + server/internal/config/config.go | 26 ++- server/internal/config/config_test.go | 85 +++++++++ server/internal/converter/converter.go | 68 +++++--- server/internal/converter/converter_test.go | 156 +++++++++++++++-- server/internal/converter/csv_parser.go | 22 ++- server/internal/converter/csv_parser_test.go | 17 ++ server/internal/converter/gpx_parser.go | 6 +- server/internal/converter/gpx_serializer.go | 14 +- server/internal/converter/marker.go | 121 +++++++++++++ server/internal/converter/marker_test.go | 140 +++++++++++++++ server/internal/converter/split.go | 47 +++++ server/internal/converter/split_test.go | 101 +++++++++++ server/internal/converter/track.go | 4 + .../002_received_forwarded_files.sql | 35 ++++ server/internal/migrations/migrations_test.go | 89 ++++++++++ server/internal/server/server.go | 162 ++++++++++++++---- server/internal/server/server_test.go | 79 ++++++++- server/internal/target/dawarich/dawarich.go | 4 +- server/main.go | 9 +- tracksync/internal/sync/sync.go | 32 ++-- tracksync/internal/sync/sync_test.go | 49 +++++- tracksync/main.go | 1 + 27 files changed, 1209 insertions(+), 123 deletions(-) create mode 100644 server/internal/converter/marker.go create mode 100644 server/internal/converter/marker_test.go create mode 100644 server/internal/converter/split.go create mode 100644 server/internal/converter/split_test.go create mode 100644 server/internal/migrations/002_received_forwarded_files.sql create mode 100644 server/internal/migrations/migrations_test.go diff --git a/.gitignore b/.gitignore index b197e47..320f75d 100644 --- a/.gitignore +++ b/.gitignore @@ -3,5 +3,9 @@ docker-compose.local.yml .env result +# IDEs +.idea/ +.vscode/ + # NixOS .direnv diff --git a/README.md b/README.md index 854e0b1..ff91d98 100644 --- a/README.md +++ b/README.md @@ -82,6 +82,9 @@ All configuration is done via environment variables. | `ACCOUNT__N__TARGET_URL` | | Yes | Target instance URL | | `ACCOUNT__N__API_KEY` | | Yes if not API_KEY_FILE | API key (inline) | | `ACCOUNT__N__API_KEY_FILE` | | Yes if not API_KEY | API key (file path) | +| `ACCOUNT__N__MARKERS` | | No | Comma-separated `marker:functionality` rules assigning behavior to point markers (see [Marker functionalities](#marker-functionalities)) | +| `ACCOUNT__N__SPLIT_MARKER_POSITION` | `start` | No | For the `split` functionality: where the marked point goes, `start` of the new track or `end` of the previous one | +| `ACCOUNT__N__SPLIT_MODE` | `tracks` | No | For the `split` functionality: `tracks` (all tracks in one file) or `files` (one upload per track) | | `CLIENT__N__ID` | | Yes | Client identifier | | `CLIENT__N__TOKEN` | | Yes if not TOKEN_FILE | Auth token (inline) | | `CLIENT__N__TOKEN_FILE` | | Yes if not TOKEN | Auth token (file path) | @@ -191,6 +194,47 @@ For example, when a Columbus P-10 Pro is configured to output CSV (which include | `columbus-csv` | Parse | Yes | Yes | Yes | No | No | | `geojson` | Serialize | Yes | Yes | Yes | Yes | Yes | +## Marker functionalities + +Points can carry a *marker*, a source-format annotation such as a manually +placed POI or waypoint. You can assign a functionality to each marker so that +tracksync acts on it. This works at the universal-track level and is +format-agnostic: every parser maps its format's native markers onto a point marker, and functionalities operate on those. + +Configure rules with `ACCOUNT__N__MARKERS` as comma-separated +`marker:functionality` pairs: + +``` +ACCOUNT__0__MARKERS=C:split +# multiple markers, each with its own functionality: +ACCOUNT__0__MARKERS=C:split,D:split +``` + +Which marker values are available depends on the source format: + +| Format | Markers | +| -------------- | ---------------------------------------------------------------------------------------- | +| `columbus-csv` | `TAG` column values other than `T`: `C` (function-key POI), `D` (second POI), `G` (automatic wake-up point, usually leave unmapped) | + +### Available functionalities + +| Functionality | Effect | +| ------------- | ---------------------------------------------------------------------------------------- | +| `split` | Start a new track at the marked point. Useful for separating legs of a journey. For example pressing a logger's function key when boarding and leaving a bus so the walking and bus legs become distinct tracks. | + +For the `split` functionality, `ACCOUNT__N__SPLIT_MARKER_POSITION` controls which +side of the split the marked point lands on: `start` (default) makes it the first +point of the new track, `end` keeps it as the last point of the previous one. + +`ACCOUNT__N__SPLIT_MODE` controls how the split tracks are delivered: + +- `tracks` (default): all tracks are written into a single file. +- `files`: each track is uploaded as its own file. Some targets, **including + Dawarich**, treat one uploaded file as a single track and rebuild their own + segmentation from the points; for those you need `files` so the split legs + actually appear as separate tracks. The output filenames are suffixed + (`track-1.geojson`, `track-2.geojson`, …). + ## Supported targets | Type | Service | Accepted formats | diff --git a/justfile b/justfile index 1978d49..3e1ef72 100644 --- a/justfile +++ b/justfile @@ -12,9 +12,13 @@ run-server: run-server-docker: cd server && docker compose up --build +clear-server-data: + rm -f server/data/state.db server/data/state.db-shm server/data/state.db-wal + # Maintenance update-vendor-hash: bash scripts/update-vendor-hash.sh -clear-local-data: - rm ~/.local/share/tracksync/state.db + +dangerously-clear-system-client-data: + rm -f ~/.local/share/tracksync/state.db ~/.local/share/tracksync/state.db-shm ~/.local/share/tracksync/state.db-wal diff --git a/nixos/tracksync.nix b/nixos/tracksync.nix index a805a76..ee48d7b 100644 --- a/nixos/tracksync.nix +++ b/nixos/tracksync.nix @@ -70,6 +70,7 @@ let ${lib.optionalString (cfg.stateDB != null) "--state-db \"${cfg.stateDB}\""}) && RC=0 || RC=$? UPLOADED=$(echo "$SUMMARY" | ${pkgs.jq}/bin/jq -r '.uploaded // 0') + FORWARDED=$(echo "$SUMMARY" | ${pkgs.jq}/bin/jq -r '.forwarded // 0') DUPLICATE=$(echo "$SUMMARY" | ${pkgs.jq}/bin/jq -r '.duplicate // 0') SKIPPED=$(echo "$SUMMARY" | ${pkgs.jq}/bin/jq -r '.skipped // 0') ERRORS=$(echo "$SUMMARY" | ${pkgs.jq}/bin/jq -r '.errors // 0') @@ -80,9 +81,9 @@ let fi if [ "$RC" = 0 ]; then - ${pkgs.libnotify}/bin/notify-send -i emblem-ok "Tracksync" "$UPLOADED uploaded, $DUPLICATE duplicate, $SKIPPED skipped" 2>/dev/null || true + ${pkgs.libnotify}/bin/notify-send -i emblem-ok "Tracksync" "$UPLOADED uploaded ($FORWARDED forwarded), $DUPLICATE duplicate, $SKIPPED skipped" 2>/dev/null || true else - ${pkgs.libnotify}/bin/notify-send -i dialog-error "Tracksync" "$UPLOADED uploaded, $DUPLICATE duplicate, $SKIPPED skipped, $ERRORS failed" 2>/dev/null || true + ${pkgs.libnotify}/bin/notify-send -i dialog-error "Tracksync" "$UPLOADED uploaded ($FORWARDED forwarded), $DUPLICATE duplicate, $SKIPPED skipped, $ERRORS failed" 2>/dev/null || true exit 1 fi ''; diff --git a/server/.env.example b/server/.env.example index 0c9966c..830d9b4 100644 --- a/server/.env.example +++ b/server/.env.example @@ -7,6 +7,10 @@ ACCOUNT__0__DEVICE_ID=my-columbus ACCOUNT__0__TARGET_URL=http://localhost:3000 ACCOUNT__0__API_KEY=your-dawarich-api-key +ACCOUNT__0__MARKERS=C:split +# ACCOUNT__0__SPLIT_MARKER_POSITION=start +# ACCOUNT__0__SPLIT_MODE=tracks + CLIENT__0__ID=my-laptop CLIENT__0__TOKEN=your-client-token CLIENT__0__ALLOWED_DEVICES=my-columbus diff --git a/server/internal/config/config.go b/server/internal/config/config.go index 3537fa8..cd54acf 100644 --- a/server/internal/config/config.go +++ b/server/internal/config/config.go @@ -7,6 +7,7 @@ import ( "strings" "time" + "github.com/Quadrubo/tracksync/server/internal/converter" "github.com/go-playground/validator/v10" "github.com/spf13/viper" ) @@ -21,8 +22,18 @@ func validateConfig(sl validator.StructLevel) { cfg := sl.Current().Interface().(Config) accountDevices := make(map[string]bool) - for _, a := range cfg.Accounts { + for i, a := range cfg.Accounts { accountDevices[a.DeviceID] = true + + if _, err := converter.ParseMarkerRules(a.Markers); err != nil { + sl.ReportError( + a.Markers, + fmt.Sprintf("Accounts[%d].Markers", i), + "Markers", + "valid_marker_rules", + err.Error(), + ) + } } for i, c := range cfg.Clients { @@ -51,11 +62,14 @@ type Config struct { } type Account struct { - DeviceID string `env:"DEVICE_ID" validate:"required"` - TargetType string `env:"TARGET_TYPE" default:"dawarich" validate:"required"` - TargetURL string `env:"TARGET_URL" validate:"required"` - APIKey string `env:"API_KEY" validate:"required_without=APIKeyFile"` - APIKeyFile string `env:"API_KEY_FILE" validate:"required_without=APIKey"` + DeviceID string `env:"DEVICE_ID" validate:"required"` + TargetType string `env:"TARGET_TYPE" default:"dawarich" validate:"required"` + TargetURL string `env:"TARGET_URL" validate:"required"` + APIKey string `env:"API_KEY" validate:"required_without=APIKeyFile"` + APIKeyFile string `env:"API_KEY_FILE" validate:"required_without=APIKey"` + Markers []string `env:"MARKERS"` + SplitMarkerPosition string `env:"SPLIT_MARKER_POSITION" default:"start" validate:"omitempty,oneof=start end"` + SplitMode string `env:"SPLIT_MODE" default:"tracks" validate:"omitempty,oneof=tracks files"` } type Client struct { diff --git a/server/internal/config/config_test.go b/server/internal/config/config_test.go index 2b5cd62..dfbeee8 100644 --- a/server/internal/config/config_test.go +++ b/server/internal/config/config_test.go @@ -41,6 +41,91 @@ func TestParseGroup_DefaultTargetType(t *testing.T) { assert.Equal(t, "dawarich", accounts[0].TargetType) } +func TestParseGroup_SplitConfig(t *testing.T) { + v := viper.New() + v.Set("ACCOUNT__0__DEVICE_ID", "dev-1") + v.Set("ACCOUNT__0__TARGET_URL", "http://localhost:3000") + v.Set("ACCOUNT__0__API_KEY", "key") + v.Set("ACCOUNT__0__MARKERS", "C:split, D:split") + v.Set("ACCOUNT__0__SPLIT_MARKER_POSITION", "end") + + accounts := parseGroup[Account](v, "ACCOUNT", "DEVICE_ID") + + require.Len(t, accounts, 1) + assert.Equal(t, []string{"C:split", "D:split"}, accounts[0].Markers) + assert.Equal(t, "end", accounts[0].SplitMarkerPosition) +} + +func TestParseGroup_SplitDefaults(t *testing.T) { + v := viper.New() + v.Set("ACCOUNT__0__DEVICE_ID", "dev-1") + v.Set("ACCOUNT__0__TARGET_URL", "http://localhost:3000") + v.Set("ACCOUNT__0__API_KEY", "key") + + accounts := parseGroup[Account](v, "ACCOUNT", "DEVICE_ID") + + require.Len(t, accounts, 1) + assert.Nil(t, accounts[0].Markers, "no markers by default") + assert.Equal(t, "start", accounts[0].SplitMarkerPosition) + assert.Equal(t, "tracks", accounts[0].SplitMode) +} + +func TestValidate_SplitMarkerPositionInvalid(t *testing.T) { + cfg := &Config{ + Accounts: []Account{{DeviceID: "d", TargetType: "dawarich", TargetURL: "http://x", APIKey: "k", SplitMarkerPosition: "middle"}}, + Clients: []Client{{ID: "c", Token: "t"}}, + } + assert.Error(t, cfg.validate(), "marker position must be start or end") +} + +func TestValidate_SplitMarkerPositionValid(t *testing.T) { + cfg := &Config{ + Accounts: []Account{{DeviceID: "d", TargetType: "dawarich", TargetURL: "http://x", APIKey: "k", SplitMarkerPosition: "end"}}, + Clients: []Client{{ID: "c", Token: "t"}}, + } + assert.NoError(t, cfg.validate()) +} + +func TestValidate_MarkersValid(t *testing.T) { + cfg := &Config{ + Accounts: []Account{{DeviceID: "d", TargetType: "dawarich", TargetURL: "http://x", APIKey: "k", Markers: []string{"C:split", "D:split"}}}, + Clients: []Client{{ID: "c", Token: "t"}}, + } + assert.NoError(t, cfg.validate()) +} + +func TestValidate_MarkersUnknownFunctionality(t *testing.T) { + cfg := &Config{ + Accounts: []Account{{DeviceID: "d", TargetType: "dawarich", TargetURL: "http://x", APIKey: "k", Markers: []string{"C:bogus"}}}, + Clients: []Client{{ID: "c", Token: "t"}}, + } + assert.Error(t, cfg.validate()) +} + +func TestValidate_MarkersMalformed(t *testing.T) { + cfg := &Config{ + Accounts: []Account{{DeviceID: "d", TargetType: "dawarich", TargetURL: "http://x", APIKey: "k", Markers: []string{"C"}}}, + Clients: []Client{{ID: "c", Token: "t"}}, + } + assert.Error(t, cfg.validate()) +} + +func TestValidate_SplitModeInvalid(t *testing.T) { + cfg := &Config{ + Accounts: []Account{{DeviceID: "d", TargetType: "dawarich", TargetURL: "http://x", APIKey: "k", SplitMode: "zip"}}, + Clients: []Client{{ID: "c", Token: "t"}}, + } + assert.Error(t, cfg.validate(), "split mode must be tracks or files") +} + +func TestValidate_SplitModeValid(t *testing.T) { + cfg := &Config{ + Accounts: []Account{{DeviceID: "d", TargetType: "dawarich", TargetURL: "http://x", APIKey: "k", SplitMode: "files"}}, + Clients: []Client{{ID: "c", Token: "t"}}, + } + assert.NoError(t, cfg.validate()) +} + func TestParseGroup_Clients(t *testing.T) { v := viper.New() v.Set("CLIENT__0__ID", "laptop") diff --git a/server/internal/converter/converter.go b/server/internal/converter/converter.go index bee5de1..78a7996 100644 --- a/server/internal/converter/converter.go +++ b/server/internal/converter/converter.go @@ -5,50 +5,66 @@ import ( "strings" ) -// Convert parses data in sourceFormat, selects the best target format from -// acceptedFormats, and serializes the tracks. Returns converted data, chosen -// format, and the new filename. +// OutputFile is a single converted file produced by Convert. +type OutputFile struct { + Data []byte + Format string + Filename string +} + +// Convert parses data in sourceFormat, applies markers, selects the best target +// format from acceptedFormats, and serializes the tracks into one or more output +// files. // -// When passthrough is true and sourceFormat matches the best target format, -// the original data is returned unchanged. Otherwise, data is always -// re-serialized to produce normalized output. -func Convert(sourceFormat string, data []byte, acceptedFormats []string, originalFilename string, passthrough bool) ([]byte, string, string, error) { +// When passthrough is true and sourceFormat matches the best target format, the +// original data is returned unchanged; otherwise it is re-serialized. With +// markers.SplitMode == "files" each track is serialized into its own file. +func Convert(sourceFormat string, data []byte, acceptedFormats []string, originalFilename string, passthrough bool, markers MarkerOptions) ([]OutputFile, error) { parser, ok := GetParser(sourceFormat) if !ok { - return nil, "", "", fmt.Errorf("no parser for format %q", sourceFormat) + return nil, fmt.Errorf("no parser for format %q", sourceFormat) } tracks, err := parser.Parse(data) if err != nil { - return nil, "", "", fmt.Errorf("parsing %s: %w", sourceFormat, err) + return nil, fmt.Errorf("parsing %s: %w", sourceFormat, err) } + result := applyMarkers(tracks, markers) + // Determine which fields the parsed tracks actually contain. - usedFields := mergeUsedFields(tracks) + usedFields := mergeUsedFields(result.Tracks()) bestFormat := selectBestFormat(usedFields, acceptedFormats) if bestFormat == "" { - return nil, "", "", fmt.Errorf("no serializer available for any accepted format: %v", acceptedFormats) + return nil, fmt.Errorf("no serializer available for any accepted format: %v", acceptedFormats) } - // Passthrough: if explicitly enabled and the best format matches the source, - // return original data without re-serializing. - if passthrough && bestFormat == sourceFormat { - return data, bestFormat, originalFilename, nil + // Passthrough only when nothing was restructured; a split rewrites the tracks. + if passthrough && bestFormat == sourceFormat && !result.Modified { + return []OutputFile{{Data: data, Format: bestFormat, Filename: originalFilename}}, nil } serializer, ok := GetSerializer(bestFormat) if !ok { - return nil, "", "", fmt.Errorf("no serializer for format %q", bestFormat) + return nil, fmt.Errorf("no serializer for format %q", bestFormat) } - out, ext, err := serializer.Serialize(tracks) - if err != nil { - return nil, "", "", fmt.Errorf("serializing to %s: %w", bestFormat, err) + // Suffix filenames only when split across multiple files. + multiFile := len(result.Files) > 1 + out := make([]OutputFile, 0, len(result.Files)) + for i, fileTracks := range result.Files { + fileData, ext, err := serializer.Serialize(fileTracks) + if err != nil { + return nil, fmt.Errorf("serializing to %s: %w", bestFormat, err) + } + filename := replaceExtension(originalFilename, ext) + if multiFile { + filename = indexedFilename(filename, i+1) + } + out = append(out, OutputFile{Data: fileData, Format: bestFormat, Filename: filename}) } - - newFilename := replaceExtension(originalFilename, ext) - return out, bestFormat, newFilename, nil + return out, nil } // mergeUsedFields combines detected fields across all tracks. @@ -69,3 +85,11 @@ func replaceExtension(filename, newExt string) string { } return filename + newExt } + +// indexedFilename inserts "-n" before the extension, e.g. track.geojson -> track-1.geojson. +func indexedFilename(filename string, n int) string { + if idx := strings.LastIndex(filename, "."); idx >= 0 { + return fmt.Sprintf("%s-%d%s", filename[:idx], n, filename[idx:]) + } + return fmt.Sprintf("%s-%d", filename, n) +} diff --git a/server/internal/converter/converter_test.go b/server/internal/converter/converter_test.go index 93aab6f..dc9610e 100644 --- a/server/internal/converter/converter_test.go +++ b/server/internal/converter/converter_test.go @@ -1,6 +1,7 @@ package converter import ( + "strings" "testing" "time" @@ -111,11 +112,51 @@ func TestConvert_GPXPassthrough(t *testing.T) { `) // GPX → GPX with passthrough enabled: return original bytes - data, format, filename, err := Convert("gpx_1.1", gpxData, []string{"gpx_1.1", "geojson"}, "track.gpx", true) + outs, err := Convert("gpx_1.1", gpxData, []string{"gpx_1.1", "geojson"}, "track.gpx", true, MarkerOptions{}) require.NoError(t, err) - assert.Equal(t, "gpx_1.1", format) - assert.Equal(t, "track.gpx", filename) - assert.Equal(t, gpxData, data, "passthrough should return original bytes") + require.Len(t, outs, 1) + assert.Equal(t, "gpx_1.1", outs[0].Format) + assert.Equal(t, "track.gpx", outs[0].Filename) + assert.Equal(t, gpxData, outs[0].Data, "passthrough should return original bytes") +} + +func TestConvert_PassthroughSkippedWhenSplitOccurs(t *testing.T) { + // A split that actually divides the tracks restructures them, so the original + // (unsplit) bytes must not be returned via passthrough even when source and + // target formats match. + csvData := []byte(`INDEX,TAG,DATE,TIME,LATITUDE N/S,LONGITUDE E/W,HEIGHT,SPEED,HEADING +1,T,260417,110529,52.0N,13.0E,38,1.4,333 +2,C,260417,110530,52.1N,13.1E,38,30.0,9 +3,T,260417,110531,52.2N,13.2E,38,30.0,13 +`) + + outs, err := Convert("columbus-csv", csvData, []string{"gpx_1.1"}, "track.csv", true, MarkerOptions{ + Rules: []MarkerRule{{Marker: "C", Functionality: MarkerSplit}}, + SplitMarkerPosition: "start", + }) + require.NoError(t, err) + require.Len(t, outs, 1) + assert.Equal(t, 2, strings.Count(string(outs[0].Data), ""), "split must produce two tracks, not passthrough") +} + +func TestConvert_PassthroughWhenConfiguredSplitDoesNotFire(t *testing.T) { + // A split rule is configured but the data carries no matching marker, so the + // tracks are unchanged and passthrough still returns the original bytes. + gpxData := []byte(` + + + 500 + +`) + + outs, err := Convert("gpx_1.1", gpxData, []string{"gpx_1.1", "geojson"}, "track.gpx", true, MarkerOptions{ + Rules: []MarkerRule{{Marker: "C", Functionality: MarkerSplit}}, + SplitMarkerPosition: "start", + }) + require.NoError(t, err) + require.Len(t, outs, 1) + assert.Equal(t, "gpx_1.1", outs[0].Format) + assert.Equal(t, gpxData, outs[0].Data, "no split fired: passthrough returns original bytes") } func TestConvert_GPXReserialized(t *testing.T) { @@ -127,11 +168,12 @@ func TestConvert_GPXReserialized(t *testing.T) { `) // GPX → GPX without passthrough: re-serialized - data, format, filename, err := Convert("gpx_1.1", gpxData, []string{"gpx_1.1", "geojson"}, "track.gpx", false) + outs, err := Convert("gpx_1.1", gpxData, []string{"gpx_1.1", "geojson"}, "track.gpx", false, MarkerOptions{}) require.NoError(t, err) - assert.Equal(t, "gpx_1.1", format) - assert.Equal(t, "track.gpx", filename) - assert.NotEqual(t, gpxData, data, "should re-serialize, not passthrough") + require.Len(t, outs, 1) + assert.Equal(t, "gpx_1.1", outs[0].Format) + assert.Equal(t, "track.gpx", outs[0].Filename) + assert.NotEqual(t, gpxData, outs[0].Data, "should re-serialize, not passthrough") } func TestConvert_GPXToGeoJSON_WhenSpeedInTarget(t *testing.T) { @@ -145,21 +187,22 @@ func TestConvert_GPXToGeoJSON_WhenSpeedInTarget(t *testing.T) { `) // GeoJSON listed first: wins the tie - data, format, filename, err := Convert("gpx_1.1", gpxData, []string{"geojson", "gpx_1.1"}, "track.gpx", false) + outs, err := Convert("gpx_1.1", gpxData, []string{"geojson", "gpx_1.1"}, "track.gpx", false, MarkerOptions{}) require.NoError(t, err) - assert.Equal(t, "geojson", format) - assert.Equal(t, "track.geojson", filename) - assert.NotEqual(t, gpxData, data) + require.Len(t, outs, 1) + assert.Equal(t, "geojson", outs[0].Format) + assert.Equal(t, "track.geojson", outs[0].Filename) + assert.NotEqual(t, gpxData, outs[0].Data) } func TestConvert_NoParser(t *testing.T) { - _, _, _, err := Convert("unknown", []byte("data"), []string{"gpx_1.1"}, "f.txt", false) + _, err := Convert("unknown", []byte("data"), []string{"gpx_1.1"}, "f.txt", false, MarkerOptions{}) assert.ErrorContains(t, err, "no parser") } func TestConvert_NoSerializer(t *testing.T) { gpxData := []byte(``) - _, _, _, err := Convert("gpx_1.1", gpxData, []string{"columbus-csv"}, "f.gpx", false) + _, err := Convert("gpx_1.1", gpxData, []string{"columbus-csv"}, "f.gpx", false, MarkerOptions{}) assert.ErrorContains(t, err, "no serializer") } @@ -167,9 +210,10 @@ func TestConvert_FilenameExtensionReplaced(t *testing.T) { gpxData := []byte(` `) - _, _, filename, err := Convert("gpx_1.1", gpxData, []string{"geojson"}, "my.track.gpx", false) + outs, err := Convert("gpx_1.1", gpxData, []string{"geojson"}, "my.track.gpx", false, MarkerOptions{}) require.NoError(t, err) - assert.Equal(t, "my.track.geojson", filename) + require.Len(t, outs, 1) + assert.Equal(t, "my.track.geojson", outs[0].Filename) } func TestReplaceExtension(t *testing.T) { @@ -178,6 +222,86 @@ func TestReplaceExtension(t *testing.T) { assert.Equal(t, "noext.geojson", replaceExtension("noext", ".geojson")) } +func TestConvert_SplitColumbusCSVToMultipleGPXTracks(t *testing.T) { + csvData := []byte(`INDEX,TAG,DATE,TIME,LATITUDE N/S,LONGITUDE E/W,HEIGHT,SPEED,HEADING +1,T,260417,110529,52.0N,13.0E,38,1.4,333 +2,C,260417,110530,52.1N,13.1E,38,30.0,9 +3,T,260417,110531,52.2N,13.2E,38,30.0,13 +`) + + // Restrict accepted formats to GPX so the split is observable as elements. + outs, err := Convert("columbus-csv", csvData, []string{"gpx_1.1"}, "track.csv", false, MarkerOptions{Rules: []MarkerRule{{Marker: "C", Functionality: MarkerSplit}}, SplitMarkerPosition: "start"}) + require.NoError(t, err) + require.Len(t, outs, 1, "tracks mode: one file") + assert.Equal(t, "gpx_1.1", outs[0].Format) + assert.Equal(t, 2, strings.Count(string(outs[0].Data), ""), "expected two tracks after split") +} + +func TestConvert_SplitFilesMode_OneFilePerTrack(t *testing.T) { + csvData := []byte(`INDEX,TAG,DATE,TIME,LATITUDE N/S,LONGITUDE E/W,HEIGHT,SPEED,HEADING +1,T,260417,110529,52.0N,13.0E,38,1.4,333 +2,C,260417,110530,52.1N,13.1E,38,30.0,9 +3,T,260417,110531,52.2N,13.2E,38,30.0,13 +`) + + outs, err := Convert("columbus-csv", csvData, []string{"gpx_1.1"}, "track.csv", false, MarkerOptions{ + Rules: []MarkerRule{{Marker: "C", Functionality: MarkerSplit}}, + SplitMarkerPosition: "start", + SplitMode: "files", + }) + require.NoError(t, err) + require.Len(t, outs, 2, "files mode: one file per track") + assert.Equal(t, "track-1.gpx", outs[0].Filename) + assert.Equal(t, "track-2.gpx", outs[1].Filename) + // Each file holds exactly one track. + assert.Equal(t, 1, strings.Count(string(outs[0].Data), "")) + assert.Equal(t, 1, strings.Count(string(outs[1].Data), "")) +} + +func TestConvert_FilesMode_SingleTrackNoSuffix(t *testing.T) { + // files mode but no split: a single track keeps the plain filename. + csvData := []byte(`INDEX,TAG,DATE,TIME,LATITUDE N/S,LONGITUDE E/W,HEIGHT,SPEED,HEADING +1,T,260417,110529,52.0N,13.0E,38,1.4,333 +2,T,260417,110530,52.1N,13.1E,38,1.4,9 +`) + outs, err := Convert("columbus-csv", csvData, []string{"gpx_1.1"}, "track.csv", false, MarkerOptions{SplitMode: "files"}) + require.NoError(t, err) + require.Len(t, outs, 1) + assert.Equal(t, "track.gpx", outs[0].Filename) +} + +func TestConvert_FilesMode_DoesNotSplitSourceNativeTracks(t *testing.T) { + // Two source-native tracks and a split rule configured, but no point carries + // the marker: files mode must not split native tracks into separate files when + // no split actually fired. + gpxData := []byte(`` + + `` + + `` + + ``) + + outs, err := Convert("gpx_1.1", gpxData, []string{"gpx_1.1"}, "track.gpx", false, MarkerOptions{ + Rules: []MarkerRule{{Marker: "C", Functionality: MarkerSplit}}, + SplitMarkerPosition: "start", + SplitMode: "files", + }) + require.NoError(t, err) + require.Len(t, outs, 1, "no split fired: a single file regardless of files mode") + assert.Equal(t, "track.gpx", outs[0].Filename) + assert.Equal(t, 2, strings.Count(string(outs[0].Data), ""), "both source tracks kept in one file") +} + +func TestConvert_NoSplitColumbusCSVSingleGPXTrack(t *testing.T) { + csvData := []byte(`INDEX,TAG,DATE,TIME,LATITUDE N/S,LONGITUDE E/W,HEIGHT,SPEED,HEADING +1,T,260417,110529,52.0N,13.0E,38,1.4,333 +2,C,260417,110530,52.1N,13.1E,38,30.0,9 +`) + + outs, err := Convert("columbus-csv", csvData, []string{"gpx_1.1"}, "track.csv", false, MarkerOptions{}) + require.NoError(t, err) + require.Len(t, outs, 1) + assert.Equal(t, 1, strings.Count(string(outs[0].Data), ""), "no split tags configured: single track") +} + func TestMergeUsedFields(t *testing.T) { spd := 5.0 ele := 100.0 diff --git a/server/internal/converter/csv_parser.go b/server/internal/converter/csv_parser.go index 6106f53..09e18c7 100644 --- a/server/internal/converter/csv_parser.go +++ b/server/internal/converter/csv_parser.go @@ -17,14 +17,15 @@ func init() { // Format: INDEX,TAG,DATE,TIME,LATITUDE N/S,LONGITUDE E/W,HEIGHT,SPEED,HEADING // Example: 1,T,260417,110529,52.4194759N,13.3076437E,62,1.4,333 // -// - DATE is yymmdd (UTC) -// - TIME is hhmmss (UTC) -// - LATITUDE: decimal degrees with N/S suffix -// - LONGITUDE: decimal degrees with E/W suffix -// - HEIGHT: meters -// - SPEED: km/h -// - HEADING: degrees -// - TAG: T=trackpoint, C=POI, D=second POI, G=wake-up point (ignored, all rows are treated as trackpoints) +// - DATE is yymmdd (UTC) +// - TIME is hhmmss (UTC) +// - LATITUDE: decimal degrees with N/S suffix +// - LONGITUDE: decimal degrees with E/W suffix +// - HEIGHT: meters +// - SPEED: km/h +// - HEADING: degrees +// - TAG: T=trackpoint, C=POI, D=second POI, G=wake-up point. Every row is kept +// as a trackpoint; tags other than T are recorded on Point.Marker. type csvParser struct{} func (p *csvParser) Parse(data []byte) ([]Track, error) { @@ -52,6 +53,11 @@ func (p *csvParser) Parse(data []byte) ([]Track, error) { return nil, err } + // Any tag other than T (trackpoint) is a marker. + if tag := record[1]; tag != "" && tag != "T" { + point.Marker = tag + } + seg.Points = append(seg.Points, point) } diff --git a/server/internal/converter/csv_parser_test.go b/server/internal/converter/csv_parser_test.go index 7788811..8163186 100644 --- a/server/internal/converter/csv_parser_test.go +++ b/server/internal/converter/csv_parser_test.go @@ -132,3 +132,20 @@ func TestParseDateTime(t *testing.T) { require.NoError(t, err) assert.Equal(t, time.Date(2026, 4, 17, 11, 5, 29, 0, time.UTC), dt) } + +func TestCSVParser_RecordsMarkers(t *testing.T) { + // Non-trackpoint tags become Point.Marker; plain T trackpoints stay empty. + data := `INDEX,TAG,DATE,TIME,LATITUDE N/S,LONGITUDE E/W,HEIGHT,SPEED,HEADING +1,T,260417,110529,52.0N,13.0E,38,1.4,333 +2,C,260417,110530,52.1N,13.1E,38,30.0,9 +3,G,260417,110531,52.2N,13.2E,38,30.0,13 +` + tracks, err := (&csvParser{}).Parse([]byte(data)) + require.NoError(t, err) + require.Len(t, tracks, 1) + pts := tracks[0].Segments[0].Points + require.Len(t, pts, 3) + assert.Equal(t, "", pts[0].Marker, "T trackpoint has no marker") + assert.Equal(t, "C", pts[1].Marker) + assert.Equal(t, "G", pts[2].Marker) +} diff --git a/server/internal/converter/gpx_parser.go b/server/internal/converter/gpx_parser.go index 7fce669..cff7fe1 100644 --- a/server/internal/converter/gpx_parser.go +++ b/server/internal/converter/gpx_parser.go @@ -22,8 +22,8 @@ type gpxFile struct { } type gpxTrk struct { - Name string `xml:"name"` - Segments []gpxTrkSeg `xml:"trkseg"` + Name string `xml:"name"` + Segments []gpxTrkSeg `xml:"trkseg"` } type gpxTrkSeg struct { @@ -46,7 +46,7 @@ type gpxTrkPt struct { // GPX 1.1 uses xsd:dateTime which allows fractional seconds and // timestamps without timezone offsets (interpreted as UTC). var gpxTimeFormats = []string{ - time.RFC3339Nano, // 2025-01-15T08:30:00.123Z / +01:00 + time.RFC3339Nano, // 2025-01-15T08:30:00.123Z / +01:00 "2006-01-02T15:04:05.999999999", // fractional seconds, no offset (UTC) "2006-01-02T15:04:05", // no fractions, no offset (UTC) } diff --git a/server/internal/converter/gpx_serializer.go b/server/internal/converter/gpx_serializer.go index 3f9e62b..05785ca 100644 --- a/server/internal/converter/gpx_serializer.go +++ b/server/internal/converter/gpx_serializer.go @@ -12,16 +12,16 @@ func init() { type gpxSerializer struct{} type gpxOutput struct { - XMLName xml.Name `xml:"gpx"` - Version string `xml:"version,attr"` - Creator string `xml:"creator,attr"` - XMLNS string `xml:"xmlns,attr"` - Tracks []gpxTrkOut `xml:"trk"` + XMLName xml.Name `xml:"gpx"` + Version string `xml:"version,attr"` + Creator string `xml:"creator,attr"` + XMLNS string `xml:"xmlns,attr"` + Tracks []gpxTrkOut `xml:"trk"` } type gpxTrkOut struct { - Name string `xml:"name,omitempty"` - Segments []gpxTrkSegOut `xml:"trkseg"` + Name string `xml:"name,omitempty"` + Segments []gpxTrkSegOut `xml:"trkseg"` } type gpxTrkSegOut struct { diff --git a/server/internal/converter/marker.go b/server/internal/converter/marker.go new file mode 100644 index 0000000..e1d4f11 --- /dev/null +++ b/server/internal/converter/marker.go @@ -0,0 +1,121 @@ +package converter + +import ( + "fmt" + "strings" +) + +// MarkerFunctionality is the behavior applied to points carrying a given marker. +type MarkerFunctionality string + +const ( + // MarkerSplit starts a new track at the marked point. + MarkerSplit MarkerFunctionality = "split" +) + +// validFunctionalities is the set of recognised marker functionalities. +var validFunctionalities = map[MarkerFunctionality]bool{ + MarkerSplit: true, +} + +// MarkerRule assigns a functionality to a point marker value, e.g. {C, split}. +type MarkerRule struct { + Marker string + Functionality MarkerFunctionality +} + +// MarkerOptions describes how parsed tracks are post-processed based on point markers. +type MarkerOptions struct { + Rules []MarkerRule + // SplitMarkerPosition places a split point at the "start" (default) of the + // new track or the "end" of the previous one. + SplitMarkerPosition string + // SplitMode emits split tracks as one file ("tracks", default) or one file + // per track ("files"). + SplitMode string +} + +// ParseMarkerRules parses "marker:functionality" entries (e.g. "C:split") into MarkerRules. +func ParseMarkerRules(entries []string) ([]MarkerRule, error) { + var rules []MarkerRule + for _, e := range entries { + marker, fn, ok := strings.Cut(e, ":") + marker = strings.TrimSpace(marker) + fn = strings.TrimSpace(fn) + if !ok || marker == "" || fn == "" { + return nil, fmt.Errorf("invalid marker rule %q: expected \"marker:functionality\"", e) + } + f := MarkerFunctionality(fn) + if !validFunctionalities[f] { + return nil, fmt.Errorf("unknown marker functionality %q in %q", fn, e) + } + rules = append(rules, MarkerRule{Marker: marker, Functionality: f}) + } + return rules, nil +} + +// MarkerResult is the outcome of applying marker rules to parsed tracks. +type MarkerResult struct { + // Files partitions tracks into output files, one serialized file each. + Files [][]Track + // Modified is true when a split restructured the tracks, ruling out passthrough. + Modified bool +} + +// Tracks returns every track across all files, in order. +func (r MarkerResult) Tracks() []Track { + var all []Track + for _, f := range r.Files { + all = append(all, f...) + } + return all +} + +// applyMarkers splits tracks per the marker rules and partitions them into +// output files. +func applyMarkers(tracks []Track, opts MarkerOptions) MarkerResult { + markerSet := make(map[string]bool) + for _, r := range opts.Rules { + if r.Functionality == MarkerSplit { + markerSet[r.Marker] = true + } + } + + if len(markerSet) == 0 { + return MarkerResult{Files: [][]Track{tracks}} + } + + markerEndsTrack := opts.SplitMarkerPosition == "end" + + legsPerTrack := make([][]Track, 0, len(tracks)) + modified := false + for _, t := range tracks { + legs := splitTrack(t, markerSet, markerEndsTrack) + if len(legs) > 1 { + modified = true + } + legsPerTrack = append(legsPerTrack, legs) + } + + // Nothing divided: keep the original tracks so passthrough stays possible. + if !modified { + return MarkerResult{Files: [][]Track{tracks}} + } + + if opts.SplitMode != "files" { + var all []Track + for _, legs := range legsPerTrack { + all = append(all, legs...) + } + return MarkerResult{Files: [][]Track{all}, Modified: true} + } + + // files mode: one file per leg, in source order. + var files [][]Track + for _, legs := range legsPerTrack { + for _, leg := range legs { + files = append(files, []Track{leg}) + } + } + return MarkerResult{Files: files, Modified: true} +} diff --git a/server/internal/converter/marker_test.go b/server/internal/converter/marker_test.go new file mode 100644 index 0000000..7eed2c9 --- /dev/null +++ b/server/internal/converter/marker_test.go @@ -0,0 +1,140 @@ +package converter + +import ( + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func TestParseMarkerRules_Valid(t *testing.T) { + rules, err := ParseMarkerRules([]string{"C:split", " D : split "}) + require.NoError(t, err) + require.Len(t, rules, 2) + assert.Equal(t, MarkerRule{Marker: "C", Functionality: MarkerSplit}, rules[0]) + assert.Equal(t, MarkerRule{Marker: "D", Functionality: MarkerSplit}, rules[1]) +} + +func TestParseMarkerRules_Empty(t *testing.T) { + rules, err := ParseMarkerRules(nil) + require.NoError(t, err) + assert.Empty(t, rules) +} + +func TestParseMarkerRules_MissingFunctionality(t *testing.T) { + _, err := ParseMarkerRules([]string{"C"}) + assert.ErrorContains(t, err, "marker:functionality") +} + +func TestParseMarkerRules_UnknownFunctionality(t *testing.T) { + _, err := ParseMarkerRules([]string{"C:teleport"}) + assert.ErrorContains(t, err, "unknown marker functionality") +} + +func TestApplyMarkers_SplitFunctionality(t *testing.T) { + in := []Track{markedTrack(mp{1, ""}, mp{2, "C"}, mp{3, ""})} + res := applyMarkers(in, MarkerOptions{ + Rules: []MarkerRule{{Marker: "C", Functionality: MarkerSplit}}, + SplitMarkerPosition: "start", + }) + assert.True(t, res.Modified) + // tracks mode (default): both resulting tracks land in a single file. + require.Len(t, res.Files, 1) + tracks := res.Tracks() + require.Len(t, tracks, 2) + assert.Len(t, tracks[0].Segments[0].Points, 1) + assert.Len(t, tracks[1].Segments[0].Points, 2) +} + +func TestApplyMarkers_SplitFilesMode(t *testing.T) { + // files mode: each resulting track becomes its own output file. + in := []Track{markedTrack(mp{1, ""}, mp{2, "C"}, mp{3, ""})} + res := applyMarkers(in, MarkerOptions{ + Rules: []MarkerRule{{Marker: "C", Functionality: MarkerSplit}}, + SplitMarkerPosition: "start", + SplitMode: "files", + }) + assert.True(t, res.Modified) + require.Len(t, res.Files, 2) + assert.Len(t, res.Files[0], 1) + assert.Len(t, res.Files[1], 1) +} + +func TestApplyMarkers_NoRulesIsNoOp(t *testing.T) { + in := []Track{markedTrack(mp{1, ""}, mp{2, "C"}, mp{3, ""})} + res := applyMarkers(in, MarkerOptions{}) + assert.False(t, res.Modified) + require.Len(t, res.Files, 1) + tracks := res.Tracks() + require.Len(t, tracks, 1) + assert.Len(t, tracks[0].Segments[0].Points, 3) +} + +func TestApplyMarkers_ConfiguredButNoMatchingMarker(t *testing.T) { + // A split rule is configured but no point carries the marker: no split. + in := []Track{markedTrack(mp{1, ""}, mp{2, "G"}, mp{3, ""})} + res := applyMarkers(in, MarkerOptions{ + Rules: []MarkerRule{{Marker: "C", Functionality: MarkerSplit}}, + SplitMarkerPosition: "start", + }) + assert.False(t, res.Modified) + require.Len(t, res.Files, 1) + tracks := res.Tracks() + require.Len(t, tracks, 1) + assert.Len(t, tracks[0].Segments[0].Points, 3) +} + +func TestApplyMarkers_FilesMode_PreservesSourceOrder(t *testing.T) { + // One file per track, in source order: an undivided track keeps its slot and a + // divided track expands into consecutive files where it sits. + in := []Track{ + markedTrack(mp{10, ""}), // untouched, before the split + markedTrack(mp{1, ""}, mp{2, "C"}, mp{3, ""}), // divided into two legs + markedTrack(mp{20, ""}), // untouched, after the split + } + res := applyMarkers(in, MarkerOptions{ + Rules: []MarkerRule{{Marker: "C", Functionality: MarkerSplit}}, + SplitMarkerPosition: "start", + SplitMode: "files", + }) + assert.True(t, res.Modified) + require.Len(t, res.Files, 4, "one file per resulting track") + // Each file holds exactly one track, ordered as in the source. + for i, f := range res.Files { + require.Lenf(t, f, 1, "file %d holds a single track", i) + } + assert.InDelta(t, 10, res.Files[0][0].Segments[0].Points[0].Lat, 0, "untouched track keeps its leading slot") + assert.InDelta(t, 1, res.Files[1][0].Segments[0].Points[0].Lat, 0, "first split leg") + assert.InDelta(t, 2, res.Files[2][0].Segments[0].Points[0].Lat, 0, "second split leg") + assert.InDelta(t, 20, res.Files[3][0].Segments[0].Points[0].Lat, 0, "trailing untouched track keeps its slot") +} + +func TestApplyMarkers_EmptySiblingDoesNotMaskSplit(t *testing.T) { + // A divided track alongside an empty source track: the empty track must not + // hide that a split happened. + in := []Track{ + markedTrack(mp{1, ""}, mp{2, "C"}, mp{3, ""}), + {}, // empty track: no segments, no points + } + res := applyMarkers(in, MarkerOptions{ + Rules: []MarkerRule{{Marker: "C", Functionality: MarkerSplit}}, + SplitMarkerPosition: "start", + }) + assert.True(t, res.Modified, "an empty source track must not mask a real split") + require.Len(t, res.Files, 1) // tracks mode + assert.Len(t, res.Tracks(), 2, "two legs, empty track contributes nothing") +} + +func TestApplyMarkers_OnlyConfiguredMarkerSplits(t *testing.T) { + // "C" is mapped to split; "G" carries no rule and is left in place. + in := []Track{markedTrack(mp{1, ""}, mp{2, "G"}, mp{3, "C"}, mp{4, ""})} + res := applyMarkers(in, MarkerOptions{ + Rules: []MarkerRule{{Marker: "C", Functionality: MarkerSplit}}, + SplitMarkerPosition: "start", + }) + assert.True(t, res.Modified) + tracks := res.Tracks() + require.Len(t, tracks, 2) + assert.Len(t, tracks[0].Segments[0].Points, 2) // pt1 + G pt2 + assert.Len(t, tracks[1].Segments[0].Points, 2) // C pt3 + pt4 +} diff --git a/server/internal/converter/split.go b/server/internal/converter/split.go new file mode 100644 index 0000000..db684ed --- /dev/null +++ b/server/internal/converter/split.go @@ -0,0 +1,47 @@ +package converter + +// splitTrack divides one track at points whose Marker is in markerSet, returning +// its legs (one leg means it was not divided). markerEndsTrack keeps the marked +// point on the previous leg ("end") rather than starting the new one ("start"). +// Segment boundaries and the track name are preserved. +func splitTrack(track Track, markerSet map[string]bool, markerEndsTrack bool) []Track { + var legs []Track + cur := Track{Name: track.Name} + seg := Segment{} + + flushSeg := func() { + if len(seg.Points) > 0 { + cur.Segments = append(cur.Segments, seg) + } + seg = Segment{} + } + flushLeg := func() { + flushSeg() + if len(cur.Segments) > 0 { + legs = append(legs, cur) + } + cur = Track{Name: track.Name} + } + + for _, origSeg := range track.Segments { + for _, pt := range origSeg.Points { + isSplit := pt.Marker != "" && markerSet[pt.Marker] + if isSplit && markerEndsTrack { + // Marker ends the current leg; the next point starts a new one. + seg.Points = append(seg.Points, pt) + flushLeg() + continue + } + if isSplit { + // Marker begins a new leg. + flushLeg() + } + seg.Points = append(seg.Points, pt) + } + // Preserve the original recording gap as a segment boundary. + flushSeg() + } + flushLeg() + + return legs +} diff --git a/server/internal/converter/split_test.go b/server/internal/converter/split_test.go new file mode 100644 index 0000000..5b42937 --- /dev/null +++ b/server/internal/converter/split_test.go @@ -0,0 +1,101 @@ +package converter + +import ( + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// markedTrack builds a single-segment track from points described as +// (lat, marker) pairs. Marker "" is an ordinary trackpoint. +func markedTrack(pts ...struct { + lat float64 + marker string +}) Track { + seg := Segment{} + for _, p := range pts { + seg.Points = append(seg.Points, Point{Lat: p.lat, Marker: p.marker}) + } + return Track{Segments: []Segment{seg}} +} + +type mp = struct { + lat float64 + marker string +} + +// markerSet builds the lookup splitTrack expects from a list of marker values. +func markerSet(markers ...string) map[string]bool { + set := make(map[string]bool, len(markers)) + for _, m := range markers { + set[m] = true + } + return set +} + +func TestSplitTrack_NoMarkersIsNoOp(t *testing.T) { + in := markedTrack(mp{1, ""}, mp{2, "C"}, mp{3, ""}) + legs := splitTrack(in, markerSet(), false) + require.Len(t, legs, 1) + assert.Len(t, legs[0].Segments[0].Points, 3) +} + +func TestSplitTrack_MarkerStart(t *testing.T) { + // walk, POI(C), bus, POI(C), walk; G is present but not configured. + in := markedTrack(mp{1, ""}, mp{2, "C"}, mp{3, ""}, mp{4, "C"}, mp{5, "G"}, mp{6, ""}) + legs := splitTrack(in, markerSet("C"), false) + require.Len(t, legs, 3) + assert.Len(t, legs[0].Segments[0].Points, 1) // walk: pt1 + require.Len(t, legs[1].Segments[0].Points, 2) // bus: marker pt2 + pt3 + assert.InDelta(t, 2, legs[1].Segments[0].Points[0].Lat, 0) // marker begins new leg + require.Len(t, legs[2].Segments[0].Points, 3) // walk: marker pt4 + G pt5 + pt6 + assert.InDelta(t, 4, legs[2].Segments[0].Points[0].Lat, 0) +} + +func TestSplitTrack_MarkerEnd(t *testing.T) { + in := markedTrack(mp{1, ""}, mp{2, "C"}, mp{3, ""}, mp{4, "C"}, mp{5, ""}) + legs := splitTrack(in, markerSet("C"), true) + require.Len(t, legs, 3) + require.Len(t, legs[0].Segments[0].Points, 2) // walk pt1 + marker pt2 (ends leg) + assert.InDelta(t, 2, legs[0].Segments[0].Points[1].Lat, 0) // marker stays at end + assert.Len(t, legs[1].Segments[0].Points, 2) // pt3 + marker pt4 + assert.Len(t, legs[2].Segments[0].Points, 1) // pt5 +} + +func TestSplitTrack_MultipleMarkers(t *testing.T) { + in := markedTrack(mp{1, ""}, mp{2, "C"}, mp{3, ""}, mp{4, "G"}, mp{5, ""}) + legs := splitTrack(in, markerSet("C", "G"), false) + require.Len(t, legs, 3) +} + +func TestSplitTrack_DropsEmptyLeadingMarker(t *testing.T) { + // A marker as the very first point must not yield an empty leading leg, so the + // track is not actually divided. + in := markedTrack(mp{1, "C"}, mp{2, ""}) + legs := splitTrack(in, markerSet("C"), false) + require.Len(t, legs, 1) + assert.Len(t, legs[0].Segments[0].Points, 2) +} + +func TestSplitTrack_PreservesNameAndSegments(t *testing.T) { + // Two recording segments; a marker in the second splits it into a new leg. + in := Track{ + Name: "trip", + Segments: []Segment{ + {Points: []Point{{Lat: 1}, {Lat: 2}}}, + {Points: []Point{{Lat: 3}, {Lat: 4, Marker: "C"}, {Lat: 5}}}, + }, + } + legs := splitTrack(in, markerSet("C"), false) + require.Len(t, legs, 2) + // First leg keeps both leading segments (gap preserved), name carried over. + assert.Equal(t, "trip", legs[0].Name) + require.Len(t, legs[0].Segments, 2) + assert.Len(t, legs[0].Segments[0].Points, 2) + assert.Len(t, legs[0].Segments[1].Points, 1) // pt3 before the marker + // Second leg starts at the marker. + assert.Equal(t, "trip", legs[1].Name) + require.Len(t, legs[1].Segments, 1) + assert.Len(t, legs[1].Segments[0].Points, 2) // marker pt4 + pt5 +} diff --git a/server/internal/converter/track.go b/server/internal/converter/track.go index d9fc35e..a0d7507 100644 --- a/server/internal/converter/track.go +++ b/server/internal/converter/track.go @@ -31,4 +31,8 @@ type Point struct { VDOP *float64 PDOP *float64 Fix *string // "none", "2d", "3d", "dgps", "pps" + + // Marker flags a special point such as a POI (e.g. a Columbus CSV tag "C"). + // Empty means an ordinary trackpoint. + Marker string } diff --git a/server/internal/migrations/002_received_forwarded_files.sql b/server/internal/migrations/002_received_forwarded_files.sql new file mode 100644 index 0000000..b90ecf5 --- /dev/null +++ b/server/internal/migrations/002_received_forwarded_files.sql @@ -0,0 +1,35 @@ +-- A payload from a client. completed_at is set once all derived files forward. +CREATE TABLE received_files ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + device_id TEXT NOT NULL, + client_id TEXT NOT NULL, + sha256 TEXT NOT NULL, + filename TEXT NOT NULL, + received_at DATETIME DEFAULT CURRENT_TIMESTAMP, + completed_at DATETIME, + UNIQUE (device_id, sha256) +); + +-- One file sent to the target; a received file can produce several. Identity is +-- (received file, filename), so a retry is deduped but distinct output files are +-- not, even with identical bytes. +CREATE TABLE forwarded_files ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + received_file_id INTEGER NOT NULL REFERENCES received_files(id) ON DELETE CASCADE, + sha256 TEXT NOT NULL, + filename TEXT NOT NULL, + forwarded_at DATETIME DEFAULT CURRENT_TIMESTAMP, + UNIQUE (received_file_id, filename) +); + +-- Carry over prior uploads, marked completed and seeded with the source file as +-- their single forwarded file. +INSERT INTO received_files (device_id, client_id, sha256, filename, received_at, completed_at) +SELECT device_id, client_id, sha256, filename, uploaded_at, uploaded_at FROM uploads; + +INSERT INTO forwarded_files (received_file_id, sha256, filename, forwarded_at) +SELECT rf.id, u.sha256, u.filename, u.uploaded_at +FROM uploads u +JOIN received_files rf ON rf.device_id = u.device_id AND rf.sha256 = u.sha256; + +DROP TABLE uploads; diff --git a/server/internal/migrations/migrations_test.go b/server/internal/migrations/migrations_test.go new file mode 100644 index 0000000..9d0ed9d --- /dev/null +++ b/server/internal/migrations/migrations_test.go @@ -0,0 +1,89 @@ +package migrations + +import ( + "database/sql" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + _ "modernc.org/sqlite" +) + +func execFile(t *testing.T, db *sql.DB, name string) { + t.Helper() + data, err := files.ReadFile(name) + require.NoError(t, err) + _, err = db.Exec(string(data)) + require.NoError(t, err) +} + +func TestRun_CreatesNewTables(t *testing.T) { + db, err := sql.Open("sqlite", ":memory:") + require.NoError(t, err) + t.Cleanup(func() { _ = db.Close() }) + + require.NoError(t, Run(db)) + + _, err = db.Exec("SELECT 1 FROM received_files") + assert.NoError(t, err, "received_files should exist") + _, err = db.Exec("SELECT 1 FROM forwarded_files") + assert.NoError(t, err, "forwarded_files should exist") + _, err = db.Exec("SELECT 1 FROM uploads") + assert.Error(t, err, "uploads should be dropped") +} + +func TestMigration002_ForwardedFilesDedupByFilename(t *testing.T) { + db, err := sql.Open("sqlite", ":memory:") + require.NoError(t, err) + t.Cleanup(func() { _ = db.Close() }) + + require.NoError(t, Run(db)) + + _, err = db.Exec("INSERT INTO received_files (id, device_id, client_id, sha256, filename) VALUES (1, 'dev', 'cli', 'src', 't.csv')") + require.NoError(t, err) + + insert := func(sha, filename string) int64 { + res, err := db.Exec("INSERT OR IGNORE INTO forwarded_files (received_file_id, sha256, filename, forwarded_at) VALUES (1, ?, ?, CURRENT_TIMESTAMP)", sha, filename) + require.NoError(t, err) + n, _ := res.RowsAffected() + return n + } + + // Identical bytes under different filenames both forward (strict per-file). + assert.Equal(t, int64(1), insert("samehash", "track-1.gpx")) + assert.Equal(t, int64(1), insert("samehash", "track-2.gpx")) + // A repeat of the same filename is recognised as already forwarded. + assert.Equal(t, int64(0), insert("samehash", "track-1.gpx")) + assert.Equal(t, int64(0), insert("differenthash", "track-1.gpx"), "filename identifies a forwarded file, regardless of content") +} + +func TestMigration002_CarriesOverExistingUploads(t *testing.T) { + db, err := sql.Open("sqlite", ":memory:") + require.NoError(t, err) + t.Cleanup(func() { _ = db.Close() }) + + // Simulate a DB already on the old schema with a processed upload. + execFile(t, db, "001_initial.sql") + _, err = db.Exec("INSERT INTO uploads (sha256, device_id, client_id, filename) VALUES ('abc', 'dev', 'cli', 't.gpx')") + require.NoError(t, err) + + execFile(t, db, "002_received_forwarded_files.sql") + + var receivedID int64 + var deviceID, sha string + var completed bool + err = db.QueryRow("SELECT id, device_id, sha256, completed_at IS NOT NULL FROM received_files").Scan(&receivedID, &deviceID, &sha, &completed) + require.NoError(t, err) + assert.Equal(t, "dev", deviceID) + assert.Equal(t, "abc", sha) + assert.True(t, completed, "carried-over upload must be marked completed so it is not re-forwarded") + + // A matching forwarded_files row is reconstructed and linked to the received file. + var linkedID int64 + var fwdSha, fwdName string + err = db.QueryRow("SELECT received_file_id, sha256, filename FROM forwarded_files").Scan(&linkedID, &fwdSha, &fwdName) + require.NoError(t, err) + assert.Equal(t, receivedID, linkedID, "forwarded file must reference the carried-over received file") + assert.Equal(t, "abc", fwdSha) + assert.Equal(t, "t.gpx", fwdName) +} diff --git a/server/internal/server/server.go b/server/internal/server/server.go index 9c0ee57..0aa4bd7 100644 --- a/server/internal/server/server.go +++ b/server/internal/server/server.go @@ -9,6 +9,7 @@ import ( "io" "log/slog" "net/http" + "strconv" "strings" "time" @@ -19,13 +20,24 @@ import ( ) type Server struct { - cfg *config.Config - db *sql.DB - targets map[string]target.Target // device ID -> target + cfg *config.Config + db *sql.DB + targets map[string]target.Target // device ID -> target + markerOpts map[string]converter.MarkerOptions // device ID -> marker config } func New(cfg *config.Config, db *sql.DB, targets map[string]target.Target) *Server { - return &Server{cfg: cfg, db: db, targets: targets} + markerOpts := make(map[string]converter.MarkerOptions, len(cfg.Accounts)) + for _, a := range cfg.Accounts { + // Validated at config load. + rules, _ := converter.ParseMarkerRules(a.Markers) + markerOpts[a.DeviceID] = converter.MarkerOptions{ + Rules: rules, + SplitMarkerPosition: a.SplitMarkerPosition, + SplitMode: a.SplitMode, + } + } + return &Server{cfg: cfg, db: db, targets: targets, markerOpts: markerOpts} } func InitDB(db *sql.DB) error { @@ -117,39 +129,31 @@ func (s *Server) handleUpload(w http.ResponseWriter, r *http.Request) { return } - // Deduplicate - claim the hash in a transaction before forwarding to - // the target. On failure the transaction is rolled back so retries work. + // Fast path: skip a fully processed payload without re-converting. hash := sha256.Sum256(data) hashHex := hex.EncodeToString(hash[:]) - tx, err := s.db.Begin() - if err != nil { - slog.Error("database error", "error", err) - http.Error(w, "internal error", http.StatusInternalServerError) - return - } - defer func() { _ = tx.Rollback() }() - - result, err := tx.Exec( - "INSERT OR IGNORE INTO uploads (sha256, device_id, client_id, filename, uploaded_at) VALUES (?, ?, ?, ?, ?)", - hashHex, deviceID, client.ID, header.Filename, time.Now().UTC(), - ) - if err != nil { + var receivedFileID int64 + var completed bool + switch err := s.db.QueryRow( + "SELECT id, completed_at IS NOT NULL FROM received_files WHERE device_id = ? AND sha256 = ?", + deviceID, hashHex, + ).Scan(&receivedFileID, &completed); { + case err == sql.ErrNoRows: + // New payload for this device. + case err != nil: slog.Error("database error", "error", err) http.Error(w, "internal error", http.StatusInternalServerError) return - } - - rows, _ := result.RowsAffected() - if rows == 0 { + case completed: slog.Info("duplicate", "file", header.Filename, "sha256", hashHex[:12], "client", client.ID) w.WriteHeader(http.StatusOK) _, _ = fmt.Fprintln(w, "duplicate") return } - convertedData, chosenFormat, newFilename, err := converter.Convert( - sourceFormat, data, t.AcceptedFormats(), header.Filename, s.cfg.PassthroughConversion, + outputs, err := converter.Convert( + sourceFormat, data, t.AcceptedFormats(), header.Filename, s.cfg.PassthroughConversion, s.markerOpts[deviceID], ) if err != nil { slog.Error("conversion failed", @@ -160,29 +164,113 @@ func (s *Server) handleUpload(w http.ResponseWriter, r *http.Request) { http.Error(w, "format conversion failed", http.StatusUnprocessableEntity) return } - - if chosenFormat != sourceFormat { + if outputs[0].Format != sourceFormat { slog.Info("converted", "file", header.Filename, "from", sourceFormat, - "to", chosenFormat, - "newFile", newFilename, + "to", outputs[0].Format, + "files", len(outputs), ) } - if err := t.Send(r.Context(), newFilename, convertedData); err != nil { - slog.Error("target failed", "target", t.Type(), "source", header.Filename, "file", newFilename, "error", err) - http.Error(w, "target forward failed", http.StatusBadGateway) - return + // Reuse the row from an earlier incomplete attempt, or claim a new one. + if receivedFileID == 0 { + receivedFileID, err = s.claimReceivedFile(deviceID, client.ID, hashHex, header.Filename) + if err != nil { + slog.Error("database error", "error", err) + http.Error(w, "internal error", http.StatusInternalServerError) + return + } } - if err := tx.Commit(); err != nil { - slog.Error("failed to commit upload", "error", err) + // Claim each file before sending and commit only after; a retry forwards + // only the files still missing. + sent := 0 + skipped := 0 + for _, out := range outputs { + outHash := sha256.Sum256(out.Data) + outHashHex := hex.EncodeToString(outHash[:]) + + tx, err := s.db.Begin() + if err != nil { + slog.Error("database error", "error", err) + http.Error(w, "internal error", http.StatusInternalServerError) + return + } + + result, err := tx.Exec( + "INSERT OR IGNORE INTO forwarded_files (received_file_id, sha256, filename, forwarded_at) VALUES (?, ?, ?, ?)", + receivedFileID, outHashHex, out.Filename, time.Now().UTC(), + ) + if err != nil { + _ = tx.Rollback() + slog.Error("database error", "error", err) + http.Error(w, "internal error", http.StatusInternalServerError) + return + } + if rows, _ := result.RowsAffected(); rows == 0 { + // Already forwarded. + _ = tx.Rollback() + skipped++ + continue + } + + if err := t.Send(r.Context(), out.Filename, out.Data); err != nil { + _ = tx.Rollback() + slog.Error("target failed", "target", t.Type(), "source", header.Filename, "file", out.Filename, "error", err) + http.Error(w, "target forward failed", http.StatusBadGateway) + return + } + if err := tx.Commit(); err != nil { + slog.Error("failed to commit forwarded file", "file", out.Filename, "error", err) + http.Error(w, "internal error", http.StatusInternalServerError) + return + } + sent++ + } + + if _, err := s.db.Exec("UPDATE received_files SET completed_at = ? WHERE id = ?", time.Now().UTC(), receivedFileID); err != nil { + slog.Error("database error", "error", err) http.Error(w, "internal error", http.StatusInternalServerError) return } - slog.Info("uploaded", "source", header.Filename, "file", newFilename, "sha256", hashHex[:12], "client", client.ID, "device", deviceID) + if sent == 0 { + // Every output file was already forwarded. + slog.Info("duplicate", "file", header.Filename, "files", len(outputs), "sha256", hashHex[:12], "client", client.ID) + w.WriteHeader(http.StatusOK) + _, _ = fmt.Fprintln(w, "duplicate") + return + } + + slog.Info("uploaded", + "source", header.Filename, + "files", len(outputs), + "sent", sent, + "duplicate", skipped, + "sha256", hashHex[:12], + "client", client.ID, + "device", deviceID, + ) + w.Header().Set("X-Tracksync-Forwarded-Files", strconv.Itoa(sent)) w.WriteHeader(http.StatusCreated) - _, _ = fmt.Fprintln(w, "uploaded") + _, _ = fmt.Fprintf(w, "uploaded %d/%d files\n", sent, len(outputs)) +} + +// claimReceivedFile inserts the payload row if absent and returns its id, +// returning the existing id on a concurrent insert of the same (device, payload). +func (s *Server) claimReceivedFile(deviceID, clientID, sha256Hex, filename string) (int64, error) { + res, err := s.db.Exec( + "INSERT OR IGNORE INTO received_files (device_id, client_id, sha256, filename, received_at) VALUES (?, ?, ?, ?, ?)", + deviceID, clientID, sha256Hex, filename, time.Now().UTC(), + ) + if err != nil { + return 0, err + } + if rows, _ := res.RowsAffected(); rows > 0 { + return res.LastInsertId() + } + var id int64 + err = s.db.QueryRow("SELECT id FROM received_files WHERE device_id = ? AND sha256 = ?", deviceID, sha256Hex).Scan(&id) + return id, err } diff --git a/server/internal/server/server_test.go b/server/internal/server/server_test.go index fb72d22..b5effba 100644 --- a/server/internal/server/server_test.go +++ b/server/internal/server/server_test.go @@ -23,15 +23,15 @@ type mockTarget struct { err error } -func (m *mockTarget) Type() string { return "mock" } -func (m *mockTarget) AcceptedFormats() []string { return []string{"gpx_1.1", "geojson"} } +func (m *mockTarget) Type() string { return "mock" } +func (m *mockTarget) AcceptedFormats() []string { return []string{"gpx_1.1", "geojson"} } func (m *mockTarget) Send(_ context.Context, filename string, data []byte) error { return m.err } type countTarget struct { count *int } -func (c *countTarget) Type() string { return "count" } +func (c *countTarget) Type() string { return "count" } func (c *countTarget) AcceptedFormats() []string { return []string{"gpx_1.1", "geojson"} } func (c *countTarget) Send(_ context.Context, filename string, data []byte) error { *c.count++ @@ -202,6 +202,79 @@ func TestUpload_TargetNotForwarded_OnDuplicate(t *testing.T) { assert.Equal(t, 1, calls, "target.Send should only be called once for duplicate content") } +func TestUpload_SplitFilesMode_SendsMultiple(t *testing.T) { + calls := 0 + db, err := sql.Open("sqlite", ":memory:") + require.NoError(t, err) + require.NoError(t, InitDB(db)) + t.Cleanup(func() { _ = db.Close() }) + + cfg := &config.Config{ + Accounts: []config.Account{{DeviceID: "dev", Markers: []string{"C:split"}, SplitMarkerPosition: "start", SplitMode: "files"}}, + Clients: []config.Client{{ID: "c", Token: "tok", AllowedDeviceIDs: []string{"dev"}}}, + } + srv := New(cfg, db, map[string]target.Target{"dev": &countTarget{count: &calls}}) + + csvData := []byte("INDEX,TAG,DATE,TIME,LATITUDE N/S,LONGITUDE E/W,HEIGHT,SPEED,HEADING\n" + + "1,T,260417,110529,52.0N,13.0E,38,1.4,333\n" + + "2,C,260417,110530,52.1N,13.1E,38,30.0,9\n" + + "3,T,260417,110531,52.2N,13.2E,38,30.0,13\n") + + rec := serveUpload(srv, "tok", "dev", "columbus-csv", "track.csv", csvData) + require.Equal(t, http.StatusCreated, rec.Code) + assert.Equal(t, 2, calls, "files mode should forward one file per split track") + assert.Equal(t, "uploaded 2/2 files\n", rec.Body.String(), "response reports the multi-file count") + assert.Equal(t, "2", rec.Header().Get("X-Tracksync-Forwarded-Files"), "fan-out count is reported to the client") +} + +// flakyTarget records every filename it accepts and fails sends for any +// filename present in failOn. +type flakyTarget struct { + sends []string + failOn map[string]bool +} + +func (f *flakyTarget) Type() string { return "flaky" } +func (f *flakyTarget) AcceptedFormats() []string { return []string{"gpx_1.1"} } +func (f *flakyTarget) Send(_ context.Context, filename string, _ []byte) error { + if f.failOn[filename] { + return fmt.Errorf("simulated failure for %s", filename) + } + f.sends = append(f.sends, filename) + return nil +} + +func TestUpload_SplitFilesMode_PartialFailureResendsOnlyMissing(t *testing.T) { + db, err := sql.Open("sqlite", ":memory:") + require.NoError(t, err) + require.NoError(t, InitDB(db)) + t.Cleanup(func() { _ = db.Close() }) + + cfg := &config.Config{ + Accounts: []config.Account{{DeviceID: "dev", Markers: []string{"C:split"}, SplitMarkerPosition: "start", SplitMode: "files"}}, + Clients: []config.Client{{ID: "c", Token: "tok", AllowedDeviceIDs: []string{"dev"}}}, + } + tgt := &flakyTarget{failOn: map[string]bool{"track-2.gpx": true}} + srv := New(cfg, db, map[string]target.Target{"dev": tgt}) + + csvData := []byte("INDEX,TAG,DATE,TIME,LATITUDE N/S,LONGITUDE E/W,HEIGHT,SPEED,HEADING\n" + + "1,T,260417,110529,52.0N,13.0E,38,1.4,333\n" + + "2,C,260417,110530,52.1N,13.1E,38,30.0,9\n" + + "3,T,260417,110531,52.2N,13.2E,38,30.0,13\n") + + // First attempt: track-1 delivered, track-2 fails -> whole upload errors. + rec1 := serveUpload(srv, "tok", "dev", "columbus-csv", "track.csv", csvData) + require.Equal(t, http.StatusBadGateway, rec1.Code) + assert.Equal(t, []string{"track-1.gpx"}, tgt.sends, "only the first file is delivered before the failure") + + // Recover the target and retry the identical upload. + tgt.failOn = nil + rec2 := serveUpload(srv, "tok", "dev", "columbus-csv", "track.csv", csvData) + require.Equal(t, http.StatusCreated, rec2.Code) + assert.Equal(t, []string{"track-1.gpx", "track-2.gpx"}, tgt.sends, "retry resends only the missing file, no duplicate of track-1") + assert.Equal(t, "uploaded 1/2 files\n", rec2.Body.String(), "retry reports one file sent, the other a duplicate") +} + func TestUpload_MissingSourceFormat(t *testing.T) { srv := setupTestServer(t, nil) rec := serveUpload(srv, "valid-token", "dev-1", "", "track.gpx", validGPX) diff --git a/server/internal/target/dawarich/dawarich.go b/server/internal/target/dawarich/dawarich.go index 0ccf2f0..d9773d9 100644 --- a/server/internal/target/dawarich/dawarich.go +++ b/server/internal/target/dawarich/dawarich.go @@ -35,8 +35,8 @@ type Dawarich struct { client *http.Client } -func (d *Dawarich) Type() string { return "dawarich" } -func (d *Dawarich) AcceptedFormats() []string { return []string{"geojson", "gpx_1.1"} } +func (d *Dawarich) Type() string { return "dawarich" } +func (d *Dawarich) AcceptedFormats() []string { return []string{"geojson", "gpx_1.1"} } func (d *Dawarich) Send(ctx context.Context, filename string, data []byte) error { apiKey, err := d.readAPIKey() diff --git a/server/main.go b/server/main.go index 43dd2f5..7b33646 100644 --- a/server/main.go +++ b/server/main.go @@ -51,7 +51,14 @@ func main() { os.Exit(1) } - db, err := sql.Open("sqlite", cfg.StateDB) + // A file: URI with an absolute path applies the foreign_keys pragma to every + // connection (the driver rejects a relative path in URI form). + dbPath, err := filepath.Abs(cfg.StateDB) + if err != nil { + slog.Error("failed to resolve database path", "error", err) + os.Exit(1) + } + db, err := sql.Open("sqlite", "file:"+dbPath+"?_pragma=foreign_keys(1)") if err != nil { slog.Error("failed to open database", "error", err) os.Exit(1) diff --git a/tracksync/internal/sync/sync.go b/tracksync/internal/sync/sync.go index 2a24ea5..5e8cd1d 100644 --- a/tracksync/internal/sync/sync.go +++ b/tracksync/internal/sync/sync.go @@ -13,6 +13,7 @@ import ( "net/http" "os" "path/filepath" + "strconv" "strings" "time" @@ -82,23 +83,25 @@ const ( StatusDuplicate // server already had this file ) -func Upload(ctx context.Context, client *http.Client, serverURL, token, deviceID, sourceFormat, filename string, data []byte) (UploadStatus, error) { +// Upload sends one source file. On success it returns how many target files the +// server forwarded (more than one when a track was split). +func Upload(ctx context.Context, client *http.Client, serverURL, token, deviceID, sourceFormat, filename string, data []byte) (UploadStatus, int, error) { var buf bytes.Buffer writer := multipart.NewWriter(&buf) part, err := writer.CreateFormFile("file", filename) if err != nil { - return 0, fmt.Errorf("creating form: %w", err) + return 0, 0, fmt.Errorf("creating form: %w", err) } if _, err := part.Write(data); err != nil { - return 0, fmt.Errorf("writing form: %w", err) + return 0, 0, fmt.Errorf("writing form: %w", err) } if err := writer.Close(); err != nil { - return 0, fmt.Errorf("closing form: %w", err) + return 0, 0, fmt.Errorf("closing form: %w", err) } req, err := http.NewRequestWithContext(ctx, "POST", serverURL+"/upload", &buf) if err != nil { - return 0, fmt.Errorf("creating request: %w", err) + return 0, 0, fmt.Errorf("creating request: %w", err) } req.Header.Set("Content-Type", writer.FormDataContentType()) req.Header.Set("Authorization", "Bearer "+token) @@ -107,7 +110,7 @@ func Upload(ctx context.Context, client *http.Client, serverURL, token, deviceID resp, err := client.Do(req) if err != nil { - return 0, fmt.Errorf("sending request: %w", err) + return 0, 0, fmt.Errorf("sending request: %w", err) } defer func() { _ = resp.Body.Close() }() @@ -116,11 +119,16 @@ func Upload(ctx context.Context, client *http.Client, serverURL, token, deviceID switch resp.StatusCode { case http.StatusCreated: - return StatusUploaded, nil + // Default to 1 when the header is missing or invalid. + forwarded, err := strconv.Atoi(resp.Header.Get("X-Tracksync-Forwarded-Files")) + if err != nil || forwarded < 1 { + forwarded = 1 + } + return StatusUploaded, forwarded, nil case http.StatusOK: - return StatusDuplicate, nil + return StatusDuplicate, 0, nil default: - return 0, fmt.Errorf("server returned %d: %s", resp.StatusCode, status) + return 0, 0, fmt.Errorf("server returned %d: %s", resp.StatusCode, status) } } @@ -130,6 +138,7 @@ type Summary struct { Duplicate int `json:"duplicate"` Skipped int `json:"skipped"` Errors int `json:"errors"` + Forwarded int `json:"forwarded"` // target files forwarded (exceeds Uploaded when tracks were split) Files []string `json:"files,omitempty"` } @@ -161,7 +170,7 @@ func SyncFiles(ctx context.Context, db *sql.DB, client *http.Client, serverURL, continue } - status, err := Upload(ctx, client, serverURL, token, deviceID, ff.Format, name, data) + status, forwarded, err := Upload(ctx, client, serverURL, token, deviceID, ff.Format, name, data) if err != nil { slog.Error("upload failed", "file", name, "error", err) summary.Errors++ @@ -176,8 +185,9 @@ func SyncFiles(ctx context.Context, db *sql.DB, client *http.Client, serverURL, switch status { case StatusUploaded: - slog.Info("uploaded", "file", name) + slog.Info("uploaded", "file", name, "forwarded", forwarded) summary.Uploaded++ + summary.Forwarded += forwarded summary.Files = append(summary.Files, name) case StatusDuplicate: slog.Info("duplicate", "file", name, "reason", "already on server") diff --git a/tracksync/internal/sync/sync_test.go b/tracksync/internal/sync/sync_test.go index 198e00a..45cfdf9 100644 --- a/tracksync/internal/sync/sync_test.go +++ b/tracksync/internal/sync/sync_test.go @@ -102,15 +102,32 @@ func TestUpload_Created(t *testing.T) { assert.Equal(t, "Bearer tok", r.Header.Get("Authorization")) assert.Equal(t, "dev-1", r.Header.Get("X-Device-ID")) assert.Equal(t, "gpx_1.1", r.Header.Get("X-Source-Format")) + w.Header().Set("X-Tracksync-Forwarded-Files", "3") w.WriteHeader(http.StatusCreated) _, _ = fmt.Fprintln(w, "uploaded") })) defer ts.Close() client := &http.Client{Timeout: 5 * time.Second} - status, err := Upload(context.Background(), client, ts.URL, "tok", "dev-1", "gpx_1.1", "track.gpx", []byte("")) + status, forwarded, err := Upload(context.Background(), client, ts.URL, "tok", "dev-1", "gpx_1.1", "track.gpx", []byte("")) require.NoError(t, err) assert.Equal(t, StatusUploaded, status) + assert.Equal(t, 3, forwarded, "forwarded-files header should be parsed") +} + +func TestUpload_CreatedNoHeaderCountsAsOne(t *testing.T) { + // A missing fan-out header counts as one forwarded file. + ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(http.StatusCreated) + _, _ = fmt.Fprintln(w, "uploaded") + })) + defer ts.Close() + + client := &http.Client{Timeout: 5 * time.Second} + status, forwarded, err := Upload(context.Background(), client, ts.URL, "tok", "dev-1", "gpx_1.1", "track.gpx", []byte("")) + require.NoError(t, err) + assert.Equal(t, StatusUploaded, status) + assert.Equal(t, 1, forwarded) } func TestUpload_Duplicate(t *testing.T) { @@ -121,9 +138,10 @@ func TestUpload_Duplicate(t *testing.T) { defer ts.Close() client := &http.Client{Timeout: 5 * time.Second} - status, err := Upload(context.Background(), client, ts.URL, "tok", "dev", "gpx_1.1", "f.gpx", []byte("data")) + status, forwarded, err := Upload(context.Background(), client, ts.URL, "tok", "dev", "gpx_1.1", "f.gpx", []byte("data")) require.NoError(t, err) assert.Equal(t, StatusDuplicate, status) + assert.Equal(t, 0, forwarded) } func TestUpload_ServerError(t *testing.T) { @@ -134,7 +152,7 @@ func TestUpload_ServerError(t *testing.T) { defer ts.Close() client := &http.Client{Timeout: 5 * time.Second} - _, err := Upload(context.Background(), client, ts.URL, "tok", "dev", "gpx_1.1", "f.gpx", []byte("data")) + _, _, err := Upload(context.Background(), client, ts.URL, "tok", "dev", "gpx_1.1", "f.gpx", []byte("data")) assert.Error(t, err) } @@ -153,12 +171,37 @@ func TestSyncFiles_Uploaded(t *testing.T) { summary := SyncFiles(context.Background(), db, &http.Client{Timeout: 5 * time.Second}, ts.URL, "tok", "dev", files) assert.Equal(t, 1, summary.Uploaded) + assert.Equal(t, 1, summary.Forwarded, "no fan-out header: one source file counts as one forwarded file") assert.Equal(t, 0, summary.Duplicate) assert.Equal(t, 0, summary.Skipped) assert.Equal(t, 0, summary.Errors) assert.Equal(t, []string{"track.gpx"}, summary.Files) } +func TestSyncFiles_AccumulatesForwardedFanOut(t *testing.T) { + // The server splits each upload into three forwarded files. + ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("X-Tracksync-Forwarded-Files", "3") + w.WriteHeader(http.StatusCreated) + _, _ = fmt.Fprintln(w, "uploaded 3/3 files") + })) + defer ts.Close() + + db := openTestDB(t) + dir := t.TempDir() + require.NoError(t, os.WriteFile(filepath.Join(dir, "a.gpx"), []byte("a"), 0644)) + require.NoError(t, os.WriteFile(filepath.Join(dir, "b.gpx"), []byte("b"), 0644)) + + files := []device.FoundFile{ + {Path: filepath.Join(dir, "a.gpx"), Format: "gpx_1.1"}, + {Path: filepath.Join(dir, "b.gpx"), Format: "gpx_1.1"}, + } + summary := SyncFiles(context.Background(), db, &http.Client{Timeout: 5 * time.Second}, ts.URL, "tok", "dev", files) + + assert.Equal(t, 2, summary.Uploaded, "two source files") + assert.Equal(t, 6, summary.Forwarded, "each source file fanned out into three target files") +} + func TestSyncFiles_Duplicate(t *testing.T) { ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { w.WriteHeader(http.StatusOK) diff --git a/tracksync/main.go b/tracksync/main.go index aada883..8cf0ff7 100644 --- a/tracksync/main.go +++ b/tracksync/main.go @@ -138,6 +138,7 @@ func main() { slog.Info("sync complete", "uploaded", summary.Uploaded, + "forwarded", summary.Forwarded, "duplicate", summary.Duplicate, "skipped", summary.Skipped, "errors", summary.Errors,