summaryrefslogtreecommitdiff
path: root/src/mongo/db/instance.cpp
diff options
context:
space:
mode:
Diffstat (limited to 'src/mongo/db/instance.cpp')
-rw-r--r--src/mongo/db/instance.cpp683
1 files changed, 376 insertions, 307 deletions
diff --git a/src/mongo/db/instance.cpp b/src/mongo/db/instance.cpp
index 926cdeecdb6..40fe6d07c33 100644
--- a/src/mongo/db/instance.cpp
+++ b/src/mongo/db/instance.cpp
@@ -14,53 +14,80 @@
*
* You should have received a copy of the GNU Affero General Public License
* along with this program. If not, see <http://www.gnu.org/licenses/>.
+*
+* As a special exception, the copyright holders give permission to link the
+* code of portions of this program with the OpenSSL library under certain
+* conditions as described in each individual source file and distribute
+* linked combinations including the program with the OpenSSL library. You
+* must comply with the GNU Affero General Public License in all respects for
+* all of the code used other than as permitted herein. If you modify file(s)
+* with this exception, you may extend this exception to your version of the
+* file(s), but you are not obligated to do so. If you do not wish to do so,
+* delete this exception statement from your version. If you delete this
+* exception statement from all source files in the program, then also delete
+* it in the license file.
*/
#include "mongo/pch.h"
+#include <boost/filesystem/operations.hpp>
#include <boost/thread/thread.hpp>
#include <fstream>
-#include <boost/filesystem/operations.hpp>
#if defined(_WIN32)
#include <io.h>
#else
#include <sys/file.h>
#endif
-#include "mongo/util/time_support.h"
#include "mongo/base/status.h"
#include "mongo/bson/util/atomic_int.h"
+#include "mongo/db/audit.h"
#include "mongo/db/auth/action_type.h"
#include "mongo/db/auth/authorization_manager.h"
+#include "mongo/db/auth/authorization_session.h"
#include "mongo/db/background.h"
-#include "mongo/db/cmdline.h"
+#include "mongo/db/clientcursor.h"
#include "mongo/db/commands/fsync.h"
#include "mongo/db/d_concurrency.h"
#include "mongo/db/db.h"
+#include "mongo/db/dbhelpers.h"
#include "mongo/db/dbmessage.h"
#include "mongo/db/dur_commitjob.h"
#include "mongo/db/dur_journal.h"
#include "mongo/db/dur_recover.h"
#include "mongo/db/instance.h"
#include "mongo/db/introspect.h"
+#include "mongo/db/jsobjmanipulator.h"
#include "mongo/db/json.h"
#include "mongo/db/kill_current_op.h"
#include "mongo/db/lasterror.h"
-#include "mongo/db/namespacestring.h"
+#include "mongo/db/matcher.h"
+#include "mongo/db/mongod_options.h"
+#include "mongo/db/namespace_string.h"
#include "mongo/db/ops/count.h"
-#include "mongo/db/ops/delete.h"
-#include "mongo/db/ops/query.h"
-#include "mongo/db/ops/update.h"
+#include "mongo/db/ops/delete_executor.h"
+#include "mongo/db/ops/delete_request.h"
+#include "mongo/db/ops/insert.h"
+#include "mongo/db/ops/update_lifecycle_impl.h"
+#include "mongo/db/ops/update_driver.h"
+#include "mongo/db/ops/update_executor.h"
+#include "mongo/db/ops/update_request.h"
#include "mongo/db/pagefault.h"
-#include "mongo/db/repl.h"
-#include "mongo/db/replutil.h"
+#include "mongo/db/query/new_find.h"
+#include "mongo/db/repl/is_master.h"
+#include "mongo/db/repl/oplog.h"
#include "mongo/db/stats/counters.h"
+#include "mongo/db/storage_options.h"
+#include "mongo/platform/process_id.h"
#include "mongo/s/d_logic.h"
#include "mongo/s/stale_exception.h" // for SendStaleConfigException
+#include "mongo/scripting/engine.h"
#include "mongo/util/fail_point_service.h"
#include "mongo/util/file_allocator.h"
+#include "mongo/util/gcov.h"
#include "mongo/util/goodies.h"
#include "mongo/util/mongoutils/str.h"
+#include "mongo/util/time_support.h"
namespace mongo {
@@ -79,8 +106,6 @@ namespace mongo {
string dbExecCommand;
- bool useHints = true;
-
KillCurrentOp killCurrentOp;
int lockFile = 0;
@@ -90,41 +115,6 @@ namespace mongo {
MONGO_FP_DECLARE(rsStopGetMore);
- /*static*/ OpTime OpTime::_now() {
- OpTime result;
- unsigned t = (unsigned) time(0);
- if ( last.secs == t ) {
- last.i++;
- result = last;
- }
- else if ( t < last.secs ) {
- result = skewed(); // separate function to keep out of the hot code path
- }
- else {
- last = OpTime(t, 1);
- result = last;
- }
- notifier.notify_all();
- return last;
- }
- OpTime OpTime::now(const mongo::mutex::scoped_lock&) {
- return _now();
- }
- OpTime OpTime::getLast(const mongo::mutex::scoped_lock&) {
- return last;
- }
- boost::condition OpTime::notifier;
- mongo::mutex OpTime::m("optime");
-
- // OpTime::now() uses mutex, thus it is in this file not in the cpp files used by drivers and such
- void BSONElementManipulator::initTimestamp() {
- massert( 10332 , "Expected CurrentTime type", _element.type() == Timestamp );
- unsigned long long &timestamp = *( reinterpret_cast< unsigned long long* >( value() ) );
- if ( timestamp == 0 ) {
- mutex::scoped_lock lk(OpTime::m);
- timestamp = OpTime::now(lk).asDate();
- }
- }
void BSONElementManipulator::SetNumber(double d) {
if ( _element.type() == NumberDouble )
*getDur().writing( reinterpret_cast< double * >( value() ) ) = d;
@@ -140,34 +130,41 @@ namespace mongo {
verify( _element.type() == NumberInt );
getDur().writingInt( *reinterpret_cast< int * >( value() ) ) = n;
}
- /* dur:: version */
- void BSONElementManipulator::ReplaceTypeAndValue( const BSONElement &e ) {
- char *d = data();
- char *v = value();
- int valsize = e.valuesize();
- int ofs = (int) (v-d);
- dassert( ofs > 0 );
- char *p = (char *) getDur().writingPtr(d, valsize + ofs);
- *p = e.type();
- memcpy( p + ofs, e.value(), valsize );
- }
void inProgCmd( Message &m, DbResponse &dbresponse ) {
+ DbMessage d(m);
+ QueryMessage q(d);
BSONObjBuilder b;
- if (!cc().getAuthorizationManager()->checkAuthorization(
- AuthorizationManager::SERVER_RESOURCE_NAME, ActionType::inprog)) {
+ const bool isAuthorized = cc().getAuthorizationSession()->isAuthorizedForActionsOnResource(
+ ResourcePattern::forClusterResource(), ActionType::inprog);
+
+ audit::logInProgAuthzCheck(
+ &cc(), q.query, isAuthorized ? ErrorCodes::OK : ErrorCodes::Unauthorized);
+
+ if (!isAuthorized) {
b.append("err", "unauthorized");
}
else {
- DbMessage d(m);
- QueryMessage q(d);
bool all = q.query["$all"].trueValue();
vector<BSONObj> vals;
{
+ BSONObj filter;
+ {
+ BSONObjBuilder b;
+ BSONObjIterator i( q.query );
+ while ( i.more() ) {
+ BSONElement e = i.next();
+ if ( str::equals( "$all", e.fieldName() ) )
+ continue;
+ b.append( e );
+ }
+ filter = b.obj();
+ }
+
Client& me = cc();
scoped_lock bl(Client::clientsMutex);
- scoped_ptr<Matcher> m(new Matcher(q.query));
+ scoped_ptr<Matcher> m(new Matcher(filter));
for( set<Client*>::iterator i = Client::clients.begin(); i != Client::clients.end(); i++ ) {
Client *c = *i;
verify( c );
@@ -195,17 +192,21 @@ namespace mongo {
}
void killOp( Message &m, DbResponse &dbresponse ) {
+ DbMessage d(m);
+ QueryMessage q(d);
BSONObj obj;
- if (!cc().getAuthorizationManager()->checkAuthorization(
- AuthorizationManager::SERVER_RESOURCE_NAME, ActionType::killop)) {
+ const bool isAuthorized = cc().getAuthorizationSession()->isAuthorizedForActionsOnResource(
+ ResourcePattern::forClusterResource(), ActionType::killop);
+ audit::logKillOpAuthzCheck(&cc(),
+ q.query,
+ isAuthorized ? ErrorCodes::OK : ErrorCodes::Unauthorized);
+ if (!isAuthorized) {
obj = fromjson("{\"err\":\"unauthorized\"}");
}
/*else if( !dbMutexInfo.isLocked() )
obj = fromjson("{\"info\":\"no op in progress/not locked\"}");
*/
else {
- DbMessage d(m);
- QueryMessage q(d);
BSONElement e = q.query.getField("op");
if( !e.isNumber() ) {
obj = fromjson("{\"err\":\"no op number field specified?\"}");
@@ -222,8 +223,11 @@ namespace mongo {
bool _unlockFsync();
void unlockFsync(const char *ns, Message& m, DbResponse &dbresponse) {
BSONObj obj;
- if (!cc().getAuthorizationManager()->checkAuthorization(
- AuthorizationManager::SERVER_RESOURCE_NAME, ActionType::unlock)) {
+ const bool isAuthorized = cc().getAuthorizationSession()->isAuthorizedForActionsOnResource(
+ ResourcePattern::forClusterResource(), ActionType::unlock);
+ audit::logFsyncUnlockAuthzCheck(
+ &cc(), isAuthorized ? ErrorCodes::OK : ErrorCodes::Unauthorized);
+ if (!isAuthorized) {
obj = fromjson("{\"err\":\"unauthorized\"}");
}
else if (strncmp(ns, "admin.", 6) != 0 ) {
@@ -254,12 +258,15 @@ namespace mongo {
shared_ptr<AssertionException> ex;
try {
- if (!NamespaceString(d.getns()).isCommand()) {
+ NamespaceString ns(d.getns());
+ if (!ns.isCommand()) {
// Auth checking for Commands happens later.
- Status status = cc().getAuthorizationManager()->checkAuthForQuery(d.getns());
- uassert(16550, status.reason(), status.isOK());
+ Client* client = &cc();
+ Status status = client->getAuthorizationSession()->checkAuthForQuery(ns, q.query);
+ audit::logQueryAuthzCheck(client, ns, q.query, status.code());
+ uassertStatusOK(status);
}
- dbresponse.exhaustNS = runQuery(m, q, op, *resp);
+ dbresponse.exhaustNS = newRunQuery(m, q, op, *resp);
verify( !resp->empty() );
}
catch ( SendStaleConfigException& e ){
@@ -324,12 +331,13 @@ namespace mongo {
return ok;
}
+ // Mongod on win32 defines a value for this function. In all other executables it is NULL.
void (*reportEventToSystem)(const char *msg) = 0;
- void mongoAbort(const char *msg) {
- if( reportEventToSystem )
+ void mongoAbort(const char *msg) {
+ if( reportEventToSystem )
reportEventToSystem(msg);
- rawOut(msg);
+ severe() << msg;
::abort();
}
@@ -339,10 +347,19 @@ namespace mongo {
// before we lock...
int op = m.operation();
bool isCommand = false;
- const char *ns = m.singleData()->_data + 4;
+
+ DbMessage dbmsg(m);
+
+ Client& c = cc();
+ if (!c.isGod())
+ c.getAuthorizationSession()->startRequest();
+
+ c.setIsWriteCmd(false);
if ( op == dbQuery ) {
- if( strstr(ns, ".$cmd") ) {
+ const char *ns = dbmsg.getns();
+
+ if (strstr(ns, ".$cmd")) {
isCommand = true;
opwrite(m);
if( strstr(ns, ".$cmd.sys.") ) {
@@ -371,10 +388,31 @@ namespace mongo {
opwrite(m);
}
- globalOpCounters.gotOp( op , isCommand );
-
- Client& c = cc();
- c.getAuthorizationManager()->startRequest();
+ // Increment op counters.
+ switch (op) {
+ case dbQuery:
+ if (!isCommand) {
+ globalOpCounters.gotQuery();
+ }
+ else {
+ // Command counting is deferred, since it is not known yet whether the command
+ // needs counting.
+ }
+ break;
+ case dbGetMore:
+ globalOpCounters.gotGetMore();
+ break;
+ case dbInsert:
+ // Insert counting is deferred, since it is not known yet whether the insert contains
+ // multiple documents (each of which needs to be counted).
+ break;
+ case dbUpdate:
+ globalOpCounters.gotUpdate();
+ break;
+ case dbDelete:
+ globalOpCounters.gotDelete();
+ break;
+ }
auto_ptr<CurOp> nestedOp;
CurOp* currentOpP = c.curop();
@@ -392,8 +430,8 @@ namespace mongo {
OpDebug& debug = currentOp.debug();
debug.op = op;
- long long logThreshold = cmdLine.slowMS;
- bool shouldLog = logLevel >= 1;
+ long long logThreshold = serverGlobalParams.slowMS;
+ bool shouldLog = logger::globalLogDomain()->shouldLog(logger::LogSeverity::Debug(1));
if ( op == dbQuery ) {
if ( handlePossibleShardedMessage( m , &dbresponse ) )
@@ -406,7 +444,8 @@ namespace mongo {
}
else if ( op == dbMsg ) {
// deprecated - replaced by commands
- char *p = m.singleData()->_data;
+ const char *p = dbmsg.getns();
+
int len = strlen(p);
if ( len > 400 )
out() << curTimeMillis64() % 10000 <<
@@ -423,8 +462,6 @@ namespace mongo {
}
else {
try {
- const NamespaceString nsString( ns );
-
// The following operations all require authorization.
// dbInsert, dbUpdate and dbDelete can be easily pre-authorized,
// here, but dbKillCursors cannot.
@@ -433,32 +470,40 @@ namespace mongo {
logThreshold = 10;
receivedKillCursors(m);
}
- else if ( !nsString.isValid() ) {
- // Only killCursors doesn't care about namespaces
- uassert( 16257, str::stream() << "Invalid ns [" << ns << "]", false );
- }
- else if ( op == dbInsert ) {
- receivedInsert(m, currentOp);
- }
- else if ( op == dbUpdate ) {
- receivedUpdate(m, currentOp);
- }
- else if ( op == dbDelete ) {
- receivedDelete(m, currentOp);
- }
- else {
+ else if (op != dbInsert && op != dbUpdate && op != dbDelete) {
mongo::log() << " operation isn't supported: " << op << endl;
currentOp.done();
shouldLog = true;
}
- }
- catch ( UserException& ue ) {
- tlog(3) << " Caught Assertion in " << opToString(op) << ", continuing "
- << ue.toString() << endl;
+ else {
+ const char* ns = dbmsg.getns();
+ const NamespaceString nsString(ns);
+
+ if (!nsString.isValid()) {
+ uassert(16257, str::stream() << "Invalid ns [" << ns << "]", false);
+ }
+ else if (op == dbInsert) {
+ receivedInsert(m, currentOp);
+ }
+ else if (op == dbUpdate) {
+ receivedUpdate(m, currentOp);
+ }
+ else if (op == dbDelete) {
+ receivedDelete(m, currentOp);
+ }
+ else {
+ invariant(false);
+ }
+ }
+ }
+ catch (const UserException& ue) {
+ setLastError(ue.getCode(), ue.getInfo().msg.c_str());
+ LOG(3) << " Caught Assertion in " << opToString(op) << ", continuing "
+ << ue.toString() << endl;
debug.exceptionInfo = ue.getInfo();
}
catch ( AssertionException& e ) {
- tlog(3) << " Caught Assertion in " << opToString(op) << ", continuing "
+ MONGO_TLOG(3) << " Caught Assertion in " << opToString(op) << ", continuing "
<< e.toString() << endl;
debug.exceptionInfo = e.getInfo();
shouldLog = true;
@@ -471,7 +516,7 @@ namespace mongo {
logThreshold += currentOp.getExpectedLatencyMs();
if ( shouldLog || debug.executionTime > logThreshold ) {
- mongo::tlog() << debug.report( currentOp ) << endl;
+ MONGO_TLOG(0) << debug.report( currentOp ) << endl;
}
if ( currentOp.shouldDBProfile( debug.executionTime ) ) {
@@ -492,22 +537,23 @@ namespace mongo {
} /* assembleResponse() */
void receivedKillCursors(Message& m) {
- int *x = (int *) m.singleData()->_data;
- x++; // reserved
- int n = *x++;
+ DbMessage dbmessage(m);
+ int n = dbmessage.pullInt();
uassert( 13659 , "sent 0 cursors to kill" , n != 0 );
massert( 13658 , str::stream() << "bad kill cursors size: " << m.dataSize() , m.dataSize() == 8 + ( 8 * n ) );
uassert( 13004 , str::stream() << "sent negative cursors to kill: " << n , n >= 1 );
if ( n > 2000 ) {
- LOG( n < 30000 ? LL_WARNING : LL_ERROR ) << "receivedKillCursors, n=" << n << endl;
+ ( n < 30000 ? warning() : error() ) << "receivedKillCursors, n=" << n << endl;
verify( n < 30000 );
}
- int found = ClientCursor::eraseIfAuthorized(n, (long long *) x);
+ const long long* cursorArray = dbmessage.getArray(n);
- if ( logLevel > 0 || found != n ) {
+ int found = CollectionCursorCache::eraseCursorGlobalIfAuthorized(n, cursorArray);
+
+ if ( logger::globalLogDomain()->shouldLog(logger::LogSeverity::Debug(1)) || found != n ) {
LOG( found == n ? 1 : 0 ) << "killcursors: found " << found << " of " << n << endl;
}
@@ -516,14 +562,14 @@ namespace mongo {
/* db - database name
path - db directory
*/
- /*static*/ void Database::closeDatabase( const char *db, const string& path ) {
+ /*static*/ void Database::closeDatabase( const string& db, const string& path ) {
verify( Lock::isW() );
Client::Context * ctx = cc().getContext();
verify( ctx );
verify( ctx->inDB( db , path ) );
Database *database = ctx->db();
- verify( database->name == db );
+ verify( database->name() == db );
oplogCheckCloseDatabase( database ); // oplog caches some things, dirty its caches
@@ -534,9 +580,6 @@ namespace mongo {
/* important: kill all open cursors on the database */
string prefix(db);
prefix += '.';
- ClientCursor::invalidate(prefix.c_str());
-
- NamespaceDetailsTransient::eraseDB( prefix );
dbHolderW().erase( db, path );
ctx->_clear();
@@ -545,8 +588,9 @@ namespace mongo {
void receivedUpdate(Message& m, CurOp& op) {
DbMessage d(m);
- const char *ns = d.getns();
- op.debug().ns = ns;
+ NamespaceString ns(d.getns());
+ uassertStatusOK( userAllowedWriteNS( ns ) );
+ op.debug().ns = ns.ns();
int flags = d.pullInt();
BSONObj query = d.nextJsObj();
@@ -560,90 +604,96 @@ namespace mongo {
bool multi = flags & UpdateOption_Multi;
bool broadcast = flags & UpdateOption_Broadcast;
- Status status = cc().getAuthorizationManager()->checkAuthForUpdate(ns, upsert);
- uassert(16538, status.reason(), status.isOK());
+ Status status = cc().getAuthorizationSession()->checkAuthForUpdate(ns,
+ query,
+ toupdate,
+ upsert);
+ audit::logUpdateAuthzCheck(&cc(), ns, query, toupdate, upsert, multi, status.code());
+ uassertStatusOK(status);
op.debug().query = query;
op.setQuery(query);
- PageFaultRetryableSection s;
- while ( 1 ) {
- try {
- Lock::DBWrite lk(ns);
-
- // void ReplSetImpl::relinquish() uses big write lock so
- // this is thus synchronized given our lock above.
- uassert( 10054 , "not master", isMasterNs( ns ) );
-
- // if this ever moves to outside of lock, need to adjust check Client::Context::_finishInit
- if ( ! broadcast && handlePossibleShardedMessage( m , 0 ) )
- return;
-
- Client::Context ctx( ns );
-
- UpdateResult res = updateObjects(ns, toupdate, query, upsert, multi, true, op.debug() );
- lastError.getSafe()->recordUpdate( res.existing , res.num , res.upserted ); // for getlasterror
- break;
- }
- catch ( PageFaultException& e ) {
- e.touch();
- }
- }
+ UpdateRequest request(ns);
+
+ request.setUpsert(upsert);
+ request.setMulti(multi);
+ request.setQuery(query);
+ request.setUpdates(toupdate);
+ request.setUpdateOpLog(); // TODO: This is wasteful if repl is not active.
+ UpdateLifecycleImpl updateLifecycle(broadcast, ns);
+ request.setLifecycle(&updateLifecycle);
+ UpdateExecutor executor(&request, &op.debug());
+ uassertStatusOK(executor.prepare());
+
+ Lock::DBWrite lk(ns.ns());
+
+ // if this ever moves to outside of lock, need to adjust check
+ // Client::Context::_finishInit
+ if ( ! broadcast && handlePossibleShardedMessage( m , 0 ) )
+ return;
+
+ Client::Context ctx( ns );
+
+ UpdateResult res = executor.execute();
+
+ // for getlasterror
+ lastError.getSafe()->recordUpdate( res.existing , res.numMatched , res.upserted );
}
void receivedDelete(Message& m, CurOp& op) {
DbMessage d(m);
- const char *ns = d.getns();
-
- Status status = cc().getAuthorizationManager()->checkAuthForDelete(ns);
- uassert(16542, status.reason(), status.isOK());
+ NamespaceString ns(d.getns());
+ uassertStatusOK( userAllowedWriteNS( ns ) );
- op.debug().ns = ns;
+ op.debug().ns = ns.ns();
int flags = d.pullInt();
bool justOne = flags & RemoveOption_JustOne;
bool broadcast = flags & RemoveOption_Broadcast;
verify( d.moreJSObjs() );
BSONObj pattern = d.nextJsObj();
-
+
+ Status status = cc().getAuthorizationSession()->checkAuthForDelete(ns, pattern);
+ audit::logDeleteAuthzCheck(&cc(), ns, pattern, status.code());
+ uassertStatusOK(status);
+
op.debug().query = pattern;
op.setQuery(pattern);
- PageFaultRetryableSection s;
- while ( 1 ) {
- try {
- Lock::DBWrite lk(ns);
-
- // writelock is used to synchronize stepdowns w/ writes
- uassert( 10056 , "not master", isMasterNs( ns ) );
-
- // if this ever moves to outside of lock, need to adjust check Client::Context::_finishInit
- if ( ! broadcast && handlePossibleShardedMessage( m , 0 ) )
- return;
-
- Client::Context ctx(ns);
-
- long long n = deleteObjects(ns, pattern, justOne, true);
- lastError.getSafe()->recordDelete( n );
- op.debug().ndeleted = n;
- break;
- }
- catch ( PageFaultException& e ) {
- LOG(2) << "recordDelete got a PageFaultException" << endl;
- e.touch();
+ {
+ PageFaultRetryableSection s;
+ while ( 1 ) {
+ try {
+ DeleteRequest request(ns);
+ request.setQuery(pattern);
+ request.setMulti(!justOne);
+ request.setUpdateOpLog(true);
+ DeleteExecutor executor(&request);
+ uassertStatusOK(executor.prepare());
+ Lock::DBWrite lk(ns.ns());
+
+ // if this ever moves to outside of lock, need to adjust check
+ // Client::Context::_finishInit
+ if ( ! broadcast && handlePossibleShardedMessage( m , 0 ) )
+ return;
+
+ Client::Context ctx(ns);
+
+ long long n = executor.execute();
+ lastError.getSafe()->recordDelete( n );
+ op.debug().ndeleted = n;
+ break;
+ }
+ catch ( PageFaultException& e ) {
+ LOG(2) << "recordDelete got a PageFaultException" << endl;
+ e.touch();
+ }
}
- }
+ } // end PageFaultRetryableSection
}
QueryResult* emptyMoreResult(long long);
- void OpTime::waitForDifferent(unsigned millis){
- mutex::scoped_lock lk(m);
- while (*this == last) {
- if (!notifier.timed_wait(lk.boost(), boost::posix_time::milliseconds(millis)))
- return; // timed out
- }
- }
-
bool receivedGetMore(DbResponse& dbresponse, Message& m, CurOp& curop ) {
bool ok = true;
@@ -669,8 +719,10 @@ namespace mongo {
const NamespaceString nsString( ns );
uassert( 16258, str::stream() << "Invalid ns [" << ns << "]", nsString.isValid() );
- Status status = cc().getAuthorizationManager()->checkAuthForGetMore(ns);
- uassert(16543, status.reason(), status.isOK());
+ Status status = cc().getAuthorizationSession()->checkAuthForGetMore(
+ nsString, cursorid);
+ audit::logGetMoreAuthzCheck(&cc(), nsString, cursorid, status.code());
+ uassertStatusOK(status);
if (str::startsWith(ns, "local.oplog.")){
while (MONGO_FAIL_POINT(rsStopGetMore)) {
@@ -686,13 +738,13 @@ namespace mongo {
}
}
- msgdata = processGetMore(ns,
- ntoreturn,
- cursorid,
- curop,
- pass,
- exhaust,
- &isCursorAuthorized);
+ msgdata = newGetMore(ns,
+ ntoreturn,
+ cursorid,
+ curop,
+ pass,
+ exhaust,
+ &isCursorAuthorized);
}
catch ( AssertionException& e ) {
if ( isCursorAuthorized ) {
@@ -701,7 +753,7 @@ namespace mongo {
// because it may now be out of sync with the client's iteration state.
// SERVER-7952
// TODO Temporary code, see SERVER-4563 for a cleanup overview.
- ClientCursor::erase( cursorid );
+ CollectionCursorCache::eraseCursorGlobal( cursorid );
}
ex.reset( new AssertionException( e.getInfo().msg, e.getCode() ) );
ok = false;
@@ -738,24 +790,16 @@ namespace mongo {
};
if (ex) {
- exhaust = false;
-
BSONObjBuilder err;
ex->getInfo().append( err );
BSONObj errObj = err.done();
- log() << errObj << endl;
-
curop.debug().exceptionInfo = ex->getInfo();
- if (ex->getCode() == 13436) {
- replyToQuery(ResultFlag_ErrSet, m, dbresponse, errObj);
- curop.debug().responseLength = dbresponse.response->header()->dataLen();
- curop.debug().nreturned = 1;
- return ok;
- }
-
- msgdata = emptyMoreResult(cursorid);
+ replyToQuery(ResultFlag_ErrSet, m, dbresponse, errObj);
+ curop.debug().responseLength = dbresponse.response->header()->dataLen();
+ curop.debug().nreturned = 1;
+ return ok;
}
Message *resp = new Message();
@@ -774,41 +818,57 @@ namespace mongo {
return ok;
}
- void checkAndInsert(const char *ns, /*modifies*/BSONObj& js) {
- uassert( 10059 , "object to insert too large", js.objsize() <= BSONObjMaxUserSize);
- {
- BSONObjIterator i( js );
- while ( i.more() ) {
- BSONElement e = i.next();
-
- // check no $ modifiers. note we only check top level.
- // (scanning deep would be quite expensive)
- uassert( 13511, "document to insert can't have $ fields", e.fieldName()[0] != '$' );
-
- // check no regexp for _id (SERVER-9502)
- if (str::equals(e.fieldName(), "_id")) {
- uassert(16824, "can't use a regex for _id", e.type() != RegEx);
- }
+ void checkAndInsert(Client::Context& ctx, const char *ns, /*modifies*/BSONObj& js,
+ PregeneratedKeys* preGen ) {
+ if ( nsToCollectionSubstring( ns ) == "system.indexes" ) {
+ string targetNS = js["ns"].String();
+ uassertStatusOK( userAllowedWriteNS( targetNS ) );
+
+ Collection* collection = ctx.db()->getCollection( targetNS );
+ if ( !collection ) {
+ // implicitly create
+ collection = ctx.db()->createCollection( targetNS );
+ verify( collection );
}
+
+ // Only permit interrupting an (index build) insert if the
+ // insert comes from a socket client request rather than a
+ // parent operation using the client interface. The parent
+ // operation might not support interrupts.
+ bool mayInterrupt = cc().curop()->parent() == NULL;
+
+ cc().curop()->setQuery(js);
+ Status status = collection->getIndexCatalog()->createIndex( js, mayInterrupt );
+
+ if ( status.code() == ErrorCodes::IndexAlreadyExists )
+ return;
+
+ uassertStatusOK( status );
+ logOp( "i", ns, js );
+ return;
+ }
+
+ StatusWith<BSONObj> fixed = fixDocumentForInsert( js );
+ uassertStatusOK( fixed.getStatus() );
+ if ( !fixed.getValue().isEmpty() )
+ js = fixed.getValue();
+
+ Collection* collection = ctx.db()->getCollection( ns );
+ if ( !collection ) {
+ collection = ctx.db()->createCollection( ns );
+ verify( collection );
}
- theDataFileMgr.insertWithObjMod(ns,
- // May be modified in the call to add an _id field.
- js,
- // Only permit interrupting an (index build) insert if the
- // insert comes from a socket client request rather than a
- // parent operation using the client interface. The parent
- // operation might not support interrupts.
- cc().curop()->parent() == NULL,
- false);
+ StatusWith<DiskLoc> status = collection->insertDocument( js, true, preGen );
+ uassertStatusOK( status.getStatus() );
logOp("i", ns, js);
}
- NOINLINE_DECL void insertMulti(bool keepGoing, const char *ns, vector<BSONObj>& objs, CurOp& op) {
+ NOINLINE_DECL void insertMulti(Client::Context& ctx, bool keepGoing, const char *ns, vector<BSONObj>& objs, CurOp& op) {
size_t i;
for (i=0; i<objs.size(); i++){
try {
- checkAndInsert(ns, objs[i]);
+ checkAndInsert(ctx, ns, objs[i], NULL);
getDur().commitIfNeeded();
} catch (const UserException&) {
if (!keepGoing || i == objs.size()-1){
@@ -828,13 +888,7 @@ namespace mongo {
const char *ns = d.getns();
op.debug().ns = ns;
- bool isIndexWrite = NamespaceString(ns).coll == "system.indexes";
-
- // Auth checking for index writes happens further down in this function.
- if (!isIndexWrite) {
- Status status = cc().getAuthorizationManager()->checkAuthForInsert(ns);
- uassert(16544, status.reason(), status.isOK());
- }
+ uassertStatusOK( userAllowedWriteNS( ns ) );
if( !d.moreJSObjs() ) {
// strange. should we complain?
@@ -845,51 +899,63 @@ namespace mongo {
while (d.moreJSObjs()){
BSONObj obj = d.nextJsObj();
multi.push_back(obj);
- if (isIndexWrite) {
- string indexNS = obj.getStringField("ns");
- uassert(16548,
- mongoutils::str::stream() << "not authorized to create index on "
- << indexNS,
- cc().getAuthorizationManager()->checkAuthorization(
- indexNS, ActionType::ensureIndex));
- }
+
+ // Check auth for insert (also handles checking if this is an index build and checks
+ // for the proper privileges in that case).
+ const NamespaceString nsString(ns);
+ Status status = cc().getAuthorizationSession()->checkAuthForInsert(nsString, obj);
+ audit::logInsertAuthzCheck(&cc(), nsString, obj, status.code());
+ uassertStatusOK(status);
}
- PageFaultRetryableSection s;
- while ( true ) {
- try {
- Lock::DBWrite lk(ns);
-
- // CONCURRENCY TODO: is being read locked in big log sufficient here?
- // writelock is used to synchronize stepdowns w/ writes
- uassert( 10058 , "not master", isMasterNs(ns) );
-
- if ( handlePossibleShardedMessage( m , 0 ) )
+ PregeneratedKeys tempHack;
+ if ( multi.size() == 1 ) {
+ StatusWith<BSONObj> fixed = fixDocumentForInsert( multi[0] );
+ uassertStatusOK( fixed.getStatus() );
+ if ( !fixed.getValue().isEmpty() )
+ multi[0] = fixed.getValue();
+
+ GeneratorHolder::getInstance()->prepare( ns, multi[0], &tempHack );
+ }
+
+ {
+ PageFaultRetryableSection s;
+ while ( true ) {
+ try {
+ Lock::DBWrite lk(ns);
+
+ // CONCURRENCY TODO: is being read locked in big log sufficient here?
+ // writelock is used to synchronize stepdowns w/ writes
+ uassert( 10058 , "not master", isMasterNs(ns) );
+
+ if ( handlePossibleShardedMessage( m , 0 ) )
+ return;
+
+ Client::Context ctx(ns);
+
+ if (multi.size() > 1) {
+ const bool keepGoing = d.reservedField() & InsertOption_ContinueOnError;
+ insertMulti(ctx, keepGoing, ns, multi, op);
+ }
+ else {
+ checkAndInsert(ctx, ns, multi[0], &tempHack);
+ globalOpCounters.incInsertInWriteLock(1);
+ op.debug().ninserted = 1;
+ }
return;
-
- Client::Context ctx(ns);
-
- if (multi.size() > 1) {
- const bool keepGoing = d.reservedField() & InsertOption_ContinueOnError;
- insertMulti(keepGoing, ns, multi, op);
- } else {
- checkAndInsert(ns, multi[0]);
- globalOpCounters.incInsertInWriteLock(1);
- op.debug().ninserted = 1;
}
- return;
- }
- catch ( PageFaultException& e ) {
- e.touch();
+ catch ( PageFaultException& e ) {
+ e.touch();
+ }
}
- }
+ } // end PageFaultRetryableSection
}
void getDatabaseNames( vector< string > &names , const string& usePath ) {
boost::filesystem::path path( usePath );
for ( boost::filesystem::directory_iterator i( path );
i != boost::filesystem::directory_iterator(); ++i ) {
- if ( directoryperdb ) {
+ if (storageGlobalParams.directoryperdb) {
boost::filesystem::path p = *i;
string dbName = p.leaf().string();
p /= ( dbName + ".ns" );
@@ -931,7 +997,21 @@ namespace mongo {
return QueryOptions(DBClientBase::_lookupAvailableOptions() & ~QueryOption_Exhaust);
}
+namespace {
+ class GodScope {
+ MONGO_DISALLOW_COPYING(GodScope);
+ public:
+ GodScope() {
+ _prev = cc().setGod(true);
+ }
+ ~GodScope() { cc().setGod(_prev); }
+ private:
+ bool _prev;
+ };
+} // namespace
+
bool DBDirectClient::call( Message &toSend, Message &response, bool assertOk , string * actualServer ) {
+ GodScope gs;
if ( lastError._get() )
lastError.startRequest( toSend, lastError._get() );
DbResponse dbResponse;
@@ -944,6 +1024,7 @@ namespace mongo {
}
void DBDirectClient::say( Message &toSend, bool isRetry, string * actualServer ) {
+ GodScope gs;
if ( lastError._get() )
lastError.startRequest( toSend, lastError._get() );
DbResponse dbResponse;
@@ -962,7 +1043,7 @@ namespace mongo {
}
void DBDirectClient::killCursor( long long id ) {
- ClientCursor::erase( id );
+ CollectionCursorCache::eraseCursorGlobal( id );
}
HostAndPort DBDirectClient::_clientHost = HostAndPort( "0.0.0.0" , 0 );
@@ -976,7 +1057,7 @@ namespace mongo {
Lock::DBRead lk( ns );
string errmsg;
int errCode;
- long long res = runCount( ns.c_str() , _countCmd( ns , query , options , limit , skip ) , errmsg, errCode );
+ long long res = runCount( ns, _countCmd( ns , query , options , limit , skip ) , errmsg, errCode );
if ( res == -1 ) {
// namespace doesn't exist
return 0;
@@ -989,6 +1070,14 @@ namespace mongo {
return new DBDirectClient();
}
+ MONGO_INITIALIZER(CreateJSDirectClient)
+ (InitializerContext* context) {
+
+ directDBClient = createDirectClient();
+
+ return Status::OK();
+ }
+
mongo::mutex exitMutex("exit");
AtomicUInt numExitCalls = 0;
@@ -996,22 +1085,6 @@ namespace mongo {
return numExitCalls > 0;
}
- void tryToOutputFatal( const string& s ) {
- try {
- rawOut( s );
- return;
- }
- catch ( ... ) {}
-
- try {
- cerr << s << endl;
- return;
- }
- catch ( ... ) {}
-
- // uh - oh, not sure there is anything else we can do...
- }
-
static void shutdownServer() {
log() << "shutdown: going to close listening sockets..." << endl;
@@ -1030,7 +1103,7 @@ namespace mongo {
log() << "shutdown: waiting for fs preallocator..." << endl;
FileAllocator::get()->waitUntilFinished();
- if( cmdLine.dur ) {
+ if (storageGlobalParams.dur) {
log() << "shutdown: lock for final commit..." << endl;
{
int n = 10;
@@ -1058,7 +1131,7 @@ namespace mongo {
MemoryMappedFile::closeAllFiles( ss3 );
log() << ss3.str() << endl;
- if( cmdLine.dur ) {
+ if (storageGlobalParams.dur) {
dur::journalCleanup(true);
}
@@ -1097,7 +1170,10 @@ namespace mongo {
/* not using log() herein in case we are already locked */
NOINLINE_DECL void dbexit( ExitCode rc, const char *why ) {
+ flushForGcov();
+
Client * c = currentClient.get();
+ audit::logShutdown(c);
{
scoped_lock lk( exitMutex );
if ( numExitCalls++ > 0 ) {
@@ -1105,25 +1181,19 @@ namespace mongo {
// this means something horrible has happened
::_exit( rc );
}
- stringstream ss;
- ss << "dbexit: " << why << "; exiting immediately";
- tryToOutputFatal( ss.str() );
+ log() << "dbexit: " << why << "; exiting immediately";
if ( c ) c->shutdown();
::_exit( rc );
}
}
- {
- stringstream ss;
- ss << "dbexit: " << why;
- tryToOutputFatal( ss.str() );
- }
+ log() << "dbexit: " << why;
try {
shutdownServer(); // gracefully shutdown instance
}
catch ( ... ) {
- tryToOutputFatal( "shutdown failed with exception" );
+ severe() << "shutdown failed with exception";
}
#if defined(_DEBUG)
@@ -1146,7 +1216,7 @@ namespace mongo {
return;
}
#endif
- tryToOutputFatal( "dbexit: really exiting now" );
+ log() << "dbexit: really exiting now";
if ( c ) c->shutdown();
::_exit(rc);
}
@@ -1154,7 +1224,7 @@ namespace mongo {
#if !defined(__sunos__)
void writePid(int fd) {
stringstream ss;
- ss << getpid() << endl;
+ ss << ProcessId::getCurrent() << endl;
string s = ss.str();
const char * data = s.c_str();
#ifdef _WIN32
@@ -1165,7 +1235,7 @@ namespace mongo {
}
void acquirePathLock(bool doingRepair) {
- string name = ( boost::filesystem::path( dbpath ) / "mongod.lock" ).string();
+ string name = (boost::filesystem::path(storageGlobalParams.dbpath) / "mongod.lock").string();
bool oldFile = false;
@@ -1214,7 +1284,7 @@ namespace mongo {
"run with --repair again.\n"
"**************";
}
- else if (cmdLine.dur) {
+ else if (storageGlobalParams.dur) {
if (!dur::haveJournalFiles(/*anyFiles=*/true)) {
// Passing anyFiles=true as we are trying to protect against starting in an
// unclean state with the journal directory unmounted. If there are any files,
@@ -1242,7 +1312,6 @@ namespace mongo {
<< "*************";
}
-
}
}
else {
@@ -1268,7 +1337,7 @@ namespace mongo {
}
// Not related to lock file, but this is where we handle unclean shutdown
- if( !cmdLine.dur && dur::haveJournalFiles() ) {
+ if (!storageGlobalParams.dur && dur::haveJournalFiles()) {
cout << "**************" << endl;
cout << "Error: journal files are present in journal directory, yet starting without journaling enabled." << endl;
cout << "It is recommended that you start with journaling enabled so that recovery may occur." << endl;
@@ -1292,7 +1361,7 @@ namespace mongo {
// TODO - this is very bad that the code above not running here.
// Not related to lock file, but this is where we handle unclean shutdown
- if( !cmdLine.dur && dur::haveJournalFiles() ) {
+ if (!storageGlobalParams.dur && dur::haveJournalFiles()) {
cout << "**************" << endl;
cout << "Error: journal files are present in journal directory, yet starting without --journal enabled." << endl;
cout << "It is recommended that you start with journaling enabled so that recovery may occur." << endl;
@@ -1310,7 +1379,7 @@ namespace mongo {
void DiagLog::openFile() {
verify( f == 0 );
stringstream ss;
- ss << dbpath << "/diaglog." << hex << time(0);
+ ss << storageGlobalParams.dbpath << "/diaglog." << hex << time(0);
string name = ss.str();
f = new ofstream(name.c_str(), ios::out | ios::binary);
if ( ! f->good() ) {