diff options
| author | Apollon Oikonomopoulos <apoikos@debian.org> | 2017-10-07 23:41:46 +0300 |
|---|---|---|
| committer | Apollon Oikonomopoulos <apoikos@debian.org> | 2017-10-07 23:41:46 +0300 |
| commit | 2a141c70bce6af0f03f734bbf298be62a8b7cdde (patch) | |
| tree | 787170dd34bcdb33ed239c52f7d9a9c9459ba870 | |
| parent | ccfd5a6e3a467d2dea89ef38a7913c286ea0b0ba (diff) | |
| parent | 930687d86a670280fe7e5f98534183a571eaecd1 (diff) | |
Updated version 3.2.17 from 'upstream/3.2.17'
with Debian dir c7d75603a464b5e168470f7eb78800fdc64cb139
108 files changed, 1528 insertions, 1274 deletions
diff --git a/etc/evergreen.yml b/etc/evergreen.yml index 9bef88f9a0e..97b0f49b755 100644 --- a/etc/evergreen.yml +++ b/etc/evergreen.yml @@ -90,7 +90,6 @@ variables: - ubuntu1404-rocksdb - debian71 - debian81 - - solaris-64-bit - osx-107-ssl - osx-107 @@ -201,7 +200,7 @@ functions: rm -rf rocksdb git clone https://github.com/facebook/rocksdb.git cd rocksdb - make static_lib + make USE_RTTI=1 static_lib fi "build new tools" : @@ -5395,89 +5394,6 @@ buildvariants: - name: push ########################################### -# Solaris buildvariants # -########################################### - -- name: solaris-64-bit - display_name: "* Solaris" - modules: - - mongo-tools - run_on: - - solaris - expansions: - push_path: sunos5 - push_bucket: downloads.mongodb.org - push_name: sunos5 - push_arch: x86_64 - gorootvars: PATH=/opt/mongodbtoolchain/v2/bin:$PATH - tooltags: -gccgoflags "-lsocket -lnsl" - compile_flags: CC=/opt/mongodbtoolchain/bin/gcc CXX=/opt/mongodbtoolchain/bin/g++ -j$(kstat cpu | sort -u | grep -c "^module") --release CCFLAGS="-m64" LINKFLAGS="-m64 -static-libstdc++ -static-libgcc" OBJCOPY=/opt/mongodbtoolchain/bin/objcopy - num_jobs_available: $(kstat cpu | sort -u | grep -c "^module") - tasks: - - name: compile - - name: aggregation - - name: aggregation_WT - - name: auth - - name: auth_WT - - name: bulk_gle_passthrough - - name: bulk_gle_passthrough_WT - - name: concurrency - - name: concurrency_WT - - name: concurrency_replication - - name: concurrency_replication_WT - - name: concurrency_sharded - - name: concurrency_sharded_WT - - name: concurrency_sharded_sccc - - name: concurrency_sharded_sccc_WT - - name: dbtest - - name: dbtest_WT - - name: disk - - name: durability - - name: failpoints - - name: httpinterface - - name: jsCore - - name: jsCore_compatibility - - name: jsCore_compatibility_WT - - name: jsCore_small_oplog - - name: jsCore_small_oplog_rs - - name: jsCore_small_oplog_rs_WT - - name: jsCore_small_oplog_WT - - name: jsCore_WT - - name: mmap - - name: mongosTest - - name: noPassthrough - - name: noPassthroughWithMongod - - name: noPassthroughWithMongod_WT - - name: noPassthrough_WT - - name: parallel - - name: parallel_compatibility - - name: parallel_compatibility_WT - - name: parallel_WT - - name: replicasets - - name: replicasets_WT - - name: replication - - name: replication_WT - - name: sharding - - name: sharding_csrs_upgrade - - name: sharded_collections_jscore_passthrough - - name: sharded_collections_jscore_passthrough_WT - - name: sharding_jscore_passthrough - - name: sharding_jscore_passthrough_WT - - name: sharding_jscore_passthrough_wire_ops_WT - - name: sharding_WT - - name: sharding_csrs_upgrade_WT - - name: slow1 - - name: slow1_WT - - name: slow2 - - name: slow2_WT - - name: tool - - name: tool_WT - - name: unittests - - name: push - distros: - - rhel70 - -########################################### # Debian buildvariants # ########################################### diff --git a/etc/longevity.yml b/etc/longevity.yml index b5cb674a354..af830baf318 100644 --- a/etc/longevity.yml +++ b/etc/longevity.yml @@ -141,15 +141,15 @@ functions: "infrastructure provisioning": - command: shell.exec - # call infrastructure-provisioning.sh. This will either create a cluster, or update tags on existing instances. + # call infrastructure-provisioning.py. This will either create a cluster, or update tags on existing instances. params: working_dir: work script: | set -e set -v source ./dsienv.sh - export PRODUCTION=true - $DSI_PATH/bin/infrastructure_provisioning.sh ${cluster} + source ./venv/bin/activate + $DSI_PATH/bin/infrastructure_provisioning.py "configure mongodb cluster": - command: shell.exec @@ -190,10 +190,11 @@ functions: set -o verbose source ./dsienv.sh # Longevity runs so rarely, we simply teardown the cluster when done. - # Note that nowadays infrastructure_teardown.sh is actually copying the terraform.tfstate into /data/infrastructure_provisioning + # Note that nowadays infrastructure_teardown.py is actually copying the terraform.tfstate into /data/infrastructure_provisioning # but as of this writing the rhel70-perf-longevity distro didn't actually use the teardown hook. source ./dsienv.sh - $DSI_PATH/bin/infrastructure_teardown.sh + source ./venv/bin/activate + $DSI_PATH/bin/infrastructure_teardown.py echo "Cluster DESTROYED." echo echo "All perf results" diff --git a/etc/system_perf.yml b/etc/system_perf.yml index 641b34bac49..e654824801c 100644 --- a/etc/system_perf.yml +++ b/etc/system_perf.yml @@ -25,8 +25,8 @@ post: params: aws_key: ${aws_key} aws_secret: ${aws_secret} - local_file: work/reports/graphs/timeseries-p1.html - remote_file: ${project}/${build_variant}/${revision}/${task_id}/${version_id}/logs/timeseries-p1-${task_name}-${build_id}.html + local_file: work/reports/graphs/timeseries-mongod.0.html + remote_file: ${project}/${build_variant}/${revision}/${task_id}/${version_id}/logs/timeseries-mongod.0-${task_name}-${build_id}.html bucket: mciuploads permissions: public-read content_type: text/html @@ -103,7 +103,6 @@ functions: ext: ${ext} script_flags : ${script_flags} dsi_rev: ${dsi_rev} - compare_task: ${compare_task} workloads_rev: ${workloads_rev} # compositions of expansions @@ -151,15 +150,15 @@ functions: "infrastructure provisioning": - command: shell.exec - # call infrastructure-provisioning.sh. This will either create a cluster, or update tags on existing instances. + # call infrastructure-provisioning.py. This will either create a cluster, or update tags on existing instances. params: working_dir: work script: | set -e set -v source ./dsienv.sh - export PRODUCTION=true - $DSI_PATH/bin/infrastructure_provisioning.sh ${cluster} + source ./venv/bin/activate + $DSI_PATH/bin/infrastructure_provisioning.py "configure mongodb cluster": - command: shell.exec @@ -237,45 +236,6 @@ functions: OVERRIDEFILE="../src/dsi/dsi/analysis/v3.2/system_perf_override.json" python -u ../src/dsi/dsi/analysis/post_run_check.py ${script_flags} --reports-analysis reports --perf-file reports/perf.json --rev ${revision} -f history.json -t tags.json --refTag $TAG --overrideFile $OVERRIDEFILE --project_id sys-perf --variant ${build_variant} --task ${task_name} - "compare": - - command: shell.exec - params: - script: | - set -o verbose - rm -rf ./src ./work - mkdir src - mkdir work - - command: manifest.load - - command: git.get_project - params: - directory: src - revisions: # for each module include revision as <module_name> : ${<module_name>_rev} - dsi: ${dsi_rev} - - command: json.get - params: - task: ${compare_task} - variant : ${variant1} - file: "work/standalone.json" - name: "perf" - - command: json.get - params: - task: ${compare_task} - variant : ${variant2} - file: "work/oplog.json" - name: "perf" - - command: shell.exec - type : test - params: - working_dir: work - script: | - set -o errexit - set -o verbose - python -u ../src/dsi/dsi/analysis/compare.py -b standalone.json -c oplog.json - - command: "json.send" - params: - name: "perf" - file: "work/perf.json" - ####################################### # Tasks # ####################################### @@ -357,6 +317,50 @@ tasks: vars: script_flags: --ycsb-throughput-analysis reports +- name: industry_benchmarks_wmajority_WT + depends_on: + - name: compile + variant: linux-standalone + commands: + - func: "prepare environment" + vars: + storageEngine: "wiredTiger" + test: "ycsb-wmajority" + - func: "infrastructure provisioning" + - func: "configure mongodb cluster" + vars: + storageEngine: "wiredTiger" + - func: "run test" + vars: + storageEngine: "wiredTiger" + test: "ycsb-wmajority" + - func: "make test log artifact" + - func: "analyze" + vars: + script_flags: --ycsb-throughput-analysis reports + +- name: industry_benchmarks_wmajority_MMAPv1 + depends_on: + - name: compile + variant: linux-standalone + commands: + - func: "prepare environment" + vars: + storageEngine: "mmapv1" + test: "ycsb-wmajority" + - func: "infrastructure provisioning" + - func: "configure mongodb cluster" + vars: + storageEngine: "mmapv1" + - func: "run test" + vars: + storageEngine: "mmapv1" + test: "ycsb-wmajority" + - func: "make test log artifact" + - func: "analyze" + vars: + script_flags: --ycsb-throughput-analysis reports + - name: core_workloads_WT depends_on: - name: compile @@ -518,70 +522,6 @@ tasks: - func: "make test log artifact" - func: "analyze" -- name: industry_benchmarks_WT_oplog_comp - depends_on: - - name: industry_benchmarks_WT - variant: linux-standalone - status : "*" - - name: industry_benchmarks_WT - variant: linux-1-node-replSet - status: "*" - commands: - - func: "compare" - vars: - compare_task: "industry_benchmarks_WT" - variant1: "linux-standalone" - variant2: "linux-1-node-replSet" - - func: "analyze" - -- name: industry_benchmarks_MMAPv1_oplog_comp - depends_on: - - name: industry_benchmarks_MMAPv1 - variant: linux-standalone - status: "*" - - name: industry_benchmarks_MMAPv1 - variant: linux-1-node-replSet - status: "*" - commands: - - func: "compare" - vars: - compare_task: "industry_benchmarks_MMAPv1" - variant1: "linux-standalone" - variant2: "linux-1-node-replSet" - - func: "analyze" - -- name: core_workloads_WT_oplog_comp - depends_on: - - name: core_workloads_WT - variant: linux-standalone - status: "*" - - name: core_workloads_WT - variant: linux-1-node-replSet - status: "*" - commands: - - func: "compare" - vars: - compare_task: "core_workloads_WT" - variant1: "linux-standalone" - variant2: "linux-1-node-replSet" - - func: "analyze" - -- name: core_workloads_MMAPv1_oplog_comp - depends_on: - - name: core_workloads_MMAPv1 - variant: linux-standalone - status: "*" - - name: core_workloads_MMAPv1 - variant: linux-1-node-replSet - status: "*" - commands: - - func: "compare" - vars: - compare_task: "core_workloads_MMAPv1" - variant1: "linux-standalone" - variant2: "linux-1-node-replSet" - - func: "analyze" - - name: initialsync_WT depends_on: - name: compile @@ -706,6 +646,8 @@ buildvariants: - name: industry_benchmarks_WT - name: core_workloads_WT - name: industry_benchmarks_MMAPv1 + - name: industry_benchmarks_wmajority_WT + - name: industry_benchmarks_wmajority_MMAPv1 - name: core_workloads_MMAPv1 - name: mongos_workloads_WT - name: mongos_workloads_MMAPv1 @@ -728,6 +670,8 @@ buildvariants: - name: industry_benchmarks_WT - name: core_workloads_WT - name: industry_benchmarks_MMAPv1 + - name: industry_benchmarks_wmajority_WT + - name: industry_benchmarks_wmajority_MMAPv1 - name: core_workloads_MMAPv1 - name: non_sharded_workloads_WT - name: non_sharded_workloads_MMAPv1 @@ -748,16 +692,3 @@ buildvariants: - name: initialsync_WT - name: initialsync_MMAPv1 -- name: linux-oplog-compare - display_name: Linux Oplog Compare - batchtime: 10080 # 7 days - modules: *modules - expansions: - project: *project - run_on: - - "rhel70-perf-single" - tasks: - - name: industry_benchmarks_WT_oplog_comp - - name: core_workloads_WT_oplog_comp - - name: industry_benchmarks_MMAPv1_oplog_comp - - name: core_workloads_MMAPv1_oplog_comp diff --git a/jstests/aggregation/bugs/server6118.js b/jstests/aggregation/bugs/server6118.js index f891135de72..3c55ae5ce33 100644 --- a/jstests/aggregation/bugs/server6118.js +++ b/jstests/aggregation/bugs/server6118.js @@ -1,12 +1,12 @@ // SERVER-6118: support for sharded sorts (function() { + 'use strict'; - var s = new ShardingTest({name: "aggregation_sort1", shards: 2, mongos: 1}); - s.stopBalancer(); + var s = new ShardingTest({shards: 2}); - s.adminCommand({enablesharding: "test"}); + assert.commandWorked(s.s0.adminCommand({enablesharding: "test"})); s.ensurePrimaryShard('test', 'shard0001'); - s.adminCommand({shardcollection: "test.data", key: {_id: 1}}); + assert.commandWorked(s.s0.adminCommand({shardcollection: "test.data", key: {_id: 1}})); var d = s.getDB("test"); @@ -20,12 +20,12 @@ bulkOp.execute(); // Split the data into 3 chunks - s.adminCommand({split: "test.data", middle: {_id: 33}}); - s.adminCommand({split: "test.data", middle: {_id: 66}}); + assert.commandWorked(s.s0.adminCommand({split: "test.data", middle: {_id: 33}})); + assert.commandWorked(s.s0.adminCommand({split: "test.data", middle: {_id: 66}})); // Migrate the middle chunk to another shard - s.adminCommand( - {movechunk: "test.data", find: {_id: 50}, to: s.getOther(s.getServer("test")).name}); + assert.commandWorked(s.s0.adminCommand( + {movechunk: "test.data", find: {_id: 50}, to: s.getOther(s.getServer("test")).name})); // Check that the results are in order. var result = d.data.aggregate({$sort: {_id: 1}}).toArray(); @@ -36,5 +36,4 @@ } s.stop(); - })(); diff --git a/jstests/aggregation/bugs/server6179.js b/jstests/aggregation/bugs/server6179.js index 1109ddaa67e..503e91a70d1 100644 --- a/jstests/aggregation/bugs/server6179.js +++ b/jstests/aggregation/bugs/server6179.js @@ -1,12 +1,12 @@ // SERVER-6179: support for two $groups in sharded agg (function() { + 'use strict'; - var s = new ShardingTest({name: "aggregation_multiple_group", shards: 2, mongos: 1}); - s.stopBalancer(); + var s = new ShardingTest({shards: 2}); - s.adminCommand({enablesharding: "test"}); + assert.commandWorked(s.s0.adminCommand({enablesharding: "test"})); s.ensurePrimaryShard('test', 'shard0001'); - s.adminCommand({shardcollection: "test.data", key: {_id: 1}}); + assert.commandWorked(s.s0.adminCommand({shardcollection: "test.data", key: {_id: 1}})); var d = s.getDB("test"); @@ -20,18 +20,18 @@ bulkOp.execute(); // Split the data into 3 chunks - s.adminCommand({split: "test.data", middle: {_id: 33}}); - s.adminCommand({split: "test.data", middle: {_id: 66}}); + assert.commandWorked(s.s0.adminCommand({split: "test.data", middle: {_id: 33}})); + assert.commandWorked(s.s0.adminCommand({split: "test.data", middle: {_id: 66}})); // Migrate the middle chunk to another shard - s.adminCommand( - {movechunk: "test.data", find: {_id: 50}, to: s.getOther(s.getServer("test")).name}); + assert.commandWorked(s.s0.adminCommand( + {movechunk: "test.data", find: {_id: 50}, to: s.getOther(s.getServer("test")).name})); // Check that we get results rather than an error var result = d.data.aggregate({$group: {_id: '$_id', i: {$first: '$i'}}}, {$group: {_id: '$i', avg_id: {$avg: '$_id'}}}, {$sort: {_id: 1}}).toArray(); - expected = [ + var expected = [ {"_id": 0, "avg_id": 45}, {"_id": 1, "avg_id": 46}, {"_id": 2, "avg_id": 47}, @@ -47,5 +47,4 @@ assert.eq(result, expected); s.stop(); - })(); diff --git a/jstests/aggregation/bugs/server7781.js b/jstests/aggregation/bugs/server7781.js index 230a8a64c9f..c3918aeb8d2 100644 --- a/jstests/aggregation/bugs/server7781.js +++ b/jstests/aggregation/bugs/server7781.js @@ -1,5 +1,6 @@ // SERVER-7781 $geoNear pipeline stage (function() { + 'use strict'; load('jstests/libs/geo_near_random.js'); load('jstests/aggregation/extras/utils.js'); @@ -59,10 +60,12 @@ shards.push(shard._id); }); - db.adminCommand({shardCollection: db[coll].getFullName(), key: {rand: 1}}); + assert.commandWorked( + db.adminCommand({shardCollection: db[coll].getFullName(), key: {rand: 1}})); for (var i = 1; i < 10; i++) { // split at 0.1, 0.2, ... 0.9 - db.adminCommand({split: db[coll].getFullName(), middle: {rand: i / 10}}); + assert.commandWorked( + db.adminCommand({split: db[coll].getFullName(), middle: {rand: i / 10}})); db.adminCommand({ moveChunk: db[coll].getFullName(), find: {rand: i / 10}, @@ -87,13 +90,13 @@ // test with defaults var queryPoint = pointMaker.mkPt(0.25); // stick to center of map - geoCmd = { + var geoCmd = { geoNear: coll, near: queryPoint, includeLocs: true, spherical: true }; - aggCmd = { + var aggCmd = { $geoNear: { near: queryPoint, includeLocs: 'stats.loc', @@ -134,7 +137,7 @@ geoCmd.num = 40; geoCmd.near = queryPoint; aggCmd.$geoNear.near = queryPoint; - aggArr = [aggCmd, {$limit: 50}, {$limit: 60}, {$limit: 40}]; + var aggArr = [aggCmd, {$limit: 50}, {$limit: 60}, {$limit: 40}]; checkOutput(db.runCommand(geoCmd), db[coll].aggregate(aggArr), 40); // Test $geoNear with an initial batchSize of 0. Regression test for SERVER-20935. @@ -157,13 +160,11 @@ test(db, false, '2dsphere'); var sharded = new ShardingTest({shards: 3, mongos: 1}); - sharded.stopBalancer(); - sharded.adminCommand({enablesharding: "test"}); + assert.commandWorked(sharded.s0.adminCommand({enablesharding: "test"})); sharded.ensurePrimaryShard('test', 'shard0001'); test(sharded.getDB('test'), true, '2d'); test(sharded.getDB('test'), true, '2dsphere'); sharded.stop(); - })(); diff --git a/jstests/aggregation/bugs/server9444.js b/jstests/aggregation/bugs/server9444.js index ad5f4b03ca6..f3dc2748b0a 100644 --- a/jstests/aggregation/bugs/server9444.js +++ b/jstests/aggregation/bugs/server9444.js @@ -1,76 +1,80 @@ // server-9444 support disk storage of intermediate results in aggregation - -var t = db.server9444; -t.drop(); - -var sharded = (typeof(RUNNING_IN_SHARDED_AGG_TEST) != 'undefined'); // see end of testshard1.js -if (sharded) { - db.adminCommand({shardcollection: t.getFullName(), key: {"_id": 'hashed'}}); -} - -var memoryLimitMB = sharded ? 200 : 100; - -function loadData() { - var bigStr = Array(1024 * 1024 + 1).toString(); // 1MB of ',' - for (var i = 0; i < memoryLimitMB + 1; i++) - t.insert({_id: i, bigStr: i + bigStr, random: Math.random()}); - - assert.gt(t.stats().size, memoryLimitMB * 1024 * 1024); -} -loadData(); - -function test(pipeline, outOfMemoryCode) { - // ensure by default we error out if exceeding memory limit - var res = t.runCommand('aggregate', {pipeline: pipeline}); - assert.commandFailed(res); - assert.eq(res.code, outOfMemoryCode); - - // ensure allowDiskUse: false does what it says - var res = t.runCommand('aggregate', {pipeline: pipeline, allowDiskUse: false}); - assert.commandFailed(res); - assert.eq(res.code, outOfMemoryCode); - - // allowDiskUse only supports bool. In particular, numbers aren't allowed. - var res = t.runCommand('aggregate', {pipeline: pipeline, allowDiskUse: 1}); - assert.commandFailed(res); - assert.eq(res.code, 16949); - - // ensure we work when allowDiskUse === true - var res = t.aggregate(pipeline, {allowDiskUse: true}); - assert.eq(res.itcount(), t.count()); // all tests output one doc per input doc -} - -var groupCode = 16945; -var sortCode = 16819; -var sortLimitCode = 16820; - -test([{$group: {_id: '$_id', bigStr: {$first: '$bigStr'}}}], groupCode); - -// sorting with _id would use index which doesn't require extsort -test([{$sort: {random: 1}}], sortCode); -test([{$sort: {bigStr: 1}}], sortCode); // big key and value - -// make sure sort + large limit won't crash the server (SERVER-10136) -test([{$sort: {bigStr: 1}}, {$limit: 1000 * 1000 * 1000}], sortLimitCode); - -// test combining two extSorts in both same and different orders -test([{$group: {_id: '$_id', bigStr: {$first: '$bigStr'}}}, {$sort: {_id: 1}}], groupCode); -test([{$group: {_id: '$_id', bigStr: {$first: '$bigStr'}}}, {$sort: {_id: -1}}], groupCode); -test([{$group: {_id: '$_id', bigStr: {$first: '$bigStr'}}}, {$sort: {random: 1}}], groupCode); -test([{$sort: {random: 1}}, {$group: {_id: '$_id', bigStr: {$first: '$bigStr'}}}], sortCode); - -var origDB = db; -if (sharded) { - // Stop balancer first before dropping so there will be no contention on the ns lock. - // It's alright to modify the global db variable since sharding tests never run in parallel. - db = db.getSiblingDB('config'); - sh.stopBalancer(); -} - -// don't leave large collection laying around -t.drop(); - -if (sharded) { - sh.startBalancer(); - db = origDB; -} +(function() { + 'use strict'; + + var t = db.server9444; + t.drop(); + + var sharded = (typeof(RUNNING_IN_SHARDED_AGG_TEST) != 'undefined'); // see end of testshard1.js + if (sharded) { + assert.commandWorked( + db.adminCommand({shardcollection: t.getFullName(), key: {"_id": 'hashed'}})); + } + + var memoryLimitMB = sharded ? 200 : 100; + + function loadData() { + var bigStr = Array(1024 * 1024 + 1).toString(); // 1MB of ',' + for (var i = 0; i < memoryLimitMB + 1; i++) + t.insert({_id: i, bigStr: i + bigStr, random: Math.random()}); + + assert.gt(t.stats().size, memoryLimitMB * 1024 * 1024); + } + loadData(); + + function test(pipeline, outOfMemoryCode) { + // ensure by default we error out if exceeding memory limit + var res = t.runCommand('aggregate', {pipeline: pipeline}); + assert.commandFailed(res); + assert.eq(res.code, outOfMemoryCode); + + // ensure allowDiskUse: false does what it says + var res = t.runCommand('aggregate', {pipeline: pipeline, allowDiskUse: false}); + assert.commandFailed(res); + assert.eq(res.code, outOfMemoryCode); + + // allowDiskUse only supports bool. In particular, numbers aren't allowed. + var res = t.runCommand('aggregate', {pipeline: pipeline, allowDiskUse: 1}); + assert.commandFailed(res); + assert.eq(res.code, 16949); + + // ensure we work when allowDiskUse === true + var res = t.aggregate(pipeline, {allowDiskUse: true}); + assert.eq(res.itcount(), t.count()); // all tests output one doc per input doc + } + + var groupCode = 16945; + var sortCode = 16819; + var sortLimitCode = 16820; + + test([{$group: {_id: '$_id', bigStr: {$first: '$bigStr'}}}], groupCode); + + // sorting with _id would use index which doesn't require extsort + test([{$sort: {random: 1}}], sortCode); + test([{$sort: {bigStr: 1}}], sortCode); // big key and value + + // make sure sort + large limit won't crash the server (SERVER-10136) + test([{$sort: {bigStr: 1}}, {$limit: 1000 * 1000 * 1000}], sortLimitCode); + + // test combining two extSorts in both same and different orders + test([{$group: {_id: '$_id', bigStr: {$first: '$bigStr'}}}, {$sort: {_id: 1}}], groupCode); + test([{$group: {_id: '$_id', bigStr: {$first: '$bigStr'}}}, {$sort: {_id: -1}}], groupCode); + test([{$group: {_id: '$_id', bigStr: {$first: '$bigStr'}}}, {$sort: {random: 1}}], groupCode); + test([{$sort: {random: 1}}, {$group: {_id: '$_id', bigStr: {$first: '$bigStr'}}}], sortCode); + + var origDB = db; + if (sharded) { + // Stop balancer first before dropping so there will be no contention on the ns lock. + // It's alright to modify the global db variable since sharding tests never run in parallel. + db = db.getSiblingDB('config'); + sh.stopBalancer(); + } + + // don't leave large collection laying around + t.drop(); + + if (sharded) { + sh.startBalancer(); + db = origDB; + } +})(); diff --git a/jstests/auth/mongos_cache_invalidation.js b/jstests/auth/mongos_cache_invalidation.js index 60700956e39..4594a18f8c4 100644 --- a/jstests/auth/mongos_cache_invalidation.js +++ b/jstests/auth/mongos_cache_invalidation.js @@ -209,13 +209,10 @@ db3.auth('spencer', 'pwd'); // s0/db1 should update its cache instantly assert.commandFailedWithCode(db1.foo.runCommand("collStats"), authzErrorCode); - // s1/db2 should update its cache in 5 seconds. - assert.soon( - function() { - return db2.foo.runCommand("collStats").code == authzErrorCode; - }, - "Mongos did not update its user cache after 5 seconds", - 6 * 1000); // Give an extra 1 second to avoid races + // s1/db2 should update its cache in 10 seconds. + assert.soon(function() { + return db2.foo.runCommand("collStats").code == authzErrorCode; + }, "Mongos did not update its user cache after 10 seconds", 10 * 1000); // We manually invalidate the cache on s2/db3. db3.adminCommand("invalidateUserCache"); diff --git a/jstests/core/group9.js b/jstests/core/group9.js new file mode 100644 index 00000000000..a7d9a0d128e --- /dev/null +++ b/jstests/core/group9.js @@ -0,0 +1,20 @@ +(function() { + 'use strict'; + var t = db.group_owned; + t.drop(); + + assert.writeOK(t.insert({_id: 1, subdoc: {id: 1}})); + assert.writeOK(t.insert({_id: 2, subdoc: {id: 2}})); + + var result = t.group({ + key: {'subdoc.id': 1}, + reduce: function(doc, value) { + value.subdoc = doc.subdoc; + return value; + }, + initial: {}, + finalize: function(res) {} + }); + + assert(result.length == 2); +}()); diff --git a/jstests/core/mr4.js b/jstests/core/mr4.js index ae5e11528af..2b8e93c3b35 100644 --- a/jstests/core/mr4.js +++ b/jstests/core/mr4.js @@ -9,7 +9,7 @@ t.save({x: 4, tags: ["b", "c"]}); m = function() { this.tags.forEach(function(z) { - emit(z, {count: xx}); + emit(z, {count: xx.val}); }); }; @@ -23,7 +23,7 @@ r = function(key, values) { }; }; -res = t.mapReduce(m, r, {out: "mr4_out", scope: {xx: 1}}); +res = t.mapReduce(m, r, {out: "mr4_out", scope: {xx: {val: 1}}}); z = res.convertToSingleObject(); assert.eq(3, Object.keySet(z).length, "A1"); @@ -33,7 +33,7 @@ assert.eq(3, z.c.count, "A4"); res.drop(); -res = t.mapReduce(m, r, {scope: {xx: 2}, out: "mr4_out"}); +res = t.mapReduce(m, r, {scope: {xx: {val: 2}}, out: "mr4_out"}); z = res.convertToSingleObject(); assert.eq(3, Object.keySet(z).length, "A1"); diff --git a/jstests/gle/gle_sharded_write.js b/jstests/gle/gle_sharded_write.js index f1feffed5b2..8d2a21cd758 100644 --- a/jstests/gle/gle_sharded_write.js +++ b/jstests/gle/gle_sharded_write.js @@ -2,192 +2,192 @@ // Ensures GLE correctly reports basic write stats and failures // Note that test should work correctly with and without write commands. // - -var st = new ShardingTest({shards: 2, mongos: 1}); -st.stopBalancer(); - -var mongos = st.s0; -var admin = mongos.getDB("admin"); -var config = mongos.getDB("config"); -var coll = mongos.getCollection(jsTestName() + ".coll"); -var shards = config.shards.find().toArray(); - -assert.commandWorked(admin.runCommand({enableSharding: coll.getDB().toString()})); -printjson(admin.runCommand({movePrimary: coll.getDB().toString(), to: shards[0]._id})); -assert.commandWorked(admin.runCommand({shardCollection: coll.toString(), key: {_id: 1}})); -assert.commandWorked(admin.runCommand({split: coll.toString(), middle: {_id: 0}})); -assert.commandWorked( - admin.runCommand({moveChunk: coll.toString(), find: {_id: 0}, to: shards[1]._id})); - -st.printShardingStatus(); - -var gle = null; - -// -// Successful insert -coll.remove({}); -coll.insert({_id: -1}); -printjson(gle = coll.getDB().runCommand({getLastError: 1})); -assert(gle.ok); -assert('err' in gle); -assert(!gle.err); -assert.eq(coll.count(), 1); - -// -// Successful update -coll.remove({}); -coll.insert({_id: 1}); -coll.update({_id: 1}, {$set: {foo: "bar"}}); -printjson(gle = coll.getDB().runCommand({getLastError: 1})); -assert(gle.ok); -assert('err' in gle); -assert(!gle.err); -assert(gle.updatedExisting); -assert.eq(gle.n, 1); -assert.eq(coll.count(), 1); - -// -// Successful multi-update -coll.remove({}); -coll.insert({_id: 1}); -coll.update({}, {$set: {foo: "bar"}}, false, true); -printjson(gle = coll.getDB().runCommand({getLastError: 1})); -assert(gle.ok); -assert('err' in gle); -assert(!gle.err); -assert(gle.updatedExisting); -assert.eq(gle.n, 1); -assert.eq(coll.count(), 1); - -// -// Successful upsert -coll.remove({}); -coll.update({_id: 1}, {_id: 1}, true); -printjson(gle = coll.getDB().runCommand({getLastError: 1})); -assert(gle.ok); -assert('err' in gle); -assert(!gle.err); -assert(!gle.updatedExisting); -assert.eq(gle.n, 1); -assert.eq(gle.upserted, 1); -assert.eq(coll.count(), 1); - -// -// Successful upserts -coll.remove({}); -coll.update({_id: -1}, {_id: -1}, true); -coll.update({_id: 1}, {_id: 1}, true); -printjson(gle = coll.getDB().runCommand({getLastError: 1})); -assert(gle.ok); -assert('err' in gle); -assert(!gle.err); -assert(!gle.updatedExisting); -assert.eq(gle.n, 1); -assert.eq(gle.upserted, 1); -assert.eq(coll.count(), 2); - -// -// Successful remove -coll.remove({}); -coll.insert({_id: 1}); -coll.remove({_id: 1}); -printjson(gle = coll.getDB().runCommand({getLastError: 1})); -assert(gle.ok); -assert('err' in gle); -assert(!gle.err); -assert.eq(gle.n, 1); -assert.eq(coll.count(), 0); - -// -// Error on one host during update -coll.remove({}); -coll.update({_id: 1}, {$invalid: "xxx"}, true); -printjson(gle = coll.getDB().runCommand({getLastError: 1})); -assert(gle.ok); -assert(gle.err); -assert(gle.code); -assert(!gle.errmsg); -assert(gle.singleShard); -assert.eq(coll.count(), 0); - -// -// Error on two hosts during remove -coll.remove({}); -coll.remove({$invalid: 'remove'}); -printjson(gle = coll.getDB().runCommand({getLastError: 1})); -assert(gle.ok); -assert(gle.err); -assert(gle.code); -assert(!gle.errmsg); -assert(gle.shards); -assert.eq(coll.count(), 0); - -// -// Repeated calls to GLE should work -coll.remove({}); -coll.update({_id: 1}, {$invalid: "xxx"}, true); -printjson(gle = coll.getDB().runCommand({getLastError: 1})); -assert(gle.ok); -assert(gle.err); -assert(gle.code); -assert(!gle.errmsg); -assert(gle.singleShard); -printjson(gle = coll.getDB().runCommand({getLastError: 1})); -assert(gle.ok); -assert(gle.err); -assert(gle.code); -assert(!gle.errmsg); -assert(gle.singleShard); -assert.eq(coll.count(), 0); - -// -// Geo $near is not supported on mongos -coll.ensureIndex({loc: "2dsphere"}); -coll.remove({}); -var query = { - loc: { - $near: { - $geometry: {type: "Point", coordinates: [0, 0]}, - $maxDistance: 1000, +(function() { + 'use strict'; + + var st = new ShardingTest({shards: 2, mongos: 1}); + + var mongos = st.s0; + var admin = mongos.getDB("admin"); + var config = mongos.getDB("config"); + var coll = mongos.getCollection(jsTestName() + ".coll"); + var shards = config.shards.find().toArray(); + + assert.commandWorked(admin.runCommand({enableSharding: coll.getDB().toString()})); + printjson(admin.runCommand({movePrimary: coll.getDB().toString(), to: shards[0]._id})); + assert.commandWorked(admin.runCommand({shardCollection: coll.toString(), key: {_id: 1}})); + assert.commandWorked(admin.runCommand({split: coll.toString(), middle: {_id: 0}})); + assert.commandWorked( + admin.runCommand({moveChunk: coll.toString(), find: {_id: 0}, to: shards[1]._id})); + + st.printShardingStatus(); + + var gle = null; + + // + // Successful insert + coll.remove({}); + coll.insert({_id: -1}); + printjson(gle = coll.getDB().runCommand({getLastError: 1})); + assert(gle.ok); + assert('err' in gle); + assert(!gle.err); + assert.eq(coll.count(), 1); + + // + // Successful update + coll.remove({}); + coll.insert({_id: 1}); + coll.update({_id: 1}, {$set: {foo: "bar"}}); + printjson(gle = coll.getDB().runCommand({getLastError: 1})); + assert(gle.ok); + assert('err' in gle); + assert(!gle.err); + assert(gle.updatedExisting); + assert.eq(gle.n, 1); + assert.eq(coll.count(), 1); + + // + // Successful multi-update + coll.remove({}); + coll.insert({_id: 1}); + coll.update({}, {$set: {foo: "bar"}}, false, true); + printjson(gle = coll.getDB().runCommand({getLastError: 1})); + assert(gle.ok); + assert('err' in gle); + assert(!gle.err); + assert(gle.updatedExisting); + assert.eq(gle.n, 1); + assert.eq(coll.count(), 1); + + // + // Successful upsert + coll.remove({}); + coll.update({_id: 1}, {_id: 1}, true); + printjson(gle = coll.getDB().runCommand({getLastError: 1})); + assert(gle.ok); + assert('err' in gle); + assert(!gle.err); + assert(!gle.updatedExisting); + assert.eq(gle.n, 1); + assert.eq(gle.upserted, 1); + assert.eq(coll.count(), 1); + + // + // Successful upserts + coll.remove({}); + coll.update({_id: -1}, {_id: -1}, true); + coll.update({_id: 1}, {_id: 1}, true); + printjson(gle = coll.getDB().runCommand({getLastError: 1})); + assert(gle.ok); + assert('err' in gle); + assert(!gle.err); + assert(!gle.updatedExisting); + assert.eq(gle.n, 1); + assert.eq(gle.upserted, 1); + assert.eq(coll.count(), 2); + + // + // Successful remove + coll.remove({}); + coll.insert({_id: 1}); + coll.remove({_id: 1}); + printjson(gle = coll.getDB().runCommand({getLastError: 1})); + assert(gle.ok); + assert('err' in gle); + assert(!gle.err); + assert.eq(gle.n, 1); + assert.eq(coll.count(), 0); + + // + // Error on one host during update + coll.remove({}); + coll.update({_id: 1}, {$invalid: "xxx"}, true); + printjson(gle = coll.getDB().runCommand({getLastError: 1})); + assert(gle.ok); + assert(gle.err); + assert(gle.code); + assert(!gle.errmsg); + assert(gle.singleShard); + assert.eq(coll.count(), 0); + + // + // Error on two hosts during remove + coll.remove({}); + coll.remove({$invalid: 'remove'}); + printjson(gle = coll.getDB().runCommand({getLastError: 1})); + assert(gle.ok); + assert(gle.err); + assert(gle.code); + assert(!gle.errmsg); + assert(gle.shards); + assert.eq(coll.count(), 0); + + // + // Repeated calls to GLE should work + coll.remove({}); + coll.update({_id: 1}, {$invalid: "xxx"}, true); + printjson(gle = coll.getDB().runCommand({getLastError: 1})); + assert(gle.ok); + assert(gle.err); + assert(gle.code); + assert(!gle.errmsg); + assert(gle.singleShard); + printjson(gle = coll.getDB().runCommand({getLastError: 1})); + assert(gle.ok); + assert(gle.err); + assert(gle.code); + assert(!gle.errmsg); + assert(gle.singleShard); + assert.eq(coll.count(), 0); + + // + // Geo $near is not supported on mongos + coll.ensureIndex({loc: "2dsphere"}); + coll.remove({}); + var query = { + loc: { + $near: { + $geometry: {type: "Point", coordinates: [0, 0]}, + $maxDistance: 1000, + } } - } -}; -printjson(coll.remove(query)); -printjson(gle = coll.getDB().runCommand({getLastError: 1})); -assert(gle.ok); -assert(gle.err); -assert(gle.code); -assert(!gle.errmsg); -assert(gle.shards); -assert.eq(coll.count(), 0); - -// -// First shard down -// - -// -// Successful bulk insert on two hosts, host dies before gle (error contacting host) -coll.remove({}); -coll.insert([{_id: 1}, {_id: -1}]); -// Wait for write to be written to shards before shutting it down. -printjson(gle = coll.getDB().runCommand({getLastError: 1})); -MongoRunner.stopMongod(st.shard0); -printjson(gle = coll.getDB().runCommand({getLastError: 1})); -// Should get an error about contacting dead host. -assert(!gle.ok); -assert(gle.errmsg); - -// -// Failed insert on two hosts, first host dead -// NOTE: This is DIFFERENT from 2.4, since we don't need to contact a host we didn't get -// successful writes from. -coll.remove({_id: 1}); -coll.insert([{_id: 1}, {_id: -1}]); -printjson(gle = coll.getDB().runCommand({getLastError: 1})); -assert(gle.ok); -assert(gle.err); -assert.eq(coll.count({_id: 1}), 1); - -jsTest.log("DONE!"); - -st.stop(); + }; + printjson(coll.remove(query)); + printjson(gle = coll.getDB().runCommand({getLastError: 1})); + assert(gle.ok); + assert(gle.err); + assert(gle.code); + assert(!gle.errmsg); + assert(gle.shards); + assert.eq(coll.count(), 0); + + // + // First shard down + // + + // + // Successful bulk insert on two hosts, host dies before gle (error contacting host) + coll.remove({}); + coll.insert([{_id: 1}, {_id: -1}]); + // Wait for write to be written to shards before shutting it down. + printjson(gle = coll.getDB().runCommand({getLastError: 1})); + MongoRunner.stopMongod(st.shard0); + printjson(gle = coll.getDB().runCommand({getLastError: 1})); + // Should get an error about contacting dead host. + assert(!gle.ok); + assert(gle.errmsg); + + // + // Failed insert on two hosts, first host dead + // NOTE: This is DIFFERENT from 2.4, since we don't need to contact a host we didn't get + // successful writes from. + coll.remove({_id: 1}); + coll.insert([{_id: 1}, {_id: -1}]); + printjson(gle = coll.getDB().runCommand({getLastError: 1})); + assert(gle.ok); + assert(gle.err); + assert.eq(coll.count({_id: 1}), 1); + + st.stop(); +})(); diff --git a/jstests/multiVersion/libs/multi_rs.js b/jstests/multiVersion/libs/multi_rs.js index 673be42d3df..da976c7f4dc 100644 --- a/jstests/multiVersion/libs/multi_rs.js +++ b/jstests/multiVersion/libs/multi_rs.js @@ -50,7 +50,12 @@ ReplSetTest.prototype.upgradeNode = function(node, opts, user, pwd) { var isMaster = node.getDB('admin').runCommand({isMaster: 1}); if (!isMaster.arbiterOnly) { - assert.commandWorked(node.adminCommand("replSetMaintenance")); + // Must retry this command, as it might return "currently running for election" and fail. + // Node might still be running for an election that will fail because it lost the election + // race with another node, at test initialization. See SERVER-23133. + assert.soon(function() { + return (node.adminCommand("replSetMaintenance").ok); + }); this.waitForState(node, ReplSetTest.State.RECOVERING); } diff --git a/jstests/noPassthrough/backup_restore.js b/jstests/noPassthrough/backup_restore.js index 95d111e2cd3..0074c8c1f8d 100644 --- a/jstests/noPassthrough/backup_restore.js +++ b/jstests/noPassthrough/backup_restore.js @@ -226,8 +226,8 @@ rst.start(secondary.nodeId, {}, true); } - // Wait up to 60 seconds until restarted node is in state secondary - rst.waitForState(rst.getSecondaries(), ReplSetTest.State.SECONDARY, 60 * 1000); + // Wait up to 5 minutes until restarted node is in state secondary. + rst.waitForState(rst.getSecondaries(), ReplSetTest.State.SECONDARY); // Add new hidden node to replSetTest var hiddenCfg = { @@ -263,8 +263,7 @@ // Wait up to 60 seconds until the new hidden node is in state RECOVERING. rst.waitForState(rst.nodes[numNodes], - [ReplSetTest.State.RECOVERING, ReplSetTest.State.SECONDARY], - 60 * 1000); + [ReplSetTest.State.RECOVERING, ReplSetTest.State.SECONDARY]); // Stop CRUD client and FSM client. assert(checkProgram(crudPid), testName + ' CRUD client was not running at end of test'); @@ -272,8 +271,8 @@ stopMongoProgramByPid(crudPid); stopMongoProgramByPid(fsmPid); - // Wait up to 60 seconds until the new hidden node is in state SECONDARY. - rst.waitForState(rst.nodes[numNodes], ReplSetTest.State.SECONDARY, 60 * 1000); + // Wait up to 5 minutes until the new hidden node is in state SECONDARY. + rst.waitForState(rst.nodes[numNodes], ReplSetTest.State.SECONDARY); // Wait for secondaries to finish catching up before shutting down. assert.writeOK(primary.getDB("test").foo.insert( diff --git a/jstests/noPassthrough/initial_sync_cloner_dups.js b/jstests/noPassthrough/initial_sync_cloner_dups.js index 1208dc8a16e..c0a3e71c0c1 100644 --- a/jstests/noPassthrough/initial_sync_cloner_dups.js +++ b/jstests/noPassthrough/initial_sync_cloner_dups.js @@ -79,8 +79,7 @@ // Wait for the secondary to get ReplSetInitiate command. replTest.waitForState( secondary, - [ReplSetTest.State.STARTUP_2, ReplSetTest.State.RECOVERING, ReplSetTest.State.SECONDARY], - 60 * 1000); + [ReplSetTest.State.STARTUP_2, ReplSetTest.State.RECOVERING, ReplSetTest.State.SECONDARY]); // This fail point will cause the first intial sync to fail, and leave an op in the buffer to // verify the fix from SERVER-17807 diff --git a/jstests/noPassthroughWithMongod/no_balance_collection.js b/jstests/noPassthroughWithMongod/no_balance_collection.js index cfec6199ca2..1c2f1aae009 100644 --- a/jstests/noPassthroughWithMongod/no_balance_collection.js +++ b/jstests/noPassthroughWithMongod/no_balance_collection.js @@ -1,14 +1,11 @@ // Tests whether the noBalance flag disables balancing for collections -var st = new ShardingTest({shards: 2, mongos: 1, verbose: 1}); +var st = new ShardingTest({shards: 2, mongos: 1}); // First, test that shell helpers require an argument assert.throws(sh.disableBalancing, [], "sh.disableBalancing requires a collection"); assert.throws(sh.enableBalancing, [], "sh.enableBalancing requires a collection"); -// Initially stop balancing -st.stopBalancer(); - var shardAName = st._shardNames[0]; var shardBName = st._shardNames[1]; @@ -70,10 +67,11 @@ jsTest.log("Chunks for " + collB + " are balanced."); // Re-disable balancing for collB sh.disableBalancing(collB); + // Wait for the balancer to fully finish the last migration and write the changelog // MUST set db var here, ugly but necessary db = st.s0.getDB("config"); -sh.waitForBalancer(true); +st.waitForBalancerRound(); // Make sure auto-migrates on insert don't move chunks var lastMigration = sh._lastMigration(collB); diff --git a/jstests/replsets/apply_ops_insert_write_conflict_nonatomic.js b/jstests/replsets/apply_ops_insert_write_conflict_nonatomic.js new file mode 100644 index 00000000000..4ef394f5682 --- /dev/null +++ b/jstests/replsets/apply_ops_insert_write_conflict_nonatomic.js @@ -0,0 +1,8 @@ +(function() { + 'use strict'; + + load("jstests/replsets/libs/apply_ops_insert_write_conflict.js"); + + new ApplyOpsInsertWriteConflictTest( + {testName: 'apply_ops_insert_write_conflict_nonatomic', atomic: false}).run(); +}()); diff --git a/jstests/replsets/election_not_blocked.js b/jstests/replsets/election_not_blocked.js index 95b53be1ebc..88d66715929 100644 --- a/jstests/replsets/election_not_blocked.js +++ b/jstests/replsets/election_not_blocked.js @@ -25,7 +25,7 @@ // so it cannot vote while fsync locked in PV1. Use PV0 explicitly here. protocolVersion: 0 }); - replTest.waitForState(replTest.nodes[0], ReplSetTest.State.PRIMARY, 60 * 1000); + replTest.waitForState(replTest.nodes[0], ReplSetTest.State.PRIMARY); var master = replTest.getPrimary(); // do a write diff --git a/jstests/replsets/initial_sync1.js b/jstests/replsets/initial_sync1.js index 55a454b765c..51d355d80c1 100644 --- a/jstests/replsets/initial_sync1.js +++ b/jstests/replsets/initial_sync1.js @@ -77,9 +77,7 @@ wait(function() { return config2.version == config.version && (config3 && config3.version == config.version); }); -replTest.waitForState(slave2, - [ReplSetTest.State.SECONDARY, ReplSetTest.State.RECOVERING], - 60 * 1000); +replTest.waitForState(slave2, [ReplSetTest.State.SECONDARY, ReplSetTest.State.RECOVERING]); print("7. Kill the secondary in the middle of syncing"); replTest.stop(slave1); @@ -91,7 +89,7 @@ replTest.waitForState(slave2, ReplSetTest.State.SECONDARY, 60 * 1000); print("9. Bring the secondary back up"); replTest.start(slave1, {}, true); reconnect(slave1); -replTest.waitForState(slave1, [ReplSetTest.State.PRIMARY, ReplSetTest.State.SECONDARY], 60 * 1000); +replTest.waitForState(slave1, [ReplSetTest.State.PRIMARY, ReplSetTest.State.SECONDARY]); print("10. Insert some stuff"); master = replTest.getPrimary(); diff --git a/jstests/replsets/initial_sync2.js b/jstests/replsets/initial_sync2.js index bab1063c072..69c91d30f04 100644 --- a/jstests/replsets/initial_sync2.js +++ b/jstests/replsets/initial_sync2.js @@ -25,7 +25,7 @@ var doTest = function() { var conns = replTest.startSet(); replTest.initiate(); - replTest.waitForState(replTest.nodes[0], ReplSetTest.State.PRIMARY, 5 * 60 * 1000); + replTest.waitForState(replTest.nodes[0], ReplSetTest.State.PRIMARY); var master = replTest.getPrimary(); var foo = master.getDB("foo"); @@ -73,27 +73,27 @@ var doTest = function() { }); admin_s2.runCommand({replSetFreeze: 999999}); - replTest.waitForState( - replTest.nodes[2], [ReplSetTest.State.SECONDARY, ReplSetTest.State.RECOVERING], 60 * 1000); + replTest.waitForState(replTest.nodes[2], + [ReplSetTest.State.SECONDARY, ReplSetTest.State.RECOVERING]); jsTest.log("7. Kill #1 in the middle of syncing"); replTest.stop(0); jsTest.log("8. Check that #3 makes it into secondary state"); - replTest.waitForState( - replTest.nodes[2], [ReplSetTest.State.PRIMARY, ReplSetTest.State.SECONDARY], 60 * 1000); + replTest.waitForState(replTest.nodes[2], + [ReplSetTest.State.PRIMARY, ReplSetTest.State.SECONDARY]); jsTest.log("9. Bring #1 back up"); replTest.start(0, {}, true); - replTest.waitForState( - replTest.nodes[0], [ReplSetTest.State.PRIMARY, ReplSetTest.State.SECONDARY], 60 * 1000); + replTest.waitForState(replTest.nodes[0], + [ReplSetTest.State.PRIMARY, ReplSetTest.State.SECONDARY]); jsTest.log("10. Initial sync should succeed"); - replTest.waitForState( - replTest.nodes[2], [ReplSetTest.State.PRIMARY, ReplSetTest.State.SECONDARY], 60 * 1000); + replTest.waitForState(replTest.nodes[2], + [ReplSetTest.State.PRIMARY, ReplSetTest.State.SECONDARY]); jsTest.log("11. Ensure #1 becomes primary"); - replTest.waitForState(replTest.nodes[0], ReplSetTest.State.PRIMARY, 60 * 1000); + replTest.waitForState(replTest.nodes[0], ReplSetTest.State.PRIMARY); jsTest.log("12. Everyone happy eventually"); replTest.awaitReplication(2 * 60 * 1000); diff --git a/jstests/replsets/libs/apply_ops_insert_write_conflict.js b/jstests/replsets/libs/apply_ops_insert_write_conflict.js new file mode 100644 index 00000000000..9bbdea7e08e --- /dev/null +++ b/jstests/replsets/libs/apply_ops_insert_write_conflict.js @@ -0,0 +1,88 @@ +/** + * Sets up a test for WriteConflictException handling in applyOps with an insert workload. + */ +var ApplyOpsInsertWriteConflictTest = function(options) { + 'use strict'; + + if (!(this instanceof ApplyOpsInsertWriteConflictTest)) { + return new ApplyOpsInsertWriteConflictTest(options); + } + + // Capture the 'this' reference + var self = this; + + self.options = options; + + /** + * Runs the test. + */ + this.run = function() { + var options = this.options; + + var replTest = new ReplSetTest({nodes: 1}); + replTest.startSet(); + replTest.initiate(); + + var primary = replTest.getPrimary(); + var primaryDB = primary.getDB('test'); + + var t = primaryDB.getCollection(options.testName); + t.drop(); + + assert.commandWorked(primaryDB.createCollection(t.getName())); + + var numOps = 1000; + var ops = Array(numOps).fill('ignored').map((unused, i) => { + return { + op: 'i', + ns: t.getFullName(), + o: {_id: i} + }; + }); + + if (!options.atomic) { + // Adding a command to the list of operations to prevent the applyOps command from + // applying + // all the operations atomically. + ops.push({ns: "test.$cmd", op: "c", o: {applyOps: []}}); + numOps++; + } + + // Probabilities for WCE are chosen based on empirical testing. + // The probability for WCE during an atomic applyOps should be much smaller than that for + // the non-atomic case because we have to attempt to re-apply the entire batch of 'numOps' + // operations on WCE in the atomic case. + var probability = (options.atomic ? 0.1 : 5.0) / numOps; + + // Set up failpoint to trigger WriteConflictException during write operations. + assert.commandWorked( + primaryDB.adminCommand({setParameter: 1, traceWriteConflictExceptions: true})); + assert.commandWorked(primaryDB.adminCommand({ + configureFailPoint: 'WTWriteConflictException', + mode: {activationProbability: probability} + })); + + // This logs each operation being applied. + var previousLogLevel = + assert.commandWorked(primaryDB.setLogLevel(3, 'replication')).was.replication.verbosity; + + var applyOpsResult = primaryDB.adminCommand({applyOps: ops}); + + // Reset log level. + primaryDB.setLogLevel(previousLogLevel, 'replication'); + + assert.eq( + numOps, + applyOpsResult.applied, + 'number of operations applied did not match list of generated insert operations. ' + + 'applyOps result: ' + tojson(applyOpsResult)); + applyOpsResult.results.forEach((operationSucceeded, i) => { + assert(operationSucceeded, + 'applyOps failed: operation with index ' + i + ' failed: operation: ' + + tojson(ops[i], '', true) + '. applyOps result: ' + tojson(applyOpsResult)); + }); + assert.commandWorked(applyOpsResult); + + replTest.stopSet(); + }; +}; diff --git a/jstests/replsets/maintenance.js b/jstests/replsets/maintenance.js index b1fe94efc0e..8b4765212b8 100644 --- a/jstests/replsets/maintenance.js +++ b/jstests/replsets/maintenance.js @@ -5,7 +5,7 @@ var conns = replTest.startSet({verbose: 1}); var config = replTest.getReplSetConfig(); config.members[0].priority = 2; replTest.initiate(config); -replTest.waitForState(replTest.nodes[0], ReplSetTest.State.PRIMARY, 60000); +replTest.waitForState(replTest.nodes[0], ReplSetTest.State.PRIMARY); // Make sure we have a master var master = replTest.getPrimary(); diff --git a/jstests/replsets/oplog_truncated_on_recovery.js b/jstests/replsets/oplog_truncated_on_recovery.js index 4d469178691..477da7b4d92 100644 --- a/jstests/replsets/oplog_truncated_on_recovery.js +++ b/jstests/replsets/oplog_truncated_on_recovery.js @@ -74,7 +74,7 @@ log(assert.commandWorked(localDB.adminCommand("replSetGetStatus"))); log("restart primary"); replTest.restart(master); - replTest.waitForState(master, ReplSetTest.State.RECOVERING, 90000); + replTest.waitForState(master, ReplSetTest.State.RECOVERING); assert.soon(function() { var mv; diff --git a/jstests/replsets/priority_takeover_one_node_higher_priority.js b/jstests/replsets/priority_takeover_one_node_higher_priority.js index 81f7717a0ee..de20f71c854 100644 --- a/jstests/replsets/priority_takeover_one_node_higher_priority.js +++ b/jstests/replsets/priority_takeover_one_node_higher_priority.js @@ -13,7 +13,7 @@ replSet.startSet(); replSet.initiate(); - replSet.waitForState(replSet.nodes[0], ReplSetTest.State.PRIMARY, 60 * 1000); + replSet.waitForState(replSet.nodes[0], ReplSetTest.State.PRIMARY); var primary = replSet.getPrimary(); replSet.awaitSecondaryNodes(); @@ -21,19 +21,19 @@ // Primary should step down long enough for election to occur on secondary. var config = assert.commandWorked(primary.adminCommand({replSetGetConfig: 1})).config; - var electionTimeoutMillis = config.settings.electionTimeoutMillis; - var stepDownGuardMillis = electionTimeoutMillis * 2; var stepDownException = assert.throws(function() { - primary.adminCommand({replSetStepDown: stepDownGuardMillis / 1000}); + primary.adminCommand({replSetStepDown: replSet.kDefaultTimeoutMS / 1000}); }); assert.neq(-1, tojson(stepDownException).indexOf('error doing query'), 'replSetStepDown did not disconnect client'); // Step down primary and wait for node 1 to be promoted to primary. - replSet.waitForState(replSet.nodes[1], ReplSetTest.State.PRIMARY, 60 * 1000); + replSet.waitForState(replSet.nodes[1], ReplSetTest.State.PRIMARY); + + // Unfreeze node 0 so it can seek election. + assert.commandWorked(primary.adminCommand({replSetFreeze: 0})); // Eventually node 0 will stand for election again because it has a higher priorty. - replSet.waitForState( - replSet.nodes[0], ReplSetTest.State.PRIMARY, stepDownGuardMillis + 60 * 1000); + replSet.waitForState(replSet.nodes[0], ReplSetTest.State.PRIMARY); })(); diff --git a/jstests/replsets/read_committed_with_catalog_changes.js b/jstests/replsets/read_committed_with_catalog_changes.js index dae14da31de..4290391bd73 100644 --- a/jstests/replsets/read_committed_with_catalog_changes.js +++ b/jstests/replsets/read_committed_with_catalog_changes.js @@ -208,13 +208,13 @@ load("jstests/replsets/rslib.js"); // For startSetIfSupportsReadMajority. // they may be passed in to a ScopedThread. function assertReadsBlock(coll) { var res = - coll.runCommand('find', {"readConcern": {"level": "majority"}, "maxTimeMS": 1000}); + coll.runCommand('find', {"readConcern": {"level": "majority"}, "maxTimeMS": 5000}); assert.commandFailedWithCode(res, ErrorCodes.ExceededTimeLimit, "Expected read of " + coll.getFullName() + " to block"); } - function assertReadsSucceed(coll, timeoutMs = 1000) { + function assertReadsSucceed(coll, timeoutMs = 20000) { var res = coll.runCommand('find', {"readConcern": {"level": "majority"}, "maxTimeMS": timeoutMs}); assert.commandWorked(res, 'reading from ' + coll.getFullName()); diff --git a/jstests/replsets/replsetadd_profile.js b/jstests/replsets/replsetadd_profile.js index 641e7ca7cfd..2e396a61eb7 100644 --- a/jstests/replsets/replsetadd_profile.js +++ b/jstests/replsets/replsetadd_profile.js @@ -19,7 +19,7 @@ masterCollection.save({a: 1}); var newNode = replTest.add(); replTest.reInitiate(); -replTest.waitForState(replTest.nodes[1], ReplSetTest.State.SECONDARY, 60 * 1000); +replTest.waitForState(replTest.nodes[1], ReplSetTest.State.SECONDARY); // Allow documents to propagate to new replica set member. replTest.awaitReplication(); diff --git a/jstests/replsets/replsetprio1.js b/jstests/replsets/replsetprio1.js index 16beb851b81..5aee0a33a92 100644 --- a/jstests/replsets/replsetprio1.js +++ b/jstests/replsets/replsetprio1.js @@ -16,10 +16,10 @@ }); // 2 should be master (give this a while to happen, as other nodes might first be elected) - replTest.waitForState(nodes[2], ReplSetTest.State.PRIMARY, 120000); + replTest.waitForState(nodes[2], ReplSetTest.State.PRIMARY); // wait for 1 to not appear to be master (we are about to make it master and need a clean slate // here) - replTest.waitForState(nodes[1], ReplSetTest.State.SECONDARY, 60000); + replTest.waitForState(nodes[1], ReplSetTest.State.SECONDARY); // Wait for election oplog entry to be replicated, to ensure 0 will vote for 1 after stopping 2. replTest.awaitReplication(); @@ -28,7 +28,7 @@ replTest.stop(2); // 1 should eventually be master - replTest.waitForState(nodes[1], ReplSetTest.State.PRIMARY, 60000); + replTest.waitForState(nodes[1], ReplSetTest.State.PRIMARY); // do some writes on 1 var master = replTest.getPrimary(); @@ -42,7 +42,7 @@ // bring 2 back up, 2 should wait until caught up and then become master replTest.restart(2); - replTest.waitForState(nodes[2], ReplSetTest.State.PRIMARY, 60000); + replTest.waitForState(nodes[2], ReplSetTest.State.PRIMARY); // make sure nothing was rolled back master = replTest.getPrimary(); diff --git a/jstests/replsets/request_primary_stepdown.js b/jstests/replsets/request_primary_stepdown.js index 02050bd55f4..3e56946397b 100644 --- a/jstests/replsets/request_primary_stepdown.js +++ b/jstests/replsets/request_primary_stepdown.js @@ -15,7 +15,7 @@ conf.protocolVersion = 0; replSet.initiate(conf); - replSet.waitForState(replSet.nodes[0], ReplSetTest.State.PRIMARY, 60 * 1000); + replSet.waitForState(replSet.nodes[0], ReplSetTest.State.PRIMARY); replSet.awaitSecondaryNodes(); replSet.awaitReplication(); var primary = replSet.getPrimary(); diff --git a/jstests/replsets/resync_with_write_load.js b/jstests/replsets/resync_with_write_load.js index 1a782ffacbe..13b6042ab58 100644 --- a/jstests/replsets/resync_with_write_load.js +++ b/jstests/replsets/resync_with_write_load.js @@ -19,7 +19,7 @@ var config = { ] }; var r = replTest.initiate(config); -replTest.waitForState(replTest.nodes[0], ReplSetTest.State.PRIMARY, 60 * 1000); +replTest.waitForState(replTest.nodes[0], ReplSetTest.State.PRIMARY); // Make sure we have a master var master = replTest.getPrimary(); var a_conn = conns[0]; diff --git a/jstests/replsets/rollback.js b/jstests/replsets/rollback.js index 56aaed37d6f..6baeb88666a 100644 --- a/jstests/replsets/rollback.js +++ b/jstests/replsets/rollback.js @@ -46,7 +46,7 @@ load("jstests/replsets/rslib.js"); }); // Make sure we have a master - replTest.waitForState(replTest.nodes[0], ReplSetTest.State.PRIMARY, 60 * 1000); + replTest.waitForState(replTest.nodes[0], ReplSetTest.State.PRIMARY); var master = replTest.getPrimary(); var a_conn = conns[0]; var A = a_conn.getDB("admin"); diff --git a/jstests/replsets/rollback2.js b/jstests/replsets/rollback2.js index cd4f5a049b7..1cbf8194880 100644 --- a/jstests/replsets/rollback2.js +++ b/jstests/replsets/rollback2.js @@ -42,7 +42,7 @@ load("jstests/replsets/rslib.js"); }); // Make sure we have a master and that that master is node A - replTest.waitForState(replTest.nodes[0], ReplSetTest.State.PRIMARY, 60 * 1000); + replTest.waitForState(replTest.nodes[0], ReplSetTest.State.PRIMARY); var master = replTest.getPrimary(); var a_conn = conns[0]; a_conn.setSlaveOk(); diff --git a/jstests/replsets/rollback3.js b/jstests/replsets/rollback3.js index 59b685c7bad..5eb9ba574d1 100644 --- a/jstests/replsets/rollback3.js +++ b/jstests/replsets/rollback3.js @@ -47,7 +47,7 @@ load("jstests/replsets/rslib.js"); }); // Make sure we have a master and that that master is node A - replTest.waitForState(replTest.nodes[0], ReplSetTest.State.PRIMARY, 60 * 1000); + replTest.waitForState(replTest.nodes[0], ReplSetTest.State.PRIMARY); var master = replTest.getPrimary(); var a_conn = conns[0]; a_conn.setSlaveOk(); diff --git a/jstests/replsets/rollback5.js b/jstests/replsets/rollback5.js index cc0007822dd..81614447607 100644 --- a/jstests/replsets/rollback5.js +++ b/jstests/replsets/rollback5.js @@ -23,7 +23,7 @@ var r = replTest.initiate({ }); // Make sure we have a master -replTest.waitForState(replTest.nodes[0], ReplSetTest.State.PRIMARY, 60 * 1000); +replTest.waitForState(replTest.nodes[0], ReplSetTest.State.PRIMARY); var master = replTest.getPrimary(); var a_conn = conns[0]; var b_conn = conns[1]; diff --git a/jstests/replsets/rollback_auth.js b/jstests/replsets/rollback_auth.js index 0c0b35b91ed..4b31a0e527e 100644 --- a/jstests/replsets/rollback_auth.js +++ b/jstests/replsets/rollback_auth.js @@ -39,7 +39,7 @@ }); // Make sure we have a master - replTest.waitForState(replTest.nodes[0], ReplSetTest.State.PRIMARY, 60 * 1000); + replTest.waitForState(replTest.nodes[0], ReplSetTest.State.PRIMARY); var master = replTest.getPrimary(); var a_conn = conns[0]; var b_conn = conns[1]; diff --git a/jstests/replsets/rollback_cmd_unrollbackable.js b/jstests/replsets/rollback_cmd_unrollbackable.js index 41b8f77f74f..01c0bc5bf38 100644 --- a/jstests/replsets/rollback_cmd_unrollbackable.js +++ b/jstests/replsets/rollback_cmd_unrollbackable.js @@ -26,7 +26,7 @@ var AID = replTest.getNodeId(a_conn); var BID = replTest.getNodeId(b_conn); // get master and do an initial write -replTest.waitForState(replTest.nodes[0], ReplSetTest.State.PRIMARY, 60 * 1000); +replTest.waitForState(replTest.nodes[0], ReplSetTest.State.PRIMARY); var master = replTest.getPrimary(); assert(master === conns[0], "conns[0] assumed to be master"); assert(a_conn.host === master.host, "a_conn assumed to be master"); diff --git a/jstests/replsets/rollback_collMod_fatal.js b/jstests/replsets/rollback_collMod_fatal.js index 61af67b15d0..76a43c8cd4c 100644 --- a/jstests/replsets/rollback_collMod_fatal.js +++ b/jstests/replsets/rollback_collMod_fatal.js @@ -25,7 +25,7 @@ var b_conn = conns[1]; var AID = replTest.getNodeId(a_conn); var BID = replTest.getNodeId(b_conn); -replTest.waitForState(replTest.nodes[0], ReplSetTest.State.PRIMARY, 60 * 1000); +replTest.waitForState(replTest.nodes[0], ReplSetTest.State.PRIMARY); // get master and do an initial write var master = replTest.getPrimary(); diff --git a/jstests/replsets/rollback_different_h.js b/jstests/replsets/rollback_different_h.js index 4b9aede1bbc..37f6f0a71cd 100644 --- a/jstests/replsets/rollback_different_h.js +++ b/jstests/replsets/rollback_different_h.js @@ -36,7 +36,7 @@ var b_conn = conns[1]; var AID = replTest.getNodeId(a_conn); var BID = replTest.getNodeId(b_conn); -replTest.waitForState(replTest.nodes[0], ReplSetTest.State.PRIMARY, 60 * 1000); +replTest.waitForState(replTest.nodes[0], ReplSetTest.State.PRIMARY); // get master and do an initial write var master = replTest.getPrimary(); diff --git a/jstests/replsets/rollback_dropdb.js b/jstests/replsets/rollback_dropdb.js index b8f7d8d09ee..5853f7c47ce 100644 --- a/jstests/replsets/rollback_dropdb.js +++ b/jstests/replsets/rollback_dropdb.js @@ -25,7 +25,7 @@ var b_conn = conns[1]; var AID = replTest.getNodeId(a_conn); var BID = replTest.getNodeId(b_conn); -replTest.waitForState(replTest.nodes[0], ReplSetTest.State.PRIMARY, 60 * 1000); +replTest.waitForState(replTest.nodes[0], ReplSetTest.State.PRIMARY); // get master and do an initial write var master = replTest.getPrimary(); diff --git a/jstests/replsets/rollback_fake_cmd.js b/jstests/replsets/rollback_fake_cmd.js index 6eba30d2c16..fcc5fbaf39b 100644 --- a/jstests/replsets/rollback_fake_cmd.js +++ b/jstests/replsets/rollback_fake_cmd.js @@ -36,7 +36,7 @@ var b_conn = conns[1]; var AID = replTest.getNodeId(a_conn); var BID = replTest.getNodeId(b_conn); -replTest.waitForState(replTest.nodes[0], ReplSetTest.State.PRIMARY, 60 * 1000); +replTest.waitForState(replTest.nodes[0], ReplSetTest.State.PRIMARY); // get master and do an initial write var master = replTest.getPrimary(); diff --git a/jstests/replsets/rollback_index.js b/jstests/replsets/rollback_index.js index d3bf747680c..66b2be66f63 100644 --- a/jstests/replsets/rollback_index.js +++ b/jstests/replsets/rollback_index.js @@ -38,7 +38,7 @@ var b_conn = conns[1]; var AID = replTest.getNodeId(a_conn); var BID = replTest.getNodeId(b_conn); -replTest.waitForState(replTest.nodes[0], ReplSetTest.State.PRIMARY, 60 * 1000); +replTest.waitForState(replTest.nodes[0], ReplSetTest.State.PRIMARY); // get master and do an initial write var master = replTest.getPrimary(); diff --git a/jstests/replsets/rslib.js b/jstests/replsets/rslib.js index b97f79c00d1..3b41531d760 100644 --- a/jstests/replsets/rslib.js +++ b/jstests/replsets/rslib.js @@ -138,7 +138,7 @@ var startSetIfSupportsReadMajority; } printjson(state); return true; - }, "not all members ready", timeout || 60000); + }, "not all members ready", timeout || 10 * 60 * 1000); print("All members are now in state PRIMARY, SECONDARY, or ARBITER"); }; diff --git a/jstests/replsets/stepdown.js b/jstests/replsets/stepdown.js index 5a8388da4d4..bed42cff3e5 100644 --- a/jstests/replsets/stepdown.js +++ b/jstests/replsets/stepdown.js @@ -20,7 +20,7 @@ var replTest = new ReplSetTest({ }); var nodes = replTest.startSet(); replTest.initiate(); -replTest.waitForState(nodes[0], ReplSetTest.State.PRIMARY, 60 * 1000); +replTest.waitForState(nodes[0], ReplSetTest.State.PRIMARY); var master = replTest.getPrimary(); // do a write diff --git a/jstests/replsets/stepdown_kill_other_ops.js b/jstests/replsets/stepdown_kill_other_ops.js index 930206046c1..d360624bc7a 100644 --- a/jstests/replsets/stepdown_kill_other_ops.js +++ b/jstests/replsets/stepdown_kill_other_ops.js @@ -15,7 +15,7 @@ ] }); - replSet.waitForState(replSet.nodes[0], ReplSetTest.State.PRIMARY, 60 * 1000); + replSet.waitForState(replSet.nodes[0], ReplSetTest.State.PRIMARY); var primary = replSet.getPrimary(); assert.eq(primary.host, nodes[0], "primary assumed to be node 0"); diff --git a/jstests/replsets/stepdown_killop.js b/jstests/replsets/stepdown_killop.js index 5c0e0ffae91..1b86d5fb8d4 100644 --- a/jstests/replsets/stepdown_killop.js +++ b/jstests/replsets/stepdown_killop.js @@ -23,7 +23,7 @@ ] }); - replSet.waitForState(replSet.nodes[0], ReplSetTest.State.PRIMARY, 60 * 1000); + replSet.waitForState(replSet.nodes[0], ReplSetTest.State.PRIMARY); var secondary = replSet.getSecondary(); jsTestLog('Disable replication on the SECONDARY ' + secondary.host); diff --git a/jstests/replsets/stepdown_long_wait_time.js b/jstests/replsets/stepdown_long_wait_time.js index 60e0fdb4247..eb1bf7007d9 100644 --- a/jstests/replsets/stepdown_long_wait_time.js +++ b/jstests/replsets/stepdown_long_wait_time.js @@ -22,7 +22,7 @@ ] }); - replSet.waitForState(replSet.nodes[0], ReplSetTest.State.PRIMARY, 60 * 1000); + replSet.waitForState(replSet.nodes[0], ReplSetTest.State.PRIMARY); var primary = replSet.getPrimary(); var secondary = replSet.getSecondary(); diff --git a/jstests/replsets/sync_passive.js b/jstests/replsets/sync_passive.js index 4899385563f..c0be375b98b 100644 --- a/jstests/replsets/sync_passive.js +++ b/jstests/replsets/sync_passive.js @@ -29,7 +29,7 @@ config.members[0].priority = 2; config.members[2].priority = 0; replTest.initiate(config); -replTest.waitForState(replTest.nodes[0], ReplSetTest.State.PRIMARY, 60 * 1000); +replTest.waitForState(replTest.nodes[0], ReplSetTest.State.PRIMARY); var master = replTest.getPrimary().getDB("test"); var server0 = master; diff --git a/jstests/replsets/two_nodes_priority_take_over.js b/jstests/replsets/two_nodes_priority_take_over.js index f6e62fe681d..ae0d53e92b6 100644 --- a/jstests/replsets/two_nodes_priority_take_over.js +++ b/jstests/replsets/two_nodes_priority_take_over.js @@ -29,7 +29,7 @@ if (false) { }); // The first node will be the primary at the beginning. - rst.waitForState(rst.nodes[0], ReplSetTest.State.PRIMARY, 60 * 1000); + rst.waitForState(rst.nodes[0], ReplSetTest.State.PRIMARY); // Get the term when replset is stable. var res = rst.getPrimary().adminCommand("replSetGetStatus"); diff --git a/jstests/sharding/balance_repl.js b/jstests/sharding/balance_repl.js index 46404646995..59cc694fc42 100644 --- a/jstests/sharding/balance_repl.js +++ b/jstests/sharding/balance_repl.js @@ -3,7 +3,7 @@ // (function() { - "use strict"; + 'use strict'; // The mongod secondaries are set to priority 0 and votes 0 to prevent the primaries // from stepping down during migrations on slow evergreen builders. @@ -27,14 +27,16 @@ } assert.writeOK(bulk.execute()); - s.adminCommand({enablesharding: "test"}); + assert.commandWorked(s.s0.adminCommand({enablesharding: "test"})); s.ensurePrimaryShard('test', 'test-rs0'); - s.adminCommand({shardcollection: "test.foo", key: {_id: 1}}); + assert.commandWorked(s.s0.adminCommand({shardcollection: "test.foo", key: {_id: 1}})); - for (i = 0; i < 20; i++) - s.adminCommand({split: "test.foo", middle: {_id: i * 100}}); + for (i = 0; i < 20; i++) { + assert.commandWorked(s.s0.adminCommand({split: "test.foo", middle: {_id: i * 100}})); + } assert.eq(2100, db.foo.find().itcount()); + var coll = db.foo; coll.setSlaveOk(); @@ -42,10 +44,9 @@ var other = s.config.shards.findOne({_id: {$ne: serverName}}); for (i = 0; i < 20; i++) { - // Needs to waitForDelete because we'll be performing a slaveOk query, - // and secondaries don't have a chunk manager so it doesn't know how to - // filter out docs it doesn't own. - assert(s.adminCommand({ + // Needs to waitForDelete because we'll be performing a slaveOk query, and secondaries don't + // have a chunk manager so it doesn't know how to filter out docs it doesn't own. + assert.commandWorked(s.s0.adminCommand({ moveChunk: "test.foo", find: {_id: i * 100}, to: other._id, @@ -53,9 +54,9 @@ writeConcern: {w: 2}, _waitForDelete: true })); + assert.eq(2100, coll.find().itcount()); } s.stop(); - }()); diff --git a/jstests/sharding/balance_tags2.js b/jstests/sharding/balance_tags2.js index e4bf370d1cd..a7f1161d6dc 100644 --- a/jstests/sharding/balance_tags2.js +++ b/jstests/sharding/balance_tags2.js @@ -1,27 +1,26 @@ // Test balancing all chunks to one shard by tagging the full shard-key range on that collection -var s = new ShardingTest( - {name: "balance_tags2", shards: 3, mongos: 1, other: {chunkSize: 1, enableBalancer: true}}); +var s = new ShardingTest({shards: 3, mongos: 1, other: {chunkSize: 1, enableBalancer: true}}); -s.adminCommand({enablesharding: "test"}); +assert.commandWorked(s.s0.adminCommand({enablesharding: "test"})); s.ensurePrimaryShard('test', 'shard0001'); var db = s.getDB("test"); var bulk = db.foo.initializeUnorderedBulkOp(); -for (i = 0; i < 21; i++) { +for (var i = 0; i < 21; i++) { bulk.insert({_id: i, x: i}); } assert.writeOK(bulk.execute()); -sh.shardCollection("test.foo", {_id: 1}); +assert.commandWorked(s.s0.adminCommand({shardCollection: "test.foo", key: {_id: 1}})); -sh.stopBalancer(); +s.stopBalancer(); -for (i = 0; i < 20; i++) { +for (var i = 0; i < 20; i++) { sh.splitAt("test.foo", {_id: i}); } -sh.startBalancer(); +s.startBalancer(); sh.status(true); diff --git a/jstests/sharding/explain_cmd.js b/jstests/sharding/explain_cmd.js index 3b8a8ef1240..b4ec0db35e9 100644 --- a/jstests/sharding/explain_cmd.js +++ b/jstests/sharding/explain_cmd.js @@ -1,174 +1,185 @@ // Tests for the mongos explain command. - -// Create a cluster with 3 shards. -var st = new ShardingTest({shards: 2}); -st.stopBalancer(); - -var db = st.s.getDB("test"); -var explain; - -// Setup a collection that will be sharded. The shard key will be 'a'. There's also an index on 'b'. -var collSharded = db.getCollection("mongos_explain_cmd"); -collSharded.drop(); -collSharded.ensureIndex({a: 1}); -collSharded.ensureIndex({b: 1}); - -// Enable sharding. -assert.commandWorked(db.adminCommand({enableSharding: db.getName()})); -st.ensurePrimaryShard(db.getName(), 'shard0001'); -db.adminCommand({shardCollection: collSharded.getFullName(), key: {a: 1}}); - -// Pre-split the collection to ensure that both shards have chunks. Explicitly -// move chunks since the balancer is disabled. -for (var i = 1; i <= 2; i++) { - assert.commandWorked(db.adminCommand({split: collSharded.getFullName(), middle: {a: i}})); - - var shardName = "shard000" + (i - 1); - printjson(db.adminCommand({moveChunk: collSharded.getFullName(), find: {a: i}, to: shardName})); -} - -// Put data on each shard. -for (var i = 0; i < 3; i++) { - collSharded.insert({_id: i, a: i, b: 1}); -} - -printjson(sh.status()); - -// Test a scatter-gather count command. -assert.eq(3, collSharded.count({b: 1})); - -// Explain the scatter-gather count. -explain = db.runCommand( - {explain: {count: collSharded.getName(), query: {b: 1}}, verbosity: "allPlansExecution"}); - -// Validate some basic properties of the result. -printjson(explain); -assert.commandWorked(explain); -assert("queryPlanner" in explain); -assert("executionStats" in explain); -assert.eq(2, explain.queryPlanner.winningPlan.shards.length); -assert.eq(2, explain.executionStats.executionStages.shards.length); - -// An explain of a command that doesn't exist should fail gracefully. -explain = db.runCommand({ - explain: {nonexistent: collSharded.getName(), query: {b: 1}}, - verbosity: "allPlansExecution" -}); -printjson(explain); -assert.commandFailed(explain); - -// ------- - -// Setup a collection that is not sharded. -var collUnsharded = db.getCollection("mongos_explain_cmd_unsharded"); -collUnsharded.drop(); -collUnsharded.ensureIndex({a: 1}); -collUnsharded.ensureIndex({b: 1}); - -for (var i = 0; i < 3; i++) { - collUnsharded.insert({_id: i, a: i, b: 1}); -} -assert.eq(3, collUnsharded.count({b: 1})); - -explain = db.runCommand({ - explain: { - group: { - ns: collUnsharded.getName(), - key: "a", - cond: "b", - $reduce: function(curr, result) {}, - initial: {} - } - }, - verbosity: "allPlansExecution" -}); - -// Basic validation: a group command can only be passed through to an unsharded collection, -// so we should confirm that the mongos stage is always SINGLE_SHARD. -printjson(explain); -assert.commandWorked(explain); -assert("queryPlanner" in explain); -assert("executionStats" in explain); -assert.eq("SINGLE_SHARD", explain.queryPlanner.winningPlan.stage); - -// The same group should fail over the sharded collection, because group is only supported -// if it is passed through to an unsharded collection. -explain = db.runCommand({ - explain: { - group: { - ns: collSharded.getName(), - key: "a", - cond: "b", - $reduce: function(curr, result) {}, - initial: {} - } - }, - verbosity: "allPlansExecution" -}); -printjson(explain); -assert.commandFailed(explain); - -// ------- - -// Explain a delete operation and verify that it hits all shards without the shard key -explain = db.runCommand({ - explain: {delete: collSharded.getName(), deletes: [{q: {b: 1}, limit: 0}]}, - verbosity: "allPlansExecution" -}); -assert.commandWorked(explain, tojson(explain)); -assert.eq(explain.queryPlanner.winningPlan.stage, "SHARD_WRITE"); -assert.eq(explain.queryPlanner.winningPlan.shards.length, 2); -assert.eq(explain.queryPlanner.winningPlan.shards[0].winningPlan.stage, "DELETE"); -assert.eq(explain.queryPlanner.winningPlan.shards[1].winningPlan.stage, "DELETE"); -// Check that the deletes didn't actually happen. -assert.eq(3, collSharded.count({b: 1})); - -// Explain a delete operation and verify that it hits only one shard with the shard key -explain = db.runCommand({ - explain: {delete: collSharded.getName(), deletes: [{q: {a: 1}, limit: 0}]}, - verbosity: "allPlansExecution" -}); -assert.commandWorked(explain, tojson(explain)); -assert.eq(explain.queryPlanner.winningPlan.shards.length, 1); -// Check that the deletes didn't actually happen. -assert.eq(3, collSharded.count({b: 1})); - -// Check that we fail gracefully if we try to do an explain of a write batch that has more -// than one operation in it. -explain = db.runCommand({ - explain: - {delete: collSharded.getName(), deletes: [{q: {a: 1}, limit: 1}, {q: {a: 2}, limit: 1}]}, - verbosity: "allPlansExecution" -}); -assert.commandFailed(explain, tojson(explain)); - -// Explain a multi upsert operation and verify that it hits all shards -explain = db.runCommand({ - explain: {update: collSharded.getName(), updates: [{q: {}, u: {$set: {b: 10}}, multi: true}]}, - verbosity: "allPlansExecution" -}); -assert.commandWorked(explain, tojson(explain)); -assert.eq(explain.queryPlanner.winningPlan.shards.length, 2); -assert.eq(explain.queryPlanner.winningPlan.stage, "SHARD_WRITE"); -assert.eq(explain.queryPlanner.winningPlan.shards.length, 2); -assert.eq(explain.queryPlanner.winningPlan.shards[0].winningPlan.stage, "UPDATE"); -assert.eq(explain.queryPlanner.winningPlan.shards[1].winningPlan.stage, "UPDATE"); -// Check that the update didn't actually happen. -assert.eq(0, collSharded.count({b: 10})); - -// Explain an upsert operation and verify that it hits only a single shard -explain = db.runCommand({ - explain: {update: collSharded.getName(), updates: [{q: {a: 10}, u: {a: 10}, upsert: true}]}, - verbosity: "allPlansExecution" -}); -assert.commandWorked(explain, tojson(explain)); -assert.eq(explain.queryPlanner.winningPlan.shards.length, 1); -// Check that the upsert didn't actually happen. -assert.eq(0, collSharded.count({a: 10})); - -// Explain an upsert operation which cannot be targeted, ensure an error is thrown -explain = db.runCommand({ - explain: {update: collSharded.getName(), updates: [{q: {b: 10}, u: {b: 10}, upsert: true}]}, - verbosity: "allPlansExecution" -}); -assert.commandFailed(explain, tojson(explain)); +(function() { + 'use strict'; + + // Create a cluster with 3 shards. + var st = new ShardingTest({shards: 2}); + + var db = st.s.getDB("test"); + var explain; + + // Setup a collection that will be sharded. The shard key will be 'a'. There's also an index on + // 'b'. + var collSharded = db.getCollection("mongos_explain_cmd"); + collSharded.drop(); + collSharded.ensureIndex({a: 1}); + collSharded.ensureIndex({b: 1}); + + // Enable sharding. + assert.commandWorked(db.adminCommand({enableSharding: db.getName()})); + st.ensurePrimaryShard(db.getName(), 'shard0001'); + db.adminCommand({shardCollection: collSharded.getFullName(), key: {a: 1}}); + + // Pre-split the collection to ensure that both shards have chunks. Explicitly + // move chunks since the balancer is disabled. + for (var i = 1; i <= 2; i++) { + assert.commandWorked(db.adminCommand({split: collSharded.getFullName(), middle: {a: i}})); + + var shardName = "shard000" + (i - 1); + printjson( + db.adminCommand({moveChunk: collSharded.getFullName(), find: {a: i}, to: shardName})); + } + + // Put data on each shard. + for (var i = 0; i < 3; i++) { + collSharded.insert({_id: i, a: i, b: 1}); + } + + st.printShardingStatus(); + + // Test a scatter-gather count command. + assert.eq(3, collSharded.count({b: 1})); + + // Explain the scatter-gather count. + explain = db.runCommand( + {explain: {count: collSharded.getName(), query: {b: 1}}, verbosity: "allPlansExecution"}); + + // Validate some basic properties of the result. + printjson(explain); + assert.commandWorked(explain); + assert("queryPlanner" in explain); + assert("executionStats" in explain); + assert.eq(2, explain.queryPlanner.winningPlan.shards.length); + assert.eq(2, explain.executionStats.executionStages.shards.length); + + // An explain of a command that doesn't exist should fail gracefully. + explain = db.runCommand({ + explain: {nonexistent: collSharded.getName(), query: {b: 1}}, + verbosity: "allPlansExecution" + }); + printjson(explain); + assert.commandFailed(explain); + + // ------- + + // Setup a collection that is not sharded. + var collUnsharded = db.getCollection("mongos_explain_cmd_unsharded"); + collUnsharded.drop(); + collUnsharded.ensureIndex({a: 1}); + collUnsharded.ensureIndex({b: 1}); + + for (var i = 0; i < 3; i++) { + collUnsharded.insert({_id: i, a: i, b: 1}); + } + assert.eq(3, collUnsharded.count({b: 1})); + + explain = db.runCommand({ + explain: { + group: { + ns: collUnsharded.getName(), + key: "a", + cond: "b", + $reduce: function(curr, result) {}, + initial: {} + } + }, + verbosity: "allPlansExecution" + }); + + // Basic validation: a group command can only be passed through to an unsharded collection, + // so we should confirm that the mongos stage is always SINGLE_SHARD. + printjson(explain); + assert.commandWorked(explain); + assert("queryPlanner" in explain); + assert("executionStats" in explain); + assert.eq("SINGLE_SHARD", explain.queryPlanner.winningPlan.stage); + + // The same group should fail over the sharded collection, because group is only supported + // if it is passed through to an unsharded collection. + explain = db.runCommand({ + explain: { + group: { + ns: collSharded.getName(), + key: "a", + cond: "b", + $reduce: function(curr, result) {}, + initial: {} + } + }, + verbosity: "allPlansExecution" + }); + printjson(explain); + assert.commandFailed(explain); + + // ------- + + // Explain a delete operation and verify that it hits all shards without the shard key + explain = db.runCommand({ + explain: {delete: collSharded.getName(), deletes: [{q: {b: 1}, limit: 0}]}, + verbosity: "allPlansExecution" + }); + assert.commandWorked(explain, tojson(explain)); + assert.eq(explain.queryPlanner.winningPlan.stage, "SHARD_WRITE"); + assert.eq(explain.queryPlanner.winningPlan.shards.length, 2); + assert.eq(explain.queryPlanner.winningPlan.shards[0].winningPlan.stage, "DELETE"); + assert.eq(explain.queryPlanner.winningPlan.shards[1].winningPlan.stage, "DELETE"); + // Check that the deletes didn't actually happen. + assert.eq(3, collSharded.count({b: 1})); + + // Explain a delete operation and verify that it hits only one shard with the shard key + explain = db.runCommand({ + explain: {delete: collSharded.getName(), deletes: [{q: {a: 1}, limit: 0}]}, + verbosity: "allPlansExecution" + }); + assert.commandWorked(explain, tojson(explain)); + assert.eq(explain.queryPlanner.winningPlan.shards.length, 1); + // Check that the deletes didn't actually happen. + assert.eq(3, collSharded.count({b: 1})); + + // Check that we fail gracefully if we try to do an explain of a write batch that has more + // than one operation in it. + explain = db.runCommand({ + explain: { + delete: collSharded.getName(), + deletes: [{q: {a: 1}, limit: 1}, {q: {a: 2}, limit: 1}] + }, + verbosity: "allPlansExecution" + }); + assert.commandFailed(explain, tojson(explain)); + + // Explain a multi upsert operation and verify that it hits all shards + explain = db.runCommand({ + explain: + {update: collSharded.getName(), updates: [{q: {}, u: {$set: {b: 10}}, multi: true}]}, + verbosity: "allPlansExecution" + }); + assert.commandWorked(explain, tojson(explain)); + assert.eq(explain.queryPlanner.winningPlan.shards.length, 2); + assert.eq(explain.queryPlanner.winningPlan.stage, "SHARD_WRITE"); + assert.eq(explain.queryPlanner.winningPlan.shards.length, 2); + assert.eq(explain.queryPlanner.winningPlan.shards[0].winningPlan.stage, "UPDATE"); + assert.eq(explain.queryPlanner.winningPlan.shards[1].winningPlan.stage, "UPDATE"); + // Check that the update didn't actually happen. + assert.eq(0, collSharded.count({b: 10})); + + // Explain an upsert operation and verify that it hits only a single shard + explain = db.runCommand({ + explain: + {update: collSharded.getName(), updates: [{q: {a: 10}, u: {a: 10}, upsert: true}]}, + verbosity: "allPlansExecution" + }); + assert.commandWorked(explain, tojson(explain)); + assert.eq(explain.queryPlanner.winningPlan.shards.length, 1); + // Check that the upsert didn't actually happen. + assert.eq(0, collSharded.count({a: 10})); + + // Explain an upsert operation which cannot be targeted, ensure an error is thrown + explain = db.runCommand({ + explain: + {update: collSharded.getName(), updates: [{q: {b: 10}, u: {b: 10}, upsert: true}]}, + verbosity: "allPlansExecution" + }); + assert.commandFailed(explain, tojson(explain)); + + st.stop(); +})(); diff --git a/jstests/sharding/explain_find_and_modify_sharded.js b/jstests/sharding/explain_find_and_modify_sharded.js index 40af14f6265..e8c69adc222 100644 --- a/jstests/sharding/explain_find_and_modify_sharded.js +++ b/jstests/sharding/explain_find_and_modify_sharded.js @@ -9,7 +9,6 @@ // Create a cluster with 2 shards. var st = new ShardingTest({shards: 2}); - st.stopBalancer(); var testDB = st.s.getDB('test'); var shardKey = { @@ -85,4 +84,5 @@ assert.commandWorked(res); assertExplainResult(res, 'executionStats', 'executionStages', 'shard0001', 'DELETE'); + st.stop(); })(); diff --git a/jstests/sharding/hash_shard_unique_compound.js b/jstests/sharding/hash_shard_unique_compound.js index 5d6d466c1f9..abaf45260b9 100644 --- a/jstests/sharding/hash_shard_unique_compound.js +++ b/jstests/sharding/hash_shard_unique_compound.js @@ -2,44 +2,42 @@ // Does 2 things and checks for consistent error: // 1.) shard collection on hashed "a", ensure unique index {a:1, b:1} // 2.) reverse order +(function() { + 'use strict'; -var s = new ShardingTest({name: jsTestName(), shards: 1, mongos: 1, verbose: 1}); -var dbName = "test"; -var collName = "foo"; -var ns = dbName + "." + collName; -var db = s.getDB(dbName); -var coll = db.getCollection(collName); + var s = new ShardingTest({shards: 1, mongos: 1}); + var dbName = "test"; + var collName = "foo"; + var ns = dbName + "." + collName; + var db = s.getDB(dbName); + var coll = db.getCollection(collName); -// Enable sharding on DB -var res = db.adminCommand({enablesharding: dbName}); + // Enable sharding on DB + assert.commandWorked(db.adminCommand({enablesharding: dbName})); -// for simplicity start by turning off balancer -var res = s.stopBalancer(); + // Shard a fresh collection using a hashed shard key + assert.commandWorked(db.adminCommand({shardcollection: ns, key: {a: "hashed"}})); -// shard a fresh collection using a hashed shard key -coll.drop(); -assert.commandWorked(db.adminCommand({shardcollection: ns, key: {a: "hashed"}})); -db.printShardingStatus(); + // Create unique index + assert.commandWorked(coll.ensureIndex({a: 1, b: 1}, {unique: true})); -// Create unique index -assert.commandWorked(coll.ensureIndex({a: 1, b: 1}, {unique: true})); + jsTest.log("------ indexes -------"); + jsTest.log(tojson(coll.getIndexes())); -jsTest.log("------ indexes -------"); -jsTest.log(tojson(coll.getIndexes())); + // Second Part + jsTest.log("------ dropping sharded collection to start part 2 -------"); + coll.drop(); -// Second Part -jsTest.log("------ dropping sharded collection to start part 2 -------"); -coll.drop(); + // Create unique index + assert.commandWorked(coll.ensureIndex({a: 1, b: 1}, {unique: true})); -// Create unique index -assert.commandWorked(coll.ensureIndex({a: 1, b: 1}, {unique: true})); + // shard a fresh collection using a hashed shard key + assert.commandWorked(db.adminCommand({shardcollection: ns, key: {a: "hashed"}}), + "shardcollection didn't worked 2"); -// shard a fresh collection using a hashed shard key -assert.commandWorked(db.adminCommand({shardcollection: ns, key: {a: "hashed"}}), - "shardcollection didn't worked 2"); + s.printShardingStatus(); + jsTest.log("------ indexes 2-------"); + jsTest.log(tojson(coll.getIndexes())); -db.printShardingStatus(); -jsTest.log("------ indexes 2-------"); -jsTest.log(tojson(coll.getIndexes())); - -s.stop(); + s.stop(); +})(); diff --git a/jstests/sharding/mapReduce_inSharded_outSharded.js b/jstests/sharding/mapReduce_inSharded_outSharded.js index d1aba2599f0..5190a1fe4ba 100644 --- a/jstests/sharding/mapReduce_inSharded_outSharded.js +++ b/jstests/sharding/mapReduce_inSharded_outSharded.js @@ -1,60 +1,70 @@ -var verifyOutput = function(out) { - printjson(out); - assert.eq(out.counts.input, 51200, "input count is wrong"); - assert.eq(out.counts.emit, 51200, "emit count is wrong"); - assert.gt(out.counts.reduce, 99, "reduce count is wrong"); - assert.eq(out.counts.output, 512, "output count is wrong"); -}; - -var st = new ShardingTest( - {shards: 2, verbose: 1, mongos: 1, other: {chunkSize: 1, enableBalancer: true}}); - -st.adminCommand({enablesharding: "mrShard"}); -st.ensurePrimaryShard('mrShard', 'shard0001'); -st.adminCommand({shardcollection: "mrShard.srcSharded", key: {"_id": 1}}); - -var db = st.getDB("mrShard"); - -var bulk = db.srcSharded.initializeUnorderedBulkOp(); -for (j = 0; j < 100; j++) { - for (i = 0; i < 512; i++) { - bulk.insert({j: j, i: i}); +(function() { + "use strict"; + + var verifyOutput = function(out) { + printjson(out); + assert.eq(out.counts.input, 51200, "input count is wrong"); + assert.eq(out.counts.emit, 51200, "emit count is wrong"); + assert.gt(out.counts.reduce, 99, "reduce count is wrong"); + assert.eq(out.counts.output, 512, "output count is wrong"); + }; + + var st = new ShardingTest( + {shards: 2, verbose: 1, mongos: 1, other: {chunkSize: 1, enableBalancer: true}}); + + var admin = st.s0.getDB('admin'); + + assert.commandWorked(admin.runCommand({enablesharding: "mrShard"})); + st.ensurePrimaryShard('mrShard', 'shard0001'); + assert.commandWorked( + admin.runCommand({shardcollection: "mrShard.srcSharded", key: {"_id": 1}})); + + var db = st.s0.getDB("mrShard"); + + var bulk = db.srcSharded.initializeUnorderedBulkOp(); + for (var j = 0; j < 100; j++) { + for (var i = 0; i < 512; i++) { + bulk.insert({j: j, i: i}); + } + } + assert.writeOK(bulk.execute()); + + function map() { + emit(this.i, 1); + } + function reduce(key, values) { + return Array.sum(values); } -} -assert.writeOK(bulk.execute()); - -function map() { - emit(this.i, 1); -} -function reduce(key, values) { - return Array.sum(values); -} - -// sharded src sharded dst -var suffix = "InShardedOutSharded"; - -var out = - db.srcSharded.mapReduce(map, reduce, {out: {replace: "mrReplace" + suffix, sharded: true}}); -verifyOutput(out); - -out = db.srcSharded.mapReduce(map, reduce, {out: {merge: "mrMerge" + suffix, sharded: true}}); -verifyOutput(out); - -out = db.srcSharded.mapReduce(map, reduce, {out: {reduce: "mrReduce" + suffix, sharded: true}}); -verifyOutput(out); - -out = db.srcSharded.mapReduce(map, reduce, {out: {inline: 1}}); -verifyOutput(out); -assert(out.results != 'undefined', "no results for inline"); - -out = db.srcSharded.mapReduce( - map, reduce, {out: {replace: "mrReplace" + suffix, db: "mrShardOtherDB", sharded: true}}); -verifyOutput(out); - -out = db.runCommand({ - mapReduce: "srcSharded", // use new name mapReduce rather than mapreduce - map: map, - reduce: reduce, - out: "mrBasic" + "srcSharded", -}); -verifyOutput(out); + + // sharded src sharded dst + var suffix = "InShardedOutSharded"; + + var out = db.srcSharded.mapReduce( + map, reduce, {out: {replace: "mrReplace" + suffix, sharded: true}}); + verifyOutput(out); + + out = db.srcSharded.mapReduce(map, reduce, {out: {merge: "mrMerge" + suffix, sharded: true}}); + verifyOutput(out); + + out = db.srcSharded.mapReduce(map, reduce, {out: {reduce: "mrReduce" + suffix, sharded: true}}); + verifyOutput(out); + + out = db.srcSharded.mapReduce(map, reduce, {out: {inline: 1}}); + verifyOutput(out); + assert(out.results != 'undefined', "no results for inline"); + + out = db.srcSharded.mapReduce( + map, reduce, {out: {replace: "mrReplace" + suffix, db: "mrShardOtherDB", sharded: true}}); + verifyOutput(out); + + out = db.runCommand({ + mapReduce: "srcSharded", // use new name mapReduce rather than mapreduce + map: map, + reduce: reduce, + out: "mrBasic" + "srcSharded", + }); + verifyOutput(out); + + st.stop(); + +})(); diff --git a/jstests/sharding/migrateBig.js b/jstests/sharding/migrateBig.js index e11782baed2..6e6be382795 100644 --- a/jstests/sharding/migrateBig.js +++ b/jstests/sharding/migrateBig.js @@ -1,64 +1,63 @@ (function() { + 'use strict'; var s = new ShardingTest({name: "migrateBig", shards: 2, other: {chunkSize: 1}}); - s.config.settings.update({_id: "balancer"}, {$set: {_waitForDelete: true}}, true); - s.adminCommand({enablesharding: "test"}); + assert.writeOK( + s.config.settings.update({_id: "balancer"}, {$set: {_waitForDelete: true}}, true)); + assert.commandWorked(s.s0.adminCommand({enablesharding: "test"})); s.ensurePrimaryShard('test', 'shard0001'); - s.adminCommand({shardcollection: "test.foo", key: {x: 1}}); + assert.commandWorked(s.s0.adminCommand({shardcollection: "test.foo", key: {x: 1}})); - db = s.getDB("test"); - coll = db.foo; + var db = s.getDB("test"); + var coll = db.foo; - big = ""; + var big = ""; while (big.length < 10000) big += "eliot"; var bulk = coll.initializeUnorderedBulkOp(); - for (x = 0; x < 100; x++) { + for (var x = 0; x < 100; x++) { bulk.insert({x: x, big: big}); } assert.writeOK(bulk.execute()); - db.printShardingStatus(); - - s.adminCommand({split: "test.foo", middle: {x: 30}}); - s.adminCommand({split: "test.foo", middle: {x: 66}}); - s.adminCommand( - {movechunk: "test.foo", find: {x: 90}, to: s.getOther(s.getServer("test")).name}); + assert.commandWorked(s.s0.adminCommand({split: "test.foo", middle: {x: 30}})); + assert.commandWorked(s.s0.adminCommand({split: "test.foo", middle: {x: 66}})); + assert.commandWorked(s.s0.adminCommand( + {movechunk: "test.foo", find: {x: 90}, to: s.getOther(s.getServer("test")).name})); db.printShardingStatus(); print("YO : " + s.getServer("test").host); - direct = new Mongo(s.getServer("test").host); + var direct = new Mongo(s.getServer("test").host); print("direct : " + direct); - directDB = direct.getDB("test"); + var directDB = direct.getDB("test"); - for (done = 0; done < 2 * 1024 * 1024; done += big.length) { + for (var done = 0; done < 2 * 1024 * 1024; done += big.length) { assert.writeOK(directDB.foo.insert({x: 50 + Math.random(), big: big})); } db.printShardingStatus(); assert.throws(function() { - s.adminCommand( - {movechunk: "test.foo", find: {x: 50}, to: s.getOther(s.getServer("test")).name}); + assert.commandWorked(s.s0.adminCommand( + {movechunk: "test.foo", find: {x: 50}, to: s.getOther(s.getServer("test")).name})); }, [], "move should fail"); - for (i = 0; i < 20; i += 2) { + for (var i = 0; i < 20; i += 2) { try { - s.adminCommand({split: "test.foo", middle: {x: i}}); + assert.commandWorked(s.s0.adminCommand({split: "test.foo", middle: {x: i}})); } catch (e) { - // we may have auto split on some of these - // which is ok + // We may have auto split on some of these, which is ok print(e); } } db.printShardingStatus(); - s.config.settings.update({_id: "balancer"}, {$set: {stopped: false}}, true); + s.startBalancer(); assert.soon(function() { var x = s.chunkDiff("foo", "test"); @@ -73,5 +72,4 @@ assert.eq(coll.count(), coll.find().itcount()); s.stop(); - })(); diff --git a/jstests/sharding/migrateBig_balancer.js b/jstests/sharding/migrateBig_balancer.js index cd44a225a62..906ed341c7c 100644 --- a/jstests/sharding/migrateBig_balancer.js +++ b/jstests/sharding/migrateBig_balancer.js @@ -1,11 +1,12 @@ (function() { + 'use strict'; var st = new ShardingTest({name: 'migrateBig_balancer', shards: 2, other: {enableBalancer: true}}); var mongos = st.s; var admin = mongos.getDB("admin"); - db = mongos.getDB("test"); + var db = mongos.getDB("test"); var coll = db.getCollection("stuff"); assert.commandWorked(admin.runCommand({enablesharding: coll.getDB().getName()})); @@ -18,7 +19,7 @@ for (var i = 0; i < nsq; i++) data += data; - dataObj = {}; + var dataObj = {}; for (var i = 0; i < n; i++) dataObj["data-" + i] = data; @@ -30,19 +31,16 @@ assert.eq(40, coll.count(), "prep1"); - printjson(coll.stats()); - - admin.printShardingStatus(); - - admin.runCommand({shardcollection: "" + coll, key: {_id: 1}}); + assert.commandWorked(admin.runCommand({shardcollection: "" + coll, key: {_id: 1}})); + st.printShardingStatus(); assert.lt( 5, mongos.getDB("config").chunks.find({ns: "test.stuff"}).count(), "not enough chunks"); assert.soon(function() { // On *extremely* slow or variable systems, we've seen migrations fail in the critical - // section and - // kill the server. Do an explicit check for this. SERVER-8781 + // section and kill the server. Do an explicit check for this. SERVER-8781 + // // TODO: Remove once we can better specify what systems to run what tests on. try { assert.commandWorked(st.shard0.getDB("admin").runCommand({ping: 1})); @@ -53,7 +51,7 @@ throw e; } - res = mongos.getDB("config").chunks.group({ + var res = mongos.getDB("config").chunks.group({ cond: {ns: "test.stuff"}, key: {shard: 1}, reduce: function(doc, out) { @@ -68,5 +66,4 @@ }, "never migrated", 10 * 60 * 1000, 1000); st.stop(); - })(); diff --git a/jstests/sharding/printShardingStatus.js b/jstests/sharding/printShardingStatus.js index 05e6eca0d4f..5bfa70c2d8f 100644 --- a/jstests/sharding/printShardingStatus.js +++ b/jstests/sharding/printShardingStatus.js @@ -3,6 +3,7 @@ // headings and the names of sharded collections and their shard keys. (function() { + 'use strict'; var st = new ShardingTest({shards: 1, mongos: 2, config: 1, other: {smallfiles: true}}); @@ -233,5 +234,4 @@ assert(mongos.getDB("test").dropDatabase()); st.stop(); - })(); diff --git a/jstests/sharding/shard3.js b/jstests/sharding/shard3.js index 6800e3f4370..290a9f79719 100644 --- a/jstests/sharding/shard3.js +++ b/jstests/sharding/shard3.js @@ -1,5 +1,4 @@ (function() { - // Include helpers for analyzing explain output. load("jstests/libs/analyze_plan.js"); @@ -17,11 +16,14 @@ } assert(sh.getBalancerState(), "A1"); - sh.setBalancerState(false); + + sh.stopBalancer(); assert(!sh.getBalancerState(), "A2"); - sh.setBalancerState(true); + + sh.startBalancer(); assert(sh.getBalancerState(), "A3"); - sh.setBalancerState(false); + + sh.stopBalancer(); assert(!sh.getBalancerState(), "A4"); s.config.databases.find().forEach(printjson); diff --git a/jstests/sharding/split_with_force_small.js b/jstests/sharding/split_with_force_small.js index 0148c924993..be21049650e 100644 --- a/jstests/sharding/split_with_force_small.js +++ b/jstests/sharding/split_with_force_small.js @@ -1,70 +1,69 @@ // // Tests autosplit locations with force : true, for small collections // +(function() { + 'use strict'; -var options = { - chunkSize: 1 // MB -}; + var st = new ShardingTest( + {shards: 1, mongos: 1, other: {chunkSize: 1, mongosOptions: {noAutoSplit: ""}}}); -var st = new ShardingTest({shards: 1, mongos: 1, other: options}); -st.stopBalancer(); + var mongos = st.s0; + var admin = mongos.getDB("admin"); + var config = mongos.getDB("config"); + var shardAdmin = st.shard0.getDB("admin"); + var coll = mongos.getCollection("foo.bar"); -var mongos = st.s0; -var admin = mongos.getDB("admin"); -var config = mongos.getDB("config"); -var shardAdmin = st.shard0.getDB("admin"); -var coll = mongos.getCollection("foo.bar"); + assert.commandWorked(admin.runCommand({enableSharding: coll.getDB() + ""})); + assert.commandWorked(admin.runCommand({shardCollection: coll + "", key: {_id: 1}})); + assert.commandWorked(admin.runCommand({split: coll + "", middle: {_id: 0}})); -assert(admin.runCommand({enableSharding: coll.getDB() + ""}).ok); -assert(admin.runCommand({shardCollection: coll + "", key: {_id: 1}}).ok); -assert(admin.runCommand({split: coll + "", middle: {_id: 0}}).ok); + jsTest.log("Insert a bunch of data into the low chunk of a collection," + + " to prevent relying on stats."); -jsTest.log("Insert a bunch of data into the low chunk of a collection," + - " to prevent relying on stats."); + var data128k = "x"; + for (var i = 0; i < 7; i++) + data128k += data128k; -var data128k = "x"; -for (var i = 0; i < 7; i++) - data128k += data128k; + var bulk = coll.initializeUnorderedBulkOp(); + for (var i = 0; i < 1024; i++) { + bulk.insert({_id: -(i + 1)}); + } + assert.writeOK(bulk.execute()); -var bulk = coll.initializeUnorderedBulkOp(); -for (var i = 0; i < 1024; i++) { - bulk.insert({_id: -(i + 1)}); -} -assert.writeOK(bulk.execute()); + jsTest.log("Insert 32 docs into the high chunk of a collection"); -jsTest.log("Insert 32 docs into the high chunk of a collection"); + bulk = coll.initializeUnorderedBulkOp(); + for (var i = 0; i < 32; i++) { + bulk.insert({_id: i}); + } + assert.writeOK(bulk.execute()); -bulk = coll.initializeUnorderedBulkOp(); -for (var i = 0; i < 32; i++) { - bulk.insert({_id: i}); -} -assert.writeOK(bulk.execute()); + jsTest.log("Split off MaxKey chunk..."); -jsTest.log("Split off MaxKey chunk..."); + assert.commandWorked(admin.runCommand({split: coll + "", middle: {_id: 32}})); -assert(admin.runCommand({split: coll + "", middle: {_id: 32}}).ok); + jsTest.log("Keep splitting chunk multiple times..."); -jsTest.log("Keep splitting chunk multiple times..."); - -st.printShardingStatus(); - -for (var i = 0; i < 5; i++) { - assert(admin.runCommand({split: coll + "", find: {_id: 0}}).ok); st.printShardingStatus(); -} -// Make sure we can't split further than 5 (2^5) times -assert(!admin.runCommand({split: coll + "", find: {_id: 0}}).ok); + for (var i = 0; i < 5; i++) { + assert.commandWorked(admin.runCommand({split: coll + "", find: {_id: 0}})); + st.printShardingStatus(); + } + + // Make sure we can't split further than 5 (2^5) times + assert.commandFailed(admin.runCommand({split: coll + "", find: {_id: 0}})); -var chunks = config.chunks.find({'min._id': {$gte: 0, $lt: 32}}).sort({min: 1}).toArray(); -printjson(chunks); + var chunks = config.chunks.find({'min._id': {$gte: 0, $lt: 32}}).sort({min: 1}).toArray(); + printjson(chunks); -// Make sure the chunks grow by 2x (except the first) -var nextSize = 1; -for (var i = 0; i < chunks.size; i++) { - assert.eq(coll.count({_id: {$gte: chunks[i].min._id, $lt: chunks[i].max._id}}), nextSize); - if (i != 0) - nextSize += nextSize; -} + // Make sure the chunks grow by 2x (except the first) + var nextSize = 1; + for (var i = 0; i < chunks.size; i++) { + assert.eq(coll.count({_id: {$gte: chunks[i].min._id, $lt: chunks[i].max._id}}), nextSize); + if (i != 0) + nextSize += nextSize; + } -st.stop(); + st.stop(); +})(); diff --git a/jstests/sharding/stale_version_write.js b/jstests/sharding/stale_version_write.js index e5885dcfa41..bd603124548 100644 --- a/jstests/sharding/stale_version_write.js +++ b/jstests/sharding/stale_version_write.js @@ -1,37 +1,37 @@ // Tests whether a reset sharding version triggers errors +(function() { + 'use strict'; -jsTest.log("Starting sharded cluster..."); + var st = new ShardingTest({shards: 1, mongos: 2}); -var st = new ShardingTest({shards: 1, mongos: 2, verbose: 2}); + var mongosA = st.s0; + var mongosB = st.s1; -st.stopBalancer(); + jsTest.log("Adding new collections..."); -var mongosA = st.s0; -var mongosB = st.s1; + var collA = mongosA.getCollection(jsTestName() + ".coll"); + assert.writeOK(collA.insert({hello: "world"})); -jsTest.log("Adding new collections..."); + var collB = mongosB.getCollection("" + collA); + assert.writeOK(collB.insert({hello: "world"})); -var collA = mongosA.getCollection(jsTestName() + ".coll"); -assert.writeOK(collA.insert({hello: "world"})); + jsTest.log("Enabling sharding..."); -var collB = mongosB.getCollection("" + collA); -assert.writeOK(collB.insert({hello: "world"})); + assert.commandWorked(mongosA.getDB("admin").adminCommand({enableSharding: "" + collA.getDB()})); + assert.commandWorked( + mongosA.getDB("admin").adminCommand({shardCollection: "" + collA, key: {_id: 1}})); -jsTest.log("Enabling sharding..."); + // MongoD doesn't know about the config shard version *until* MongoS tells it + collA.findOne(); -printjson(mongosA.getDB("admin").runCommand({enableSharding: "" + collA.getDB()})); -printjson(mongosA.getDB("admin").runCommand({shardCollection: "" + collA, key: {_id: 1}})); + jsTest.log("Trigger shard version mismatch..."); -// MongoD doesn't know about the config shard version *until* MongoS tells it -collA.findOne(); + assert.writeOK(collB.insert({goodbye: "world"})); -jsTest.log("Trigger shard version mismatch..."); + print("Inserted..."); -assert.writeOK(collB.insert({goodbye: "world"})); + assert.eq(3, collA.find().itcount()); + assert.eq(3, collB.find().itcount()); -print("Inserted..."); - -assert.eq(3, collA.find().itcount()); -assert.eq(3, collB.find().itcount()); - -st.stop(); + st.stop(); +})(); diff --git a/jstests/slow2/mr_during_migrate.js b/jstests/slow2/mr_during_migrate.js index cb439aeb241..1b3f55721f4 100644 --- a/jstests/slow2/mr_during_migrate.js +++ b/jstests/slow2/mr_during_migrate.js @@ -1,113 +1,112 @@ // Do parallel ops with migrates occurring +(function() { + 'use strict'; -var st = new ShardingTest({shards: 10, mongos: 2, verbose: 2}); + var st = new ShardingTest({shards: 10, mongos: 2, verbose: 2}); -jsTest.log("Doing parallel operations..."); + var mongos = st.s0; + var admin = mongos.getDB("admin"); + var coll = st.s.getCollection(jsTest.name() + ".coll"); -// Stop balancer, since it'll just get in the way of these -st.stopBalancer(); + var numDocs = 1024 * 1024; + var dataSize = 1024; // bytes, must be power of 2 -var mongos = st.s0; -var admin = mongos.getDB("admin"); -var coll = st.s.getCollection(jsTest.name() + ".coll"); + var data = "x"; + while (data.length < dataSize) + data += data; -var numDocs = 1024 * 1024; -var dataSize = 1024; // bytes, must be power of 2 + var bulk = coll.initializeUnorderedBulkOp(); + for (var i = 0; i < numDocs; i++) { + bulk.insert({_id: i, data: data}); + } + assert.writeOK(bulk.execute()); -var data = "x"; -while (data.length < dataSize) - data += data; + // Make sure everything got inserted + assert.eq(numDocs, coll.find().itcount()); -var bulk = coll.initializeUnorderedBulkOp(); -for (var i = 0; i < numDocs; i++) { - bulk.insert({_id: i, data: data}); -} -assert.writeOK(bulk.execute()); + jsTest.log("Inserted " + sh._dataFormat(dataSize * numDocs) + " of data."); -// Make sure everything got inserted -assert.eq(numDocs, coll.find().itcount()); + // Shard collection + st.shardColl(coll, {_id: 1}, false); -jsTest.log("Inserted " + sh._dataFormat(dataSize * numDocs) + " of data."); + st.printShardingStatus(); -// Shard collection -st.shardColl(coll, {_id: 1}, false); + jsTest.log("Sharded collection now initialized, starting migrations..."); -st.printShardingStatus(); - -jsTest.log("Sharded collection now initialized, starting migrations..."); + var checkMigrate = function() { + print("Result of migrate : "); + printjson(this); + }; -var checkMigrate = function() { - print("Result of migrate : "); - printjson(this); -}; + // Creates a number of migrations of random chunks to diff shard servers + var ops = []; + for (var i = 0; i < st._connections.length; i++) { + ops.push({ + op: "command", + ns: "admin", + command: { + moveChunk: "" + coll, + find: {_id: {"#RAND_INT": [0, numDocs]}}, + to: st._connections[i].shardName, + _waitForDelete: true + }, + showResult: true + }); + } -// Creates a number of migrations of random chunks to diff shard servers -var ops = []; -for (var i = 0; i < st._connections.length; i++) { - ops.push({ - op: "command", - ns: "admin", - command: { - moveChunk: "" + coll, - find: {_id: {"#RAND_INT": [0, numDocs]}}, - to: st._connections[i].shardName, - _waitForDelete: true - }, - showResult: true - }); -} + // TODO: Also migrate output collection -// TODO: Also migrate output collection + jsTest.log("Starting migrations now..."); -jsTest.log("Starting migrations now..."); + var bid = benchStart({ops: ops, host: st.s.host, parallel: 1, handleErrors: false}); -var bid = benchStart({ops: ops, host: st.s.host, parallel: 1, handleErrors: false}); + //####################### + // Tests during migration -//####################### -// Tests during migration + var numTests = 5; -var numTests = 5; + for (var t = 0; t < numTests; t++) { + jsTest.log("Test #" + t); -for (var t = 0; t < numTests; t++) { - jsTest.log("Test #" + t); + var mongos = st.s1; // use other mongos so we get stale shard versions + var coll = mongos.getCollection(coll + ""); + var outputColl = mongos.getCollection(coll + "_output"); - var mongos = st.s1; // use other mongos so we get stale shard versions - var coll = mongos.getCollection(coll + ""); - var outputColl = mongos.getCollection(coll + "_output"); + var numTypes = 32; + var map = function() { + emit(this._id % 32 /* must be hardcoded */, {c: 1}); + }; - var numTypes = 32; - var map = function() { - emit(this._id % 32 /* must be hardcoded */, {c: 1}); - }; - var reduce = function(k, vals) { - var total = 0; - for (var i = 0; i < vals.length; i++) - total += vals[i].c; - return { - c: total + var reduce = function(k, vals) { + var total = 0; + for (var i = 0; i < vals.length; i++) + total += vals[i].c; + return { + c: total + }; }; - }; - printjson(coll.find({_id: 0}).itcount()); + printjson(coll.find({_id: 0}).itcount()); - jsTest.log("Starting new mapReduce run #" + t); + jsTest.log("Starting new mapReduce run #" + t); - // assert.eq( coll.find().itcount(), numDocs ) + // assert.eq( coll.find().itcount(), numDocs ) - coll.getMongo().getDB("admin").runCommand({setParameter: 1, traceExceptions: true}); + coll.getMongo().getDB("admin").runCommand({setParameter: 1, traceExceptions: true}); - printjson(coll.mapReduce( - map, reduce, {out: {replace: outputColl.getName(), db: outputColl.getDB() + ""}})); + printjson(coll.mapReduce( + map, reduce, {out: {replace: outputColl.getName(), db: outputColl.getDB() + ""}})); - jsTest.log("MapReduce run #" + t + " finished."); + jsTest.log("MapReduce run #" + t + " finished."); - assert.eq(outputColl.find().itcount(), numTypes); + assert.eq(outputColl.find().itcount(), numTypes); - outputColl.find().forEach(function(x) { - assert.eq(x.value.c, numDocs / numTypes); - }); -} + outputColl.find().forEach(function(x) { + assert.eq(x.value.c, numDocs / numTypes); + }); + } -printjson(benchFinish(bid)); + printjson(benchFinish(bid)); -st.stop(); + st.stop(); +})(); diff --git a/src/mongo/base/validate_locale.cpp b/src/mongo/base/validate_locale.cpp index 5a4320c34c5..81207d4d89d 100644 --- a/src/mongo/base/validate_locale.cpp +++ b/src/mongo/base/validate_locale.cpp @@ -30,6 +30,7 @@ #include <boost/filesystem/operations.hpp> #include "mongo/base/init.h" +#include "mongo/util/mongoutils/str.h" namespace mongo { @@ -38,13 +39,15 @@ MONGO_INITIALIZER_GENERAL(ValidateLocale, MONGO_NO_PREREQUISITES, MONGO_DEFAULT_ try { // Validate that boost can correctly load the user's locale boost::filesystem::path("/").has_root_directory(); - } catch (const std::runtime_error&) { - return Status(ErrorCodes::BadValue, - "Invalid or no user locale set." + } catch (const std::runtime_error& e) { + return Status( + ErrorCodes::BadValue, + str::stream() + << "Invalid or no user locale set. " #ifndef _WIN32 - " Please ensure LANG and/or LC_* environment variables are set correctly." + << " Please ensure LANG and/or LC_* environment variables are set correctly. " #endif - ); + << e.what()); } return Status::OK(); } diff --git a/src/mongo/db/commands/mr.cpp b/src/mongo/db/commands/mr.cpp index f93dee2b554..ced9f2bbdb7 100644 --- a/src/mongo/db/commands/mr.cpp +++ b/src/mongo/db/commands/mr.cpp @@ -308,7 +308,7 @@ Config::Config(const string& _dbname, const BSONObj& cmdObj) { // scope and code if (cmdObj["scope"].type() == Object) - scopeSetup = cmdObj["scope"].embeddedObjectUserCheck(); + scopeSetup = cmdObj["scope"].embeddedObjectUserCheck().getOwned(); mapper.reset(new JSMapper(cmdObj["map"])); reducer.reset(new JSReducer(cmdObj["reduce"])); @@ -316,7 +316,7 @@ Config::Config(const string& _dbname, const BSONObj& cmdObj) { finalizer.reset(new JSFinalizer(cmdObj["finalize"])); if (cmdObj["mapparams"].type() == Array) { - mapParams = cmdObj["mapparams"].embeddedObjectUserCheck(); + mapParams = cmdObj["mapparams"].embeddedObjectUserCheck().getOwned(); } } @@ -787,6 +787,7 @@ void State::init() { AuthorizationSession::get(ClientBasic::getCurrent())->getAuthenticatedUserNamesToken(); _scope.reset(globalScriptEngine->newScopeForCurrentThread()); _scope->registerOperation(_txn); + _scope->requireOwnedObjects(); _scope->setLocalDB(_config.dbname); _scope->loadStored(_txn, true); @@ -1435,6 +1436,7 @@ public: BSONObj o; PlanExecutor::ExecState execState; while (PlanExecutor::ADVANCED == (execState = exec->getNext(&o, NULL))) { + o = o.getOwned(); // we will be accessing outside of the lock // check to see if this is a new object we don't own yet // because of a chunk migration if (collMetadata) { diff --git a/src/mongo/db/dbhelpers.cpp b/src/mongo/db/dbhelpers.cpp index d3c62ea0ef8..ced194addbf 100644 --- a/src/mongo/db/dbhelpers.cpp +++ b/src/mongo/db/dbhelpers.cpp @@ -448,11 +448,12 @@ long long Helpers::removeRange(OperationContext* txn, txn, repl::ReplClientInfo::forClient(txn->getClient()).getLastOp(), writeConcern); - if (replStatus.status.code() == ErrorCodes::ExceededTimeLimit) { + if (replStatus.status.code() == ErrorCodes::ExceededTimeLimit || + replStatus.status.code() == ErrorCodes::WriteConcernFailed) { warning(LogComponent::kSharding) << "replication to secondaries for removeRange at " "least 60 seconds behind"; } else { - massertStatusOK(replStatus.status); + uassertStatusOK(replStatus.status); } millisWaitingForReplication += replStatus.duration; } diff --git a/src/mongo/db/exec/group.cpp b/src/mongo/db/exec/group.cpp index 4ce1d4031db..31a72df4772 100644 --- a/src/mongo/db/exec/group.cpp +++ b/src/mongo/db/exec/group.cpp @@ -148,7 +148,8 @@ Status GroupStage::processObject(const BSONObj& obj) { } } - _scope->setObject("obj", obj, true); + BSONObj objCopy = obj.getOwned(); + _scope->setObject("obj", objCopy, true); _scope->setNumber("n", n - 1); try { diff --git a/src/mongo/db/ops/update_driver.cpp b/src/mongo/db/ops/update_driver.cpp index c0c11082a2e..a99376174a8 100644 --- a/src/mongo/db/ops/update_driver.cpp +++ b/src/mongo/db/ops/update_driver.cpp @@ -183,6 +183,13 @@ Status UpdateDriver::populateDocumentWithQueryFields(const BSONObj& query, return populateDocumentWithQueryFields(*cq, immutablePaths, doc); } +namespace { + +const FieldRef idPath("_id"); +const vector<FieldRef*> emptyImmutablePaths; + +} // namespace + Status UpdateDriver::populateDocumentWithQueryFields(const CanonicalQuery& query, const vector<FieldRef*>* immutablePathsPtr, mutablebson::Document& doc) const { @@ -192,13 +199,12 @@ Status UpdateDriver::populateDocumentWithQueryFields(const CanonicalQuery& query if (isDocReplacement()) { FieldRefSet pathsToExtract; - // TODO: Refactor update logic, make _id just another immutable field - static const FieldRef idPath("_id"); - static const vector<FieldRef*> emptyImmutablePaths; const vector<FieldRef*>& immutablePaths = immutablePathsPtr ? *immutablePathsPtr : emptyImmutablePaths; pathsToExtract.fillFrom(immutablePaths); + + // TODO: Refactor update logic, make _id just another immutable field pathsToExtract.insert(&idPath); // Extract only immutable fields from replacement-style diff --git a/src/mongo/db/repl/oplog.cpp b/src/mongo/db/repl/oplog.cpp index 63858050014..2c764baf9e3 100644 --- a/src/mongo/db/repl/oplog.cpp +++ b/src/mongo/db/repl/oplog.cpp @@ -825,11 +825,7 @@ Status applyOperation_inlock(OperationContext* txn, Status status{ErrorCodes::NotYetInitialized, ""}; { WriteUnitOfWork wuow(txn); - try { - status = collection->insertDocument(txn, o, true); - } catch (DBException dbe) { - status = dbe.toStatus(); - } + status = collection->insertDocument(txn, o, true); if (status.isOK()) { wuow.commit(); } diff --git a/src/mongo/dbtests/jstests.cpp b/src/mongo/dbtests/jstests.cpp index 53cead63a3a..04a9180dbbd 100644 --- a/src/mongo/dbtests/jstests.cpp +++ b/src/mongo/dbtests/jstests.cpp @@ -2361,6 +2361,49 @@ public: } }; +class RequiresOwnedObjects { +public: + void run() { + char buf[] = {5, 0, 0, 0, 0}; + BSONObj unowned(buf); + BSONObj owned = unowned.getOwned(); + + ASSERT(!unowned.isOwned()); + ASSERT(owned.isOwned()); + + // Ensure that by default we can bind owned and unowned + { + unique_ptr<Scope> s(globalScriptEngine->newScope()); + s->setObject("unowned", unowned, true); + s->setObject("owned", owned, true); + } + + // After we set the flag, we should only be able to set owned + { + unique_ptr<Scope> s(globalScriptEngine->newScope()); + s->requireOwnedObjects(); + s->setObject("owned", owned, true); + + bool threwException = false; + try { + s->setObject("unowned", unowned, true); + } catch (...) { + threwException = true; + + auto status = exceptionToStatus(); + + ASSERT_EQUALS(status.code(), ErrorCodes::BadValue); + } + + ASSERT(threwException); + + // after resetting, we can set unowned's again + s->reset(); + s->setObject("unowned", unowned, true); + } + } +}; + class All : public Suite { public: All() : Suite("js") { @@ -2422,6 +2465,7 @@ public: add<RecursiveInvoke>(); add<ErrorCodeFromInvoke>(); + add<RequiresOwnedObjects>(); add<RoundTripTests::DBRefTest>(); add<RoundTripTests::DBPointerTest>(); diff --git a/src/mongo/s/SConscript b/src/mongo/s/SConscript index 34efb20cf88..cf8d04827a6 100644 --- a/src/mongo/s/SConscript +++ b/src/mongo/s/SConscript @@ -75,9 +75,9 @@ env.Library( '$BUILD_DIR/mongo/s/catalog/forwarding_catalog_manager', '$BUILD_DIR/mongo/s/catalog/replset/catalog_manager_replica_set', '$BUILD_DIR/mongo/s/coreshard', - '$BUILD_DIR/mongo/s/mongoscore', '$BUILD_DIR/mongo/util/clock_source_mock', '$BUILD_DIR/mongo/util/net/message_port_mock', + 'mongoscore', ], LIBDEPS_TAGS=[ # Depends on coreshard, but that would be circular diff --git a/src/mongo/s/catalog/catalog_manager.h b/src/mongo/s/catalog/catalog_manager.h index 5257f53f7db..2f75ba2da69 100644 --- a/src/mongo/s/catalog/catalog_manager.h +++ b/src/mongo/s/catalog/catalog_manager.h @@ -451,9 +451,6 @@ public: */ virtual bool isMetadataConsistentFromLastCheck(OperationContext* txn) = 0; -protected: - CatalogManager() = default; - /** * Obtains a reference to the distributed lock manager instance to use for synchronizing * system-wide changes. @@ -462,6 +459,9 @@ protected: * be cached. */ virtual DistLockManager* getDistLockManager() = 0; + +protected: + CatalogManager() = default; }; } // namespace mongo diff --git a/src/mongo/s/chunk_manager.cpp b/src/mongo/s/chunk_manager.cpp index f785e748c54..29e596948d4 100644 --- a/src/mongo/s/chunk_manager.cpp +++ b/src/mongo/s/chunk_manager.cpp @@ -98,6 +98,9 @@ public: string shardFor(OperationContext* txn, const string& hostName) const final { const auto shard = grid.shardRegistry()->getShard(txn, hostName); + uassert(ErrorCodes::ShardNotFound, + str::stream() << "Shard " << hostName << " not found.", + shard); return shard->getId(); } diff --git a/src/mongo/scripting/deadline_monitor.cpp b/src/mongo/scripting/deadline_monitor.cpp index 5eb0f52e5de..f75c34a0cc2 100644 --- a/src/mongo/scripting/deadline_monitor.cpp +++ b/src/mongo/scripting/deadline_monitor.cpp @@ -34,7 +34,7 @@ namespace mongo { -MONGO_EXPORT_SERVER_PARAMETER(scriptingEngineInterruptIntervalMS, int, 1000); +MONGO_EXPORT_SERVER_PARAMETER(scriptingEngineInterruptIntervalMS, int, 0); int getScriptingEngineInterruptInterval() { return scriptingEngineInterruptIntervalMS.load(); diff --git a/src/mongo/scripting/deadline_monitor.h b/src/mongo/scripting/deadline_monitor.h index b1c2855dd46..05acd5c349c 100644 --- a/src/mongo/scripting/deadline_monitor.h +++ b/src/mongo/scripting/deadline_monitor.h @@ -138,7 +138,7 @@ private: const Date_t now = Date_t::now(); const auto interruptInterval = Milliseconds{getScriptingEngineInterruptInterval()}; - if (now - lastInterruptCycle > interruptInterval) { + if ((interruptInterval.count() > 0) && (now - lastInterruptCycle > interruptInterval)) { for (const auto& task : _tasks) { if (task.second > now) task.first->interrupt(); @@ -148,9 +148,13 @@ private: // wait for a task to be added or a deadline to expire if (_nearestDeadlineWallclock > now) { - if (_nearestDeadlineWallclock == Date_t::max() || - _nearestDeadlineWallclock - now > interruptInterval) { - _newDeadlineAvailable.wait_for(lk, interruptInterval); + if (_nearestDeadlineWallclock == Date_t::max()) { + if ((interruptInterval.count() > 0) && + (_nearestDeadlineWallclock - now > interruptInterval)) { + _newDeadlineAvailable.wait_for(lk, interruptInterval); + } else { + _newDeadlineAvailable.wait(lk); + } } else { _newDeadlineAvailable.wait_until(lk, _nearestDeadlineWallclock.toSystemTimePoint()); diff --git a/src/mongo/scripting/engine.cpp b/src/mongo/scripting/engine.cpp index 7e6628344ab..71c42f2fe74 100644 --- a/src/mongo/scripting/engine.cpp +++ b/src/mongo/scripting/engine.cpp @@ -422,6 +422,9 @@ public: void advanceGeneration() { _real->advanceGeneration(); } + void requireOwnedObjects() override { + _real->requireOwnedObjects(); + } bool isKillPending() const { return _real->isKillPending(); } diff --git a/src/mongo/scripting/engine.h b/src/mongo/scripting/engine.h index d31311f0114..2fbae5c347e 100644 --- a/src/mongo/scripting/engine.h +++ b/src/mongo/scripting/engine.h @@ -105,6 +105,8 @@ public: virtual void advanceGeneration() = 0; + virtual void requireOwnedObjects() = 0; + virtual ScriptingFunction createFunction(const char* code); /** diff --git a/src/mongo/scripting/mozjs/bson.cpp b/src/mongo/scripting/mozjs/bson.cpp index 1e7e4b7c190..3783bc0755a 100644 --- a/src/mongo/scripting/mozjs/bson.cpp +++ b/src/mongo/scripting/mozjs/bson.cpp @@ -59,13 +59,17 @@ namespace { * the appearance of mutable state on the read/write versions. */ struct BSONHolder { - BSONHolder(const BSONObj& obj, const BSONObj* parent, std::size_t generation, bool ro) + BSONHolder(const BSONObj& obj, const BSONObj* parent, const MozJSImplScope* scope, bool ro) : _obj(obj), - _generation(generation), + _generation(scope->getGeneration()), _isOwned(obj.isOwned() || (parent && parent->isOwned())), _resolved(false), _readOnly(ro), _altered(false) { + uassert( + ErrorCodes::BadValue, + "Attempt to bind an unowned BSON Object to a JS scope marked as requiring ownership", + _isOwned || (!scope->requiresOwnedObjects())); if (parent) { _parent.emplace(*parent); } @@ -107,7 +111,7 @@ void BSONInfo::make( auto scope = getScope(cx); scope->getProto<BSONInfo>().newObject(obj); - JS_SetPrivate(obj, new BSONHolder(bson, parent, scope->getGeneration(), ro)); + JS_SetPrivate(obj, new BSONHolder(bson, parent, scope, ro)); } void BSONInfo::finalize(JSFreeOp* fop, JSObject* obj) { diff --git a/src/mongo/scripting/mozjs/implscope.cpp b/src/mongo/scripting/mozjs/implscope.cpp index 3b260c1558c..e85efd97505 100644 --- a/src/mongo/scripting/mozjs/implscope.cpp +++ b/src/mongo/scripting/mozjs/implscope.cpp @@ -322,6 +322,7 @@ MozJSImplScope::MozJSImplScope(MozJSScriptEngine* engine) _status(Status::OK()), _quickExit(false), _generation(0), + _requireOwnedObjects(false), _hasOutOfMemoryException(false), _binDataProto(_context), _bsonProto(_context), @@ -749,6 +750,7 @@ void MozJSImplScope::reset() { unregisterOperation(); _pendingKill.store(false); _pendingGC.store(false); + _requireOwnedObjects = false; advanceGeneration(); } @@ -869,6 +871,14 @@ void MozJSImplScope::advanceGeneration() { _generation++; } +void MozJSImplScope::requireOwnedObjects() { + _requireOwnedObjects = true; +} + +bool MozJSImplScope::requiresOwnedObjects() const { + return _requireOwnedObjects; +} + const std::string& MozJSImplScope::getParentStack() const { return _parentStack; } diff --git a/src/mongo/scripting/mozjs/implscope.h b/src/mongo/scripting/mozjs/implscope.h index bfff70a93b0..edfc236e90c 100644 --- a/src/mongo/scripting/mozjs/implscope.h +++ b/src/mongo/scripting/mozjs/implscope.h @@ -305,6 +305,10 @@ public: void advanceGeneration() override; + void requireOwnedObjects() override; + + bool requiresOwnedObjects() const; + JS::HandleId getInternedStringId(InternedString name) { return _internedStrings.getInternedString(name); } @@ -373,6 +377,7 @@ private: bool _quickExit; std::string _parentStack; std::size_t _generation; + bool _requireOwnedObjects; bool _hasOutOfMemoryException; WrapType<BinDataInfo> _binDataProto; diff --git a/src/mongo/scripting/mozjs/objectwrapper.cpp b/src/mongo/scripting/mozjs/objectwrapper.cpp index 77ea9a1eca8..8781f792484 100644 --- a/src/mongo/scripting/mozjs/objectwrapper.cpp +++ b/src/mongo/scripting/mozjs/objectwrapper.cpp @@ -564,7 +564,12 @@ ObjectWrapper::WriteFieldRecursionFrame::WriteFieldRecursionFrame(JSContext* cx, ids.infallibleAppend(rid); } } else { - JS::AutoIdArray rids(cx, JS_Enumerate(cx, thisv)); + auto ridArrayPtr = JS_Enumerate(cx, thisv); + if (!ridArrayPtr) { + throwCurrentJSException( + cx, ErrorCodes::JSInterpreterFailure, "Failure to enumerate object"); + } + JS::AutoIdArray rids(cx, ridArrayPtr); if (!ids.reserve(rids.length())) { throwCurrentJSException( diff --git a/src/mongo/scripting/mozjs/proxyscope.cpp b/src/mongo/scripting/mozjs/proxyscope.cpp index bb4cd5c06ff..9bb11ce400b 100644 --- a/src/mongo/scripting/mozjs/proxyscope.cpp +++ b/src/mongo/scripting/mozjs/proxyscope.cpp @@ -120,6 +120,10 @@ void MozJSProxyScope::advanceGeneration() { run([&] { _implScope->advanceGeneration(); }); } +void MozJSProxyScope::requireOwnedObjects() { + run([&] { _implScope->requireOwnedObjects(); }); +} + double MozJSProxyScope::getNumber(const char* field) { double out; run([&] { out = _implScope->getNumber(field); }); diff --git a/src/mongo/scripting/mozjs/proxyscope.h b/src/mongo/scripting/mozjs/proxyscope.h index 451981330a1..4dd69a3ebe9 100644 --- a/src/mongo/scripting/mozjs/proxyscope.h +++ b/src/mongo/scripting/mozjs/proxyscope.h @@ -129,6 +129,8 @@ public: void advanceGeneration() override; + void requireOwnedObjects() override; + double getNumber(const char* field) override; int getNumberInt(const char* field) override; long long getNumberLongLong(const char* field) override; diff --git a/src/mongo/shell/replsettest.js b/src/mongo/shell/replsettest.js index 4a64f5e50a0..89ab91da580 100644 --- a/src/mongo/shell/replsettest.js +++ b/src/mongo/shell/replsettest.js @@ -89,7 +89,7 @@ var ReplSetTest = function(opts) { var _unbridgedPorts; var _unbridgedNodes; - this.kDefaultTimeoutMS = 5 * 60 * 1000; + this.kDefaultTimeoutMS = 10 * 60 * 1000; // Publicly exposed variables diff --git a/src/third_party/wiredtiger/dist/s_string.ok b/src/third_party/wiredtiger/dist/s_string.ok index 1f7f7d9fd3a..f3852d00ac8 100644 --- a/src/third_party/wiredtiger/dist/s_string.ok +++ b/src/third_party/wiredtiger/dist/s_string.ok @@ -731,6 +731,7 @@ fsyncLock fsyncs ftruncate func +fvisibility gcc gdb ge diff --git a/src/third_party/wiredtiger/dist/stat_data.py b/src/third_party/wiredtiger/dist/stat_data.py index ac79ffd029a..512892eb44d 100644 --- a/src/third_party/wiredtiger/dist/stat_data.py +++ b/src/third_party/wiredtiger/dist/stat_data.py @@ -150,6 +150,7 @@ connection_stats = [ ConnStat('read_io', 'total read I/Os'), ConnStat('rwlock_read', 'pthread mutex shared lock read-lock calls'), ConnStat('rwlock_write', 'pthread mutex shared lock write-lock calls'), + ConnStat('time_travel', 'detected system time went backwards'), ConnStat('write_io', 'total write I/Os'), ########################################## diff --git a/src/third_party/wiredtiger/import.data b/src/third_party/wiredtiger/import.data index abf3fe5cb9c..8bd00db3aa2 100644 --- a/src/third_party/wiredtiger/import.data +++ b/src/third_party/wiredtiger/import.data @@ -1,5 +1,5 @@ { - "commit": "b8f590dea0400666ef26e21adf11c5997bb5ef1b", + "commit": "827b48a34227243c809d41fac3dc909ed46b0c5e", "github": "wiredtiger/wiredtiger.git", "vendor": "wiredtiger", "branch": "mongodb-3.2" diff --git a/src/third_party/wiredtiger/src/btree/row_key.c b/src/third_party/wiredtiger/src/btree/row_key.c index 032fdf7d897..5bb09832eed 100644 --- a/src/third_party/wiredtiger/src/btree/row_key.c +++ b/src/third_party/wiredtiger/src/btree/row_key.c @@ -471,6 +471,8 @@ __wt_row_ikey_alloc(WT_SESSION_IMPL *session, { WT_IKEY *ikey; + WT_ASSERT(session, key != NULL); /* quiet clang scan-build */ + /* * Allocate memory for the WT_IKEY structure and the key, then copy * the key into place. diff --git a/src/third_party/wiredtiger/src/docs/programming.dox b/src/third_party/wiredtiger/src/docs/programming.dox index aa76bef4614..205e7544c6c 100644 --- a/src/third_party/wiredtiger/src/docs/programming.dox +++ b/src/third_party/wiredtiger/src/docs/programming.dox @@ -65,19 +65,20 @@ each of which is ordered by one or more columns. - @subpage_single wtperf - @subpage_single wtstats <p> -- @subpage_single tune_memory_allocator -- @subpage_single tune_page_size_and_comp -- @subpage_single tune_cache +- @subpage_single tune_build_options - @subpage_single tune_bulk_load +- @subpage_single tune_cache +- @subpage_single tune_checksum +- @subpage_single tune_close - @subpage_single tune_cursor_persist -- @subpage_single tune_read_only - @subpage_single tune_durability -- @subpage_single tune_checksum - @subpage_single tune_file_alloc +- @subpage_single tune_memory_allocator +- @subpage_single tune_mutex +- @subpage_single tune_page_size_and_comp +- @subpage_single tune_read_only - @subpage_single tune_system_buffer_cache - @subpage_single tune_transparent_huge_pages -- @subpage_single tune_close -- @subpage_single tune_mutex - @subpage_single tune_zone_reclaim */ diff --git a/src/third_party/wiredtiger/src/docs/spell.ok b/src/third_party/wiredtiger/src/docs/spell.ok index bc2e16b1122..5d629f4c49f 100644 --- a/src/third_party/wiredtiger/src/docs/spell.ok +++ b/src/third_party/wiredtiger/src/docs/spell.ok @@ -237,6 +237,7 @@ fput freelist fsync ftruncate +fvisibility gcc gdbm ge diff --git a/src/third_party/wiredtiger/src/docs/tune-build-options.dox b/src/third_party/wiredtiger/src/docs/tune-build-options.dox new file mode 100644 index 00000000000..79cd60b1105 --- /dev/null +++ b/src/third_party/wiredtiger/src/docs/tune-build-options.dox @@ -0,0 +1,9 @@ +/*! @page tune_build_options gcc/clang build options + +WiredTiger can be built using the gcc/clang \c -fvisibility=hidden flag, +which may significantly reduce the size and load time of the WiredTiger +library when built as a dynamic shared object, and allow the optimizer +to produce better code (for example, by eliminating most lookups in the +procedure linkage table). + + */ diff --git a/src/third_party/wiredtiger/src/evict/evict_lru.c b/src/third_party/wiredtiger/src/evict/evict_lru.c index 26bbf9f679b..cc3c5a5c824 100644 --- a/src/third_party/wiredtiger/src/evict/evict_lru.c +++ b/src/third_party/wiredtiger/src/evict/evict_lru.c @@ -941,6 +941,13 @@ __evict_tune_workers(WT_SESSION_IMPL *session) conn = S2C(session); cache = conn->cache; + /* + * If we have a fixed number of eviction threads, there is no value in + * calculating if we should do any tuning. + */ + if (conn->evict_threads_max == conn->evict_threads_min) + return (0); + WT_ASSERT(session, conn->evict_threads.threads[0]->session == session); pgs_evicted_cur = pgs_evicted_persec_cur = 0; diff --git a/src/third_party/wiredtiger/src/include/extern.h b/src/third_party/wiredtiger/src/include/extern.h index bf3279d0f94..dfd2d03707f 100644 --- a/src/third_party/wiredtiger/src/include/extern.h +++ b/src/third_party/wiredtiger/src/include/extern.h @@ -570,6 +570,7 @@ extern int __wt_schema_destroy_index(WT_SESSION_IMPL *session, WT_INDEX **idxp) extern int __wt_schema_destroy_table(WT_SESSION_IMPL *session, WT_TABLE **tablep) WT_GCC_FUNC_DECL_ATTRIBUTE((warn_unused_result)); extern int __wt_schema_remove_table(WT_SESSION_IMPL *session, WT_TABLE *table) WT_GCC_FUNC_DECL_ATTRIBUTE((warn_unused_result)); extern int __wt_schema_close_tables(WT_SESSION_IMPL *session) WT_GCC_FUNC_DECL_ATTRIBUTE((warn_unused_result)); +extern int __wt_schema_sweep_tables(WT_SESSION_IMPL *session) WT_GCC_FUNC_DECL_ATTRIBUTE((warn_unused_result)); extern int __wt_schema_colgroup_name(WT_SESSION_IMPL *session, WT_TABLE *table, const char *cgname, size_t len, WT_ITEM *buf) WT_GCC_FUNC_DECL_ATTRIBUTE((warn_unused_result)); extern int __wt_schema_open_colgroups(WT_SESSION_IMPL *session, WT_TABLE *table) WT_GCC_FUNC_DECL_ATTRIBUTE((warn_unused_result)); extern int __wt_schema_open_index(WT_SESSION_IMPL *session, WT_TABLE *table, const char *idxname, size_t len, WT_INDEX **indexp) WT_GCC_FUNC_DECL_ATTRIBUTE((warn_unused_result)); diff --git a/src/third_party/wiredtiger/src/include/gcc.h b/src/third_party/wiredtiger/src/include/gcc.h index 22d78fc165a..6c712813d29 100644 --- a/src/third_party/wiredtiger/src/include/gcc.h +++ b/src/third_party/wiredtiger/src/include/gcc.h @@ -9,7 +9,7 @@ #define WT_PTRDIFFT_FMT "td" /* ptrdiff_t format string */ #define WT_SIZET_FMT "zu" /* size_t format string */ -/* Add GCC-specific attributes to types and function declarations. */ +/* GCC-specific attributes. */ #define WT_PACKED_STRUCT_BEGIN(name) \ struct __attribute__ ((__packed__)) name { #define WT_PACKED_STRUCT_END \ diff --git a/src/third_party/wiredtiger/src/include/lint.h b/src/third_party/wiredtiger/src/include/lint.h index 2d0f47988b7..813b9182683 100644 --- a/src/third_party/wiredtiger/src/include/lint.h +++ b/src/third_party/wiredtiger/src/include/lint.h @@ -9,6 +9,7 @@ #define WT_PTRDIFFT_FMT "td" /* ptrdiff_t format string */ #define WT_SIZET_FMT "zu" /* size_t format string */ +/* Lint-specific attributes. */ #define WT_PACKED_STRUCT_BEGIN(name) \ struct name { #define WT_PACKED_STRUCT_END \ diff --git a/src/third_party/wiredtiger/src/include/misc.i b/src/third_party/wiredtiger/src/include/misc.i index 7040886cf82..fad10f01103 100644 --- a/src/third_party/wiredtiger/src/include/misc.i +++ b/src/third_party/wiredtiger/src/include/misc.i @@ -55,6 +55,31 @@ __wt_seconds(WT_SESSION_IMPL *session, time_t *timep) } /* + * __wt_time_check_monotonic -- + * Check and prevent time running backward. If we detect that it has, we + * set the time structure to the previous values, making time stand still + * until we see a time in the future of the highest value seen so far. + */ +static inline void +__wt_time_check_monotonic(WT_SESSION_IMPL *session, struct timespec *tsp) +{ + /* + * Detect time going backward. If so, use the last + * saved timestamp. + */ + if (session == NULL) + return; + + if (tsp->tv_sec < session->last_epoch.tv_sec || + (tsp->tv_sec == session->last_epoch.tv_sec && + tsp->tv_nsec < session->last_epoch.tv_nsec)) { + WT_STAT_CONN_INCR(session, time_travel); + *tsp = session->last_epoch; + } else + session->last_epoch = *tsp; +} + +/* * __wt_verbose -- * Verbose message. * diff --git a/src/third_party/wiredtiger/src/include/msvc.h b/src/third_party/wiredtiger/src/include/msvc.h index 6c5c8b67647..c9399be3185 100644 --- a/src/third_party/wiredtiger/src/include/msvc.h +++ b/src/third_party/wiredtiger/src/include/msvc.h @@ -16,9 +16,7 @@ #define WT_PTRDIFFT_FMT "Id" /* ptrdiff_t format string */ #define WT_SIZET_FMT "Iu" /* size_t format string */ -/* - * Add MSVC-specific attributes and pragmas to types and function declarations. - */ +/* MSVC-specific attributes. */ #define WT_PACKED_STRUCT_BEGIN(name) \ __pragma(pack(push,1)) \ struct name { diff --git a/src/third_party/wiredtiger/src/include/session.h b/src/third_party/wiredtiger/src/include/session.h index 1b2dfd1ed2b..d05dee68641 100644 --- a/src/third_party/wiredtiger/src/include/session.h +++ b/src/third_party/wiredtiger/src/include/session.h @@ -66,6 +66,7 @@ struct __wt_session_impl { /* Session handle reference list */ TAILQ_HEAD(__dhandles, __wt_data_handle_cache) dhandles; time_t last_sweep; /* Last sweep for dead handles */ + struct timespec last_epoch; /* Last epoch time returned */ /* Cursors closed with the session */ TAILQ_HEAD(__cursors, __wt_cursor) cursors; @@ -97,6 +98,12 @@ struct __wt_session_impl { */ TAILQ_HEAD(__tables, __wt_table) tables; + /* + * Updated when the table cache is swept of all tables older than the + * current schema generation. + */ + uint64_t table_sweep_gen; + /* Current rwlock for callback. */ WT_RWLOCK *current_rwlock; uint8_t current_rwticket; diff --git a/src/third_party/wiredtiger/src/include/stat.h b/src/third_party/wiredtiger/src/include/stat.h index 6c274484bcb..db48a841571 100644 --- a/src/third_party/wiredtiger/src/include/stat.h +++ b/src/third_party/wiredtiger/src/include/stat.h @@ -361,6 +361,7 @@ struct __wt_connection_stats { int64_t cache_eviction_clean; int64_t cond_auto_wait_reset; int64_t cond_auto_wait; + int64_t time_travel; int64_t file_open; int64_t memory_allocation; int64_t memory_free; diff --git a/src/third_party/wiredtiger/src/include/txn.h b/src/third_party/wiredtiger/src/include/txn.h index 7e802c188ab..fdf9c714afa 100644 --- a/src/third_party/wiredtiger/src/include/txn.h +++ b/src/third_party/wiredtiger/src/include/txn.h @@ -93,6 +93,8 @@ struct __wt_txn_global { * the global transaction state. */ WT_RWLOCK scan_rwlock; + /* Protects logging, checkpoints and transaction visibility. */ + WT_RWLOCK visibility_rwlock; /* * Track information about the running checkpoint. The transaction diff --git a/src/third_party/wiredtiger/src/include/wiredtiger.in b/src/third_party/wiredtiger/src/include/wiredtiger.in index ddecb2ac765..821efdf5fa1 100644 --- a/src/third_party/wiredtiger/src/include/wiredtiger.in +++ b/src/third_party/wiredtiger/src/include/wiredtiger.in @@ -39,6 +39,16 @@ extern "C" { #define __F(func) (*(func)) #endif +/* + * We support configuring WiredTiger with the gcc/clang -fvisibility=hidden + * flags, but that requires public APIs be specifically marked. + */ +#if defined(DOXYGEN) || defined(SWIG) || !defined(__GNUC__) +#define WT_ATTRIBUTE_LIBRARY_VISIBLE +#else +#define WT_ATTRIBUTE_LIBRARY_VISIBLE __attribute__((visibility("default"))) +#endif + #ifdef SWIG %{ #include <wiredtiger.h> @@ -2553,7 +2563,7 @@ struct __wt_connection { */ int wiredtiger_open(const char *home, WT_EVENT_HANDLER *errhandler, const char *config, - WT_CONNECTION **connectionp); + WT_CONNECTION **connectionp) WT_ATTRIBUTE_LIBRARY_VISIBLE; /*! * Return information about a WiredTiger error as a string (see @@ -2564,7 +2574,7 @@ int wiredtiger_open(const char *home, * @param error a return value from a WiredTiger, ISO C, or POSIX standard API * @returns a string representation of the error */ -const char *wiredtiger_strerror(int error); +const char *wiredtiger_strerror(int error) WT_ATTRIBUTE_LIBRARY_VISIBLE; #if !defined(SWIG) /*! @@ -2701,7 +2711,8 @@ struct __wt_event_handler { * @errors */ int wiredtiger_struct_pack(WT_SESSION *session, - void *buffer, size_t size, const char *format, ...); + void *buffer, size_t size, const char *format, ...) + WT_ATTRIBUTE_LIBRARY_VISIBLE; /*! * Calculate the size required to pack a structure. @@ -2719,7 +2730,7 @@ int wiredtiger_struct_pack(WT_SESSION *session, * @errors */ int wiredtiger_struct_size(WT_SESSION *session, - size_t *sizep, const char *format, ...); + size_t *sizep, const char *format, ...) WT_ATTRIBUTE_LIBRARY_VISIBLE; /*! * Unpack a structure from a buffer. @@ -2736,7 +2747,8 @@ int wiredtiger_struct_size(WT_SESSION *session, * @errors */ int wiredtiger_struct_unpack(WT_SESSION *session, - const void *buffer, size_t size, const char *format, ...); + const void *buffer, size_t size, const char *format, ...) + WT_ATTRIBUTE_LIBRARY_VISIBLE; #if !defined(SWIG) @@ -2763,7 +2775,8 @@ typedef struct __wt_pack_stream WT_PACK_STREAM; * @errors */ int wiredtiger_pack_start(WT_SESSION *session, - const char *format, void *buffer, size_t size, WT_PACK_STREAM **psp); + const char *format, void *buffer, size_t size, WT_PACK_STREAM **psp) + WT_ATTRIBUTE_LIBRARY_VISIBLE; /*! * Start an unpacking operation from a buffer with the given format string. @@ -2779,7 +2792,8 @@ int wiredtiger_pack_start(WT_SESSION *session, * @errors */ int wiredtiger_unpack_start(WT_SESSION *session, - const char *format, const void *buffer, size_t size, WT_PACK_STREAM **psp); + const char *format, const void *buffer, size_t size, WT_PACK_STREAM **psp) + WT_ATTRIBUTE_LIBRARY_VISIBLE; /*! * Close a packing stream. @@ -2788,7 +2802,8 @@ int wiredtiger_unpack_start(WT_SESSION *session, * @param[out] usedp the number of bytes in the buffer used by the stream * @errors */ -int wiredtiger_pack_close(WT_PACK_STREAM *ps, size_t *usedp); +int wiredtiger_pack_close(WT_PACK_STREAM *ps, size_t *usedp) + WT_ATTRIBUTE_LIBRARY_VISIBLE; /*! * Pack an item into a packing stream. @@ -2797,7 +2812,8 @@ int wiredtiger_pack_close(WT_PACK_STREAM *ps, size_t *usedp); * @param item an item to pack * @errors */ -int wiredtiger_pack_item(WT_PACK_STREAM *ps, WT_ITEM *item); +int wiredtiger_pack_item(WT_PACK_STREAM *ps, WT_ITEM *item) + WT_ATTRIBUTE_LIBRARY_VISIBLE; /*! * Pack a signed integer into a packing stream. @@ -2806,7 +2822,8 @@ int wiredtiger_pack_item(WT_PACK_STREAM *ps, WT_ITEM *item); * @param i a signed integer to pack * @errors */ -int wiredtiger_pack_int(WT_PACK_STREAM *ps, int64_t i); +int wiredtiger_pack_int(WT_PACK_STREAM *ps, int64_t i) + WT_ATTRIBUTE_LIBRARY_VISIBLE; /*! * Pack a string into a packing stream. @@ -2815,7 +2832,8 @@ int wiredtiger_pack_int(WT_PACK_STREAM *ps, int64_t i); * @param s a string to pack * @errors */ -int wiredtiger_pack_str(WT_PACK_STREAM *ps, const char *s); +int wiredtiger_pack_str(WT_PACK_STREAM *ps, const char *s) + WT_ATTRIBUTE_LIBRARY_VISIBLE; /*! * Pack an unsigned integer into a packing stream. @@ -2824,7 +2842,8 @@ int wiredtiger_pack_str(WT_PACK_STREAM *ps, const char *s); * @param u an unsigned integer to pack * @errors */ -int wiredtiger_pack_uint(WT_PACK_STREAM *ps, uint64_t u); +int wiredtiger_pack_uint(WT_PACK_STREAM *ps, uint64_t u) + WT_ATTRIBUTE_LIBRARY_VISIBLE; /*! * Unpack an item from a packing stream. @@ -2833,7 +2852,8 @@ int wiredtiger_pack_uint(WT_PACK_STREAM *ps, uint64_t u); * @param item an item to unpack * @errors */ -int wiredtiger_unpack_item(WT_PACK_STREAM *ps, WT_ITEM *item); +int wiredtiger_unpack_item(WT_PACK_STREAM *ps, WT_ITEM *item) + WT_ATTRIBUTE_LIBRARY_VISIBLE; /*! * Unpack a signed integer from a packing stream. @@ -2842,7 +2862,8 @@ int wiredtiger_unpack_item(WT_PACK_STREAM *ps, WT_ITEM *item); * @param[out] ip the unpacked signed integer * @errors */ -int wiredtiger_unpack_int(WT_PACK_STREAM *ps, int64_t *ip); +int wiredtiger_unpack_int(WT_PACK_STREAM *ps, int64_t *ip) + WT_ATTRIBUTE_LIBRARY_VISIBLE; /*! * Unpack a string from a packing stream. @@ -2851,7 +2872,8 @@ int wiredtiger_unpack_int(WT_PACK_STREAM *ps, int64_t *ip); * @param[out] sp the unpacked string * @errors */ -int wiredtiger_unpack_str(WT_PACK_STREAM *ps, const char **sp); +int wiredtiger_unpack_str(WT_PACK_STREAM *ps, const char **sp) + WT_ATTRIBUTE_LIBRARY_VISIBLE; /*! * Unpack an unsigned integer from a packing stream. @@ -2860,7 +2882,8 @@ int wiredtiger_unpack_str(WT_PACK_STREAM *ps, const char **sp); * @param[out] up the unpacked unsigned integer * @errors */ -int wiredtiger_unpack_uint(WT_PACK_STREAM *ps, uint64_t *up); +int wiredtiger_unpack_uint(WT_PACK_STREAM *ps, uint64_t *up) + WT_ATTRIBUTE_LIBRARY_VISIBLE; /*! @} */ /*! @@ -2938,7 +2961,8 @@ struct __wt_config_item { * @snippet ex_all.c Validate a configuration string */ int wiredtiger_config_validate(WT_SESSION *session, - WT_EVENT_HANDLER *errhandler, const char *name, const char *config); + WT_EVENT_HANDLER *errhandler, const char *name, const char *config) + WT_ATTRIBUTE_LIBRARY_VISIBLE; #endif /*! @@ -2958,7 +2982,8 @@ int wiredtiger_config_validate(WT_SESSION *session, * @snippet ex_config_parse.c Create a configuration parser */ int wiredtiger_config_parser_open(WT_SESSION *session, - const char *config, size_t len, WT_CONFIG_PARSER **config_parserp); + const char *config, size_t len, WT_CONFIG_PARSER **config_parserp) + WT_ATTRIBUTE_LIBRARY_VISIBLE; /*! * A handle that can be used to search and traverse configuration strings @@ -3047,7 +3072,8 @@ struct __wt_config_parser { * @param patchp a location where the patch version number is returned * @returns a string representation of the version */ -const char *wiredtiger_version(int *majorp, int *minorp, int *patchp); +const char *wiredtiger_version(int *majorp, int *minorp, int *patchp) + WT_ATTRIBUTE_LIBRARY_VISIBLE; /******************************************* * Error returns @@ -4546,304 +4572,306 @@ extern int wiredtiger_extension_terminate(WT_CONNECTION *connection); #define WT_STAT_CONN_COND_AUTO_WAIT_RESET 1102 /*! connection: auto adjusting condition wait calls */ #define WT_STAT_CONN_COND_AUTO_WAIT 1103 +/*! connection: detected system time went backwards */ +#define WT_STAT_CONN_TIME_TRAVEL 1104 /*! connection: files currently open */ -#define WT_STAT_CONN_FILE_OPEN 1104 +#define WT_STAT_CONN_FILE_OPEN 1105 /*! connection: memory allocations */ -#define WT_STAT_CONN_MEMORY_ALLOCATION 1105 +#define WT_STAT_CONN_MEMORY_ALLOCATION 1106 /*! connection: memory frees */ -#define WT_STAT_CONN_MEMORY_FREE 1106 +#define WT_STAT_CONN_MEMORY_FREE 1107 /*! connection: memory re-allocations */ -#define WT_STAT_CONN_MEMORY_GROW 1107 +#define WT_STAT_CONN_MEMORY_GROW 1108 /*! connection: pthread mutex condition wait calls */ -#define WT_STAT_CONN_COND_WAIT 1108 +#define WT_STAT_CONN_COND_WAIT 1109 /*! connection: pthread mutex shared lock read-lock calls */ -#define WT_STAT_CONN_RWLOCK_READ 1109 +#define WT_STAT_CONN_RWLOCK_READ 1110 /*! connection: pthread mutex shared lock write-lock calls */ -#define WT_STAT_CONN_RWLOCK_WRITE 1110 +#define WT_STAT_CONN_RWLOCK_WRITE 1111 /*! connection: total fsync I/Os */ -#define WT_STAT_CONN_FSYNC_IO 1111 +#define WT_STAT_CONN_FSYNC_IO 1112 /*! connection: total read I/Os */ -#define WT_STAT_CONN_READ_IO 1112 +#define WT_STAT_CONN_READ_IO 1113 /*! connection: total write I/Os */ -#define WT_STAT_CONN_WRITE_IO 1113 +#define WT_STAT_CONN_WRITE_IO 1114 /*! cursor: cursor create calls */ -#define WT_STAT_CONN_CURSOR_CREATE 1114 +#define WT_STAT_CONN_CURSOR_CREATE 1115 /*! cursor: cursor insert calls */ -#define WT_STAT_CONN_CURSOR_INSERT 1115 +#define WT_STAT_CONN_CURSOR_INSERT 1116 /*! cursor: cursor next calls */ -#define WT_STAT_CONN_CURSOR_NEXT 1116 +#define WT_STAT_CONN_CURSOR_NEXT 1117 /*! cursor: cursor prev calls */ -#define WT_STAT_CONN_CURSOR_PREV 1117 +#define WT_STAT_CONN_CURSOR_PREV 1118 /*! cursor: cursor remove calls */ -#define WT_STAT_CONN_CURSOR_REMOVE 1118 +#define WT_STAT_CONN_CURSOR_REMOVE 1119 /*! cursor: cursor reset calls */ -#define WT_STAT_CONN_CURSOR_RESET 1119 +#define WT_STAT_CONN_CURSOR_RESET 1120 /*! cursor: cursor restarted searches */ -#define WT_STAT_CONN_CURSOR_RESTART 1120 +#define WT_STAT_CONN_CURSOR_RESTART 1121 /*! cursor: cursor search calls */ -#define WT_STAT_CONN_CURSOR_SEARCH 1121 +#define WT_STAT_CONN_CURSOR_SEARCH 1122 /*! cursor: cursor search near calls */ -#define WT_STAT_CONN_CURSOR_SEARCH_NEAR 1122 +#define WT_STAT_CONN_CURSOR_SEARCH_NEAR 1123 /*! cursor: cursor update calls */ -#define WT_STAT_CONN_CURSOR_UPDATE 1123 +#define WT_STAT_CONN_CURSOR_UPDATE 1124 /*! cursor: truncate calls */ -#define WT_STAT_CONN_CURSOR_TRUNCATE 1124 +#define WT_STAT_CONN_CURSOR_TRUNCATE 1125 /*! data-handle: connection data handles currently active */ -#define WT_STAT_CONN_DH_CONN_HANDLE_COUNT 1125 +#define WT_STAT_CONN_DH_CONN_HANDLE_COUNT 1126 /*! data-handle: connection sweep candidate became referenced */ -#define WT_STAT_CONN_DH_SWEEP_REF 1126 +#define WT_STAT_CONN_DH_SWEEP_REF 1127 /*! data-handle: connection sweep dhandles closed */ -#define WT_STAT_CONN_DH_SWEEP_CLOSE 1127 +#define WT_STAT_CONN_DH_SWEEP_CLOSE 1128 /*! data-handle: connection sweep dhandles removed from hash list */ -#define WT_STAT_CONN_DH_SWEEP_REMOVE 1128 +#define WT_STAT_CONN_DH_SWEEP_REMOVE 1129 /*! data-handle: connection sweep time-of-death sets */ -#define WT_STAT_CONN_DH_SWEEP_TOD 1129 +#define WT_STAT_CONN_DH_SWEEP_TOD 1130 /*! data-handle: connection sweeps */ -#define WT_STAT_CONN_DH_SWEEPS 1130 +#define WT_STAT_CONN_DH_SWEEPS 1131 /*! data-handle: session dhandles swept */ -#define WT_STAT_CONN_DH_SESSION_HANDLES 1131 +#define WT_STAT_CONN_DH_SESSION_HANDLES 1132 /*! data-handle: session sweep attempts */ -#define WT_STAT_CONN_DH_SESSION_SWEEPS 1132 +#define WT_STAT_CONN_DH_SESSION_SWEEPS 1133 /*! lock: checkpoint lock acquisitions */ -#define WT_STAT_CONN_LOCK_CHECKPOINT_COUNT 1133 +#define WT_STAT_CONN_LOCK_CHECKPOINT_COUNT 1134 /*! lock: checkpoint lock application thread wait time (usecs) */ -#define WT_STAT_CONN_LOCK_CHECKPOINT_WAIT_APPLICATION 1134 +#define WT_STAT_CONN_LOCK_CHECKPOINT_WAIT_APPLICATION 1135 /*! lock: checkpoint lock internal thread wait time (usecs) */ -#define WT_STAT_CONN_LOCK_CHECKPOINT_WAIT_INTERNAL 1135 +#define WT_STAT_CONN_LOCK_CHECKPOINT_WAIT_INTERNAL 1136 /*! lock: handle-list lock eviction thread wait time (usecs) */ -#define WT_STAT_CONN_LOCK_HANDLE_LIST_WAIT_EVICTION 1136 +#define WT_STAT_CONN_LOCK_HANDLE_LIST_WAIT_EVICTION 1137 /*! lock: metadata lock acquisitions */ -#define WT_STAT_CONN_LOCK_METADATA_COUNT 1137 +#define WT_STAT_CONN_LOCK_METADATA_COUNT 1138 /*! lock: metadata lock application thread wait time (usecs) */ -#define WT_STAT_CONN_LOCK_METADATA_WAIT_APPLICATION 1138 +#define WT_STAT_CONN_LOCK_METADATA_WAIT_APPLICATION 1139 /*! lock: metadata lock internal thread wait time (usecs) */ -#define WT_STAT_CONN_LOCK_METADATA_WAIT_INTERNAL 1139 +#define WT_STAT_CONN_LOCK_METADATA_WAIT_INTERNAL 1140 /*! lock: schema lock acquisitions */ -#define WT_STAT_CONN_LOCK_SCHEMA_COUNT 1140 +#define WT_STAT_CONN_LOCK_SCHEMA_COUNT 1141 /*! lock: schema lock application thread wait time (usecs) */ -#define WT_STAT_CONN_LOCK_SCHEMA_WAIT_APPLICATION 1141 +#define WT_STAT_CONN_LOCK_SCHEMA_WAIT_APPLICATION 1142 /*! lock: schema lock internal thread wait time (usecs) */ -#define WT_STAT_CONN_LOCK_SCHEMA_WAIT_INTERNAL 1142 +#define WT_STAT_CONN_LOCK_SCHEMA_WAIT_INTERNAL 1143 /*! lock: table lock acquisitions */ -#define WT_STAT_CONN_LOCK_TABLE_COUNT 1143 +#define WT_STAT_CONN_LOCK_TABLE_COUNT 1144 /*! * lock: table lock application thread time waiting for the table lock * (usecs) */ -#define WT_STAT_CONN_LOCK_TABLE_WAIT_APPLICATION 1144 +#define WT_STAT_CONN_LOCK_TABLE_WAIT_APPLICATION 1145 /*! * lock: table lock internal thread time waiting for the table lock * (usecs) */ -#define WT_STAT_CONN_LOCK_TABLE_WAIT_INTERNAL 1145 +#define WT_STAT_CONN_LOCK_TABLE_WAIT_INTERNAL 1146 /*! log: busy returns attempting to switch slots */ -#define WT_STAT_CONN_LOG_SLOT_SWITCH_BUSY 1146 +#define WT_STAT_CONN_LOG_SLOT_SWITCH_BUSY 1147 /*! log: consolidated slot closures */ -#define WT_STAT_CONN_LOG_SLOT_CLOSES 1147 +#define WT_STAT_CONN_LOG_SLOT_CLOSES 1148 /*! log: consolidated slot join active slot closed */ -#define WT_STAT_CONN_LOG_SLOT_ACTIVE_CLOSED 1148 +#define WT_STAT_CONN_LOG_SLOT_ACTIVE_CLOSED 1149 /*! log: consolidated slot join races */ -#define WT_STAT_CONN_LOG_SLOT_RACES 1149 +#define WT_STAT_CONN_LOG_SLOT_RACES 1150 /*! log: consolidated slot join transitions */ -#define WT_STAT_CONN_LOG_SLOT_TRANSITIONS 1150 +#define WT_STAT_CONN_LOG_SLOT_TRANSITIONS 1151 /*! log: consolidated slot joins */ -#define WT_STAT_CONN_LOG_SLOT_JOINS 1151 +#define WT_STAT_CONN_LOG_SLOT_JOINS 1152 /*! log: consolidated slot transitions unable to find free slot */ -#define WT_STAT_CONN_LOG_SLOT_NO_FREE_SLOTS 1152 +#define WT_STAT_CONN_LOG_SLOT_NO_FREE_SLOTS 1153 /*! log: consolidated slot unbuffered writes */ -#define WT_STAT_CONN_LOG_SLOT_UNBUFFERED 1153 +#define WT_STAT_CONN_LOG_SLOT_UNBUFFERED 1154 /*! log: log bytes of payload data */ -#define WT_STAT_CONN_LOG_BYTES_PAYLOAD 1154 +#define WT_STAT_CONN_LOG_BYTES_PAYLOAD 1155 /*! log: log bytes written */ -#define WT_STAT_CONN_LOG_BYTES_WRITTEN 1155 +#define WT_STAT_CONN_LOG_BYTES_WRITTEN 1156 /*! log: log files manually zero-filled */ -#define WT_STAT_CONN_LOG_ZERO_FILLS 1156 +#define WT_STAT_CONN_LOG_ZERO_FILLS 1157 /*! log: log flush operations */ -#define WT_STAT_CONN_LOG_FLUSH 1157 +#define WT_STAT_CONN_LOG_FLUSH 1158 /*! log: log force write operations */ -#define WT_STAT_CONN_LOG_FORCE_WRITE 1158 +#define WT_STAT_CONN_LOG_FORCE_WRITE 1159 /*! log: log force write operations skipped */ -#define WT_STAT_CONN_LOG_FORCE_WRITE_SKIP 1159 +#define WT_STAT_CONN_LOG_FORCE_WRITE_SKIP 1160 /*! log: log records compressed */ -#define WT_STAT_CONN_LOG_COMPRESS_WRITES 1160 +#define WT_STAT_CONN_LOG_COMPRESS_WRITES 1161 /*! log: log records not compressed */ -#define WT_STAT_CONN_LOG_COMPRESS_WRITE_FAILS 1161 +#define WT_STAT_CONN_LOG_COMPRESS_WRITE_FAILS 1162 /*! log: log records too small to compress */ -#define WT_STAT_CONN_LOG_COMPRESS_SMALL 1162 +#define WT_STAT_CONN_LOG_COMPRESS_SMALL 1163 /*! log: log release advances write LSN */ -#define WT_STAT_CONN_LOG_RELEASE_WRITE_LSN 1163 +#define WT_STAT_CONN_LOG_RELEASE_WRITE_LSN 1164 /*! log: log scan operations */ -#define WT_STAT_CONN_LOG_SCANS 1164 +#define WT_STAT_CONN_LOG_SCANS 1165 /*! log: log scan records requiring two reads */ -#define WT_STAT_CONN_LOG_SCAN_REREADS 1165 +#define WT_STAT_CONN_LOG_SCAN_REREADS 1166 /*! log: log server thread advances write LSN */ -#define WT_STAT_CONN_LOG_WRITE_LSN 1166 +#define WT_STAT_CONN_LOG_WRITE_LSN 1167 /*! log: log server thread write LSN walk skipped */ -#define WT_STAT_CONN_LOG_WRITE_LSN_SKIP 1167 +#define WT_STAT_CONN_LOG_WRITE_LSN_SKIP 1168 /*! log: log sync operations */ -#define WT_STAT_CONN_LOG_SYNC 1168 +#define WT_STAT_CONN_LOG_SYNC 1169 /*! log: log sync time duration (usecs) */ -#define WT_STAT_CONN_LOG_SYNC_DURATION 1169 +#define WT_STAT_CONN_LOG_SYNC_DURATION 1170 /*! log: log sync_dir operations */ -#define WT_STAT_CONN_LOG_SYNC_DIR 1170 +#define WT_STAT_CONN_LOG_SYNC_DIR 1171 /*! log: log sync_dir time duration (usecs) */ -#define WT_STAT_CONN_LOG_SYNC_DIR_DURATION 1171 +#define WT_STAT_CONN_LOG_SYNC_DIR_DURATION 1172 /*! log: log write operations */ -#define WT_STAT_CONN_LOG_WRITES 1172 +#define WT_STAT_CONN_LOG_WRITES 1173 /*! log: logging bytes consolidated */ -#define WT_STAT_CONN_LOG_SLOT_CONSOLIDATED 1173 +#define WT_STAT_CONN_LOG_SLOT_CONSOLIDATED 1174 /*! log: maximum log file size */ -#define WT_STAT_CONN_LOG_MAX_FILESIZE 1174 +#define WT_STAT_CONN_LOG_MAX_FILESIZE 1175 /*! log: number of pre-allocated log files to create */ -#define WT_STAT_CONN_LOG_PREALLOC_MAX 1175 +#define WT_STAT_CONN_LOG_PREALLOC_MAX 1176 /*! log: pre-allocated log files not ready and missed */ -#define WT_STAT_CONN_LOG_PREALLOC_MISSED 1176 +#define WT_STAT_CONN_LOG_PREALLOC_MISSED 1177 /*! log: pre-allocated log files prepared */ -#define WT_STAT_CONN_LOG_PREALLOC_FILES 1177 +#define WT_STAT_CONN_LOG_PREALLOC_FILES 1178 /*! log: pre-allocated log files used */ -#define WT_STAT_CONN_LOG_PREALLOC_USED 1178 +#define WT_STAT_CONN_LOG_PREALLOC_USED 1179 /*! log: records processed by log scan */ -#define WT_STAT_CONN_LOG_SCAN_RECORDS 1179 +#define WT_STAT_CONN_LOG_SCAN_RECORDS 1180 /*! log: total in-memory size of compressed records */ -#define WT_STAT_CONN_LOG_COMPRESS_MEM 1180 +#define WT_STAT_CONN_LOG_COMPRESS_MEM 1181 /*! log: total log buffer size */ -#define WT_STAT_CONN_LOG_BUFFER_SIZE 1181 +#define WT_STAT_CONN_LOG_BUFFER_SIZE 1182 /*! log: total size of compressed records */ -#define WT_STAT_CONN_LOG_COMPRESS_LEN 1182 +#define WT_STAT_CONN_LOG_COMPRESS_LEN 1183 /*! log: written slots coalesced */ -#define WT_STAT_CONN_LOG_SLOT_COALESCED 1183 +#define WT_STAT_CONN_LOG_SLOT_COALESCED 1184 /*! log: yields waiting for previous log file close */ -#define WT_STAT_CONN_LOG_CLOSE_YIELDS 1184 +#define WT_STAT_CONN_LOG_CLOSE_YIELDS 1185 /*! reconciliation: fast-path pages deleted */ -#define WT_STAT_CONN_REC_PAGE_DELETE_FAST 1185 +#define WT_STAT_CONN_REC_PAGE_DELETE_FAST 1186 /*! reconciliation: page reconciliation calls */ -#define WT_STAT_CONN_REC_PAGES 1186 +#define WT_STAT_CONN_REC_PAGES 1187 /*! reconciliation: page reconciliation calls for eviction */ -#define WT_STAT_CONN_REC_PAGES_EVICTION 1187 +#define WT_STAT_CONN_REC_PAGES_EVICTION 1188 /*! reconciliation: pages deleted */ -#define WT_STAT_CONN_REC_PAGE_DELETE 1188 +#define WT_STAT_CONN_REC_PAGE_DELETE 1189 /*! reconciliation: split bytes currently awaiting free */ -#define WT_STAT_CONN_REC_SPLIT_STASHED_BYTES 1189 +#define WT_STAT_CONN_REC_SPLIT_STASHED_BYTES 1190 /*! reconciliation: split objects currently awaiting free */ -#define WT_STAT_CONN_REC_SPLIT_STASHED_OBJECTS 1190 +#define WT_STAT_CONN_REC_SPLIT_STASHED_OBJECTS 1191 /*! session: open cursor count */ -#define WT_STAT_CONN_SESSION_CURSOR_OPEN 1191 +#define WT_STAT_CONN_SESSION_CURSOR_OPEN 1192 /*! session: open session count */ -#define WT_STAT_CONN_SESSION_OPEN 1192 +#define WT_STAT_CONN_SESSION_OPEN 1193 /*! session: table alter failed calls */ -#define WT_STAT_CONN_SESSION_TABLE_ALTER_FAIL 1193 +#define WT_STAT_CONN_SESSION_TABLE_ALTER_FAIL 1194 /*! session: table alter successful calls */ -#define WT_STAT_CONN_SESSION_TABLE_ALTER_SUCCESS 1194 +#define WT_STAT_CONN_SESSION_TABLE_ALTER_SUCCESS 1195 /*! session: table alter unchanged and skipped */ -#define WT_STAT_CONN_SESSION_TABLE_ALTER_SKIP 1195 +#define WT_STAT_CONN_SESSION_TABLE_ALTER_SKIP 1196 /*! session: table compact failed calls */ -#define WT_STAT_CONN_SESSION_TABLE_COMPACT_FAIL 1196 +#define WT_STAT_CONN_SESSION_TABLE_COMPACT_FAIL 1197 /*! session: table compact successful calls */ -#define WT_STAT_CONN_SESSION_TABLE_COMPACT_SUCCESS 1197 +#define WT_STAT_CONN_SESSION_TABLE_COMPACT_SUCCESS 1198 /*! session: table create failed calls */ -#define WT_STAT_CONN_SESSION_TABLE_CREATE_FAIL 1198 +#define WT_STAT_CONN_SESSION_TABLE_CREATE_FAIL 1199 /*! session: table create successful calls */ -#define WT_STAT_CONN_SESSION_TABLE_CREATE_SUCCESS 1199 +#define WT_STAT_CONN_SESSION_TABLE_CREATE_SUCCESS 1200 /*! session: table drop failed calls */ -#define WT_STAT_CONN_SESSION_TABLE_DROP_FAIL 1200 +#define WT_STAT_CONN_SESSION_TABLE_DROP_FAIL 1201 /*! session: table drop successful calls */ -#define WT_STAT_CONN_SESSION_TABLE_DROP_SUCCESS 1201 +#define WT_STAT_CONN_SESSION_TABLE_DROP_SUCCESS 1202 /*! session: table rebalance failed calls */ -#define WT_STAT_CONN_SESSION_TABLE_REBALANCE_FAIL 1202 +#define WT_STAT_CONN_SESSION_TABLE_REBALANCE_FAIL 1203 /*! session: table rebalance successful calls */ -#define WT_STAT_CONN_SESSION_TABLE_REBALANCE_SUCCESS 1203 +#define WT_STAT_CONN_SESSION_TABLE_REBALANCE_SUCCESS 1204 /*! session: table rename failed calls */ -#define WT_STAT_CONN_SESSION_TABLE_RENAME_FAIL 1204 +#define WT_STAT_CONN_SESSION_TABLE_RENAME_FAIL 1205 /*! session: table rename successful calls */ -#define WT_STAT_CONN_SESSION_TABLE_RENAME_SUCCESS 1205 +#define WT_STAT_CONN_SESSION_TABLE_RENAME_SUCCESS 1206 /*! session: table salvage failed calls */ -#define WT_STAT_CONN_SESSION_TABLE_SALVAGE_FAIL 1206 +#define WT_STAT_CONN_SESSION_TABLE_SALVAGE_FAIL 1207 /*! session: table salvage successful calls */ -#define WT_STAT_CONN_SESSION_TABLE_SALVAGE_SUCCESS 1207 +#define WT_STAT_CONN_SESSION_TABLE_SALVAGE_SUCCESS 1208 /*! session: table truncate failed calls */ -#define WT_STAT_CONN_SESSION_TABLE_TRUNCATE_FAIL 1208 +#define WT_STAT_CONN_SESSION_TABLE_TRUNCATE_FAIL 1209 /*! session: table truncate successful calls */ -#define WT_STAT_CONN_SESSION_TABLE_TRUNCATE_SUCCESS 1209 +#define WT_STAT_CONN_SESSION_TABLE_TRUNCATE_SUCCESS 1210 /*! session: table verify failed calls */ -#define WT_STAT_CONN_SESSION_TABLE_VERIFY_FAIL 1210 +#define WT_STAT_CONN_SESSION_TABLE_VERIFY_FAIL 1211 /*! session: table verify successful calls */ -#define WT_STAT_CONN_SESSION_TABLE_VERIFY_SUCCESS 1211 +#define WT_STAT_CONN_SESSION_TABLE_VERIFY_SUCCESS 1212 /*! thread-state: active filesystem fsync calls */ -#define WT_STAT_CONN_THREAD_FSYNC_ACTIVE 1212 +#define WT_STAT_CONN_THREAD_FSYNC_ACTIVE 1213 /*! thread-state: active filesystem read calls */ -#define WT_STAT_CONN_THREAD_READ_ACTIVE 1213 +#define WT_STAT_CONN_THREAD_READ_ACTIVE 1214 /*! thread-state: active filesystem write calls */ -#define WT_STAT_CONN_THREAD_WRITE_ACTIVE 1214 +#define WT_STAT_CONN_THREAD_WRITE_ACTIVE 1215 /*! thread-yield: application thread time evicting (usecs) */ -#define WT_STAT_CONN_APPLICATION_EVICT_TIME 1215 +#define WT_STAT_CONN_APPLICATION_EVICT_TIME 1216 /*! thread-yield: application thread time waiting for cache (usecs) */ -#define WT_STAT_CONN_APPLICATION_CACHE_TIME 1216 +#define WT_STAT_CONN_APPLICATION_CACHE_TIME 1217 /*! thread-yield: page acquire busy blocked */ -#define WT_STAT_CONN_PAGE_BUSY_BLOCKED 1217 +#define WT_STAT_CONN_PAGE_BUSY_BLOCKED 1218 /*! thread-yield: page acquire eviction blocked */ -#define WT_STAT_CONN_PAGE_FORCIBLE_EVICT_BLOCKED 1218 +#define WT_STAT_CONN_PAGE_FORCIBLE_EVICT_BLOCKED 1219 /*! thread-yield: page acquire locked blocked */ -#define WT_STAT_CONN_PAGE_LOCKED_BLOCKED 1219 +#define WT_STAT_CONN_PAGE_LOCKED_BLOCKED 1220 /*! thread-yield: page acquire read blocked */ -#define WT_STAT_CONN_PAGE_READ_BLOCKED 1220 +#define WT_STAT_CONN_PAGE_READ_BLOCKED 1221 /*! thread-yield: page acquire time sleeping (usecs) */ -#define WT_STAT_CONN_PAGE_SLEEP 1221 +#define WT_STAT_CONN_PAGE_SLEEP 1222 /*! transaction: number of named snapshots created */ -#define WT_STAT_CONN_TXN_SNAPSHOTS_CREATED 1222 +#define WT_STAT_CONN_TXN_SNAPSHOTS_CREATED 1223 /*! transaction: number of named snapshots dropped */ -#define WT_STAT_CONN_TXN_SNAPSHOTS_DROPPED 1223 +#define WT_STAT_CONN_TXN_SNAPSHOTS_DROPPED 1224 /*! transaction: transaction begins */ -#define WT_STAT_CONN_TXN_BEGIN 1224 +#define WT_STAT_CONN_TXN_BEGIN 1225 /*! transaction: transaction checkpoint currently running */ -#define WT_STAT_CONN_TXN_CHECKPOINT_RUNNING 1225 +#define WT_STAT_CONN_TXN_CHECKPOINT_RUNNING 1226 /*! transaction: transaction checkpoint generation */ -#define WT_STAT_CONN_TXN_CHECKPOINT_GENERATION 1226 +#define WT_STAT_CONN_TXN_CHECKPOINT_GENERATION 1227 /*! transaction: transaction checkpoint max time (msecs) */ -#define WT_STAT_CONN_TXN_CHECKPOINT_TIME_MAX 1227 +#define WT_STAT_CONN_TXN_CHECKPOINT_TIME_MAX 1228 /*! transaction: transaction checkpoint min time (msecs) */ -#define WT_STAT_CONN_TXN_CHECKPOINT_TIME_MIN 1228 +#define WT_STAT_CONN_TXN_CHECKPOINT_TIME_MIN 1229 /*! transaction: transaction checkpoint most recent time (msecs) */ -#define WT_STAT_CONN_TXN_CHECKPOINT_TIME_RECENT 1229 +#define WT_STAT_CONN_TXN_CHECKPOINT_TIME_RECENT 1230 /*! transaction: transaction checkpoint scrub dirty target */ -#define WT_STAT_CONN_TXN_CHECKPOINT_SCRUB_TARGET 1230 +#define WT_STAT_CONN_TXN_CHECKPOINT_SCRUB_TARGET 1231 /*! transaction: transaction checkpoint scrub time (msecs) */ -#define WT_STAT_CONN_TXN_CHECKPOINT_SCRUB_TIME 1231 +#define WT_STAT_CONN_TXN_CHECKPOINT_SCRUB_TIME 1232 /*! transaction: transaction checkpoint total time (msecs) */ -#define WT_STAT_CONN_TXN_CHECKPOINT_TIME_TOTAL 1232 +#define WT_STAT_CONN_TXN_CHECKPOINT_TIME_TOTAL 1233 /*! transaction: transaction checkpoints */ -#define WT_STAT_CONN_TXN_CHECKPOINT 1233 +#define WT_STAT_CONN_TXN_CHECKPOINT 1234 /*! * transaction: transaction checkpoints skipped because database was * clean */ -#define WT_STAT_CONN_TXN_CHECKPOINT_SKIPPED 1234 +#define WT_STAT_CONN_TXN_CHECKPOINT_SKIPPED 1235 /*! transaction: transaction failures due to cache overflow */ -#define WT_STAT_CONN_TXN_FAIL_CACHE 1235 +#define WT_STAT_CONN_TXN_FAIL_CACHE 1236 /*! * transaction: transaction fsync calls for checkpoint after allocating * the transaction ID */ -#define WT_STAT_CONN_TXN_CHECKPOINT_FSYNC_POST 1236 +#define WT_STAT_CONN_TXN_CHECKPOINT_FSYNC_POST 1237 /*! * transaction: transaction fsync duration for checkpoint after * allocating the transaction ID (usecs) */ -#define WT_STAT_CONN_TXN_CHECKPOINT_FSYNC_POST_DURATION 1237 +#define WT_STAT_CONN_TXN_CHECKPOINT_FSYNC_POST_DURATION 1238 /*! transaction: transaction range of IDs currently pinned */ -#define WT_STAT_CONN_TXN_PINNED_RANGE 1238 +#define WT_STAT_CONN_TXN_PINNED_RANGE 1239 /*! transaction: transaction range of IDs currently pinned by a checkpoint */ -#define WT_STAT_CONN_TXN_PINNED_CHECKPOINT_RANGE 1239 +#define WT_STAT_CONN_TXN_PINNED_CHECKPOINT_RANGE 1240 /*! * transaction: transaction range of IDs currently pinned by named * snapshots */ -#define WT_STAT_CONN_TXN_PINNED_SNAPSHOT_RANGE 1240 +#define WT_STAT_CONN_TXN_PINNED_SNAPSHOT_RANGE 1241 /*! transaction: transaction sync calls */ -#define WT_STAT_CONN_TXN_SYNC 1241 +#define WT_STAT_CONN_TXN_SYNC 1242 /*! transaction: transactions committed */ -#define WT_STAT_CONN_TXN_COMMIT 1242 +#define WT_STAT_CONN_TXN_COMMIT 1243 /*! transaction: transactions rolled back */ -#define WT_STAT_CONN_TXN_ROLLBACK 1243 +#define WT_STAT_CONN_TXN_ROLLBACK 1244 /*! * @} diff --git a/src/third_party/wiredtiger/src/os_common/os_alloc.c b/src/third_party/wiredtiger/src/os_common/os_alloc.c index ef96ed09ea7..c54bcc718f2 100644 --- a/src/third_party/wiredtiger/src/os_common/os_alloc.c +++ b/src/third_party/wiredtiger/src/os_common/os_alloc.c @@ -266,6 +266,8 @@ __wt_strndup(WT_SESSION_IMPL *session, const void *str, size_t len, void *retp) WT_RET(__wt_malloc(session, len + 1, &p)); + WT_ASSERT(session, p != NULL); /* quiet clang scan-build */ + /* * Don't change this to strncpy, we rely on this function to duplicate * "strings" that contain nul bytes. diff --git a/src/third_party/wiredtiger/src/os_common/os_getopt.c b/src/third_party/wiredtiger/src/os_common/os_getopt.c index 960776c3999..fa21123ba0e 100644 --- a/src/third_party/wiredtiger/src/os_common/os_getopt.c +++ b/src/third_party/wiredtiger/src/os_common/os_getopt.c @@ -59,13 +59,17 @@ #include "wt_internal.h" -extern int __wt_opterr, __wt_optind, __wt_optopt, __wt_optreset; +extern int __wt_opterr WT_ATTRIBUTE_LIBRARY_VISIBLE; +extern int __wt_optind WT_ATTRIBUTE_LIBRARY_VISIBLE; +extern int __wt_optopt WT_ATTRIBUTE_LIBRARY_VISIBLE; +extern int __wt_optreset WT_ATTRIBUTE_LIBRARY_VISIBLE; + int __wt_opterr = 1, /* if error message should be printed */ __wt_optind = 1, /* index into parent argv vector */ __wt_optopt, /* character checked for validity */ __wt_optreset; /* reset getopt */ -extern char *__wt_optarg; +extern char *__wt_optarg WT_ATTRIBUTE_LIBRARY_VISIBLE; char *__wt_optarg; /* argument associated with option */ #define BADCH (int)'?' diff --git a/src/third_party/wiredtiger/src/os_posix/os_dir.c b/src/third_party/wiredtiger/src/os_posix/os_dir.c index 627278540d1..b1b6571e4ba 100644 --- a/src/third_party/wiredtiger/src/os_posix/os_dir.c +++ b/src/third_party/wiredtiger/src/os_posix/os_dir.c @@ -37,7 +37,13 @@ __wt_posix_directory_list(WT_FILE_SYSTEM *file_system, dirallocsz = 0; entries = NULL; + /* + * If opendir fails, we should have a NULL pointer with an error value, + * but various static analysis programs remain unconvinced, check both. + */ WT_SYSCALL_RETRY(((dirp = opendir(directory)) == NULL ? -1 : 0), ret); + if (dirp == NULL && ret == 0) + ret = EINVAL; if (ret != 0) WT_RET_MSG(session, ret, "%s: directory-list: opendir", directory); diff --git a/src/third_party/wiredtiger/src/os_posix/os_time.c b/src/third_party/wiredtiger/src/os_posix/os_time.c index 6f150ee8ffe..fe337fea7cf 100644 --- a/src/third_party/wiredtiger/src/os_posix/os_time.c +++ b/src/third_party/wiredtiger/src/os_posix/os_time.c @@ -16,6 +16,7 @@ void __wt_epoch(WT_SESSION_IMPL *session, struct timespec *tsp) WT_GCC_FUNC_ATTRIBUTE((visibility("default"))) { + struct timespec tmp; WT_DECL_RET; /* @@ -27,21 +28,34 @@ __wt_epoch(WT_SESSION_IMPL *session, struct timespec *tsp) tsp->tv_sec = 0; tsp->tv_nsec = 0; + /* + * Read into a local variable so that we're comparing the correct + * value when we check for monotonic increasing time. There are + * many places we read into an unlocked global variable. + */ #if defined(HAVE_CLOCK_GETTIME) - WT_SYSCALL_RETRY(clock_gettime(CLOCK_REALTIME, tsp), ret); - if (ret == 0) + WT_SYSCALL_RETRY(clock_gettime(CLOCK_REALTIME, &tmp), ret); + if (ret == 0) { + __wt_time_check_monotonic(session, &tmp); + tsp->tv_sec = tmp.tv_sec; + tsp->tv_nsec = tmp.tv_nsec; return; + } WT_PANIC_MSG(session, ret, "clock_gettime"); #elif defined(HAVE_GETTIMEOFDAY) + { struct timeval v; WT_SYSCALL_RETRY(gettimeofday(&v, NULL), ret); if (ret == 0) { - tsp->tv_sec = v.tv_sec; - tsp->tv_nsec = v.tv_usec * WT_THOUSAND; + tmp.tv_sec = v.tv_sec; + tmp.tv_nsec = v.tv_usec * WT_THOUSAND; + __wt_time_check_monotonic(session, &tmp); + *tsp = tmp; return; } WT_PANIC_MSG(session, ret, "gettimeofday"); + } #else NO TIME-OF-DAY IMPLEMENTATION: see src/os_posix/os_time.c #endif diff --git a/src/third_party/wiredtiger/src/os_win/os_time.c b/src/third_party/wiredtiger/src/os_win/os_time.c index 6aa5b3719f6..ba71341ab22 100644 --- a/src/third_party/wiredtiger/src/os_win/os_time.c +++ b/src/third_party/wiredtiger/src/os_win/os_time.c @@ -15,17 +15,18 @@ void __wt_epoch(WT_SESSION_IMPL *session, struct timespec *tsp) { + struct timespec tmp; FILETIME time; uint64_t ns100; - WT_UNUSED(session); - GetSystemTimeAsFileTime(&time); ns100 = (((int64_t)time.dwHighDateTime << 32) + time.dwLowDateTime) - 116444736000000000LL; - tsp->tv_sec = ns100 / 10000000; - tsp->tv_nsec = (long)((ns100 % 10000000) * 100); + tmp.tv_sec = ns100 / 10000000; + tmp.tv_nsec = (long)((ns100 % 10000000) * 100); + __wt_time_check_monotonic(session, &tmp); + *tsp = tmp; } /* diff --git a/src/third_party/wiredtiger/src/schema/schema_list.c b/src/third_party/wiredtiger/src/schema/schema_list.c index 74ef5135a4a..3dc51b6cb43 100644 --- a/src/third_party/wiredtiger/src/schema/schema_list.c +++ b/src/third_party/wiredtiger/src/schema/schema_list.c @@ -249,3 +249,34 @@ __wt_schema_close_tables(WT_SESSION_IMPL *session) WT_TRET(__wt_schema_remove_table(session, table)); return (ret); } + +/* + * __wt_schema_sweep_tables -- + * Close all idle, obsolete tables in a session. + */ +int +__wt_schema_sweep_tables(WT_SESSION_IMPL *session) +{ + WT_TABLE *table, *next; + uint64_t schema_gen; + bool old_table_busy; + + WT_ORDERED_READ(schema_gen, S2C(session)->schema_gen); + if (schema_gen == session->table_sweep_gen) + return (0); + + old_table_busy = false; + TAILQ_FOREACH_SAFE(table, &session->tables, q, next) + if (table->schema_gen != schema_gen) { + if (table->refcnt == 0) + WT_RET(__wt_schema_remove_table( + session, table)); + else + old_table_busy = true; + } + + if (!old_table_busy) + session->table_sweep_gen = schema_gen; + + return (0); +} diff --git a/src/third_party/wiredtiger/src/session/session_api.c b/src/third_party/wiredtiger/src/session/session_api.c index b7daf0e2e02..5ce6135cfca 100644 --- a/src/third_party/wiredtiger/src/session/session_api.c +++ b/src/third_party/wiredtiger/src/session/session_api.c @@ -818,6 +818,8 @@ __session_reset(WT_SESSION *wt_session) WT_TRET(__wt_session_reset_cursors(session, true)); + WT_TRET(__wt_schema_sweep_tables(session)); + /* Release common session resources. */ WT_TRET(__wt_session_release_resources(session)); @@ -1105,7 +1107,6 @@ int __wt_session_range_truncate(WT_SESSION_IMPL *session, const char *uri, WT_CURSOR *start, WT_CURSOR *stop) { - WT_CURSOR *cursor; WT_DECL_RET; int cmp; bool local_start; @@ -1134,12 +1135,13 @@ __wt_session_range_truncate(WT_SESSION_IMPL *session, } /* - * Cursor truncate is only supported for some objects, check for the - * supporting methods we need, range_truncate and compare. + * Cursor truncate is only supported for some objects, check for a + * supporting compare method. */ - cursor = start == NULL ? stop : start; - if (cursor->compare == NULL) - WT_ERR(__wt_bad_object_type(session, cursor->uri)); + if (start != NULL && start->compare == NULL) + WT_ERR(__wt_bad_object_type(session, start->uri)); + if (stop != NULL && stop->compare == NULL) + WT_ERR(__wt_bad_object_type(session, stop->uri)); /* * If both cursors set, check they're correctly ordered with respect to @@ -1150,6 +1152,9 @@ __wt_session_range_truncate(WT_SESSION_IMPL *session, * reference the same object and the keys are set. */ if (start != NULL && stop != NULL) { + /* quiet clang scan-build */ + WT_ASSERT(session, start->compare != NULL); + WT_ERR(start->compare(start, stop, &cmp)); if (cmp > 0) WT_ERR_MSG(session, EINVAL, diff --git a/src/third_party/wiredtiger/src/support/stat.c b/src/third_party/wiredtiger/src/support/stat.c index 2c2217f8c20..8b72e653658 100644 --- a/src/third_party/wiredtiger/src/support/stat.c +++ b/src/third_party/wiredtiger/src/support/stat.c @@ -728,6 +728,7 @@ static const char * const __stats_connection_desc[] = { "cache: unmodified pages evicted", "connection: auto adjusting condition resets", "connection: auto adjusting condition wait calls", + "connection: detected system time went backwards", "connection: files currently open", "connection: memory allocations", "connection: memory frees", @@ -1014,6 +1015,7 @@ __wt_stat_connection_clear_single(WT_CONNECTION_STATS *stats) stats->cache_eviction_clean = 0; stats->cond_auto_wait_reset = 0; stats->cond_auto_wait = 0; + stats->time_travel = 0; /* not clearing file_open */ stats->memory_allocation = 0; stats->memory_free = 0; @@ -1320,6 +1322,7 @@ __wt_stat_connection_aggregate( to->cache_eviction_clean += WT_STAT_READ(from, cache_eviction_clean); to->cond_auto_wait_reset += WT_STAT_READ(from, cond_auto_wait_reset); to->cond_auto_wait += WT_STAT_READ(from, cond_auto_wait); + to->time_travel += WT_STAT_READ(from, time_travel); to->file_open += WT_STAT_READ(from, file_open); to->memory_allocation += WT_STAT_READ(from, memory_allocation); to->memory_free += WT_STAT_READ(from, memory_free); diff --git a/src/third_party/wiredtiger/src/txn/txn.c b/src/third_party/wiredtiger/src/txn/txn.c index ea7faa2e966..76fdf71e715 100644 --- a/src/third_party/wiredtiger/src/txn/txn.c +++ b/src/third_party/wiredtiger/src/txn/txn.c @@ -503,13 +503,17 @@ __wt_txn_commit(WT_SESSION_IMPL *session, const char *cfg[]) WT_CONNECTION_IMPL *conn; WT_DECL_RET; WT_TXN *txn; + WT_TXN_GLOBAL *txn_global; WT_TXN_OP *op; u_int i; - bool did_update; + bool did_update, locked; txn = &session->txn; conn = S2C(session); + txn_global = &conn->txn_global; did_update = txn->mod_count != 0; + locked = false; + WT_ASSERT(session, !F_ISSET(txn, WT_TXN_ERROR) || !did_update); if (!F_ISSET(txn, WT_TXN_RUNNING)) @@ -580,6 +584,14 @@ __wt_txn_commit(WT_SESSION_IMPL *session, const char *cfg[]) * This is particularly important for checkpoints. */ __wt_txn_release_snapshot(session); + /* + * We hold the visibility lock for reading from the time + * we write our log record until the time we release our + * transaction so that the LSN any checkpoint gets will + * always reflect visible data. + */ + __wt_readlock(session, &txn_global->visibility_rwlock); + locked = true; ret = __wt_txn_log_commit(session, cfg); } @@ -590,8 +602,12 @@ __wt_txn_commit(WT_SESSION_IMPL *session, const char *cfg[]) * Nothing can fail after this point. */ if (ret != 0) { + if (locked) + __wt_readunlock(session, + &txn_global->visibility_rwlock); WT_TRET(__wt_txn_rollback(session, cfg)); return (ret); + } /* Free memory associated with updates. */ @@ -600,6 +616,8 @@ __wt_txn_commit(WT_SESSION_IMPL *session, const char *cfg[]) txn->mod_count = 0; __wt_txn_release(session); + if (locked) + __wt_readunlock(session, &txn_global->visibility_rwlock); return (0); } @@ -770,6 +788,7 @@ __wt_txn_global_init(WT_SESSION_IMPL *session, const char *cfg[]) &txn_global->id_lock, "transaction id lock")); WT_RET(__wt_rwlock_init(session, &txn_global->scan_rwlock)); WT_RET(__wt_rwlock_init(session, &txn_global->nsnap_rwlock)); + WT_RET(__wt_rwlock_init(session, &txn_global->visibility_rwlock)); txn_global->nsnap_oldest_id = WT_TXN_NONE; TAILQ_INIT(&txn_global->nsnaph); @@ -801,6 +820,7 @@ __wt_txn_global_destroy(WT_SESSION_IMPL *session) __wt_spin_destroy(session, &txn_global->id_lock); __wt_rwlock_destroy(session, &txn_global->scan_rwlock); __wt_rwlock_destroy(session, &txn_global->nsnap_rwlock); + __wt_rwlock_destroy(session, &txn_global->visibility_rwlock); __wt_free(session, txn_global->states); } diff --git a/src/third_party/wiredtiger/src/txn/txn_log.c b/src/third_party/wiredtiger/src/txn/txn_log.c index 2931dc1ce82..cb3b3436786 100644 --- a/src/third_party/wiredtiger/src/txn/txn_log.c +++ b/src/third_party/wiredtiger/src/txn/txn_log.c @@ -294,11 +294,13 @@ __wt_txn_checkpoint_log( WT_ITEM *ckpt_snapshot, empty; WT_LSN *ckpt_lsn; WT_TXN *txn; + WT_TXN_GLOBAL *txn_global; uint8_t *end, *p; size_t recsize; uint32_t i, rectype = WT_LOGREC_CHECKPOINT; const char *fmt = WT_UNCHECKED_STRING(IIIIu); + txn_global = &S2C(session)->txn_global; txn = &session->txn; ckpt_lsn = &txn->ckpt_lsn; @@ -320,6 +322,15 @@ __wt_txn_checkpoint_log( txn->full_ckpt = true; WT_ERR(__wt_log_flush_lsn(session, ckpt_lsn, true)); /* + * We take and immediately release the visibility lock. + * Acquiring the write lock guarantees that any transaction + * that has written to the log has also made its transaction + * visible at this time. + */ + __wt_writelock(session, &txn_global->visibility_rwlock); + __wt_writeunlock(session, &txn_global->visibility_rwlock); + + /* * We need to make sure that the log records in the checkpoint * LSN are on disk. In particular to make sure that the * current log file exists. diff --git a/version.json b/version.json index ad5f47b99da..84ff3f9e40e 100644 --- a/version.json +++ b/version.json @@ -1,4 +1,4 @@ { - "githash": "056bf45128114e44c5358c7a8776fb582363e094", - "version": "3.2.16" + "githash": "186656d79574f7dfe0831a7e7821292ab380f667", + "version": "3.2.17" }
\ No newline at end of file |
