| 110 | } |
| 111 | |
| 112 | int rocksdb_wrapper::write_batch_put_ctx(const db_write_context &ctx, |
| 113 | dsn::string_view raw_key, |
| 114 | dsn::string_view value, |
| 115 | uint32_t expire_sec) |
| 116 | { |
| 117 | FAIL_POINT_INJECT_F("db_write_batch_put", |
| 118 | [](dsn::string_view) -> int { return FAIL_DB_WRITE_BATCH_PUT; }); |
| 119 | |
| 120 | uint64_t new_timetag = ctx.remote_timetag; |
| 121 | if (!ctx.is_duplicated_write()) { // local write |
| 122 | new_timetag = generate_timetag(ctx.timestamp, get_cluster_id_if_exists(), false); |
| 123 | } |
| 124 | |
| 125 | if (ctx.verify_timetag && // needs read-before-write |
| 126 | _pegasus_data_version >= 1 && // data version 0 doesn't support timetag. |
| 127 | !raw_key.empty()) { // not an empty write |
| 128 | |
| 129 | db_get_context get_ctx; |
| 130 | int err = get(raw_key, &get_ctx); |
| 131 | if (dsn_unlikely(err != rocksdb::Status::kOk)) { |
| 132 | return err; |
| 133 | } |
| 134 | // if record exists and is not expired. |
| 135 | if (get_ctx.found && !get_ctx.expired) { |
| 136 | uint64_t local_timetag = |
| 137 | pegasus_extract_timetag(_pegasus_data_version, get_ctx.raw_value); |
| 138 | |
| 139 | if (local_timetag >= new_timetag) { |
| 140 | // ignore this stale update with lower timetag, |
| 141 | // and write an empty record instead |
| 142 | raw_key = value = dsn::string_view(); |
| 143 | } |
| 144 | } |
| 145 | } |
| 146 | |
| 147 | rocksdb::Slice skey = utils::to_rocksdb_slice(raw_key); |
| 148 | rocksdb::SliceParts skey_parts(&skey, 1); |
| 149 | rocksdb::SliceParts svalue = _value_generator->generate_value( |
| 150 | _pegasus_data_version, value, db_expire_ts(expire_sec), new_timetag); |
| 151 | rocksdb::Status s = _write_batch->Put(skey_parts, svalue); |
| 152 | if (dsn_unlikely(!s.ok())) { |
| 153 | ::dsn::blob hash_key, sort_key; |
| 154 | pegasus_restore_key(::dsn::blob(raw_key.data(), 0, raw_key.size()), hash_key, sort_key); |
| 155 | LOG_ERROR_ROCKSDB("WriteBatchPut", |
| 156 | s.ToString(), |
| 157 | "decree: {}, hash_key: {}, sort_key: {}, expire_ts: {}", |
| 158 | ctx.decree, |
| 159 | utils::c_escape_string(hash_key), |
| 160 | utils::c_escape_string(sort_key), |
| 161 | expire_sec); |
| 162 | } |
| 163 | return s.code(); |
| 164 | } |
| 165 | |
| 166 | int rocksdb_wrapper::write(int64_t decree) |
| 167 | { |