diff --git a/mongo_connector/connector.py b/mongo_connector/connector.py index 7465b6d4..bb516598 100644 --- a/mongo_connector/connector.py +++ b/mongo_connector/connector.py @@ -124,6 +124,9 @@ def __init__(self, mongo_address, doc_managers=None, **kwargs): # Timezone awareness self.tz_aware = kwargs.get('tz_aware', False) + # If oplog access is via the proxy + self.is_oplog_proxy= kwargs.get('is_oplog_proxy', False) + # SSL keyword arguments to MongoClient. ssl_certfile = kwargs.pop('ssl_certfile', None) ssl_ca_certs = kwargs.pop('ssl_ca_certs', None) @@ -215,7 +218,8 @@ def from_config(cls, config): ssl_keyfile=config['ssl.sslKeyfile'], ssl_ca_certs=config['ssl.sslCACerts'], ssl_cert_reqs=config['ssl.sslCertificatePolicy'], - tz_aware=config['timezoneAware'] + tz_aware=config['timezoneAware'], + is_oplog_proxy=config['isOplogProxy'] ) return connector @@ -374,12 +378,19 @@ def run(self): ) return - # Establish a connection to the replica set as a whole - self.main_conn.close() - self.main_conn = self.create_authed_client( - replicaSet=is_master['setName']) - - self.update_version_from_client(self.main_conn) + if not self.is_oplog_proxy: + # Establish a connection to the replica set as a whole + self.main_conn.close() + self.main_conn = self.create_authed_client( + replicaSet=is_master['setName']) + self.update_version_from_client(self.main_conn) + else: + try: + # Check if local.oplog.rs is readable + self.main_conn.local.oplog.rs.find_one() + except pymongo.errors.OperationFailure: + LOG.error('Could not read local.oplog.rs!') + sys.exit(1) # non sharded configuration oplog = OplogThread( @@ -486,6 +497,16 @@ def add_option(*args, **kwargs): " would be a valid argument to `-m`. Don't use" " quotes around the address.") + is_oplog_proxy = add_option( + config_key="isOplogProxy", + default=False, + type=bool) + + is_oplog_proxy.add_cli( + "--is_oplog_proxy", dest="is_oplog_proxy", help= + "True if passed uri is a proxy with access to mongo" + "oplog.rs.") + oplog_file = add_option( config_key="oplogFile", default="oplog.timestamp",