diff --git a/.github/workflows/build-docker.yml b/.github/workflows/build-docker.yml index a3ebe75cb..a24ed39a9 100644 --- a/.github/workflows/build-docker.yml +++ b/.github/workflows/build-docker.yml @@ -3,7 +3,7 @@ name: Create and publish a Docker image on: push: - branches: ['dev', 'master', 'dev-flex', 'mtc-deploy'] + branches: ['dev', 'master', 'mtc-deploy'] env: REGISTRY: ghcr.io diff --git a/pom.xml b/pom.xml index c162661f4..4230d9f57 100644 --- a/pom.xml +++ b/pom.xml @@ -296,7 +296,7 @@ com.github.ibi-group gtfs-lib - d39534084f9df5a798ae8d2b6f61cd8a6ab4947d + 065d51c957dfdd6231b3f10e9fddef4dd867146a @@ -424,9 +424,9 @@ - com.github.conveyal + com.github.ibi-group java-snapshot-matcher - 3495b32f7b4d3f82590e0a2284029214070b6984 + master test @@ -505,6 +505,11 @@ gson 2.8.9 + + org.apache.commons + commons-lang3 + 3.17.0 + commons-cli commons-cli diff --git a/src/main/java/com/conveyal/datatools/manager/DataManager.java b/src/main/java/com/conveyal/datatools/manager/DataManager.java index b1f4c2335..a711c3209 100644 --- a/src/main/java/com/conveyal/datatools/manager/DataManager.java +++ b/src/main/java/com/conveyal/datatools/manager/DataManager.java @@ -227,6 +227,13 @@ static void registerRoutes() throws IOException { new EditorControllerImpl(EDITOR_API_PREFIX, Table.FARE_PRODUCTS, DataManager.GTFS_DATA_SOURCE); new EditorControllerImpl(EDITOR_API_PREFIX, Table.FARE_TRANSFER_RULES, DataManager.GTFS_DATA_SOURCE); new EditorControllerImpl(EDITOR_API_PREFIX, Table.FEED_INFO, DataManager.GTFS_DATA_SOURCE); + + // NOTE: Booking rules, locations and location shapes are GTFS Flex additions. + new EditorControllerImpl(EDITOR_API_PREFIX, Table.BOOKING_RULES, DataManager.GTFS_DATA_SOURCE); + new EditorControllerImpl(EDITOR_API_PREFIX, Table.LOCATIONS, DataManager.GTFS_DATA_SOURCE); + new EditorControllerImpl(EDITOR_API_PREFIX, Table.LOCATION_GROUP, DataManager.GTFS_DATA_SOURCE); + new EditorControllerImpl(EDITOR_API_PREFIX, Table.LOCATION_GROUP_STOPS, DataManager.GTFS_DATA_SOURCE); + new EditorControllerImpl(EDITOR_API_PREFIX, Table.LOCATION_SHAPES, DataManager.GTFS_DATA_SOURCE); new EditorControllerImpl(EDITOR_API_PREFIX, Table.NETWORKS, DataManager.GTFS_DATA_SOURCE); new EditorControllerImpl(EDITOR_API_PREFIX, Table.RIDER_CATEGORIES, DataManager.GTFS_DATA_SOURCE); new EditorControllerImpl(EDITOR_API_PREFIX, Table.ROUTES, DataManager.GTFS_DATA_SOURCE); diff --git a/src/main/java/com/conveyal/datatools/manager/jobs/MergeFeedsJob.java b/src/main/java/com/conveyal/datatools/manager/jobs/MergeFeedsJob.java index 6c382ed16..5b0754bf5 100644 --- a/src/main/java/com/conveyal/datatools/manager/jobs/MergeFeedsJob.java +++ b/src/main/java/com/conveyal/datatools/manager/jobs/MergeFeedsJob.java @@ -18,8 +18,12 @@ import com.conveyal.datatools.manager.persistence.Persistence; import com.conveyal.datatools.manager.utils.ErrorUtils; import com.conveyal.gtfs.loader.Feed; +import com.conveyal.gtfs.loader.JdbcGtfsExporter; import com.conveyal.gtfs.loader.Table; +import com.conveyal.gtfs.model.Location; +import com.conveyal.gtfs.model.LocationShape; import com.conveyal.gtfs.model.StopTime; +import com.conveyal.gtfs.util.GeoJsonUtil; import com.fasterxml.jackson.annotation.JsonIgnore; import com.fasterxml.jackson.databind.ObjectMapper; import com.google.common.collect.Lists; @@ -39,13 +43,17 @@ import java.util.List; import java.util.Set; import java.util.stream.Collectors; +import java.util.zip.ZipEntry; import java.util.zip.ZipOutputStream; import static com.conveyal.datatools.manager.jobs.feedmerge.MergeFeedsType.SERVICE_PERIOD; import static com.conveyal.datatools.manager.jobs.feedmerge.MergeFeedsType.REGIONAL; import static com.conveyal.datatools.manager.jobs.feedmerge.MergeStrategy.CHECK_STOP_TIMES; import static com.conveyal.datatools.manager.models.FeedRetrievalMethod.REGIONAL_MERGE; -import static com.conveyal.datatools.manager.utils.MergeFeedUtils.*; +import static com.conveyal.datatools.manager.utils.MergeFeedUtils.getMergedVersion; +import static com.conveyal.datatools.manager.utils.MergeFeedUtils.stopTimesMatchSimplified; +import static com.conveyal.datatools.manager.utils.StringUtils.getCleanName; +import static com.conveyal.gtfs.loader.Table.LOCATION_GEO_JSON_FILE_NAME; /** * This job handles merging two or more feed versions according to logic specific to the specified merge type. @@ -205,6 +213,8 @@ public void jobLogic() { } } + mergeLocations(out); + // Loop over GTFS tables and merge each feed one table at a time. for (int i = 0; i < numberOfTables; i++) { Table table = tablesToMerge.get(i); @@ -226,7 +236,7 @@ public void jobLogic() { status.fail(message, e); } finally { try { - feedMergeContext.close(); + if (feedMergeContext != null) feedMergeContext.close(); } catch (IOException e) { logAndReportToBugsnag(e, "Error closing FeedMergeContext object"); } @@ -246,6 +256,39 @@ public void jobLogic() { } } + /** + * Merge locations.geojson files. These files are not compatible with the CSV merging strategy. Instead, the + * location.geojson file is flattened into locations and locations shapes. The location id is then updated with the + * scope id to keep feed locations unique, converted back into geojson and written to the zip file. + * + * Locations must be processed prior to other CSV files so the location ids are available for foreign reference + * checks. + */ + void mergeLocations(ZipOutputStream out) throws IOException { + Set mergedLocations = new HashSet<>(); + Set mergedLocationShapes = new HashSet<>(); + for (FeedToMerge feed : feedMergeContext.feedsToMerge) { + ZipEntry locationGeoJsonFile = feed.zipFile.getEntry(LOCATION_GEO_JSON_FILE_NAME); + if (locationGeoJsonFile != null) { + String idScope = getCleanName(feed.version.parentFeedSource().name) + feed.version.version; + List locations = GeoJsonUtil.getLocationsFromGeoJson(feed.zipFile, locationGeoJsonFile, null); + for (Location location : locations) { + location.location_id = String.join(":", idScope, location.location_id); + mergedLocations.add(location); + feedMergeContext.locationIds.add(location.location_id); + } + List locationShapes = GeoJsonUtil.getLocationShapesFromGeoJson(feed.zipFile, locationGeoJsonFile, null); + for (LocationShape locationShape : locationShapes) { + locationShape.location_id = String.join(":", idScope, locationShape.location_id); + mergedLocationShapes.add(locationShape); + } + } + } + if (!mergedLocations.isEmpty()) { + JdbcGtfsExporter.writeLocationsToFile(out, new ArrayList<>(mergedLocations), new ArrayList<>(mergedLocationShapes)); + } + } + /** * Obtains trip ids whose entries in the stop_times table differ between the active and future feed. */ @@ -282,6 +325,10 @@ private boolean shouldSkipTable(String tableName) { LOG.warn("Skipping editor-only table {}.", tableName); return true; } + if (tableName.equals(Table.LOCATIONS.name) || tableName.equals(Table.LOCATION_SHAPES.name)) { + LOG.warn("{} detected. Skipping traditional merge in favour of bespoke merge.", LOCATION_GEO_JSON_FILE_NAME); + return true; + } return false; } diff --git a/src/main/java/com/conveyal/datatools/manager/jobs/ProcessSingleFeedJob.java b/src/main/java/com/conveyal/datatools/manager/jobs/ProcessSingleFeedJob.java index 25a3fb215..a3540f663 100644 --- a/src/main/java/com/conveyal/datatools/manager/jobs/ProcessSingleFeedJob.java +++ b/src/main/java/com/conveyal/datatools/manager/jobs/ProcessSingleFeedJob.java @@ -14,6 +14,7 @@ import com.conveyal.datatools.manager.models.transform.FeedTransformZipTarget; import com.conveyal.datatools.manager.models.transform.RemoveNonRevenueTripsTransformation; import com.conveyal.datatools.manager.models.transform.ZipTransformation; +import com.conveyal.gtfs.validator.ValidationResult; import com.fasterxml.jackson.annotation.JsonIgnore; import com.fasterxml.jackson.annotation.JsonProperty; import org.slf4j.Logger; @@ -227,15 +228,20 @@ private String getErrorReasonMessage() { public String getNotificationMessage() { StringBuilder message = new StringBuilder(); if (!status.error) { - message.append(String.format("New feed version created for %s (valid from %s - %s). ", - feedSource.name, - feedVersion.validationResult.firstCalendarDate, - feedVersion.validationResult.lastCalendarDate)); - if (feedVersion.validationResult.errorCount > 0) { - message.append(String.format("During validation, we found %s issue(s)", - feedVersion.validationResult.errorCount)); + ValidationResult validationResult = feedVersion.validationResult; + if (validationResult != null) { + message.append(String.format("New feed version created for %s (valid from %s - %s).", + feedSource.name, + validationResult.firstCalendarDate, + validationResult.lastCalendarDate + )); + if (validationResult.errorCount > 0) { + message.append(String.format(" During validation, we found %s issue(s)", validationResult.errorCount)); + } else { + message.append(" The validation check found no issues with this new dataset!"); + } } else { - message.append("The validation check found no issues with this new dataset!"); + message.append(String.format("New feed version created for %s.", feedSource.name)); } } else { // Processing did not complete. Depending on which sub-task this occurred in, diff --git a/src/main/java/com/conveyal/datatools/manager/jobs/feedmerge/FeedMergeContext.java b/src/main/java/com/conveyal/datatools/manager/jobs/feedmerge/FeedMergeContext.java index 371112026..da5974856 100644 --- a/src/main/java/com/conveyal/datatools/manager/jobs/feedmerge/FeedMergeContext.java +++ b/src/main/java/com/conveyal/datatools/manager/jobs/feedmerge/FeedMergeContext.java @@ -9,6 +9,7 @@ import java.io.Closeable; import java.io.IOException; import java.time.LocalDate; +import java.util.HashSet; import java.util.List; import java.util.Set; @@ -23,6 +24,8 @@ public class FeedMergeContext implements Closeable { public final boolean tripIdsMatch; public final LocalDate futureFirstCalendarStartDate; public final Set sharedTripIds; + public Set locationIds = new HashSet<>(); + public FeedMergeContext(Set feedVersions, Auth0UserProfile owner) throws IOException { feedsToMerge = MergeFeedUtils.collectAndSortFeeds(feedVersions, owner); diff --git a/src/main/java/com/conveyal/datatools/manager/jobs/feedmerge/MergeFeedsResult.java b/src/main/java/com/conveyal/datatools/manager/jobs/feedmerge/MergeFeedsResult.java index 881c78051..b9f848b4c 100644 --- a/src/main/java/com/conveyal/datatools/manager/jobs/feedmerge/MergeFeedsResult.java +++ b/src/main/java/com/conveyal/datatools/manager/jobs/feedmerge/MergeFeedsResult.java @@ -36,6 +36,12 @@ public class MergeFeedsResult implements Serializable { */ public Set calendarDatesServiceIds = new HashSet<>(); + /** + * Track various table ids for resolving foreign references. + */ + public Set stopIds = new HashSet<>(); + public Set locationGroupStopIds = new HashSet<>(); + /** * Track the set of route IDs to end up in the merged feed in order to determine which route_attributes * records should be retained in the merged result. diff --git a/src/main/java/com/conveyal/datatools/manager/jobs/feedmerge/MergeLineContext.java b/src/main/java/com/conveyal/datatools/manager/jobs/feedmerge/MergeLineContext.java index 2a275e2d7..63fb7b243 100644 --- a/src/main/java/com/conveyal/datatools/manager/jobs/feedmerge/MergeLineContext.java +++ b/src/main/java/com/conveyal/datatools/manager/jobs/feedmerge/MergeLineContext.java @@ -240,6 +240,7 @@ public void startNewRow() throws IOException { // For this table, use all fields found in the feeds to merge. sharedSpecFields = List.copyOf(allFields); } else { + setForeignReferenceKeyValues(); // Get the spec fields and custom/proprietary fields to export List specFields = table.specFields(); // Filter the spec fields on the set of fields found in all feeds to be merged. @@ -250,8 +251,24 @@ public void startNewRow() throws IOException { } /** - * Determine which reference table to use. If there is only one reference use this. If there are multiple references - * determine the context and then the correct reference table to use. + * Build a list of table key id values to be used in foreign key field look-ups. + */ + private void setForeignReferenceKeyValues() { + switch (table.name) { + case "stops": + mergeFeedsResult.stopIds.add(getIdWithScope(keyValue)); + break; + case "location_group_stops": + mergeFeedsResult.locationGroupStopIds.add(getIdWithScope(keyValue)); + break; + default: + // nothing. + } + } + + /** + * Determine which reference table to use. If there is only one reference use this. If there are multiple + * references, determine the context and then the correct reference table to use. */ private Table getReferenceTable(FieldContext fieldContext, Field field) { if (field.referenceTables.size() == 1) { @@ -266,7 +283,12 @@ private Table getReferenceTable(FieldContext fieldContext, Field field) { getTableScopedValue(Table.CALENDAR, fieldContext.getValue()), getTableScopedValue(Table.CALENDAR_DATES, fieldContext.getValue()) ); - // Include other cases as multiple references are added e.g. flex!. + case STOP_TIMES_STOP_ID_KEY: + case LOCATION_GROUP_STOPS_STOP_ID_KEY: + return ReferenceTableDiscovery.getStopReferenceTable( + fieldContext.getValueToWrite(), + mergeFeedsResult + ); default: return null; } diff --git a/src/main/java/com/conveyal/datatools/manager/jobs/feedmerge/ReferenceTableDiscovery.java b/src/main/java/com/conveyal/datatools/manager/jobs/feedmerge/ReferenceTableDiscovery.java index a972dbc25..285a8912d 100644 --- a/src/main/java/com/conveyal/datatools/manager/jobs/feedmerge/ReferenceTableDiscovery.java +++ b/src/main/java/com/conveyal/datatools/manager/jobs/feedmerge/ReferenceTableDiscovery.java @@ -11,6 +11,9 @@ public class ReferenceTableDiscovery { public static final String REF_TABLE_SEPARATOR = "#~#"; + /** + * Tables that have two or more foreign references. + */ public enum ReferenceTableKey { TRIP_SERVICE_ID_KEY( @@ -22,6 +25,22 @@ public enum ReferenceTableKey { Table.CALENDAR_DATES.name, Table.SCHEDULE_EXCEPTIONS.name ) + ), + LOCATION_GROUP_STOPS_STOP_ID_KEY( + String.join( + REF_TABLE_SEPARATOR, + Table.LOCATION_GROUP_STOPS.name, + "stop_id", + Table.STOPS.name + ) + ), + STOP_TIMES_STOP_ID_KEY( + String.join( + REF_TABLE_SEPARATOR, + Table.STOP_TIMES.name, + "stop_id", + Table.STOPS.name + ) ); private final String value; @@ -87,4 +106,14 @@ public static Table getTripServiceIdReferenceTable( } return null; } + + /** + * Defines the reference table as a stop if the field value matches a stop id. + */ + public static Table getStopReferenceTable(String fieldValue, MergeFeedsResult mergeFeedsResult) { + if (mergeFeedsResult.stopIds.contains(fieldValue)) { + return Table.STOPS; + } + return null; + } } diff --git a/src/main/java/com/conveyal/datatools/manager/models/FeedSource.java b/src/main/java/com/conveyal/datatools/manager/models/FeedSource.java index 3723f1a91..7d71229ed 100644 --- a/src/main/java/com/conveyal/datatools/manager/models/FeedSource.java +++ b/src/main/java/com/conveyal/datatools/manager/models/FeedSource.java @@ -18,6 +18,7 @@ import com.conveyal.datatools.manager.utils.connections.ConnectionResponse; import com.conveyal.datatools.manager.utils.connections.HttpURLConnectionResponse; import com.conveyal.gtfs.GTFS; +import com.conveyal.gtfs.loader.FeedLoadResult; import com.fasterxml.jackson.annotation.JsonIgnore; import com.fasterxml.jackson.annotation.JsonIgnoreProperties; import com.fasterxml.jackson.annotation.JsonInclude; @@ -180,6 +181,14 @@ public String organizationId () { public String editorNamespace; + /** + * If this feed source references a feed version that is GTFS Flex, this value will permanently be set to true. + * This is evaluated every time a feed version is retrieved as different feed versions may or may not be flex. + */ + public boolean flex; + + public boolean flexUIFeaturesEnabled; + /** * Create a new feed. */ @@ -388,6 +397,7 @@ public FeedVersion processFetchResponse( String message = String.format("Fetch complete for %s", this.name); LOG.info(message); status.completeSuccessfully(message); + presetFlex(version); return version; } } @@ -429,10 +439,16 @@ public String toString () { */ @JsonIgnore public FeedVersion retrieveLatest() { - return Persistence.feedVersions.getOneFiltered( + FeedVersion newestVersion = Persistence.feedVersions.getOneFiltered( eq("feedSourceId", this.id), Sorts.descending("version") ); + if (newestVersion == null) { + // Is this what happens if there are none? + return null; + } + presetFlex(newestVersion); + return newestVersion; } /** @@ -448,6 +464,7 @@ public FeedVersion retrievePublishedVersion() { // Is this what happens if there are none? return null; } + presetFlex(publishedVersion); return publishedVersion; } @@ -809,4 +826,19 @@ public List getActiveTransformations(FeedVersi .flatMap(Collection::stream) .collect(Collectors.toList()); } + + /** + * If a flex feed version is loaded, set this feed source to permanently be a flex feed source. Subsequent feed + * versions which are not flex will have no affected. + */ + public void presetFlex(FeedVersion feedVersion) { + boolean isFlexFeedVersion = false; + if (feedVersion.feedLoadResult != null) isFlexFeedVersion = feedVersion.feedLoadResult.isGTFSFlex(); + if (!flex && isFlexFeedVersion) { + // If the feed version is flex (and the feed source was not flex) permanently enable flex. + flex = flexUIFeaturesEnabled = true; + Persistence.feedSources.updateField(this.id, "flexUIFeaturesEnabled", true); + Persistence.feedSources.updateField(this.id, "flex", true); + } + } } diff --git a/src/main/java/com/conveyal/datatools/manager/models/FeedVersion.java b/src/main/java/com/conveyal/datatools/manager/models/FeedVersion.java index 6f8586bc8..d3508b28e 100644 --- a/src/main/java/com/conveyal/datatools/manager/models/FeedVersion.java +++ b/src/main/java/com/conveyal/datatools/manager/models/FeedVersion.java @@ -306,6 +306,7 @@ public void load(MonitorableJob.Status status, boolean isNewVersion) { assignGtfsFileAttributes(gtfsFile); String gtfsFilePath = gtfsFile.getPath(); this.feedLoadResult = GTFS.load(gtfsFilePath, DataManager.GTFS_DATA_SOURCE); + if (this.feedLoadResult.fatalException != null) { status.fail("Could not load feed due to " + feedLoadResult.fatalException); return; @@ -449,7 +450,6 @@ public void validateMobility(MonitorableJob.Status status) { // Wait for the file to be entirely copied into the directory. // 5 seconds + ~1 second per 10mb Thread.sleep(5000 + (this.fileSize / 10000)); - File gtfsZip = this.retrieveGtfsFile(); // Namespace based folders avoid clash for validation being run on multiple versions of a feed. // TODO: do we know that there will always be a namespace? String validatorOutputDirectory = "/tmp/datatools_gtfs/" + this.namespace + "/"; @@ -457,7 +457,7 @@ public void validateMobility(MonitorableJob.Status status) { status.update("MobilityData Analysis...", 20); // Set up MobilityData validator. ValidationRunnerConfig.Builder builder = ValidationRunnerConfig.builder(); - builder.setGtfsSource(gtfsZip.toURI()); + builder.setGtfsSource(this.retrieveGtfsFile().toURI()); builder.setOutputDirectory(Path.of(validatorOutputDirectory)); ValidationRunnerConfig mbValidatorConfig = builder.build(); @@ -656,6 +656,7 @@ public void delete() { deleteFeedVersionFile(); deleteDBSchema(namespace); + // Remove this FeedVersion from all Deployments associated with this FeedVersion's FeedSource's Project // TODO TEST THOROUGHLY THAT THIS UPDATE EXPRESSION IS CORRECT // Although outright deleting the feedVersion from deployments could be surprising and shouldn't be done anyway. diff --git a/src/main/java/com/conveyal/datatools/manager/models/transform/FeedTransformation.java b/src/main/java/com/conveyal/datatools/manager/models/transform/FeedTransformation.java index 694fd063c..d119693ba 100644 --- a/src/main/java/com/conveyal/datatools/manager/models/transform/FeedTransformation.java +++ b/src/main/java/com/conveyal/datatools/manager/models/transform/FeedTransformation.java @@ -2,6 +2,7 @@ import com.conveyal.datatools.common.status.MonitorableJob; import com.conveyal.datatools.manager.utils.GtfsUtils; +import com.conveyal.gtfs.loader.Table; import com.fasterxml.jackson.annotation.JsonIgnoreProperties; import com.fasterxml.jackson.annotation.JsonSubTypes; import com.fasterxml.jackson.annotation.JsonTypeInfo; @@ -102,8 +103,16 @@ public void doTransform(FeedTransformTarget target, MonitorableJob.Status status */ protected void validateTableName(MonitorableJob.Status status) { // Validate fields before running transform. - if (GtfsUtils.getGtfsTable(table) == null) { - status.fail("Table must be valid GTFS spec table name (without .txt)."); + if (GtfsUtils.getGtfsTable(table) == null && !Table.LOCATION_GEO_JSON_FILE_NAME.equals(table)) { + status.fail(String.format("Table must be valid GTFS spec table name (without .txt) or %s.", Table.LOCATION_GEO_JSON_FILE_NAME)); } } + + protected String getTableName() { + return Table.LOCATION_GEO_JSON_FILE_NAME.equals(table) ? table : table + ".txt"; + } + + protected String getTableSuffix() { + return Table.LOCATION_GEO_JSON_FILE_NAME.equals(table) ? ".geojson" : ".txt"; + } } diff --git a/src/main/java/com/conveyal/datatools/manager/models/transform/NormalizeFieldTransformation.java b/src/main/java/com/conveyal/datatools/manager/models/transform/NormalizeFieldTransformation.java index fd3bb6266..4a69c42b9 100644 --- a/src/main/java/com/conveyal/datatools/manager/models/transform/NormalizeFieldTransformation.java +++ b/src/main/java/com/conveyal/datatools/manager/models/transform/NormalizeFieldTransformation.java @@ -182,7 +182,12 @@ public static List getInvalidSubstitutionPatterns(List sub @Override public void transform(FeedTransformZipTarget zipTarget, MonitorableJob.Status status) { - String tableName = table + ".txt"; + if (Table.LOCATION_GEO_JSON_FILE_NAME.equals(table)) { + // It's not possible to select the locations.geojson file from the normalize field transformation list. This + // is here in case that changes. + throw new UnsupportedOperationException("It is not possible to normalize geo json fields."); + } + String tableName = getTableName(); try( // Hold output before writing to ZIP StringWriter stringWriter = new StringWriter(); diff --git a/src/main/java/com/conveyal/datatools/manager/models/transform/PreserveCustomFieldsTransformation.java b/src/main/java/com/conveyal/datatools/manager/models/transform/PreserveCustomFieldsTransformation.java index e04db2407..976ebe0f2 100644 --- a/src/main/java/com/conveyal/datatools/manager/models/transform/PreserveCustomFieldsTransformation.java +++ b/src/main/java/com/conveyal/datatools/manager/models/transform/PreserveCustomFieldsTransformation.java @@ -60,14 +60,14 @@ private static HashMap> createCsvHashMap(CsvMapReade @Override public void transform(FeedTransformZipTarget zipTarget, MonitorableJob.Status status) throws Exception{ - String tableName = table + ".txt"; + String tableName = getTableName(); Path targetZipPath = Paths.get(zipTarget.gtfsFile.getAbsolutePath()); Optional streamResult = Arrays.stream(Table.tablesInOrder) .filter(t -> t.name.equals(table)) .findFirst(); - if (!streamResult.isPresent()) { - throw new Exception(String.format("could not find specTable for table %s", table)); + if (streamResult.isEmpty()) { + throw new IOException(String.format("could not find specTable for table %s", table)); } Table specTable = streamResult.get(); @@ -77,8 +77,8 @@ public void transform(FeedTransformZipTarget zipTarget, MonitorableJob.Status st Path targetTxtFilePath = getTablePathInZip(tableName, targetZipFs); - final File tempFile = File.createTempFile(tableName + "-temp", ".txt"); - File output = File.createTempFile(tableName + "-output-temp", ".txt"); + final File tempFile = File.createTempFile(tableName + "-temp", getTableSuffix()); + File output = File.createTempFile(tableName + "-output-temp", getTableSuffix()); int rowsModified = 0; List customFields; diff --git a/src/main/java/com/conveyal/datatools/manager/models/transform/RemoveNonRevenueTripsTransformation.java b/src/main/java/com/conveyal/datatools/manager/models/transform/RemoveNonRevenueTripsTransformation.java index b3a05cfb1..f68feaf67 100644 --- a/src/main/java/com/conveyal/datatools/manager/models/transform/RemoveNonRevenueTripsTransformation.java +++ b/src/main/java/com/conveyal/datatools/manager/models/transform/RemoveNonRevenueTripsTransformation.java @@ -59,6 +59,9 @@ public void transform(FeedTransformZipTarget zipTarget, MonitorableJob.Status st Files.copy(originalZipPath, tempZipPath, StandardCopyOption.REPLACE_EXISTING); Table gtfsTable = GtfsUtils.getGtfsTable("stop_times"); + if (gtfsTable == null) { + return; + } CsvReader csvReaderForStopTimes = CsvReaderUtil.getCsvReaderAccordingToFileName( gtfsTable, new ZipFile(tempZipPath.toAbsolutePath().toString()), @@ -78,6 +81,9 @@ public void transform(FeedTransformZipTarget zipTarget, MonitorableJob.Status st ); gtfsTable = GtfsUtils.getGtfsTable("trips"); + if (gtfsTable == null) { + return; + } CsvReader csvReaderForTrips = CsvReaderUtil.getCsvReaderAccordingToFileName( gtfsTable, new ZipFile(tempZipPath.toAbsolutePath().toString()), diff --git a/src/main/java/com/conveyal/datatools/manager/models/transform/ReplaceFileFromVersionTransformation.java b/src/main/java/com/conveyal/datatools/manager/models/transform/ReplaceFileFromVersionTransformation.java index 53963fb16..76234ae3e 100644 --- a/src/main/java/com/conveyal/datatools/manager/models/transform/ReplaceFileFromVersionTransformation.java +++ b/src/main/java/com/conveyal/datatools/manager/models/transform/ReplaceFileFromVersionTransformation.java @@ -42,7 +42,7 @@ public void validateParameters(MonitorableJob.Status status) { @Override public void transform(FeedTransformZipTarget zipTarget, MonitorableJob.Status status) { FeedVersion sourceVersion = getSourceVersion(); - String tableName = table + ".txt"; + String tableName = getTableName(); // Run the replace transformation Path sourceZipPath = Paths.get(sourceVersion.retrieveGtfsFile().getAbsolutePath()); try (FileSystem sourceZipFs = FileSystems.newFileSystem(sourceZipPath, (ClassLoader) null)) { diff --git a/src/main/java/com/conveyal/datatools/manager/models/transform/StringTransformation.java b/src/main/java/com/conveyal/datatools/manager/models/transform/StringTransformation.java index c845d4f37..f822eb6b9 100644 --- a/src/main/java/com/conveyal/datatools/manager/models/transform/StringTransformation.java +++ b/src/main/java/com/conveyal/datatools/manager/models/transform/StringTransformation.java @@ -32,7 +32,7 @@ public void validateParameters(MonitorableJob.Status status) { @Override public void transform(FeedTransformZipTarget zipTarget, MonitorableJob.Status status) { - String tableName = table + ".txt"; + String tableName = getTableName(); Path targetZipPath = Paths.get(zipTarget.gtfsFile.getAbsolutePath()); try ( FileSystem targetZipFs = FileSystems.newFileSystem(targetZipPath, (ClassLoader) null); diff --git a/src/main/java/com/conveyal/datatools/manager/utils/MergeFeedUtils.java b/src/main/java/com/conveyal/datatools/manager/utils/MergeFeedUtils.java index 1d412fe1c..9817235b9 100644 --- a/src/main/java/com/conveyal/datatools/manager/utils/MergeFeedUtils.java +++ b/src/main/java/com/conveyal/datatools/manager/utils/MergeFeedUtils.java @@ -155,8 +155,8 @@ public static boolean hasDuplicateError(Set errors) { } /** - * Checks whether the future and active stop_times for a particular trip_id are an exact match, - * using these criteria only: arrival_time, departure_time, stop_id, and stop_sequence + * Checks whether the future and active stop_times for a particular trip_id are an exact match, using these criteria + * only: arrival_time, departure_time, stop_sequence, stop_id, location_group_id and location_id * instead of StopTime::equals (Revised MTC feed merge requirement). */ public static boolean stopTimesMatchSimplified(List futureStopTimes, List activeStopTimes) { @@ -171,7 +171,9 @@ public static boolean stopTimesMatchSimplified(List futureStopTimes, L activeTime.arrival_time != futureTime.arrival_time || activeTime.departure_time != futureTime.departure_time || activeTime.stop_sequence != futureTime.stop_sequence || - !activeTime.stop_id.equals(futureTime.stop_id) + !activeTime.stop_id.equals(futureTime.stop_id) || + !Objects.equals(activeTime.location_group_id, futureTime.location_group_id) || + !Objects.equals(activeTime.location_id, futureTime.location_id) ) { return false; } diff --git a/src/test/java/com/conveyal/datatools/manager/jobs/FlexMergeFeedsJobTest.java b/src/test/java/com/conveyal/datatools/manager/jobs/FlexMergeFeedsJobTest.java new file mode 100644 index 000000000..0d04589dc --- /dev/null +++ b/src/test/java/com/conveyal/datatools/manager/jobs/FlexMergeFeedsJobTest.java @@ -0,0 +1,150 @@ +package com.conveyal.datatools.manager.jobs; + +import com.conveyal.datatools.DatatoolsTest; +import com.conveyal.datatools.TestUtils; +import com.conveyal.datatools.UnitTest; +import com.conveyal.datatools.manager.auth.Auth0Connection; +import com.conveyal.datatools.manager.auth.Auth0UserProfile; +import com.conveyal.datatools.manager.jobs.feedmerge.MergeFeedsType; +import com.conveyal.datatools.manager.models.FeedSource; +import com.conveyal.datatools.manager.models.FeedVersion; +import com.conveyal.datatools.manager.models.Project; +import com.conveyal.datatools.manager.persistence.Persistence; +import com.conveyal.gtfs.error.NewGTFSErrorType; +import com.conveyal.gtfs.loader.FeedLoadResult; +import com.conveyal.gtfs.loader.TableLoadResult; +import org.junit.jupiter.api.AfterAll; +import org.junit.jupiter.api.BeforeAll; +import org.junit.jupiter.api.Test; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.io.IOException; +import java.sql.SQLException; +import java.util.Date; +import java.util.HashSet; +import java.util.Set; + +import static com.conveyal.datatools.TestUtils.assertThatFeedHasNoErrorsOfType; +import static com.conveyal.datatools.manager.models.FeedRetrievalMethod.MANUALLY_UPLOADED; +import static org.junit.jupiter.api.Assertions.assertEquals; + +/** + * Tests for the various {@link MergeFeedsJob} merge types. + */ +class FlexMergeFeedsJobTest extends UnitTest { + private static final Logger LOG = LoggerFactory.getLogger(FlexMergeFeedsJobTest.class); + private static final Auth0UserProfile user = Auth0UserProfile.createTestAdminUser(); + private static FeedVersion fakeAgencyWithFlexVersion1; + private static FeedVersion fakeAgencyWithFlexVersion2; + private static Project project; + + /** + * Prepare and start a testing-specific web server + */ + @BeforeAll + static void setUp() throws IOException { + // start server if it isn't already running + DatatoolsTest.setUp(); + Auth0Connection.setAuthDisabled(true); + + // Create a project, feed sources, and feed versions to merge. + project = new Project(); + project.name = String.format("Test %s", new Date()); + Persistence.projects.create(project); + + FeedSource flexAgencyA = new FeedSource("FLEX-AGENCY-A", project.id, MANUALLY_UPLOADED); + Persistence.feedSources.create(flexAgencyA); + fakeAgencyWithFlexVersion1 = TestUtils.createFeedVersion(flexAgencyA, TestUtils.zipFolderFiles("fake-agency-with-flex-version-1")); + + FeedSource flexAgencyB = new FeedSource("FLEX-AGENCY-B", project.id, MANUALLY_UPLOADED); + Persistence.feedSources.create(flexAgencyB); + fakeAgencyWithFlexVersion2 = TestUtils.createFeedVersion(flexAgencyB, TestUtils.zipFolderFiles("fake-agency-with-flex-version-2")); + } + + /** + * Delete project on tear down (feed sources/versions will also be deleted). + */ + @AfterAll + static void tearDown() { + if (project != null) { + project.delete(); + } + Auth0Connection.setAuthDisabled(Auth0Connection.getDefaultAuthDisabled()); + } + + /** + * Ensures that a regional feed merge will produce a feed that includes all entities from each feed. + */ + @Test + void canMergeRegional() throws SQLException { + // Set up list of feed versions to merge. + Set versions = new HashSet<>(); + versions.add(fakeAgencyWithFlexVersion1); + versions.add(fakeAgencyWithFlexVersion2); + FeedVersion mergedVersion = regionallyMergeVersions(versions); + + FeedLoadResult r1 = fakeAgencyWithFlexVersion1.feedLoadResult; + FeedLoadResult r2 = fakeAgencyWithFlexVersion2.feedLoadResult; + FeedLoadResult merged = mergedVersion.feedLoadResult; + + // Ensure the feed has the row counts we expect. + assertRowCount(r1.agency, r2.agency, merged.agency, "Agency"); + assertRowCount(r1.attributions, r2.attributions, merged.attributions, "Attributions"); + assertRowCount(r1.bookingRules, r2.bookingRules, merged.bookingRules, "Booking rules"); + assertRowCount(r1.calendar, r2.calendar, merged.calendar, "Calendar"); + assertRowCount(r1.calendarDates, r2.calendarDates, merged.calendarDates, "Calendar dates"); + assertRowCount(r1.fareAttributes, r2.fareAttributes, merged.fareAttributes, "Fare attributes"); + assertRowCount(r1.fareRules, r2.fareRules, merged.fareRules, "Fare rules"); + assertRowCount(r1.frequencies, r2.frequencies, merged.frequencies, "Frequencies"); + // For GeoJSON locations, the change from https://github.com/ibi-group/gtfs-lib/pull/28 + // results in an additional two locations registered into feed versions 1 and 2 each. + // These additional locations are flagged with an error because the location type is not polygon or linestring, + // and they are not stored or exported during merge, so we subtract these locations from the expected count. + // The expected locations in the merged feed are unchanged. + assertEquals( + r1.locations.rowCount - r1.locations.errorCount + r2.locations.rowCount - r2.locations.errorCount, + merged.locations.rowCount + ); + + assertRowCount(r1.locationGroup, r2.locationGroup, merged.locationGroup, "Location Groups"); + assertRowCount(r1.locationGroupStops, r2.locationGroupStops, merged.locationGroupStops, "Location Group Stops"); + assertRowCount(r1.locationShapes, r2.locationShapes, merged.locationShapes, "Location shapes"); + assertRowCount(r1.routes, r2.routes, merged.routes, "Routes"); + assertRowCount(r1.shapes, r2.shapes, merged.shapes, "Shapes"); + assertRowCount(r1.stops, r2.stops, merged.stops, "Stops"); + assertRowCount(r1.stopTimes, r2.stopTimes, merged.stopTimes, "Stop times"); + assertRowCount(r1.trips, r2.trips, merged.trips, "Trips"); + assertRowCount(r1.translations, r2.translations, merged.translations, "Translations"); + + // Ensure there are no referential integrity errors, duplicate ID, or wrong number of fields errors. + assertThatFeedHasNoErrorsOfType( + mergedVersion.namespace, + NewGTFSErrorType.REFERENTIAL_INTEGRITY.toString(), + NewGTFSErrorType.DUPLICATE_ID.toString(), + NewGTFSErrorType.WRONG_NUMBER_OF_FIELDS.toString() + ); + } + + /** + * Merges a set of FeedVersions and then creates a new FeedSource and FeedVersion of the merged feed. + */ + private FeedVersion regionallyMergeVersions(Set versions) { + MergeFeedsJob mergeFeedsJob = new MergeFeedsJob(user, versions, project.id, MergeFeedsType.REGIONAL); + // Run the job in this thread (we're not concerned about concurrency here). + mergeFeedsJob.run(); + LOG.info("Regional merged file: {}", mergeFeedsJob.mergedVersion.retrieveGtfsFile().getAbsolutePath()); + return mergeFeedsJob.mergedVersion; + } + + /** + * Helper method to confirm that the sum of the two feed table rows match the merged feed table rows. + */ + private void assertRowCount(TableLoadResult feedOne, TableLoadResult feedTwo, TableLoadResult feedMerged, String entity) { + assertEquals( + feedOne.rowCount + feedTwo.rowCount, + feedMerged.rowCount, + String.format("%s count for merged feed should equal the sum for the versions merged.", entity) + ); + } +} diff --git a/src/test/java/com/conveyal/datatools/manager/jobs/NormalizeFieldTransformJobTest.java b/src/test/java/com/conveyal/datatools/manager/jobs/NormalizeFieldTransformJobTest.java index 6a35edb56..a5e7c8863 100644 --- a/src/test/java/com/conveyal/datatools/manager/jobs/NormalizeFieldTransformJobTest.java +++ b/src/test/java/com/conveyal/datatools/manager/jobs/NormalizeFieldTransformJobTest.java @@ -19,12 +19,18 @@ import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.BeforeAll; import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.Arguments; +import org.junit.jupiter.params.provider.MethodSource; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; import java.io.BufferedReader; import java.io.IOException; import java.io.InputStream; import java.io.InputStreamReader; import java.util.Date; +import java.util.stream.Stream; import java.util.zip.ZipEntry; import java.util.zip.ZipFile; @@ -35,13 +41,11 @@ public class NormalizeFieldTransformJobTest extends DatatoolsTest { - private static final String TABLE_NAME = "routes"; - private static final String FIELD_NAME = "route_long_name"; + private static final Logger LOG = LoggerFactory.getLogger(NormalizeFieldTransformJobTest.class); private static Project project; private static FeedSource feedSource; private FeedVersion targetVersion; - /** * Initialize Data Tools and set up a simple feed source and project. */ @@ -78,42 +82,50 @@ public void tearDownTest() { * Test that a {@link NormalizeFieldTransformation} will successfully complete. * FIXME: On certain Windows machines, this test fails. */ - @Test - public void canNormalizeField() throws IOException { - // Create transform. - // In this test, as an illustration, replace "Route" with the "Rte" abbreviation in routes.txt. - FeedTransformation transformation = NormalizeFieldTransformationTest.createTransformation( - TABLE_NAME, FIELD_NAME, null, Lists.newArrayList( - new Substitution("Route", "Rte") - ) - ); - initializeFeedSource(transformation); + @ParameterizedTest + @MethodSource("createNormalizedFieldCases") + void canNormalizeField(TransformationCase transformationCase) throws IOException { + initializeFeedSource(transformationCase.table, createTransformation(transformationCase)); // Create target version that the transform will operate on. targetVersion = createFeedVersion( feedSource, - zipFolderFiles("fake-agency-with-only-calendar") + zipFolderFiles("fake-agency-for-field-normalizing") ); try (ZipFile zip = new ZipFile(targetVersion.retrieveGtfsFile())) { - // Check that new version has routes table modified. - ZipEntry entry = zip.getEntry(TABLE_NAME + ".txt"); - assertNotNull(entry); - - // Scan the first data row in routes.txt and check that the substitution - // that was defined in setUp was done. - try ( - InputStream stream = zip.getInputStream(entry); - InputStreamReader streamReader = new InputStreamReader(stream); - BufferedReader reader = new BufferedReader(streamReader) - ) { - String[] columns = reader.readLine().split(","); - int fieldIndex = ArrayUtils.indexOf(columns, FIELD_NAME); - - String row1 = reader.readLine(); - String[] row1Fields = row1.split(","); - assertTrue(row1Fields[fieldIndex].startsWith("Rte "), row1); - } + // Check that new version has expected modifications. + checkTableForModification(zip, transformationCase); + } + } + + private static Stream createNormalizedFieldCases() { + return Stream.of( + Arguments.of(new TransformationCase("routes", "route_long_name", "Route", "Rte")), + Arguments.of(new TransformationCase("booking_rules", "pickup_message", "Message", "Msg")) + ); + } + + private void checkTableForModification(ZipFile zip, TransformationCase transformationCase) throws IOException { + String tableName = transformationCase.table + ".txt"; + LOG.info("Getting table {} from zip {}", tableName, zip.getName()); + // Check that the new version has been modified. + ZipEntry entry = zip.getEntry(tableName); + assertNotNull(entry); + + // Scan the first data row and check that the substitution that was defined in the set-up was done. + try ( + InputStream stream = zip.getInputStream(entry); + InputStreamReader streamReader = new InputStreamReader(stream); + BufferedReader reader = new BufferedReader(streamReader) + ) { + String[] columns = reader.readLine().split(","); + int fieldIndex = ArrayUtils.indexOf(columns, transformationCase.fieldName); + + String rowOne = reader.readLine(); + assertNotNull(rowOne, String.format("First row in table %s is null!", transformationCase.table)); + String[] row1Fields = rowOne.split(","); + assertTrue(row1Fields[fieldIndex].contains(transformationCase.replacement), rowOne); } } @@ -121,16 +133,16 @@ public void canNormalizeField() throws IOException { * Test that a {@link NormalizeFieldTransformation} will fail if invalid substitution patterns are provided. */ @Test - public void canHandleInvalidSubstitutionPatterns() throws IOException { + void canHandleInvalidSubstitutionPatterns() throws IOException { // Create transform. // In this test, we provide an invalid pattern '\Cir\b' (instead of '\bCir\b'), // when trying to replace e.g. 'Piedmont Cir' with 'Piedmont Circle'. FeedTransformation transformation = NormalizeFieldTransformationTest.createTransformation( - TABLE_NAME, FIELD_NAME, null, Lists.newArrayList( + "routes", "route_long_name", null, Lists.newArrayList( new Substitution("\\Cir\\b", "Circle") ) ); - initializeFeedSource(transformation); + initializeFeedSource("routes", transformation); // Create target version that the transform will operate on. targetVersion = createFeedVersion( @@ -143,16 +155,27 @@ public void canHandleInvalidSubstitutionPatterns() throws IOException { assertTrue(targetVersion.hasCriticalErrors()); } + private FeedTransformation createTransformation(TransformationCase transformationCase) { + return NormalizeFieldTransformationTest + .createTransformation( + transformationCase.table, + transformationCase.fieldName, + null, + Lists.newArrayList( + new Substitution(transformationCase.pattern, transformationCase.replacement) + ) + ); + } + /** * Create and persist a feed source using the given transformation. */ - private static void initializeFeedSource(FeedTransformation transformation) { - FeedTransformRules transformRules = new FeedTransformRules(transformation); + private void initializeFeedSource(String table, FeedTransformation transformation) { // Create feed source with above transform. - feedSource = new FeedSource("Normalize Field Test Feed", project.id, FeedRetrievalMethod.MANUALLY_UPLOADED); + feedSource = new FeedSource(table + " Normalize Field Test Feed", project.id, FeedRetrievalMethod.MANUALLY_UPLOADED); feedSource.deployable = false; - feedSource.transformRules.add(transformRules); + feedSource.transformRules.add(new FeedTransformRules(transformation)); Persistence.feedSources.create(feedSource); } @@ -160,8 +183,22 @@ private static void initializeFeedSource(FeedTransformation