Initial commit
This commit is contained in:
@@ -0,0 +1,423 @@
|
||||
require 'carto/importer/table_setup'
|
||||
|
||||
module CartoDB
|
||||
module Synchronization
|
||||
class Adapter
|
||||
STATEMENT_TIMEOUT = (1.hour * 1000).freeze
|
||||
DESTINATION_SCHEMA = 'public'.freeze
|
||||
THE_GEOM = 'the_geom'.freeze
|
||||
OVERWRITE_ERROR = 2013
|
||||
|
||||
def initialize(table_name, runner, database, user, overviews_creator, synchronization_id)
|
||||
@table_name = table_name
|
||||
@runner = runner
|
||||
@database = database
|
||||
@user = user
|
||||
@overviews_creator = overviews_creator
|
||||
@failed = false
|
||||
@table_setup = ::Carto::Importer::TableSetup.new(
|
||||
user: user,
|
||||
overviews_creator: overviews_creator,
|
||||
log: runner.log
|
||||
)
|
||||
@error_code = nil
|
||||
@synchronization_id = synchronization_id
|
||||
end
|
||||
|
||||
def run(&tracker)
|
||||
runner.run(&tracker)
|
||||
result = runner.results.select(&:success?).first
|
||||
|
||||
if runner.remote_data_updated?
|
||||
if result.nil?
|
||||
data_for_exception = "Expecting success data for table '#{table_name}'\nResults:#{runner.results.to_s}\n"
|
||||
data_for_exception << "1st result:#{runner.results.first.inspect}"
|
||||
raise data_for_exception
|
||||
end
|
||||
|
||||
Carto::GhostTablesManager.run_synchronized(
|
||||
user.id, attempts: 10, timeout: 3000,
|
||||
message: "Couldn't acquire bolt to register. Registering sync without bolt",
|
||||
user: user,
|
||||
synchronization_id: @synchronization_id
|
||||
) do
|
||||
move_to_schema(result)
|
||||
geo_type = fix_the_geom_type!(user.database_schema, result.table_name)
|
||||
import_cleanup(user.database_schema, result.table_name)
|
||||
@table_setup.cartodbfy(result.table_name)
|
||||
overwrite(user.database_schema, table_name, result, geo_type)
|
||||
setup_table(table_name, geo_type)
|
||||
@table_setup.recreate_overviews(table_name)
|
||||
end
|
||||
end
|
||||
self
|
||||
rescue => exception
|
||||
@failed = true
|
||||
puts '=================='
|
||||
puts exception.to_s
|
||||
puts exception.backtrace
|
||||
puts '=================='
|
||||
drop(result.table_name) if result && exists?(result.table_name)
|
||||
raise exception
|
||||
end
|
||||
|
||||
def user
|
||||
@user
|
||||
end
|
||||
|
||||
def overwrite(schema, table_name, result, geo_type)
|
||||
# Determine what kind of overwrite to perform
|
||||
# overwrite_replace substitutes the existing table by the new one,
|
||||
# so any modifications since previous import/sync are lost.
|
||||
# overwrite_sync will preserve columns added since the import/sync,
|
||||
# and the geometry as well if the new table doesn't have it (nil geo_type)
|
||||
|
||||
# For the time being the latter method will only be used with tables
|
||||
# that have had the geocoder analysis applied, resulting in an column
|
||||
# named carto_geocode_hash being present.
|
||||
# TODO: we could perform the sync if there's any column named `_carto_*`
|
||||
# (carto_geocode_hash would need be renamed as _carto_geocode_hash)
|
||||
sync = has_column(schema, table_name, 'carto_geocode_hash')
|
||||
|
||||
if sync
|
||||
overwrite_sync(schema, table_name, result, geo_type)
|
||||
else
|
||||
overwrite_replace(schema, table_name, result)
|
||||
end
|
||||
end
|
||||
|
||||
def overwrite_sync(schema, table_name, result, geo_type)
|
||||
return false unless runner.remote_data_updated?
|
||||
|
||||
# NOTE the import table is already moved to the user schema;
|
||||
# this was done (#7543) because the cartodbfication performs
|
||||
# queries on CDB_UserQuotaSize and other functions expected
|
||||
# to exist in the schema of the table.
|
||||
qualified_result_table_name = %{"#{schema}"."#{result.table_name}"}
|
||||
skip_columns = '{the_geom, the_geom_webmercator}'
|
||||
|
||||
database.transaction do
|
||||
if geo_type.nil?
|
||||
# If there's no geometry in the result table, not worth
|
||||
# syncing. Maybe those were added via geocoding
|
||||
database.execute(%{
|
||||
SELECT cartodb.CDB_SyncTable(
|
||||
'#{qualified_result_table_name}',
|
||||
'#{schema}', '#{table_name}',
|
||||
'#{skip_columns}'
|
||||
)})
|
||||
else
|
||||
database.execute(%{
|
||||
SELECT cartodb.CDB_SyncTable(
|
||||
'#{qualified_result_table_name}',
|
||||
'#{schema}', '#{table_name}'
|
||||
)})
|
||||
end
|
||||
end
|
||||
|
||||
drop(result.table_name) if exists?(result.table_name)
|
||||
|
||||
# TODO not sure whether these two are needed
|
||||
@table_setup.fix_oid(table_name)
|
||||
@table_setup.update_cdb_tablemetadata(table_name)
|
||||
rescue => exception
|
||||
@error_code = OVERWRITE_ERROR
|
||||
puts "Sync overwrite ERROR: #{exception.message}: #{exception.backtrace.join}"
|
||||
|
||||
# Gets all attributes in the result except for 'log_trace', as it is too long for Rollbar
|
||||
result_hash = CartoDB::Importer2::Result::ATTRIBUTES.map { |m| [m, result.send(m)] if m != 'log_trace' }
|
||||
.compact.to_h
|
||||
CartoDB::Logger.error(message: 'Error in sync overwrite',
|
||||
exception: exception,
|
||||
user: user,
|
||||
table: table_name,
|
||||
result: result_hash)
|
||||
drop(result.table_name) if exists?(result.table_name)
|
||||
raise exception
|
||||
end
|
||||
|
||||
def overwrite_replace(schema, table_name, result)
|
||||
return false unless runner.remote_data_updated?
|
||||
|
||||
@table_setup.copy_privileges(schema, table_name, schema, result.table_name)
|
||||
index_statements = @table_setup.generate_index_statements(schema, table_name)
|
||||
|
||||
temporary_name = temporary_name_for(result.table_name)
|
||||
database.transaction do
|
||||
rename(table_name, temporary_name) if exists?(table_name)
|
||||
drop(temporary_name) if exists?(temporary_name)
|
||||
rename(result.table_name, table_name)
|
||||
end
|
||||
@table_setup.fix_oid(table_name)
|
||||
@table_setup.update_cdb_tablemetadata(table_name)
|
||||
@table_setup.run_index_statements(index_statements, @database)
|
||||
rescue => exception
|
||||
@error_code = OVERWRITE_ERROR
|
||||
puts "Sync overwrite ERROR: #{exception.message}: #{exception.backtrace.join}"
|
||||
|
||||
# Gets all attributes in the result except for 'log_trace', as it is too long for Rollbar
|
||||
result_hash = CartoDB::Importer2::Result::ATTRIBUTES.map { |m| [m, result.send(m)] if m != 'log_trace' }
|
||||
.compact.to_h
|
||||
CartoDB::Logger.error(message: 'Error in sync overwrite',
|
||||
exception: exception,
|
||||
user: user,
|
||||
table: table_name,
|
||||
result: result_hash)
|
||||
drop(result.table_name) if exists?(result.table_name)
|
||||
raise exception
|
||||
end
|
||||
|
||||
def setup_table(table_name, geo_type)
|
||||
table = Carto::UserTable.find(user.tables.where(name: table_name).first.id).service
|
||||
|
||||
table.force_schema = true
|
||||
|
||||
table.import_to_cartodb(table_name)
|
||||
table.schema(reload: true)
|
||||
table.reload
|
||||
|
||||
# We send the detected geometry type to avoid manipulating geoms twice
|
||||
# set_the_geom_column! should just edit the metadata with the specified type
|
||||
table.send :set_the_geom_column!, geo_type
|
||||
table.save
|
||||
rescue => exception
|
||||
CartoDB::Logger.error(message: 'Error in setup cartodbfy',
|
||||
exception: exception,
|
||||
user: user,
|
||||
table: table_name)
|
||||
ensure
|
||||
@table_setup.fix_oid(table_name)
|
||||
end
|
||||
|
||||
def has_column(schema_name, table_name, column_name)
|
||||
qualified_table_name = "\"#{schema_name}\".#{table_name}"
|
||||
sql = %{
|
||||
SELECT TRUE as has_column FROM pg_catalog.pg_attribute a
|
||||
WHERE
|
||||
a.attname = '#{column_name}'
|
||||
AND a.attnum > 0
|
||||
AND NOT a.attisdropped
|
||||
AND a.attrelid = '#{qualified_table_name}'::regclass::oid
|
||||
LIMIT 1
|
||||
}
|
||||
result = user.in_database[sql].first
|
||||
result && result[:has_column]
|
||||
end
|
||||
|
||||
# From Table#get_the_geom_type!, adapted to unregistered tables
|
||||
# returns type to run Table#get_the_geom_type! afterwards again, which
|
||||
# saves the type in table metadata
|
||||
def fix_the_geom_type!(schema_name, table_name)
|
||||
qualified_table_name = "\"#{schema_name}\".#{table_name}"
|
||||
|
||||
type = nil
|
||||
the_geom_data = user.in_database[%Q{
|
||||
SELECT a.attname, t.typname
|
||||
FROM pg_attribute a, pg_type t
|
||||
WHERE attrelid = '#{qualified_table_name}'::regclass
|
||||
AND attname = '#{THE_GEOM}'
|
||||
AND a.atttypid = t.oid
|
||||
AND a.attstattarget < 0
|
||||
LIMIT 1
|
||||
}].first
|
||||
return nil unless the_geom_data
|
||||
|
||||
if the_geom_data[:typname] != 'geometry'
|
||||
user.in_database.execute %{
|
||||
ALTER TABLE #{qualified_table_name} RENAME COLUMN "#{THE_GEOM}" TO "the_geom_str"
|
||||
}
|
||||
return nil
|
||||
end
|
||||
|
||||
geom_type = user.in_database[%Q{
|
||||
SELECT GeometryType(#{THE_GEOM})
|
||||
FROM #{qualified_table_name}
|
||||
WHERE #{THE_GEOM} IS NOT null
|
||||
LIMIT 1
|
||||
}].first
|
||||
|
||||
type = geom_type[:geometrytype].to_s.downcase if geom_type
|
||||
|
||||
# if the geometry is MULTIPOINT we convert it to POINT
|
||||
if type == 'multipoint'
|
||||
user.db_service.in_database_direct_connection(statement_timeout: STATEMENT_TIMEOUT) do |user_database|
|
||||
user_database.run("SELECT public.AddGeometryColumn('#{schema_name}', '#{table_name}','the_geom_simple',4326, 'GEOMETRY', 2);")
|
||||
user_database.run(%Q{UPDATE #{qualified_table_name} SET the_geom_simple = ST_GeometryN(the_geom,1);})
|
||||
user_database.run("SELECT DropGeometryColumn('#{schema_name}', '#{table_name}','the_geom');")
|
||||
user_database.run(%Q{ALTER TABLE #{qualified_table_name} RENAME COLUMN the_geom_simple TO the_geom;})
|
||||
end
|
||||
type = 'point'
|
||||
end
|
||||
|
||||
# if the geometry is LINESTRING or POLYGON we convert it to MULTILINESTRING or MULTIPOLYGON
|
||||
if %w(linestring polygon).include?(type)
|
||||
user.db_service.in_database_direct_connection(statement_timeout: STATEMENT_TIMEOUT) do |user_database|
|
||||
user_database.run("SELECT public.AddGeometryColumn('#{schema_name}', '#{table_name}','the_geom_simple',4326, 'GEOMETRY', 2);")
|
||||
user_database.run(%Q{UPDATE #{qualified_table_name} SET the_geom_simple = ST_Multi(the_geom);})
|
||||
user_database.run("SELECT DropGeometryColumn('#{schema_name}', '#{table_name}','the_geom');")
|
||||
user_database.run(%Q{ALTER TABLE #{qualified_table_name} RENAME COLUMN the_geom_simple TO the_geom;})
|
||||
|
||||
type = user_database[%Q{
|
||||
SELECT GeometryType(#{THE_GEOM})
|
||||
FROM #{qualified_table_name}
|
||||
WHERE #{THE_GEOM} IS NOT null
|
||||
LIMIT 1
|
||||
}].first[:geometrytype]
|
||||
end
|
||||
end
|
||||
|
||||
type
|
||||
end
|
||||
|
||||
# From Table#import_cleanup, with column schema checks adapted to unregistered tables
|
||||
def import_cleanup(schema_name, table_name)
|
||||
qualified_table_name = "\"#{schema_name}\".#{table_name}"
|
||||
|
||||
user.db_service.in_database_direct_connection(statement_timeout: STATEMENT_TIMEOUT) do |user_database|
|
||||
|
||||
# For consistency with regular imports, also eases testing
|
||||
Table.sanitize_columns(table_name, {database_schema: schema_name, connection: user_database})
|
||||
|
||||
# When tables are created using ogr2ogr they are added a ogc_fid or gid primary key
|
||||
# In that case:
|
||||
# - If cartodb_id already exists, remove ogc_fid
|
||||
# - If cartodb_id does not exist, treat this field as the auxiliary column
|
||||
aux_cartodb_id_column = [:ogc_fid, :gid].find do |col|
|
||||
valid_cartodb_id_candidate?(user, table_name, qualified_table_name, col)
|
||||
end
|
||||
|
||||
# Remove primary key
|
||||
existing_pk = user_database[%Q{
|
||||
SELECT c.conname AS pk_name
|
||||
FROM pg_class r, pg_constraint c, pg_namespace n
|
||||
WHERE r.oid = c.conrelid AND contype='p' AND relname = '#{table_name}'
|
||||
AND r.relnamespace = n.oid and n.nspname= '#{schema_name}'
|
||||
}].first
|
||||
existing_pk = existing_pk[:pk_name] unless existing_pk.nil?
|
||||
user_database.run(%Q{
|
||||
ALTER TABLE #{qualified_table_name} DROP CONSTRAINT "#{existing_pk}"
|
||||
}) unless existing_pk.nil?
|
||||
|
||||
# All normal fields casted to text
|
||||
varchar_columns = user_database[%Q{
|
||||
SELECT a.attname, t.typname
|
||||
FROM pg_attribute a, pg_type t
|
||||
WHERE attrelid = '#{qualified_table_name}'::regclass
|
||||
AND typname = 'varchar'
|
||||
AND a.atttypid = t.oid
|
||||
}].all
|
||||
|
||||
varchar_columns.each do |column|
|
||||
user_database.run(%Q{ALTER TABLE #{qualified_table_name} ALTER COLUMN "#{column[:attname]}" TYPE text})
|
||||
end
|
||||
|
||||
# If there's an auxiliary column, copy to cartodb_id and restart the sequence to the max(cartodb_id)+1
|
||||
if aux_cartodb_id_column.present?
|
||||
begin
|
||||
already_had_cartodb_id = false
|
||||
user_database.run(%Q{ALTER TABLE #{qualified_table_name} ADD COLUMN cartodb_id SERIAL})
|
||||
rescue
|
||||
already_had_cartodb_id = true
|
||||
end
|
||||
unless already_had_cartodb_id
|
||||
user_database.run(%Q{UPDATE #{qualified_table_name} SET cartodb_id = CAST(#{aux_cartodb_id_column} AS INTEGER)})
|
||||
cartodb_id_sequence_name = user_database["SELECT pg_get_serial_sequence('#{schema_name}.#{table_name}', 'cartodb_id')"].first[:pg_get_serial_sequence]
|
||||
max_cartodb_id = user_database[%Q{SELECT max(cartodb_id) FROM #{qualified_table_name}}].first[:max]
|
||||
# only reset the sequence on real imports.
|
||||
|
||||
if max_cartodb_id
|
||||
user_database.run("ALTER SEQUENCE #{cartodb_id_sequence_name} RESTART WITH #{max_cartodb_id + 1}")
|
||||
end
|
||||
end
|
||||
user_database.run(%Q{ALTER TABLE #{qualified_table_name} DROP COLUMN #{aux_cartodb_id_column}})
|
||||
end
|
||||
|
||||
end
|
||||
end
|
||||
|
||||
def success?
|
||||
(!@failed && runner.success?)
|
||||
end
|
||||
|
||||
def etag
|
||||
runner.etag
|
||||
end
|
||||
|
||||
def last_modified
|
||||
runner.last_modified
|
||||
end
|
||||
|
||||
def checksum
|
||||
runner.checksum
|
||||
end
|
||||
|
||||
def move_to_schema(result, schema=DESTINATION_SCHEMA)
|
||||
# The new table to sync is moved to user schema to allow CartoDBfication.
|
||||
# This temporary table should not be registered (check ghost_tables_manager.rb)
|
||||
return self if schema == result.schema
|
||||
database.execute(%Q{
|
||||
ALTER TABLE "#{result.schema}"."#{result.table_name}"
|
||||
SET SCHEMA "#{user.database_schema}"
|
||||
})
|
||||
end
|
||||
|
||||
def rename(current_name, new_name)
|
||||
database.execute(%Q{
|
||||
ALTER TABLE "#{user.database_schema}"."#{current_name}"
|
||||
RENAME TO #{new_name}
|
||||
})
|
||||
end
|
||||
|
||||
def drop(table_name)
|
||||
database.execute(%Q(DROP TABLE "#{user.database_schema}"."#{table_name}"))
|
||||
end
|
||||
|
||||
def exists?(table_name)
|
||||
database.table_exists?(table_name)
|
||||
end
|
||||
|
||||
def results
|
||||
runner.results
|
||||
end
|
||||
|
||||
def error_code
|
||||
@error_code || runner.results.map(&:error_code).compact.first
|
||||
end
|
||||
|
||||
def runner_log_trace
|
||||
runner.results.map(&:log_trace).compact.first
|
||||
end
|
||||
|
||||
def error_message
|
||||
''
|
||||
end
|
||||
|
||||
def temporary_name_for(table_name)
|
||||
"#{table_name}_to_be_deleted"
|
||||
end
|
||||
|
||||
private
|
||||
|
||||
def valid_cartodb_id_candidate?(user, table_name, qualified_table_name, col_name)
|
||||
return false unless column_names(user, table_name).include?(col_name)
|
||||
user.transaction_with_timeout(statement_timeout: STATEMENT_TIMEOUT, as: :superuser) do |db|
|
||||
return db["SELECT 1 FROM #{qualified_table_name} WHERE #{col_name} IS NULL LIMIT 1"].first.nil?
|
||||
end
|
||||
end
|
||||
|
||||
def column_names(user, table_name)
|
||||
user.in_database.schema(table_name, schema: user.database_schema).map { |row| row[0] }
|
||||
rescue => e
|
||||
CartoDB::Logger.error(
|
||||
message: 'Error in column_names from sync adapter',
|
||||
exception: e,
|
||||
user: user,
|
||||
table: table_name
|
||||
)
|
||||
[]
|
||||
end
|
||||
|
||||
attr_reader :table_name, :runner, :database, :user
|
||||
end
|
||||
end
|
||||
end
|
||||
@@ -0,0 +1,75 @@
|
||||
require 'set'
|
||||
require_relative './member'
|
||||
require_relative '../../../services/data-repository/structures/collection'
|
||||
|
||||
module CartoDB
|
||||
module Synchronization
|
||||
SIGNATURE = 'synchronizations'
|
||||
|
||||
class << self
|
||||
attr_accessor :repository
|
||||
end
|
||||
|
||||
class Collection
|
||||
def initialize(attributes={}, options={})
|
||||
@collection = DataRepository::Collection.new(
|
||||
signature: SIGNATURE,
|
||||
repository: options.fetch(:repository, Synchronization.repository),
|
||||
member_class: Member
|
||||
)
|
||||
end #initialize
|
||||
|
||||
DataRepository::Collection::INTERFACE.each do |method_name|
|
||||
define_method(method_name) do |*arguments, &block|
|
||||
result = collection.send(method_name, *arguments, &block)
|
||||
return self if result.is_a?(DataRepository::Collection)
|
||||
result
|
||||
end
|
||||
end
|
||||
|
||||
def fetch(filters={})
|
||||
per_page_filter = filters.delete(:per_page)
|
||||
dataset = repository.collection(filters, [])
|
||||
self.total_entries = dataset.count
|
||||
dataset = repository.paginate(dataset, per_page_filter.present? ? {per_page:per_page_filter} : {})
|
||||
|
||||
collection.storage = Set.new(dataset.map { |attributes|
|
||||
Synchronization::Member.new(attributes)
|
||||
})
|
||||
|
||||
self
|
||||
end #fetch
|
||||
|
||||
def store
|
||||
#map { |member| member.fetch.store }
|
||||
self
|
||||
end #store
|
||||
|
||||
def destroy
|
||||
map(&:delete)
|
||||
self
|
||||
end #destroy
|
||||
|
||||
def to_poro
|
||||
map { |member| member.to_hash }
|
||||
end
|
||||
|
||||
attr_reader :total_entries
|
||||
|
||||
private
|
||||
|
||||
attr_reader :collection
|
||||
attr_writer :total_entries
|
||||
|
||||
def order(dataset, criteria={})
|
||||
return dataset if criteria.nil? || criteria.empty?
|
||||
dataset.order(*order_params_from(criteria))
|
||||
end #order
|
||||
|
||||
def order_params_from(criteria)
|
||||
criteria.map { |key, order| Sequel.send(order.to_sym, key.to_sym) }
|
||||
end #order_params_from
|
||||
end # Collection
|
||||
end # Synchronization
|
||||
end # CartoDB
|
||||
|
||||
@@ -0,0 +1,574 @@
|
||||
require 'virtus'
|
||||
require_relative 'adapter'
|
||||
require_relative '../../../services/importer/lib/importer'
|
||||
require_relative '../visualization/collection'
|
||||
require_relative '../../../services/importer/lib/importer/datasource_downloader'
|
||||
require_relative '../../../services/datasources/lib/datasources'
|
||||
require_relative '../log'
|
||||
require_relative '../../../services/importer/lib/importer/unp'
|
||||
require_relative '../../../services/importer/lib/importer/post_import_handler'
|
||||
require_relative '../../../services/importer/lib/importer/overviews'
|
||||
require_relative '../../../lib/cartodb/errors'
|
||||
require_relative '../../../lib/cartodb/import_error_codes'
|
||||
require_relative '../../../services/platform-limits/platform_limits'
|
||||
|
||||
require_dependency 'carto/configuration'
|
||||
|
||||
include CartoDB::Datasources
|
||||
|
||||
module CartoDB
|
||||
module Synchronization
|
||||
|
||||
class << self
|
||||
attr_accessor :repository
|
||||
end
|
||||
|
||||
class Member
|
||||
include Carto::Configuration
|
||||
include Virtus.model
|
||||
|
||||
MAX_RETRIES = 10
|
||||
MIN_INTERVAL_SECONDS = 15 * 60
|
||||
|
||||
# Seconds required between manual sync now
|
||||
SYNC_NOW_TIMESPAN = 900
|
||||
|
||||
STATE_CREATED = 'created'
|
||||
# Already at resque, waiting for slot
|
||||
STATE_QUEUED = 'queued'
|
||||
# Actually syncing
|
||||
STATE_SYNCING = 'syncing'
|
||||
STATE_SUCCESS = 'success'
|
||||
STATE_FAILURE = 'failure'
|
||||
|
||||
attribute :id, String
|
||||
attribute :name, String
|
||||
attribute :interval, Integer, default: 3600
|
||||
attribute :url, String
|
||||
attribute :state, String, default: STATE_CREATED
|
||||
attribute :user_id, String
|
||||
attribute :created_at, Time
|
||||
attribute :updated_at, Time
|
||||
attribute :run_at, Time
|
||||
attribute :ran_at, Time
|
||||
attribute :modified_at, Time
|
||||
attribute :etag, String
|
||||
attribute :checksum, String
|
||||
attribute :log_id, String
|
||||
attribute :error_code, Integer
|
||||
attribute :error_message, String
|
||||
attribute :retried_times, Integer, default: 0
|
||||
attribute :service_name, String
|
||||
attribute :service_item_id, String
|
||||
attribute :type_guessing, Boolean, default: true
|
||||
attribute :quoted_fields_guessing, Boolean, default: true
|
||||
attribute :content_guessing, Boolean, default: false
|
||||
attribute :visualization_id, String
|
||||
|
||||
def initialize(attributes={}, repository=Synchronization.repository)
|
||||
super(attributes)
|
||||
|
||||
@log = nil
|
||||
|
||||
self.log_trace = nil
|
||||
@repository = repository
|
||||
self.id ||= @repository.next_id
|
||||
self.state ||= STATE_CREATED
|
||||
self.ran_at ||= Time.now
|
||||
self.interval ||= 3600
|
||||
self.run_at ||= Time.now + interval
|
||||
self.retried_times ||= 0
|
||||
self.log_id ||= log.id unless log.nil?
|
||||
self.service_name ||= nil
|
||||
self.service_item_id ||= nil
|
||||
self.checksum ||= ''
|
||||
|
||||
raise InvalidInterval.new unless self.interval >= MIN_INTERVAL_SECONDS
|
||||
|
||||
end
|
||||
|
||||
def to_s
|
||||
"<CartoDB::Synchronization::Member id:\"#{@id}\" name:\"#{@name}\" ran_at:\"#{@ran_at}\" run_at:\"#{@run_at}\" " \
|
||||
"interval:\"#{@interval}\" state:\"#{@state}\" retried_times:\"#{@retried_times}\" log_id:\"#{self.log_id}\" " \
|
||||
"service_name:\"#{@service_name}\" service_item_id:\"#{@service_item_id}\" checksum:\"#{@checksum}\" " \
|
||||
"url:\"#{@url}\" error_code:\"#{@error_code}\" error_message:\"#{@error_message}\" modified_at:\"#{@modified_at}\" " \
|
||||
" user_id:\"#{@user_id}\" type_guessing:\"#{@type_guessing}\" " \
|
||||
"quoted_fields_guessing:\"#{@quoted_fields_guessing}\" visualization_id:\"#{@visualization_id}\">"
|
||||
end
|
||||
|
||||
def synchronizations_logger
|
||||
@@synchronizations_logger ||= CartoDB.unformatted_logger(log_file_path("synchronizations.log"))
|
||||
end
|
||||
|
||||
def interval=(seconds=3600)
|
||||
super(seconds.to_i)
|
||||
if seconds
|
||||
self.run_at = Time.now + (seconds.to_i || 3600)
|
||||
end
|
||||
seconds
|
||||
end
|
||||
|
||||
def store
|
||||
raise CartoDB::InvalidMember unless self.valid?
|
||||
set_timestamps
|
||||
repository.store(id, attributes.to_hash)
|
||||
self
|
||||
end
|
||||
|
||||
def fetch
|
||||
data = repository.fetch(id)
|
||||
raise KeyError if data.nil?
|
||||
self.attributes = data
|
||||
self
|
||||
end
|
||||
|
||||
def fetch_by_visualization_id
|
||||
data = repository.fetch(visualization_id, "visualization_id")
|
||||
raise KeyError if data.nil?
|
||||
self.attributes = data
|
||||
self
|
||||
end
|
||||
|
||||
def delete
|
||||
repository.delete(id)
|
||||
self.attributes.keys.each { |key| self.send("#{key}=", nil) }
|
||||
self
|
||||
end
|
||||
|
||||
def enqueue
|
||||
Resque.enqueue(Resque::SynchronizationJobs, job_id: id)
|
||||
self.error_code = nil
|
||||
self.error_message = nil
|
||||
self.state = CartoDB::Synchronization::Member::STATE_QUEUED
|
||||
self.store
|
||||
end
|
||||
|
||||
# @return bool
|
||||
def can_manually_sync?
|
||||
# Last sync ok, last sync failed, or too much time in queued state
|
||||
( self.state == STATE_SUCCESS ||
|
||||
self.state == STATE_FAILURE ||
|
||||
(self.state == STATE_QUEUED && self.updated_at + SYNC_NOW_TIMESPAN < Time.now)
|
||||
) && (self.ran_at + SYNC_NOW_TIMESPAN < Time.now)
|
||||
end
|
||||
|
||||
# @return bool
|
||||
def should_auto_sync?
|
||||
self.state == STATE_SUCCESS && (self.run_at < Time.now)
|
||||
end
|
||||
|
||||
# This should be joined with data_import to stop the madness of duplicated code
|
||||
def run
|
||||
importer = nil
|
||||
self.state = STATE_SYNCING
|
||||
self.store
|
||||
|
||||
# TODO: See if we can remove this code
|
||||
# First import is a "normal import" so still has no id, then run gets called and will get log first time
|
||||
# but we need this to fix old logs
|
||||
if log.nil?
|
||||
@log = CartoDB::Log.new(type: CartoDB::Log::TYPE_SYNCHRONIZATION, user_id: user.id)
|
||||
@log.store
|
||||
self.log_id = @log.id
|
||||
store
|
||||
else
|
||||
@log.type = CartoDB::Log::TYPE_SYNCHRONIZATION
|
||||
@log.clear
|
||||
@log.store
|
||||
end
|
||||
|
||||
if user.nil?
|
||||
raise "Couldn't instantiate synchronization user. Data: #{to_s}"
|
||||
end
|
||||
|
||||
raise "Can't run a synchronization for inactive user: #{user.username}" unless user.reload.active?
|
||||
|
||||
if !authorize?(user)
|
||||
raise CartoDB::Datasources::AuthError.new('User is not authorized to sync tables')
|
||||
end
|
||||
|
||||
runner = service_name == 'connector' ? get_connector : get_runner
|
||||
|
||||
database = user.in_database
|
||||
overviews_creator = CartoDB::Importer2::Overviews.new(runner, user)
|
||||
importer = CartoDB::Synchronization::Adapter.new(name, runner, database, user, overviews_creator, id)
|
||||
|
||||
importer.run
|
||||
self.ran_at = Time.now
|
||||
self.run_at = Time.now + interval
|
||||
|
||||
if importer.success?
|
||||
set_success_state_from(importer)
|
||||
else
|
||||
set_failure_state_from(importer)
|
||||
end
|
||||
|
||||
store
|
||||
|
||||
notify
|
||||
|
||||
rescue => exception
|
||||
if exception.is_a? CartoDB::Datasources::NotFoundDownloadError
|
||||
CartoDB::Logger.debug(exception: exception, sync_id: id)
|
||||
else
|
||||
CartoDB::Logger.error(exception: exception, sync_id: id)
|
||||
end
|
||||
log.append_and_store exception.message, truncate = false
|
||||
log.append exception.backtrace.join("\n"), truncate = false
|
||||
|
||||
if importer.nil?
|
||||
if exception.is_a?(NotFoundDownloadError)
|
||||
set_general_failure_state_from(exception, 1017, 'File not found, you must import it again')
|
||||
elsif exception.is_a?(CartoDB::Importer2::FileTooBigError)
|
||||
set_general_failure_state_from(exception, exception.error_code,
|
||||
CartoDB::IMPORTER_ERROR_CODES[exception.error_code][:title])
|
||||
elsif exception.is_a?(AuthError)
|
||||
set_general_failure_state_from(exception, 1011, 'Unauthorized')
|
||||
else
|
||||
set_general_failure_state_from(exception)
|
||||
end
|
||||
else
|
||||
set_failure_state_from(importer)
|
||||
end
|
||||
|
||||
store
|
||||
|
||||
if exception.is_a?(TokenExpiredOrInvalidError)
|
||||
begin
|
||||
user.oauths.remove(exception.service_name)
|
||||
rescue => ex
|
||||
log.append "Exception removing OAuth: #{ex.message}"
|
||||
log.append ex.backtrace
|
||||
end
|
||||
end
|
||||
notify
|
||||
self
|
||||
ensure
|
||||
CartoDB::PlatformLimits::Importer::UserConcurrentSyncsAmount.new(
|
||||
user: user, redis: { db: $users_metadata }
|
||||
).decrement!
|
||||
end
|
||||
|
||||
def get_runner
|
||||
downloader = get_downloader
|
||||
|
||||
post_import_handler = CartoDB::Importer2::PostImportHandler.new
|
||||
unless downloader.datasource.nil?
|
||||
case downloader.datasource.class::DATASOURCE_NAME
|
||||
when Url::ArcGIS::DATASOURCE_NAME
|
||||
post_import_handler.add_fix_geometries_task
|
||||
when Search::Twitter::DATASOURCE_NAME
|
||||
post_import_handler.add_transform_geojson_geom_column
|
||||
end
|
||||
end
|
||||
|
||||
runner = CartoDB::Importer2::Runner.new(
|
||||
pg: pg_options,
|
||||
downloader: downloader,
|
||||
log: log,
|
||||
user: user,
|
||||
unpacker: CartoDB::Importer2::Unp.new(Cartodb.config[:importer], Cartodb.config[:ogr2ogr]),
|
||||
post_import_handler: post_import_handler,
|
||||
importer_config: Cartodb.config[:importer]
|
||||
)
|
||||
runner.loader_options = ogr2ogr_options.merge content_guessing_options
|
||||
|
||||
runner.include_additional_errors_mapping(
|
||||
AuthError => 1011,
|
||||
DataDownloadError => 1011,
|
||||
TokenExpiredOrInvalidError => 1012,
|
||||
DatasourceBaseError => 1012,
|
||||
InvalidServiceError => 1012,
|
||||
MissingConfigurationError => 1012,
|
||||
UninitializedError => 1012
|
||||
)
|
||||
|
||||
runner
|
||||
end
|
||||
|
||||
def get_connector
|
||||
CartoDB::Importer2::ConnectorRunner.check_availability!(user)
|
||||
CartoDB::Importer2::ConnectorRunner.new(
|
||||
service_item_id,
|
||||
user: user,
|
||||
pg: pg_options,
|
||||
log: log
|
||||
)
|
||||
end
|
||||
|
||||
def notify
|
||||
sync_log = {
|
||||
'name' => self.name,
|
||||
'sync_time' => self.updated_at - self.created_at,
|
||||
'sync_timestamp' => Time.now,
|
||||
'user' => user.username,
|
||||
'queue_server' => `hostname`.strip,
|
||||
'resque_ppid' => Process.ppid,
|
||||
'state' => self.state,
|
||||
'user_timeout' => ::DataImport.http_timeout_for(user)
|
||||
}
|
||||
synchronizations_logger.info(sync_log.to_json)
|
||||
end
|
||||
|
||||
def get_downloader
|
||||
datasource_name = (service_name.nil? || service_name.size == 0) ? Url::PublicUrl::DATASOURCE_NAME : service_name
|
||||
if service_item_id.nil? || service_item_id.size == 0
|
||||
self.service_item_id = url
|
||||
end
|
||||
|
||||
datasource_provider = get_datasource(datasource_name)
|
||||
if datasource_provider.nil?
|
||||
raise CartoDB::DataSourceError.new("Datasource #{datasource_name} could not be instantiated")
|
||||
end
|
||||
|
||||
if service_item_id.nil?
|
||||
raise CartoDB::DataSourceError.new("Datasource #{datasource_name} without item id")
|
||||
end
|
||||
|
||||
log.append "Fetching datasource #{datasource_provider.to_s} metadata for item id #{service_item_id} from user #{user.id}"
|
||||
metadata = datasource_provider.get_resource_metadata(service_item_id)
|
||||
|
||||
if hit_platform_limit?(datasource_provider, metadata, user)
|
||||
raise CartoDB::Importer2::FileTooBigError.new(metadata.inspect)
|
||||
end
|
||||
|
||||
if datasource_provider.providers_download_url?
|
||||
metadata_url = metadata[:url]
|
||||
resource_url = (metadata_url.present? && datasource_provider.providers_download_url?) ? metadata_url : url
|
||||
|
||||
raise CartoDB::DataSourceError.new("Missing resource URL to download. Data:#{self}") unless resource_url
|
||||
|
||||
log.append "File will be downloaded from #{resource_url}"
|
||||
|
||||
http_options = {
|
||||
http_timeout: DataImport.http_timeout_for(user),
|
||||
etag: etag,
|
||||
last_modified: modified_at,
|
||||
checksum: checksum,
|
||||
verify_ssl_cert: false
|
||||
}
|
||||
|
||||
CartoDB::Importer2::Downloader.new(user.id,
|
||||
resource_url,
|
||||
http_options,
|
||||
importer_config: Cartodb.config[:importer])
|
||||
else
|
||||
log.append 'Downloading file data from datasource'
|
||||
|
||||
options = {
|
||||
http_timeout: DataImport.http_timeout_for(user),
|
||||
checksum: checksum,
|
||||
importer_config: Cartodb.config[:importer],
|
||||
user_id: user.id
|
||||
}
|
||||
|
||||
CartoDB::Importer2::DatasourceDownloader.new(datasource_provider, metadata, options, log)
|
||||
end
|
||||
end
|
||||
|
||||
def hit_platform_limit?(datasource, metadata, user)
|
||||
if datasource.has_resource_size?(metadata)
|
||||
CartoDB::PlatformLimits::Importer::InputFileSize.new({ user: user })
|
||||
.is_over_limit!(metadata[:size])
|
||||
else
|
||||
false
|
||||
end
|
||||
end
|
||||
|
||||
def set_success_state_from(importer)
|
||||
log.append_and_store '******** synchronization succeeded ********'
|
||||
self.log_trace = importer.runner_log_trace
|
||||
self.state = STATE_SUCCESS
|
||||
self.etag = importer.etag
|
||||
self.checksum = importer.checksum
|
||||
self.error_code = nil
|
||||
self.error_message = nil
|
||||
self.retried_times = 0
|
||||
self.run_at = Time.now + interval
|
||||
self.modified_at = importer.last_modified
|
||||
geocode_table
|
||||
rescue => exception
|
||||
CartoDB::Logger.error(exception: exception,
|
||||
message: 'Error updating state for sync table',
|
||||
sync_id: self.id)
|
||||
self
|
||||
end
|
||||
|
||||
def set_failure_state_from(importer)
|
||||
log.append_and_store '******** synchronization failed ********'
|
||||
self.log_trace = importer.runner_log_trace
|
||||
log.append "*** Runner log: #{self.log_trace} \n***" unless self.log_trace.nil?
|
||||
self.state = STATE_FAILURE
|
||||
self.error_code = importer.error_code.blank? ? 9999 : importer.error_code
|
||||
self.error_message = importer.error_message
|
||||
# Try to fill empty messages with the list
|
||||
if self.error_message.blank? && !self.error_code.nil?
|
||||
default_message = CartoDB::IMPORTER_ERROR_CODES.fetch(self.error_code, {})
|
||||
self.error_message = default_message.fetch(:title, '')
|
||||
end
|
||||
if self.retried_times < MAX_RETRIES
|
||||
self.retried_times += 1
|
||||
self.run_at = Time.now + interval
|
||||
else
|
||||
::Resque.enqueue(::Resque::UserJobs::Mail::Sync::MaxRetriesReached, self.user_id,
|
||||
self.visualization_id, self.name, self.error_code, self.error_message)
|
||||
end
|
||||
end
|
||||
|
||||
def set_general_failure_state_from(exception, error_code = 99999, error_message = 'Unknown error, please try again')
|
||||
log.append_and_store '******** synchronization raised exception ********'
|
||||
self.log_trace = exception.message + ' ' + exception.backtrace.join("\n")
|
||||
self.state = STATE_FAILURE
|
||||
self.error_code = error_code
|
||||
self.error_message = error_message
|
||||
if self.retried_times < MAX_RETRIES
|
||||
self.retried_times += 1
|
||||
self.run_at = Time.now + interval
|
||||
else
|
||||
::Resque.enqueue(::Resque::UserJobs::Mail::Sync::MaxRetriesReached, self.user_id,
|
||||
self.visualization_id, self.name, self.error_code, self.error_message)
|
||||
end
|
||||
rescue => e
|
||||
CartoDB.notify_exception(e,
|
||||
{
|
||||
error_code: error_code,
|
||||
error_message: error_message,
|
||||
retried_times: self.retried_times
|
||||
}
|
||||
)
|
||||
end
|
||||
|
||||
# Tries to run automatic geocoding if present
|
||||
def geocode_table
|
||||
return unless table && table.automatic_geocoding
|
||||
log.append 'Running automatic geocoding...'
|
||||
table.automatic_geocoding.run
|
||||
rescue => e
|
||||
log.append "Error while running automatic geocoding: #{e.message}"
|
||||
end # geocode_table
|
||||
|
||||
def to_hash
|
||||
attributes.merge({from_external_source: from_external_source?}).to_hash
|
||||
end
|
||||
|
||||
def to_json(*args)
|
||||
attributes.merge({from_external_source: from_external_source?}).to_json(*args)
|
||||
end
|
||||
|
||||
def valid?
|
||||
true
|
||||
end
|
||||
|
||||
def set_timestamps
|
||||
self.created_at ||= Time.now
|
||||
self.updated_at = Time.now
|
||||
self
|
||||
end
|
||||
|
||||
def user
|
||||
@user ||= ::User.where(id: user_id).first
|
||||
end
|
||||
|
||||
def table
|
||||
if @table.nil?
|
||||
@table = ::Table.new(name: name, user_id: user.id)
|
||||
end
|
||||
@table
|
||||
end
|
||||
|
||||
def visualization
|
||||
@visualization ||= CartoDB::Visualization::Member.new(id: @visualization_id).fetch
|
||||
rescue KeyErrror
|
||||
@visualization = nil
|
||||
end
|
||||
|
||||
def authorize?(user)
|
||||
user.id == user_id && (!!user.sync_tables_enabled || from_external_source?)
|
||||
end
|
||||
|
||||
def pg_options
|
||||
SequelRails.configuration.environment_for(Rails.env)
|
||||
.merge(
|
||||
username: user.database_username,
|
||||
password: user.database_password,
|
||||
database: user.database_name,
|
||||
host: user.user_database_host
|
||||
)
|
||||
end
|
||||
|
||||
def ogr2ogr_options
|
||||
options = Cartodb.config.fetch(:ogr2ogr, {})
|
||||
if options['binary'].nil? || options['csv_guessing'].nil?
|
||||
{}
|
||||
else
|
||||
ogr_options = {
|
||||
ogr2ogr_binary: options['binary'],
|
||||
ogr2ogr_csv_guessing: options['csv_guessing'] && @type_guessing,
|
||||
quoted_fields_guessing: @quoted_fields_guessing,
|
||||
}
|
||||
if options['memory_limit'].present?
|
||||
ogr_options.merge!(ogr2ogr_memory_limit: options['memory_limit'])
|
||||
end
|
||||
return ogr_options
|
||||
end
|
||||
end
|
||||
|
||||
# TODO code duplicated from data_import.rb, refactor
|
||||
def content_guessing_options
|
||||
guessing_config = Cartodb.config.fetch(:importer, {}).deep_symbolize_keys.fetch(:content_guessing, {})
|
||||
geocoder_config = Cartodb.config.fetch(:geocoder, {}).deep_symbolize_keys
|
||||
if guessing_config[:enabled] and self.content_guessing and geocoder_config
|
||||
{ guessing: guessing_config, geocoder: geocoder_config }
|
||||
else
|
||||
{ guessing: { enabled: false } }
|
||||
end
|
||||
end
|
||||
|
||||
def log
|
||||
return @log unless @log.nil?
|
||||
|
||||
log_attributes = {
|
||||
id: self.log_id
|
||||
}
|
||||
|
||||
log_attributes.merge(user_id: user.id) if user
|
||||
|
||||
@log = CartoDB::Log.where(log_attributes).first
|
||||
end
|
||||
|
||||
def valid_uuid?(text)
|
||||
!!UUIDTools::UUID.parse(text)
|
||||
rescue TypeError
|
||||
false
|
||||
rescue ArgumentError
|
||||
false
|
||||
end
|
||||
|
||||
def get_datasource(datasource_name)
|
||||
begin
|
||||
datasource = DatasourcesFactory.get_datasource(datasource_name, user, {
|
||||
http_timeout: ::DataImport.http_timeout_for(user)
|
||||
})
|
||||
if datasource.kind_of? BaseOAuth
|
||||
oauth = user.oauths.select(datasource_name)
|
||||
datasource.token = oauth.token unless oauth.nil?
|
||||
end
|
||||
rescue => ex
|
||||
log.append "Exception: #{ex.message}"
|
||||
log.append ex.backtrace
|
||||
datasource = nil
|
||||
end
|
||||
datasource
|
||||
end
|
||||
|
||||
def from_external_source?
|
||||
::ExternalDataImport.where(synchronization_id: self.id).first != nil
|
||||
end
|
||||
|
||||
attr_reader :repository
|
||||
|
||||
attr_accessor :log_trace, :service_name, :service_item_id
|
||||
|
||||
end
|
||||
end
|
||||
end
|
||||
@@ -0,0 +1,58 @@
|
||||
require_relative '../../../services/datasources/lib/datasources'
|
||||
|
||||
# @see DB table 'synchronization_oauths'
|
||||
class SynchronizationOauth < Sequel::Model
|
||||
|
||||
many_to_one :user
|
||||
|
||||
PUBLIC_ATTRIBUTES = [
|
||||
:id,
|
||||
:user_id,
|
||||
:service,
|
||||
:token
|
||||
]
|
||||
|
||||
def public_values
|
||||
Hash[PUBLIC_ATTRIBUTES.map{ |k| [k, (self.send(k) rescue self[k].to_s)] }]
|
||||
end #public_values
|
||||
|
||||
def validate
|
||||
super
|
||||
|
||||
validates_presence :token
|
||||
|
||||
if new?
|
||||
existing_oauth = SynchronizationOauth.filter(
|
||||
user_id: user_id,
|
||||
service: service
|
||||
).first
|
||||
errors.add(:user_id, " already has an oauth token created for service #{:service}") unless existing_oauth.nil?
|
||||
else
|
||||
existing_oauth = SynchronizationOauth.filter(
|
||||
id: id
|
||||
).first
|
||||
errors.add(:id, ' cannot change user or service, only token') unless (existing_oauth.service == service && existing_oauth.user_id == user_id)
|
||||
end
|
||||
end #validate
|
||||
|
||||
def before_save
|
||||
super
|
||||
self.updated_at = Time.now
|
||||
end #before_save
|
||||
|
||||
def ==(oauth_object)
|
||||
return false unless oauth_object
|
||||
self.id == oauth_object.id
|
||||
end #==
|
||||
|
||||
|
||||
def get_service_datasource
|
||||
user = ::User.where(id: user_id).first
|
||||
datasource = CartoDB::Datasources::DatasourcesFactory.get_datasource(service, user, {
|
||||
http_timeout: ::DataImport.http_timeout_for(user)
|
||||
})
|
||||
datasource.token = token unless datasource.nil?
|
||||
datasource
|
||||
end #get_service_datasource
|
||||
|
||||
end
|
||||
Reference in New Issue
Block a user