diff --git a/daemon/Makefile b/daemon/Makefile index 1294871a4..9e8722a43 100644 --- a/daemon/Makefile +++ b/daemon/Makefile @@ -80,7 +80,7 @@ endif CFLAGS += $(CFLAGS_MQTT) LDLIBS += $(LDLIBS_MQTT) -SRCS := main.c kernel.c helpers.c control_tcp.c call.c control_udp.c redis.c \ +SRCS := main.c kernel.c helpers.c control_tcp.c call.c call_checkpoint.c control_udp.c redis.c \ cookie_cache.c udp_listener.c control_ng_flags_parser.c control_ng.c sdp.strhash.c stun.c rtcp.c \ crypto.c rtp.c call_interfaces.c dtls.c log.c cli.strhash.c graphite.c ice.c \ media_socket.c homer.c recording.c statistics.c cdr.c ssrc.c iptables.c tcp_listener.c \ diff --git a/daemon/call.c b/daemon/call.c index 25adfb8d7..1092b3a52 100644 --- a/daemon/call.c +++ b/daemon/call.c @@ -1,4 +1,5 @@ #include "call.h" +#include "call_checkpoint.h" #include #include @@ -69,7 +70,6 @@ static int64_t add_ongoing_calls_dur_in_interval(int64_t interval_start, int64_t static void __call_free(call_t *p); static void __call_cleanup(call_t *c); static void __monologue_stop(struct call_monologue *ml); -static void media_stop(struct call_media *m); __attribute__((nonnull(1, 2, 4))) static struct media_subscription *__subscribe_medias_both_ways(struct call_media * a, struct call_media * b, bool is_offer, medias_q *); @@ -5182,7 +5182,7 @@ static void __call_cleanup(call_t *c) { for (__auto_type l = c->medias.head; l; l = l->next) { struct call_media *md = l->data; ice_shutdown(&md->ice_agent); - media_stop(md); + call_media_stop(md); t38_gateway_put(&md->t38_gateway); audio_player_free(md); mutex_destroy(&md->dtmf_lock); @@ -5498,6 +5498,7 @@ static void __call_free(call_t *c) { //ilog(LOG_DEBUG, "freeing main call struct"); + call_checkpoint_free_all(c); obj_release(c->dtls_cert); mqtt_timer_stop(&c->mqtt_timer); @@ -6418,7 +6419,7 @@ int call_get_mono_dialogue(struct call_monologue *monologues[2], return call_get_dialogue(monologues, call, callid, fromtag, totag, viabranch, flags, ep); } -static void media_stop(struct call_media *m) { +void call_media_stop(struct call_media *m) { if (!m) return; t38_gateway_stop(m->t38_gateway); @@ -6443,7 +6444,7 @@ static void monologue_stop(struct call_monologue *ml, bool stop_media_subscriber __monologue_stop(ml); for (unsigned int i = 0; i < ml->medias->len; i++) { - media_stop(ml->medias->pdata[i]); + call_media_stop(ml->medias->pdata[i]); } /* monologue's subscribers */ if (stop_media_subscribers) { @@ -6453,7 +6454,7 @@ static void monologue_stop(struct call_monologue *ml, bool stop_media_subscriber if (!media) continue; IQUEUE_FOREACH(&media->media_subscribers, ms) { - media_stop(ms->media); + call_media_stop(ms->media); __monologue_stop(ms->monologue); } } diff --git a/daemon/call_checkpoint.c b/daemon/call_checkpoint.c new file mode 100644 index 000000000..3ed5e1134 --- /dev/null +++ b/daemon/call_checkpoint.c @@ -0,0 +1,1206 @@ +#include "call_checkpoint.h" + +#include "call.h" +#include "codec.h" +#include "dtls.h" +#include "ice.h" +#include "log_d.h" +#include "main.h" +#include "redis.h" + +#include + +struct checkpoint_stream { + struct checkpoint_stream *next; + struct packet_stream *stream; + stream_fd_q sfds; + GArray *sfd_local_endpoints; + stream_fd *selected_sfd; + bool selected_sfd_set; + bool redis_restored; + struct endpoint endpoint; + struct endpoint advertised_endpoint; + struct endpoint learned_endpoint; + struct endpoint detected_endpoints[4]; + endpoint_t last_local_endpoint; + int64_t ep_detect_signal; + enum endpoint_learning el_flags; + uint64_t flags; +}; + +struct checkpoint_media { + struct checkpoint_media *next; + struct call_media *media; + struct endpoint_map *endpoint_map; + const struct transport_protocol *protocol; + str protocol_str; + str format_str; + sockfamily_t *desired_family; + struct logical_intf *logical_intf; + uint64_t flags; + sdes_q sdes_in; + sdes_q sdes_out; + struct dtls_fingerprint fingerprint; + const struct dtls_hash_func *fp_hash_func; + str tls_id; + struct codec_store codecs; + struct codec_store offered_codecs; + candidate_q ice_candidates; + str ice_ufrag[2]; + str ice_pwd[2]; + bool had_ice; + int ptime; + int maxptime; + struct session_bandwidth bandwidth; + struct checkpoint_stream *streams; +}; + +struct checkpoint_monologue { + struct call_monologue *monologue; + sockfamily_t *desired_family; + struct logical_intf *logical_intf; + uint64_t flags; + struct session_bandwidth bandwidth; + sdp_origin sdp_orig_in; + sdp_origin sdp_orig_out; + str session_name; + str session_timing; + GString *last_out_sdp; + unsigned int medias_len; +}; + +struct call_checkpoint { + struct call_checkpoint *next; + struct call_monologue *offerer; + struct call_monologue *answerer; + uint64_t generation; + uint64_t pending_generation; + bool pending; + struct checkpoint_monologue monologues[2]; + struct checkpoint_media *medias; +}; + +static void checkpoint_stream_free(struct checkpoint_stream *); +static void checkpoint_media_free(struct checkpoint_media *); +static void checkpoint_free(struct call_checkpoint *); + +static void checkpoint_json_str(JsonBuilder *b, const char *name, const str *value) { + json_builder_set_member_name(b, name); + if (!value || !value->s) { + json_builder_add_string_value(b, ""); + return; + } + g_autofree char *string = g_strndup(value->s, value->len); + json_builder_add_string_value(b, string); +} + +static void checkpoint_json_endpoint(JsonBuilder *b, const char *name, const endpoint_t *ep) { + json_builder_set_member_name(b, name); + json_builder_add_string_value(b, ep->address.family ? endpoint_print_buf(ep) : ""); +} + +static void checkpoint_json_bandwidth(JsonBuilder *b, const struct session_bandwidth *bw) { + json_builder_set_member_name(b, "bandwidth"); + json_builder_begin_object(b); + json_builder_set_member_name(b, "as"); + json_builder_add_int_value(b, bw->as); + json_builder_set_member_name(b, "ct"); + json_builder_add_int_value(b, bw->ct); + json_builder_set_member_name(b, "rr"); + json_builder_add_int_value(b, bw->rr); + json_builder_set_member_name(b, "rs"); + json_builder_add_int_value(b, bw->rs); + json_builder_set_member_name(b, "tias"); + json_builder_add_int_value(b, bw->tias); + json_builder_end_object(b); +} + +static void checkpoint_json_origin(JsonBuilder *b, const char *name, const sdp_origin *o) { + json_builder_set_member_name(b, name); + json_builder_begin_object(b); + json_builder_set_member_name(b, "parsed"); + json_builder_add_boolean_value(b, o->parsed); + json_builder_set_member_name(b, "version"); + json_builder_add_int_value(b, o->version_num); + checkpoint_json_str(b, "username", &o->username); + checkpoint_json_str(b, "session-id", &o->session_id); + checkpoint_json_str(b, "network-type", &o->address.network_type); + checkpoint_json_str(b, "address-type", &o->address.address_type); + checkpoint_json_str(b, "address", &o->address.address); + json_builder_end_object(b); +} + +static void checkpoint_json_stream(JsonBuilder *b, const struct checkpoint_stream *s) { + json_builder_begin_object(b); + json_builder_set_member_name(b, "stream-id"); + json_builder_add_int_value(b, s->stream->unique_id); + json_builder_set_member_name(b, "selected-sfd-id"); + json_builder_add_int_value(b, s->selected_sfd ? s->selected_sfd->unique_id : -1); + json_builder_set_member_name(b, "selected-interface-id"); + json_builder_add_int_value(b, s->selected_sfd ? s->selected_sfd->local_intf->unique_id : -1); + json_builder_set_member_name(b, "flags"); + json_builder_add_int_value(b, s->flags); + json_builder_set_member_name(b, "ep-detect-signal"); + json_builder_add_int_value(b, s->ep_detect_signal); + json_builder_set_member_name(b, "el-flags"); + json_builder_add_int_value(b, s->el_flags); + checkpoint_json_endpoint(b, "endpoint", &s->endpoint); + checkpoint_json_endpoint(b, "advertised-endpoint", &s->advertised_endpoint); + checkpoint_json_endpoint(b, "learned-endpoint", &s->learned_endpoint); + checkpoint_json_endpoint(b, "last-local-endpoint", &s->last_local_endpoint); + json_builder_set_member_name(b, "detected-endpoints"); + json_builder_begin_array(b); + for (unsigned int i = 0; i < G_N_ELEMENTS(s->detected_endpoints); i++) + json_builder_add_string_value(b, s->detected_endpoints[i].address.family + ? endpoint_print_buf(&s->detected_endpoints[i]) : ""); + json_builder_end_array(b); + json_builder_set_member_name(b, "sfd-ids"); + json_builder_begin_array(b); + for (__auto_type l = s->sfds.head; l; l = l->next) + json_builder_add_int_value(b, ((stream_fd *) l->data)->unique_id); + json_builder_end_array(b); + json_builder_set_member_name(b, "sfd-local-endpoints"); + json_builder_begin_array(b); + if (s->sfd_local_endpoints) + for (unsigned int i = 0; i < s->sfd_local_endpoints->len; i++) { + const endpoint_t *ep = &g_array_index(s->sfd_local_endpoints, endpoint_t, i); + json_builder_add_string_value(b, ep->address.family ? endpoint_print_buf(ep) : ""); + } + json_builder_end_array(b); + json_builder_end_object(b); +} + +static void checkpoint_json_monologue(JsonBuilder *b, const struct checkpoint_monologue *m) { + json_builder_begin_object(b); + json_builder_set_member_name(b, "monologue-id"); + json_builder_add_int_value(b, m->monologue->unique_id); + json_builder_set_member_name(b, "flags"); + json_builder_add_int_value(b, m->flags); + json_builder_set_member_name(b, "medias-len"); + json_builder_add_int_value(b, m->medias_len); + checkpoint_json_str(b, "desired-family", + m->desired_family ? STR_PTR(m->desired_family->rfc_name) : NULL); + checkpoint_json_str(b, "logical-interface", + m->logical_intf ? &m->logical_intf->name : NULL); + checkpoint_json_str(b, "session-name", &m->session_name); + checkpoint_json_str(b, "session-timing", &m->session_timing); + str last_sdp = m->last_out_sdp ? STR_LEN(m->last_out_sdp->str, m->last_out_sdp->len) : STR_NULL; + checkpoint_json_str(b, "last-out-sdp", &last_sdp); + checkpoint_json_bandwidth(b, &m->bandwidth); + checkpoint_json_origin(b, "origin-in", &m->sdp_orig_in); + checkpoint_json_origin(b, "origin-out", &m->sdp_orig_out); + json_builder_end_object(b); +} + +static void checkpoint_json_candidates(JsonBuilder *b, const candidate_q *q) { + json_builder_set_member_name(b, "ice-candidates"); + json_builder_begin_array(b); + for (__auto_type l = q->head; l; l = l->next) { + const struct ice_candidate *c = l->data; + json_builder_begin_object(b); + checkpoint_json_str(b, "foundation", &c->foundation); + json_builder_set_member_name(b, "component"); + json_builder_add_int_value(b, c->component_id); + checkpoint_json_str(b, "transport", c->transport ? STR_PTR(c->transport->name) : NULL); + json_builder_set_member_name(b, "priority"); + json_builder_add_int_value(b, c->priority); + checkpoint_json_endpoint(b, "endpoint", &c->endpoint); + json_builder_set_member_name(b, "type"); + json_builder_add_int_value(b, c->type); + checkpoint_json_endpoint(b, "related", &c->related); + checkpoint_json_str(b, "ufrag", &c->ufrag); + json_builder_end_object(b); + } + json_builder_end_array(b); +} + +static void checkpoint_json_media(JsonBuilder *b, const struct checkpoint_media *m) { + json_builder_begin_object(b); + json_builder_set_member_name(b, "media-id"); + json_builder_add_int_value(b, m->media->unique_id); + json_builder_set_member_name(b, "endpoint-map-id"); + json_builder_add_int_value(b, m->endpoint_map ? m->endpoint_map->unique_id : -1); + json_builder_set_member_name(b, "flags"); + json_builder_add_int_value(b, m->flags); + json_builder_set_member_name(b, "ptime"); + json_builder_add_int_value(b, m->ptime); + json_builder_set_member_name(b, "maxptime"); + json_builder_add_int_value(b, m->maxptime); + json_builder_set_member_name(b, "had-ice"); + json_builder_add_boolean_value(b, m->had_ice); + checkpoint_json_str(b, "protocol", m->protocol ? STR_PTR(m->protocol->name) : &m->protocol_str); + checkpoint_json_str(b, "protocol-string", &m->protocol_str); + checkpoint_json_str(b, "format-string", &m->format_str); + checkpoint_json_str(b, "desired-family", + m->desired_family ? STR_PTR(m->desired_family->rfc_name) : NULL); + checkpoint_json_str(b, "logical-interface", + m->logical_intf ? &m->logical_intf->name : NULL); + checkpoint_json_str(b, "tls-id", &m->tls_id); + checkpoint_json_str(b, "ice-ufrag-local", &m->ice_ufrag[0]); + checkpoint_json_str(b, "ice-ufrag-remote", &m->ice_ufrag[1]); + checkpoint_json_str(b, "ice-pwd-local", &m->ice_pwd[0]); + checkpoint_json_str(b, "ice-pwd-remote", &m->ice_pwd[1]); + checkpoint_json_bandwidth(b, &m->bandwidth); + checkpoint_json_candidates(b, &m->ice_candidates); + ng_parser_ctx_t crypto_ctx; + ng_parser_json.init(&crypto_ctx, NULL); + parser_arg crypto = ng_parser_json.dict(&crypto_ctx); + redis_encode_sdes_params(&ng_parser_json, crypto, "sdes_in", &m->sdes_in); + redis_encode_sdes_params(&ng_parser_json, crypto, "sdes_out", &m->sdes_out); + redis_encode_dtls_fingerprint(&ng_parser_json, crypto, &m->fingerprint); + if (m->fp_hash_func) + ng_parser_json.dict_add_string(crypto, "preferred_hash_func", m->fp_hash_func->name); + json_builder_set_member_name(b, "crypto"); + json_builder_add_value(b, crypto.json); + parser_arg codecs = ng_parser_json.list(&crypto_ctx); + redis_encode_codec_store(&ng_parser_json, codecs, &m->codecs); + json_builder_set_member_name(b, "codecs"); + json_builder_add_value(b, codecs.json); + parser_arg offered_codecs = ng_parser_json.list(&crypto_ctx); + redis_encode_codec_store(&ng_parser_json, offered_codecs, &m->offered_codecs); + json_builder_set_member_name(b, "offered-codecs"); + json_builder_add_value(b, offered_codecs.json); + json_builder_set_member_name(b, "streams"); + json_builder_begin_array(b); + for (const struct checkpoint_stream *s = m->streams; s; s = s->next) + checkpoint_json_stream(b, s); + json_builder_end_array(b); + json_builder_end_object(b); +} + +str call_checkpoint_serialize(call_t *call, void **to_free) { + *to_free = NULL; + if (!call->checkpoints) + return STR_NULL; + + JsonBuilder *b = json_builder_new(); + json_builder_begin_object(b); + json_builder_set_member_name(b, "version"); + json_builder_add_int_value(b, 1); + json_builder_set_member_name(b, "checkpoints"); + json_builder_begin_array(b); + for (const struct call_checkpoint *cp = call->checkpoints; cp; cp = cp->next) { + json_builder_begin_object(b); + json_builder_set_member_name(b, "offerer-id"); + json_builder_add_int_value(b, cp->offerer->unique_id); + json_builder_set_member_name(b, "answerer-id"); + json_builder_add_int_value(b, cp->answerer->unique_id); + json_builder_set_member_name(b, "generation"); + json_builder_add_int_value(b, cp->generation); + json_builder_set_member_name(b, "pending-generation"); + json_builder_add_int_value(b, cp->pending_generation); + json_builder_set_member_name(b, "pending"); + json_builder_add_boolean_value(b, cp->pending); + if (cp->pending) { + json_builder_set_member_name(b, "monologues"); + json_builder_begin_array(b); + for (unsigned int i = 0; i < G_N_ELEMENTS(cp->monologues); i++) + checkpoint_json_monologue(b, &cp->monologues[i]); + json_builder_end_array(b); + json_builder_set_member_name(b, "medias"); + json_builder_begin_array(b); + for (const struct checkpoint_media *m = cp->medias; m; m = m->next) + checkpoint_json_media(b, m); + json_builder_end_array(b); + } + json_builder_end_object(b); + } + json_builder_end_array(b); + json_builder_end_object(b); + + JsonGenerator *g = json_generator_new(); + JsonNode *root = json_builder_get_root(b); + json_generator_set_root(g, root); + gsize len = 0; + char *data = json_generator_to_data(g, &len); + json_node_free(root); + g_object_unref(g); + g_object_unref(b); + *to_free = data; + return STR_LEN(data, len); +} + +static struct call_monologue *checkpoint_find_monologue_id(call_t *call, unsigned int id) { + for (__auto_type l = call->monologues.head; l; l = l->next) { + struct call_monologue *ml = l->data; + if (ml->unique_id == id) + return ml; + } + return NULL; +} + +static struct call_media *checkpoint_find_media_id(call_t *call, unsigned int id) { + for (__auto_type l = call->medias.head; l; l = l->next) { + struct call_media *m = l->data; + if (m->unique_id == id) + return m; + } + return NULL; +} + +static struct packet_stream *checkpoint_find_stream_id(call_t *call, unsigned int id) { + for (__auto_type l = call->streams.head; l; l = l->next) { + struct packet_stream *s = l->data; + if (s->unique_id == id) + return s; + } + return NULL; +} + +static stream_fd *checkpoint_find_sfd_id(call_t *call, unsigned int id) { + for (__auto_type l = call->stream_fds.head; l; l = l->next) { + stream_fd *sfd = l->data; + if (sfd->unique_id == id) + return sfd; + } + return NULL; +} + +static struct endpoint_map *checkpoint_find_endpoint_map_id(call_t *call, unsigned int id) { + for (__auto_type l = call->endpoint_maps.head; l; l = l->next) { + struct endpoint_map *map = l->data; + if (map->unique_id == id) + return map; + } + return NULL; +} + +static bool checkpoint_json_value_is(JsonObject *o, const char *name, GType type) { + JsonNode *node = o ? json_object_get_member(o, name) : NULL; + return node && JSON_NODE_HOLDS_VALUE(node) && json_node_get_value_type(node) == type; +} + +static bool checkpoint_json_int(JsonObject *o, const char *name) { + return checkpoint_json_value_is(o, name, G_TYPE_INT64); +} + +static bool checkpoint_json_bool(JsonObject *o, const char *name) { + return checkpoint_json_value_is(o, name, G_TYPE_BOOLEAN); +} + +static bool checkpoint_json_string(JsonObject *o, const char *name) { + return checkpoint_json_value_is(o, name, G_TYPE_STRING); +} + +static bool checkpoint_json_object(JsonObject *o, const char *name) { + JsonNode *node = o ? json_object_get_member(o, name) : NULL; + return node && JSON_NODE_HOLDS_OBJECT(node); +} + +static bool checkpoint_json_array(JsonObject *o, const char *name) { + JsonNode *node = o ? json_object_get_member(o, name) : NULL; + return node && JSON_NODE_HOLDS_ARRAY(node); +} + +static JsonObject *checkpoint_json_array_object(JsonArray *array, unsigned int index) { + JsonNode *node = array ? json_array_get_element(array, index) : NULL; + return node && JSON_NODE_HOLDS_OBJECT(node) ? json_node_get_object(node) : NULL; +} + +static str checkpoint_json_call_str(call_t *call, JsonObject *o, const char *name) { + const char *s = json_object_get_string_member(o, name); + str value = STR(s); + return call_str_cpy(&value); +} + +static int checkpoint_json_get_endpoint(endpoint_t *ep, JsonObject *o, const char *name) { + if (!checkpoint_json_string(o, name)) + return -1; + const char *value = json_object_get_string_member(o, name); + if (!value[0]) { + ZERO(*ep); + return 0; + } + return endpoint_parse_any(ep, value) ? 0 : -1; +} + +static int checkpoint_json_parse_endpoint(endpoint_t *ep, const char *value) { + if (!value[0]) { + ZERO(*ep); + return 0; + } + return endpoint_parse_any(ep, value) ? 0 : -1; +} + +static int checkpoint_json_get_bandwidth(struct session_bandwidth *bw, JsonObject *o) { + if (!checkpoint_json_object(o, "bandwidth")) + return -1; + JsonObject *b = json_object_get_object_member(o, "bandwidth"); + if (!checkpoint_json_int(b, "as") || !checkpoint_json_int(b, "ct") + || !checkpoint_json_int(b, "rr") || !checkpoint_json_int(b, "rs") + || !checkpoint_json_int(b, "tias")) + return -1; + bw->as = json_object_get_int_member(b, "as"); + bw->ct = json_object_get_int_member(b, "ct"); + bw->rr = json_object_get_int_member(b, "rr"); + bw->rs = json_object_get_int_member(b, "rs"); + bw->tias = json_object_get_int_member(b, "tias"); + return 0; +} + +static int checkpoint_json_get_origin(call_t *call, sdp_origin *origin, JsonObject *o, + const char *name) +{ + if (!checkpoint_json_object(o, name)) + return -1; + JsonObject *v = json_object_get_object_member(o, name); + if (!checkpoint_json_bool(v, "parsed") || !checkpoint_json_int(v, "version") + || !checkpoint_json_string(v, "username") + || !checkpoint_json_string(v, "session-id") + || !checkpoint_json_string(v, "network-type") + || !checkpoint_json_string(v, "address-type") + || !checkpoint_json_string(v, "address")) + return -1; + origin->parsed = json_object_get_boolean_member(v, "parsed"); + origin->version_num = json_object_get_int_member(v, "version"); + origin->username = checkpoint_json_call_str(call, v, "username"); + origin->session_id = checkpoint_json_call_str(call, v, "session-id"); + origin->address.network_type = checkpoint_json_call_str(call, v, "network-type"); + origin->address.address_type = checkpoint_json_call_str(call, v, "address-type"); + origin->address.address = checkpoint_json_call_str(call, v, "address"); + return 0; +} + +static stream_fd *checkpoint_reopen_sfd(call_t *call, stream_fd *old, const endpoint_t *local, + struct endpoint_map *map) +{ + if (!local->address.family || !local->port || old->socket.local.port) + return old; + stream_fd *existing = stream_fd_lookup(local); + stream_fd *replacement; + if (existing) { + obj_release(existing); + replacement = existing; + } + else { + struct socket_port_link spl = get_specific_port(local->port, old->local_intf->spec, + &call->callid); + if (!spl.socket.family) + return NULL; + set_tos(&spl.socket, call->tos); + replacement = stream_fd_new(&spl, call, old->local_intf); + } + for (__auto_type l = map ? map->intf_sfds.head : NULL; l; l = l->next) { + struct sfd_intf_list *il = l->data; + for (__auto_type k = il->list.head; k; k = k->next) + if (k->data == old) { + stream_fd_inc(replacement); + k->data = replacement; + stream_fd_dec(old); + } + } + return replacement; +} + +static struct checkpoint_stream *checkpoint_json_get_stream(call_t *call, JsonObject *o, + struct endpoint_map *map) +{ + const char *integers[] = { "stream-id", "selected-sfd-id", "selected-interface-id", + "flags", "ep-detect-signal", "el-flags" }; + for (unsigned int i = 0; i < G_N_ELEMENTS(integers); i++) + if (!checkpoint_json_int(o, integers[i])) + return NULL; + if (!checkpoint_json_array(o, "detected-endpoints") + || !checkpoint_json_array(o, "sfd-ids") + || !checkpoint_json_array(o, "sfd-local-endpoints")) + return NULL; + + struct packet_stream *ps = checkpoint_find_stream_id(call, + json_object_get_int_member(o, "stream-id")); + if (!ps) + return NULL; + struct checkpoint_stream *s = g_new0(__typeof(*s), 1); + s->stream = ps; + s->redis_restored = true; + s->flags = json_object_get_int_member(o, "flags"); + s->ep_detect_signal = json_object_get_int_member(o, "ep-detect-signal"); + s->el_flags = json_object_get_int_member(o, "el-flags"); + if (checkpoint_json_get_endpoint(&s->endpoint, o, "endpoint") + || checkpoint_json_get_endpoint(&s->advertised_endpoint, o, "advertised-endpoint") + || checkpoint_json_get_endpoint(&s->learned_endpoint, o, "learned-endpoint") + || checkpoint_json_get_endpoint(&s->last_local_endpoint, o, "last-local-endpoint")) + goto err; + + JsonArray *detected = json_object_get_array_member(o, "detected-endpoints"); + if (!detected || json_array_get_length(detected) != G_N_ELEMENTS(s->detected_endpoints)) + goto err; + for (unsigned int i = 0; i < G_N_ELEMENTS(s->detected_endpoints); i++) { + JsonNode *element = json_array_get_element(detected, i); + if (!element || !JSON_NODE_HOLDS_VALUE(element) + || json_node_get_value_type(element) != G_TYPE_STRING + || checkpoint_json_parse_endpoint(&s->detected_endpoints[i], + json_node_get_string(element))) + goto err; + } + + JsonArray *sfds = json_object_get_array_member(o, "sfd-ids"); + JsonArray *locals = json_object_get_array_member(o, "sfd-local-endpoints"); + if (!sfds || !locals || json_array_get_length(sfds) != json_array_get_length(locals)) + goto err; + gint64 selected = json_object_get_int_member(o, "selected-sfd-id"); + s->selected_sfd_set = selected >= 0; + s->sfd_local_endpoints = g_array_new(FALSE, FALSE, sizeof(endpoint_t)); + for (unsigned int i = 0; i < json_array_get_length(sfds); i++) { + JsonNode *sfd_node = json_array_get_element(sfds, i); + JsonNode *local_node = json_array_get_element(locals, i); + if (!sfd_node || !JSON_NODE_HOLDS_VALUE(sfd_node) + || json_node_get_value_type(sfd_node) != G_TYPE_INT64 + || !local_node || !JSON_NODE_HOLDS_VALUE(local_node) + || json_node_get_value_type(local_node) != G_TYPE_STRING) + goto err; + gint64 sfd_id = json_node_get_int(sfd_node); + stream_fd *old = checkpoint_find_sfd_id(call, sfd_id); + endpoint_t local; + if (!old || checkpoint_json_parse_endpoint(&local, + json_node_get_string(local_node))) + goto err; + stream_fd *sfd = checkpoint_reopen_sfd(call, old, &local, map); + if (!sfd) + goto err; + stream_fd_inc(sfd); + t_queue_push_tail(&s->sfds, sfd); + g_array_append_val(s->sfd_local_endpoints, local); + if (sfd_id == selected) + s->selected_sfd = sfd; + } + if (selected >= 0 && !s->selected_sfd) + goto err; + return s; + +err: + checkpoint_stream_free(s); + return NULL; +} + +static int checkpoint_json_get_monologue(call_t *call, struct checkpoint_monologue *m, + JsonObject *o) +{ + const char *integers[] = { "monologue-id", "flags", "medias-len" }; + const char *strings[] = { "desired-family", "logical-interface", "session-name", + "session-timing", "last-out-sdp" }; + for (unsigned int i = 0; i < G_N_ELEMENTS(integers); i++) + if (!checkpoint_json_int(o, integers[i])) + return -1; + for (unsigned int i = 0; i < G_N_ELEMENTS(strings); i++) + if (!checkpoint_json_string(o, strings[i])) + return -1; + m->monologue = checkpoint_find_monologue_id(call, + json_object_get_int_member(o, "monologue-id")); + if (!m->monologue) + return -1; + m->flags = json_object_get_int_member(o, "flags"); + m->medias_len = json_object_get_int_member(o, "medias-len"); + str family = STR(json_object_get_string_member(o, "desired-family")); + m->desired_family = family.len ? get_socket_family_rfc(&family) : NULL; + str intf = STR(json_object_get_string_member(o, "logical-interface")); + m->logical_intf = intf.len ? get_logical_interface(&intf, m->desired_family, 0) : NULL; + if ((family.len && !m->desired_family) || (intf.len && !m->logical_intf)) + return -1; + m->session_name = checkpoint_json_call_str(call, o, "session-name"); + m->session_timing = checkpoint_json_call_str(call, o, "session-timing"); + const char *last_sdp = json_object_get_string_member(o, "last-out-sdp"); + if (*last_sdp) + m->last_out_sdp = g_string_new(last_sdp); + if (checkpoint_json_get_bandwidth(&m->bandwidth, o) + || checkpoint_json_get_origin(call, &m->sdp_orig_in, o, "origin-in") + || checkpoint_json_get_origin(call, &m->sdp_orig_out, o, "origin-out")) + return -1; + return 0; +} + +static int checkpoint_json_get_candidates(call_t *call, candidate_q *q, JsonObject *o) { + if (!checkpoint_json_array(o, "ice-candidates")) + return -1; + JsonArray *array = json_object_get_array_member(o, "ice-candidates"); + if (!array) + return -1; + for (unsigned int i = 0; i < json_array_get_length(array); i++) { + JsonObject *v = checkpoint_json_array_object(array, i); + const char *integers[] = { "component", "priority", "type" }; + const char *strings[] = { "foundation", "transport", "endpoint", "related", "ufrag" }; + for (unsigned int j = 0; j < G_N_ELEMENTS(integers); j++) + if (!checkpoint_json_int(v, integers[j])) + goto err; + for (unsigned int j = 0; j < G_N_ELEMENTS(strings); j++) + if (!checkpoint_json_string(v, strings[j])) + goto err; + struct ice_candidate *c = g_new0(__typeof(*c), 1); + c->foundation = checkpoint_json_call_str(call, v, "foundation"); + c->component_id = json_object_get_int_member(v, "component"); + str transport = STR(json_object_get_string_member(v, "transport")); + c->transport = transport.len ? get_socket_type(&transport) : NULL; + c->priority = json_object_get_int_member(v, "priority"); + gint64 type = json_object_get_int_member(v, "type"); + if (!c->transport || type <= ICT_UNKNOWN || type >= __ICT_LAST + || checkpoint_json_get_endpoint(&c->endpoint, v, "endpoint") + || checkpoint_json_get_endpoint(&c->related, v, "related")) { + g_free(c); + goto err; + } + c->type = type; + c->ufrag = checkpoint_json_call_str(call, v, "ufrag"); + t_queue_push_tail(q, c); + } + return 0; + +err: + ice_candidates_free(q); + return -1; +} + +static struct checkpoint_media *checkpoint_json_get_media(call_t *call, JsonObject *o) { + const char *stage = "required fields"; + const char *integers[] = { "media-id", "endpoint-map-id", "flags", "ptime", "maxptime" }; + const char *strings[] = { "protocol", "protocol-string", "format-string", "desired-family", + "logical-interface", "tls-id", "ice-ufrag-local", "ice-ufrag-remote", + "ice-pwd-local", "ice-pwd-remote" }; + for (unsigned int i = 0; i < G_N_ELEMENTS(integers); i++) + if (!checkpoint_json_int(o, integers[i])) + return NULL; + for (unsigned int i = 0; i < G_N_ELEMENTS(strings); i++) + if (!checkpoint_json_string(o, strings[i])) + return NULL; + if (!checkpoint_json_bool(o, "had-ice") || !checkpoint_json_object(o, "crypto") + || !checkpoint_json_array(o, "codecs") + || !checkpoint_json_array(o, "offered-codecs") + || !checkpoint_json_array(o, "streams")) + return NULL; + struct call_media *live = checkpoint_find_media_id(call, + json_object_get_int_member(o, "media-id")); + if (!live) + return NULL; + struct checkpoint_media *m = g_new0(__typeof(*m), 1); + m->media = live; + gint64 endpoint_map_id = json_object_get_int_member(o, "endpoint-map-id"); + if (endpoint_map_id >= 0 + && !(m->endpoint_map = checkpoint_find_endpoint_map_id(call, endpoint_map_id))) + goto err; + codec_store_init(&m->codecs, live); + codec_store_init(&m->offered_codecs, live); + m->flags = json_object_get_int_member(o, "flags"); + m->ptime = json_object_get_int_member(o, "ptime"); + m->maxptime = json_object_get_int_member(o, "maxptime"); + m->had_ice = json_object_get_boolean_member(o, "had-ice"); + str protocol = STR(json_object_get_string_member(o, "protocol")); + m->protocol = protocol.len ? transport_protocol(&protocol) : NULL; + if (protocol.len && !m->protocol) + goto err; + m->protocol_str = checkpoint_json_call_str(call, o, "protocol-string"); + m->format_str = checkpoint_json_call_str(call, o, "format-string"); + str family = STR(json_object_get_string_member(o, "desired-family")); + m->desired_family = family.len ? get_socket_family_rfc(&family) : NULL; + str intf = STR(json_object_get_string_member(o, "logical-interface")); + m->logical_intf = intf.len ? get_logical_interface(&intf, m->desired_family, 0) : NULL; + if ((family.len && !m->desired_family) || (intf.len && !m->logical_intf)) + goto err; + m->tls_id = checkpoint_json_call_str(call, o, "tls-id"); + m->ice_ufrag[0] = checkpoint_json_call_str(call, o, "ice-ufrag-local"); + m->ice_ufrag[1] = checkpoint_json_call_str(call, o, "ice-ufrag-remote"); + m->ice_pwd[0] = checkpoint_json_call_str(call, o, "ice-pwd-local"); + m->ice_pwd[1] = checkpoint_json_call_str(call, o, "ice-pwd-remote"); + stage = "bandwidth"; + if (checkpoint_json_get_bandwidth(&m->bandwidth, o)) + goto err; + stage = "ICE candidates"; + if (checkpoint_json_get_candidates(call, &m->ice_candidates, o)) + goto err; + stage = "crypto object"; + JsonNode *crypto_node = json_object_get_member(o, "crypto"); + struct redis_hash crypto = {0}; + parser_arg crypto_arg = { .json = crypto_node }; + if (redis_hash_from_parser(&crypto, &ng_parser_json, crypto_arg)) + goto err; + stage = "crypto parameters"; + int crypto_ret = redis_decode_sdes_params(&m->sdes_in, &crypto, "sdes_in") + || redis_decode_sdes_params(&m->sdes_out, &crypto, "sdes_out") + || redis_decode_dtls_fingerprint(&m->fingerprint, &crypto); + str *preferred = g_hash_table_lookup(crypto.ht, "preferred_hash_func"); + if (preferred) + m->fp_hash_func = dtls_find_hash_func(preferred); + if (preferred && !m->fp_hash_func) + crypto_ret = -1; + redis_hash_destroy(&crypto); + if (crypto_ret) + goto err; + stage = "codec objects"; + JsonNode *codecs_node = json_object_get_member(o, "codecs"); + JsonNode *offered_node = json_object_get_member(o, "offered-codecs"); + parser_arg codecs_arg = { .json = codecs_node }; + parser_arg offered_arg = { .json = offered_node }; + stage = "codec stores"; + if (redis_decode_codec_store(&ng_parser_json, codecs_arg, &m->codecs) + || redis_decode_codec_store(&ng_parser_json, offered_arg, &m->offered_codecs)) + goto err; + JsonArray *streams = json_object_get_array_member(o, "streams"); + stage = "streams array"; + stage = "stream state"; + for (unsigned int i = 0; i < json_array_get_length(streams); i++) { + struct checkpoint_stream *s = checkpoint_json_get_stream(call, + checkpoint_json_array_object(streams, i), m->endpoint_map); + if (!s) + goto err; + s->next = m->streams; + m->streams = s; + } + return m; + +err: + ilog(LOG_WARNING, "Failed to restore checkpoint media at %s", stage); + checkpoint_media_free(m); + return NULL; +} + +int call_checkpoint_deserialize(call_t *call, const str *data) { + const char *stage = "JSON document"; + JsonParser *parser = json_parser_new(); + GError *error = NULL; + struct call_checkpoint *head = NULL; + if (!json_parser_load_from_data(parser, data->s, data->len, &error)) + goto err; + JsonNode *node = json_parser_get_root(parser); + if (!node || !JSON_NODE_HOLDS_OBJECT(node)) + goto err; + JsonObject *root = json_node_get_object(node); + stage = "root object"; + if (!checkpoint_json_int(root, "version") + || json_object_get_int_member(root, "version") != 1 + || !checkpoint_json_array(root, "checkpoints")) + goto err; + JsonArray *checkpoints = json_object_get_array_member(root, "checkpoints"); + if (!checkpoints) + goto err; + for (unsigned int i = 0; i < json_array_get_length(checkpoints); i++) { + stage = "checkpoint fields"; + JsonObject *o = checkpoint_json_array_object(checkpoints, i); + const char *integers[] = { "offerer-id", "answerer-id", "generation", + "pending-generation" }; + for (unsigned int j = 0; j < G_N_ELEMENTS(integers); j++) + if (!checkpoint_json_int(o, integers[j])) + goto err; + if (!checkpoint_json_bool(o, "pending")) + goto err; + struct call_checkpoint *cp = g_new0(__typeof(*cp), 1); + cp->offerer = checkpoint_find_monologue_id(call, + json_object_get_int_member(o, "offerer-id")); + cp->answerer = checkpoint_find_monologue_id(call, + json_object_get_int_member(o, "answerer-id")); + if (!cp->offerer || !cp->answerer) { + stage = "checkpoint monologue IDs"; + g_free(cp); + goto err; + } + cp->generation = json_object_get_int_member(o, "generation"); + cp->pending_generation = json_object_get_int_member(o, "pending-generation"); + cp->pending = json_object_get_boolean_member(o, "pending"); + cp->next = head; + head = cp; + if (!cp->pending) + continue; + stage = "pending checkpoint fields"; + if (!checkpoint_json_array(o, "monologues") + || !checkpoint_json_array(o, "medias")) + goto err; + JsonArray *monologues = json_object_get_array_member(o, "monologues"); + if (!monologues || json_array_get_length(monologues) != G_N_ELEMENTS(cp->monologues)) + goto err; + for (unsigned int j = 0; j < G_N_ELEMENTS(cp->monologues); j++) + if (checkpoint_json_get_monologue(call, &cp->monologues[j], + checkpoint_json_array_object(monologues, j))) { + stage = "monologue state"; + goto err; + } + stage = "media array"; + JsonArray *medias = json_object_get_array_member(o, "medias"); + if (!medias) + goto err; + for (unsigned int j = 0; j < json_array_get_length(medias); j++) { + stage = "media state"; + struct checkpoint_media *m = checkpoint_json_get_media(call, + checkpoint_json_array_object(medias, j)); + if (!m) + goto err; + m->next = cp->medias; + cp->medias = m; + } + } + call_checkpoint_free_all(call); + call->checkpoints = head; + if (error) + g_error_free(error); + g_object_unref(parser); + return 0; + +err: + ilog(LOG_WARNING, "Failed to deserialize checkpoint data at %s", stage); + while (head) { + struct call_checkpoint *cp = head; + head = cp->next; + checkpoint_free(cp); + } + if (error) + g_error_free(error); + g_object_unref(parser); + return -1; +} + +static void checkpoint_candidates_copy(candidate_q *dst, const candidate_q *src) { + for (__auto_type l = src->head; l; l = l->next) { + struct ice_candidate *copy = g_new(__typeof(*copy), 1); + *copy = *(struct ice_candidate *) l->data; + t_queue_push_tail(dst, copy); + } +} +static void checkpoint_stream_free(struct checkpoint_stream *stream) { + t_queue_clear_full(&stream->sfds, stream_fd_dec); + if (stream->sfd_local_endpoints) + g_array_free(stream->sfd_local_endpoints, TRUE); + g_free(stream); +} + +static void checkpoint_media_free(struct checkpoint_media *media) { + while (media->streams) { + struct checkpoint_stream *stream = media->streams; + media->streams = stream->next; + checkpoint_stream_free(stream); + } + crypto_params_sdes_queue_clear(&media->sdes_in); + crypto_params_sdes_queue_clear(&media->sdes_out); + codec_store_cleanup(&media->codecs); + codec_store_cleanup(&media->offered_codecs); + ice_candidates_free(&media->ice_candidates); + g_free(media); +} + +static void checkpoint_clear_snapshot(struct call_checkpoint *cp) { + while (cp->medias) { + struct checkpoint_media *media = cp->medias; + cp->medias = media->next; + checkpoint_media_free(media); + } + for (unsigned int i = 0; i < G_N_ELEMENTS(cp->monologues); i++) { + if (cp->monologues[i].last_out_sdp) + g_string_free(cp->monologues[i].last_out_sdp, TRUE); + ZERO(cp->monologues[i]); + } + cp->pending = false; + cp->pending_generation = 0; +} + +static void checkpoint_free(struct call_checkpoint *cp) { + checkpoint_clear_snapshot(cp); + g_free(cp); +} + +static struct call_checkpoint *checkpoint_find(call_t *call, struct call_monologue *a, + struct call_monologue *b) +{ + for (struct call_checkpoint *cp = call->checkpoints; cp; cp = cp->next) { + if ((cp->offerer == a && cp->answerer == b) || (cp->offerer == b && cp->answerer == a)) + return cp; + } + return NULL; +} + +static void checkpoint_take_stream(struct checkpoint_media *media, struct packet_stream *ps) { + struct checkpoint_stream *stream = g_new0(__typeof(*stream), 1); + stream->stream = ps; + stream->endpoint = ps->endpoint; + stream->advertised_endpoint = ps->advertised_endpoint; + stream->learned_endpoint = ps->learned_endpoint; + memcpy(stream->detected_endpoints, ps->detected_endpoints, sizeof(stream->detected_endpoints)); + stream->last_local_endpoint = ps->last_local_endpoint; + stream->ep_detect_signal = ps->ep_detect_signal; + stream->el_flags = ps->el_flags; + stream->flags = atomic64_get_na(&ps->ps_flags); + stream->selected_sfd = ps->selected_sfd; + stream->selected_sfd_set = ps->selected_sfd != NULL; + stream->sfd_local_endpoints = g_array_new(FALSE, FALSE, sizeof(endpoint_t)); + for (__auto_type l = ps->sfds.head; l; l = l->next) { + stream_fd *sfd = l->data; + stream_fd_inc(sfd); + t_queue_push_tail(&stream->sfds, sfd); + g_array_append_val(stream->sfd_local_endpoints, sfd->socket.local); + } + stream->next = media->streams; + media->streams = stream; +} + +static void checkpoint_take_media(struct call_checkpoint *cp, struct call_media *m) { + struct checkpoint_media *media = g_new0(__typeof(*media), 1); + media->media = m; + media->endpoint_map = m->endpoint_map; + media->protocol = m->protocol; + media->protocol_str = m->protocol_str; + media->format_str = m->format_str; + media->desired_family = m->desired_family; + media->logical_intf = m->logical_intf; + media->flags = atomic64_get_na(&m->media_flags); + crypto_params_sdes_queue_copy(&media->sdes_in, &m->sdes_in); + crypto_params_sdes_queue_copy(&media->sdes_out, &m->sdes_out); + media->fingerprint = m->fingerprint; + media->fp_hash_func = m->fp_hash_func; + media->tls_id = m->tls_id; + codec_store_init(&media->codecs, m); + codec_store_init(&media->offered_codecs, m); + codec_store_copy(&media->codecs, &m->codecs); + codec_store_copy(&media->offered_codecs, &m->offered_codecs); + checkpoint_candidates_copy(&media->ice_candidates, &m->ice_candidates); + if (m->ice_agent) { + media->had_ice = true; + memcpy(media->ice_ufrag, m->ice_agent->ufrag, sizeof(media->ice_ufrag)); + memcpy(media->ice_pwd, m->ice_agent->pwd, sizeof(media->ice_pwd)); + } + media->ptime = m->ptime; + media->maxptime = m->maxptime; + media->bandwidth = m->sdp_media_bandwidth; + for (__auto_type l = m->streams.head; l; l = l->next) + checkpoint_take_stream(media, l->data); + media->next = cp->medias; + cp->medias = media; +} + +static void checkpoint_take_monologue(struct call_checkpoint *cp, unsigned int idx, + struct call_monologue *ml) +{ + struct checkpoint_monologue *snap = &cp->monologues[idx]; + snap->monologue = ml; + snap->desired_family = ml->desired_family; + snap->logical_intf = ml->logical_intf; + snap->flags = atomic64_get_na(&ml->ml_flags); + snap->bandwidth = ml->sdp_session_bandwidth; + snap->sdp_orig_in = ml->sdp_orig_in; + snap->sdp_orig_out = ml->sdp_orig_out; + snap->session_name = ml->sdp_session_name; + snap->session_timing = ml->sdp_session_timing; + if (ml->last_out_sdp) + snap->last_out_sdp = g_string_new_len(ml->last_out_sdp->str, ml->last_out_sdp->len); + snap->medias_len = ml->medias->len; + for (unsigned int i = 0; i < ml->medias->len; i++) { + struct call_media *media = ml->medias->pdata[i]; + if (media) + checkpoint_take_media(cp, media); + } +} + +uint64_t call_checkpoint_offer(call_t *call, struct call_monologue *offerer, + struct call_monologue *answerer, bool enable) +{ + struct call_checkpoint *cp = checkpoint_find(call, offerer, answerer); + if (!cp && !enable) + return 0; + if (!cp) { + cp = g_new0(__typeof(*cp), 1); + cp->next = call->checkpoints; + call->checkpoints = cp; + } + /* Multiple offers can be outstanding before either an answer or rollback + * arrives. They all belong to the same uncommitted exchange, so retain the + * snapshot and generation of the last committed state. Replacing either + * here would make a later rollback restore state from an earlier rejected + * offer instead. */ + if (cp->pending) + return cp->pending_generation; + checkpoint_clear_snapshot(cp); + cp->offerer = offerer; + cp->answerer = answerer; + cp->pending_generation = cp->generation + 1; + checkpoint_take_monologue(cp, 0, offerer); + checkpoint_take_monologue(cp, 1, answerer); + cp->pending = true; + return cp->pending_generation; +} + +uint64_t call_checkpoint_answer(call_t *call, struct call_monologue *a, + struct call_monologue *b, bool *enabled) +{ + struct call_checkpoint *cp = checkpoint_find(call, a, b); + *enabled = cp != NULL; + if (!cp) + return 0; + if (cp->pending) { + cp->generation = cp->pending_generation; + checkpoint_clear_snapshot(cp); + } + return cp->generation; +} + +static void checkpoint_restore_stream(struct checkpoint_stream *snap) { + struct packet_stream *ps = snap->stream; + dtls_shutdown(ps); + /* A Redis-restored snapshot can refer to an SFD that existed on the old + * instance but could not be rebound there, while normal call restoration + * has already selected a usable local socket on this instance. Keep that + * live binding. For in-process rollback, retain it only when it is the + * socket captured by the snapshot. */ + bool preserve_live_binding = ps->selected_sfd && ps->selected_sfd->socket.local.port + && (ps->selected_sfd == snap->selected_sfd + || (snap->selected_sfd_set + && (!snap->selected_sfd || !snap->selected_sfd->socket.local.port))); + if (!preserve_live_binding) { + t_queue_clear_full(&ps->sfds, stream_fd_dec); + for (__auto_type l = snap->sfds.head; l; l = l->next) { + stream_fd *sfd = l->data; + stream_fd_inc(sfd); + t_queue_push_tail(&ps->sfds, sfd); + } + ps->selected_sfd = snap->selected_sfd; + } + ps->endpoint = snap->endpoint; + ps->advertised_endpoint = snap->advertised_endpoint; + ps->learned_endpoint = snap->learned_endpoint; + memcpy(ps->detected_endpoints, snap->detected_endpoints, sizeof(ps->detected_endpoints)); + ps->last_local_endpoint = snap->last_local_endpoint; + ps->ep_detect_signal = snap->ep_detect_signal; + ps->el_flags = snap->el_flags; + atomic64_set_na(&ps->ps_flags, snap->flags); +} + +static bool checkpoint_endpoint_map_is_bound(const struct endpoint_map *map) { + if (!map) + return false; + for (__auto_type l = map->intf_sfds.head; l; l = l->next) { + struct sfd_intf_list *il = l->data; + for (__auto_type k = il->list.head; k; k = k->next) { + stream_fd *sfd = k->data; + if (sfd->socket.local.port) + return true; + } + } + return false; +} + +static struct endpoint_map *checkpoint_live_endpoint_map(struct call_media *media) { + for (__auto_type l = media->endpoint_maps.head; l; l = l->next) { + struct endpoint_map *map = l->data; + for (__auto_type k = map->intf_sfds.head; k; k = k->next) { + struct sfd_intf_list *il = k->data; + for (__auto_type n = il->list.head; n; n = n->next) { + stream_fd *sfd = n->data; + for (__auto_type p = media->streams.head; p; p = p->next) + if (((struct packet_stream *) p->data)->selected_sfd == sfd) + return map; + } + } + } + return NULL; +} + +static void checkpoint_prepare_live_bindings(struct checkpoint_media *snap) { + struct endpoint_map *live_map = checkpoint_live_endpoint_map(snap->media); + bool used_live_binding = false; + for (struct checkpoint_stream *stream = snap->streams; stream; stream = stream->next) { + struct packet_stream *ps = stream->stream; + if (!stream->redis_restored || !stream->selected_sfd_set + || (stream->selected_sfd && stream->selected_sfd->socket.local.port) + || !ps->selected_sfd || !ps->selected_sfd->socket.local.port) + continue; + t_queue_clear_full(&stream->sfds, stream_fd_dec); + for (__auto_type l = ps->sfds.head; l; l = l->next) { + stream_fd *sfd = l->data; + stream_fd_inc(sfd); + t_queue_push_tail(&stream->sfds, sfd); + } + stream->selected_sfd = ps->selected_sfd; + used_live_binding = true; + } + if (used_live_binding && live_map) + snap->endpoint_map = live_map; +} + +static void checkpoint_restore_media(struct checkpoint_media *snap) { + struct call_media *m = snap->media; + struct endpoint_map *live_map = checkpoint_live_endpoint_map(m); + m->endpoint_map = checkpoint_endpoint_map_is_bound(snap->endpoint_map) + ? snap->endpoint_map : live_map; + m->protocol = snap->protocol; + m->protocol_str = snap->protocol_str; + m->format_str = snap->format_str; + m->desired_family = snap->desired_family; + m->logical_intf = snap->logical_intf; + atomic64_set_na(&m->media_flags, snap->flags); + crypto_params_sdes_queue_clear(&m->sdes_in); + crypto_params_sdes_queue_clear(&m->sdes_out); + crypto_params_sdes_queue_copy(&m->sdes_in, &snap->sdes_in); + crypto_params_sdes_queue_copy(&m->sdes_out, &snap->sdes_out); + m->fingerprint = snap->fingerprint; + m->fp_hash_func = snap->fp_hash_func; + m->tls_id = snap->tls_id; + codec_store_copy(&m->codecs, &snap->codecs); + codec_store_copy(&m->offered_codecs, &snap->offered_codecs); + m->ptime = snap->ptime; + m->maxptime = snap->maxptime; + m->sdp_media_bandwidth = snap->bandwidth; + ice_candidates_free(&m->ice_candidates); + checkpoint_candidates_copy(&m->ice_candidates, &snap->ice_candidates); + for (struct checkpoint_stream *stream = snap->streams; stream; stream = stream->next) + checkpoint_restore_stream(stream); + codec_handlers_free(m); + if (snap->had_ice) { + ice_agent_init(&m->ice_agent, m); + ice_rollback(m->ice_agent, snap->ice_ufrag, snap->ice_pwd, &snap->ice_candidates); + } + else + ice_shutdown(&m->ice_agent); +} + +static void checkpoint_restore_monologue(struct checkpoint_monologue *snap) { + struct call_monologue *ml = snap->monologue; + for (unsigned int i = snap->medias_len; i < ml->medias->len; i++) { + struct call_media *media = ml->medias->pdata[i]; + if (media) + call_media_stop(media); + } + t_ptr_array_set_size(ml->medias, snap->medias_len); + ml->desired_family = snap->desired_family; + ml->logical_intf = snap->logical_intf; + atomic64_set_na(&ml->ml_flags, snap->flags); + ml->sdp_session_bandwidth = snap->bandwidth; + ml->sdp_orig_in = snap->sdp_orig_in; + ml->sdp_orig_out = snap->sdp_orig_out; + ml->sdp_session_name = snap->session_name; + ml->sdp_session_timing = snap->session_timing; + if (ml->last_out_sdp) + g_string_free(ml->last_out_sdp, TRUE); + ml->last_out_sdp = snap->last_out_sdp + ? g_string_new_len(snap->last_out_sdp->str, snap->last_out_sdp->len) : NULL; +} + +int call_checkpoint_rollback(call_t *call, struct call_monologue *a, struct call_monologue *b, + uint64_t generation, bool generation_given, uint64_t *current_generation) +{ + struct call_checkpoint *cp = checkpoint_find(call, a, b); + if (!cp) { + *current_generation = 0; + return 0; + } + *current_generation = cp->generation; + if (!cp->pending || (generation_given && generation != cp->pending_generation)) + return 0; + + for (struct checkpoint_media *media = cp->medias; media; media = media->next) + checkpoint_prepare_live_bindings(media); + for (struct checkpoint_media *media = cp->medias; media; media = media->next) + checkpoint_restore_media(media); + for (unsigned int i = 0; i < G_N_ELEMENTS(cp->monologues); i++) + checkpoint_restore_monologue(&cp->monologues[i]); + update_init_monologue_subscribers(cp->offerer, OP_OFFER); + update_init_monologue_subscribers(cp->answerer, OP_ANSWER); + for (struct checkpoint_media *media = cp->medias; media; media = media->next) + for (struct checkpoint_stream *stream = media->streams; stream; stream = stream->next) + { + checkpoint_restore_stream(stream); + __init_stream(stream->stream); + } + checkpoint_clear_snapshot(cp); + call->last_signal_us = rtpe_now; + return 1; +} + + +void call_checkpoint_free_all(call_t *call) { + while (call->checkpoints) { + struct call_checkpoint *cp = call->checkpoints; + call->checkpoints = cp->next; + checkpoint_free(cp); + } +} diff --git a/daemon/call_flags.c b/daemon/call_flags.c index 6b6cdd4dd..5ba708638 100644 --- a/daemon/call_flags.c +++ b/daemon/call_flags.c @@ -526,6 +526,8 @@ static const char *call_ng_flags_supports(str *s, unsigned int idx, helper_arg a sdp_ng_flags *out = arg.flags; if (!str_cmp(s, "load limit")) out->supports_load_limit = true; + else if (!str_cmp(s, "rollback")) + out->supports_rollback = true; else ilog(LOG_INFO | LOG_FLAG_LIMIT, "Optional feature '" STR_FORMAT "' not supported", STR_FMT(s)); @@ -909,6 +911,10 @@ const char *call_ng_flags_flags(str *s, unsigned int idx, helper_arg arg) { case CSH_LOOKUP("reset"): out->reset = true; break; + case CSH_LOOKUP("track-state"): + case CSH_LOOKUP("track state"): + out->track_state = true; + break; case CSH_LOOKUP("single-codec"): case CSH_LOOKUP("single codec"): out->single_codec = true; diff --git a/daemon/call_interfaces.c b/daemon/call_interfaces.c index ab004a095..16a5c93a3 100644 --- a/daemon/call_interfaces.c +++ b/daemon/call_interfaces.c @@ -1,4 +1,5 @@ #include "call_interfaces.h" +#include "call_checkpoint.h" #include #include @@ -557,6 +558,8 @@ static const char *call_offer_answer_ng(ng_command_ctx_t *ctx, const char *addr) parser_arg output = ctx->resp; const ng_parser_t *parser = ctx->parser_ctx.parser; g_auto(str) sdp_out = STR_NULL; + uint64_t generation = 0; + bool checkpointing = false; call_ng_process_flags_RETURN(&flags, ctx); @@ -652,10 +655,22 @@ static const char *call_offer_answer_ng(ng_command_ctx_t *ctx, const char *addr) t_hash_table_insert(call->endpoints, memory_arena_objdup(streams.head->data->rtp_endpoint), from_ml); + if (flags.opmode == OP_OFFER) { + /* The checkpoint is opened before SDP processing so it captures the + * pre-offer state. If monologue_offer_answer() rejects the offer, the + * pending checkpoint deliberately remains: it is still a valid snapshot + * of the committed state, although the caller never received its + * generation. A later tracked offer reuses that pending generation. */ + generation = call_checkpoint_offer(call, from_ml, to_ml, flags.track_state); + checkpointing = generation != 0; + } + struct recording *recording = NULL; /* offer/answer model processing */ if ((ret = monologue_offer_answer(monologues, &streams, &flags)) == 0) { + if (flags.opmode == OP_ANSWER) + generation = call_checkpoint_answer(call, from_ml, to_ml, &checkpointing); update_metadata_monologue(from_ml, &flags); detect_setup_recording(call, &flags); @@ -681,6 +696,12 @@ static const char *call_offer_answer_ng(ng_command_ctx_t *ctx, const char *addr) /* place return output SDP */ ctx->ngbuf->sdp_out = sdp_out.s; ctx->parser_ctx.parser->dict_add_str(output, "sdp", &sdp_out); + if (checkpointing) + parser->dict_add_int(output, "generation", generation); + if (flags.supports_rollback) { + parser_arg supported = parser->dict_add_list(output, "supported"); + parser->list_add_string(supported, "rollback"); + } meta_write_sdp_after(recording, &sdp_out, from_ml, flags.opmode); @@ -734,6 +755,66 @@ const char *call_answer_ng(ng_command_ctx_t *ctx) { return call_offer_answer_ng(ctx, NULL); } +static bool monologue_has_tag(const struct call_monologue *ml, const str *tag) { + if (!str_cmp_str(&ml->tag, tag)) + return true; + for (__auto_type l = ml->tag_aliases.head; l; l = l->next) { + if (!str_cmp_str(l->data, tag)) + return true; + } + return false; +} + +const char *call_rollback_ng(ng_command_ctx_t *ctx) { + const ng_parser_t *parser = ctx->parser_ctx.parser; + str call_id = parser->dict_get_str(ctx->req, "call-id"); + str from_tag = parser->dict_get_str(ctx->req, "from-tag"); + str to_tag = parser->dict_get_str(ctx->req, "to-tag"); + str via_branch = parser->dict_get_str(ctx->req, "via-branch"); + long long generation_value = parser->dict_get_int_str(ctx->req, "generation", -1); + + if (!call_id.len) + return "No call-id in message"; + if (!from_tag.len) + return "No from-tag in message"; + if (!to_tag.len) + return "No to-tag in message"; + if (generation_value < -1) + return "Invalid generation"; + + call_t *call = call_get(&call_id); + if (!call) + return "Unknown call-id"; + + struct call_monologue *from_ml = call_get_monologue(call, &from_tag); + struct call_monologue *to_ml = via_branch.len + ? t_hash_table_lookup(call->viabranches, &via_branch) + : call_get_monologue(call, &to_tag); + if (!from_ml || !to_ml || from_ml == to_ml + || !monologue_has_tag(from_ml, &from_tag) + || !monologue_has_tag(to_ml, &to_tag) + || !g_hash_table_contains(from_ml->associated_tags, to_ml)) + { + rwlock_unlock_w(&call->master_lock); + obj_put(call); + return "Unknown dialogue"; + } + if (rtpe_config.active_switchover && IS_FOREIGN_CALL(call)) + call_make_own_foreign(call, false); + + uint64_t current_generation; + int rolled_back = call_checkpoint_rollback(call, from_ml, to_ml, + generation_value < 0 ? 0 : (uint64_t) generation_value, + generation_value >= 0, ¤t_generation); + parser->dict_add_int(ctx->resp, "rolled-back", rolled_back); + if (current_generation) + parser->dict_add_int(ctx->resp, "generation", current_generation); + rwlock_unlock_w(&call->master_lock); + redis_update_onekey(call, rtpe_redis_write); + obj_put(call); + return NULL; +} + const char *call_delete_ng(ng_command_ctx_t *ctx) { g_auto(sdp_ng_flags) rtpp_flags; parser_arg output = ctx->resp; diff --git a/daemon/ice.c b/daemon/ice.c index dde134897..66c158765 100644 --- a/daemon/ice.c +++ b/daemon/ice.c @@ -53,6 +53,7 @@ static void __agent_schedule(struct ice_agent *ag, int64_t); static void __agent_schedule_abs(struct ice_agent *ag, int64_t tv); static void __agent_deschedule(struct ice_agent *ag); static void __ice_agent_free_components(struct ice_agent *ag); +static void __ice_pairings(struct ice_agent *ag); static void __agent_shutdown(struct ice_agent *ag); static void ice_agents_timer_run(void *); @@ -358,7 +359,7 @@ TYPED_GHASHTABLE_IMPL(foundation_ht, __found_hash, __found_equal, NULL, NULL) TYPED_GHASHTABLE_IMPL(priority_ht, g_direct_hash, g_direct_equal, NULL, NULL) TYPED_GHASHTABLE_IMPL(transaction_ht, __trans_hash, __trans_equal, NULL, NULL) -static void __ice_agent_initialize(struct ice_agent *ag) { +static void __ice_agent_initialize(struct ice_agent *ag, bool generate_credentials) { struct call_media *media = ag->media; call_t *call = ag->call; @@ -377,8 +378,10 @@ static void __ice_agent_initialize(struct ice_agent *ag) { ag->succeeded_pairs = g_tree_new(__pair_prio_cmp); ag->all_pairs = g_tree_new(__pair_prio_cmp); - create_random_ice_string(call, &ag->ufrag[1], 8); - create_random_ice_string(call, &ag->pwd[1], 26); + if (generate_credentials) { + create_random_ice_string(call, &ag->ufrag[1], 8); + create_random_ice_string(call, &ag->pwd[1], 26); + } atomic64_set_na(&ag->last_activity, rtpe_now); } @@ -394,7 +397,7 @@ static struct ice_agent *__ice_agent_new(struct call_media *media) { ag->media = media; mutex_init(&ag->lock); - __ice_agent_initialize(ag); + __ice_agent_initialize(ag, true); return ag; } @@ -423,7 +426,38 @@ static void __ice_reset(struct ice_agent *ag) { ZERO(ag->active_components); ag->start_nominating = 0; ag->tt_obj.last_run = 0; - __ice_agent_initialize(ag); + __ice_agent_initialize(ag, true); +} + +/* Restore the credentials of a completed exchange. Candidate-pair and + * nomination state is deliberately rebuilt rather than snapshotted. + * Called with the call lock held in W, hence agent doesn't need to be locked. */ +void ice_rollback(struct ice_agent *ag, const str ufrag[2], const str pwd[2], + const candidate_q *candidates) +{ + if (!ag) + return; + + __agent_deschedule(ag); + __ice_agent_free_components(ag); + ZERO(ag->active_components); + ag->start_nominating = 0; + ag->tt_obj.last_run = 0; + __ice_agent_initialize(ag, false); + memcpy(ag->ufrag, ufrag, sizeof(ag->ufrag)); + memcpy(ag->pwd, pwd, sizeof(ag->pwd)); + + for (__auto_type l = candidates->head; l; l = l->next) { + struct ice_candidate *copy = g_new(__typeof(*copy), 1); + *copy = *(struct ice_candidate *) l->data; + t_hash_table_insert(ag->candidate_hash, copy, copy); + t_hash_table_insert(ag->cand_prio_hash, GUINT_TO_POINTER(copy->priority), copy); + t_hash_table_insert(ag->foundation_hash, copy, copy); + t_queue_push_tail(&ag->remote_candidates, copy); + ag->active_components = MAX(ag->active_components, copy->component_id); + } + __ice_pairings(ag); + ice_start(ag); } /* if the other side did a restart */ diff --git a/daemon/redis.c b/daemon/redis.c index c4547afe4..6a8440ce1 100644 --- a/daemon/redis.c +++ b/daemon/redis.c @@ -21,6 +21,7 @@ #include "compat.h" #include "helpers.h" #include "call.h" +#include "call_checkpoint.h" #include "log_d.h" #include "str.h" #include "crypto.h" @@ -1080,6 +1081,14 @@ static const char *json_get_hash_iter(const ng_parser_t *parser, str *key, parse return NULL; } +int redis_hash_from_parser(struct redis_hash *out, const ng_parser_t *parser, parser_arg dict) { + out->ht = g_hash_table_new_full(g_str_hash, g_str_equal, free, free); + if (!out->ht) + return -1; + parser->dict_iter(parser, dict, json_get_hash_iter, out->ht); + return 0; +} + static int json_get_hash(struct redis_hash *out, const char *key, unsigned int id, parser_arg root) { @@ -1103,16 +1112,10 @@ static int json_get_hash(struct redis_hash *out, return -1; } - out->ht = g_hash_table_new_full(g_str_hash, g_str_equal, free, free); - if (!out->ht) - return -1; - - redis_parser->dict_iter(redis_parser, dict, json_get_hash_iter, out->ht); - - return 0; + return redis_hash_from_parser(out, redis_parser, dict); } -static void json_destroy_hash(struct redis_hash *rh) { +void redis_hash_destroy(struct redis_hash *rh) { g_hash_table_destroy(rh->ht); } @@ -1120,7 +1123,7 @@ static void json_destroy_list(struct redis_list *rl) { unsigned int i; for (i = 0; i < rl->len; i++) { - json_destroy_hash(&rl->rh[i]); + redis_hash_destroy(&rl->rh[i]); } free(rl->rh); free(rl->ptrs); @@ -1369,7 +1372,7 @@ static int json_get_list_hash(struct redis_list *out, free(out->ptrs); while (i) { i--; - json_destroy_hash(&out->rh[i]); + redis_hash_destroy(&out->rh[i]); } err1: free(out->rh); @@ -1418,7 +1421,7 @@ static int redis_hash_get_sdes_params1(struct crypto_params *out, const struct r rlog(LOG_ERR, "Crypto params error: %s", err); return -1; } -static int redis_hash_get_sdes_params(sdes_q *out, const struct redis_hash *h, const char *k) { +int redis_decode_sdes_params(sdes_q *out, const struct redis_hash *h, const char *k) { char key[32], tagkey[64]; const char *kk = k; unsigned int tag; @@ -1445,6 +1448,17 @@ static int redis_hash_get_sdes_params(sdes_q *out, const struct redis_hash *h, c return 0; } +int redis_decode_dtls_fingerprint(struct dtls_fingerprint *out, const struct redis_hash *h) { + str hash; + if (redis_hash_get_str(&hash, h, "hash_func")) + return 0; + out->hash_func = dtls_find_hash_func(&hash); + if (!out->hash_func || redis_hash_get_c_buf_f(out->digest, h, "fingerprint")) + return -1; + out->digest_len = out->hash_func->num_bytes; + return 0; +} + static int redis_sfds(call_t *c, struct redis_list *sfds) { unsigned int i; str family, intf_name; @@ -1653,10 +1667,24 @@ static rtp_payload_type *rbl_cb_plts_g(str *s, struct redis_list *list, void *pt return pt; } -static int rbl_cb_plts_r(str *s, callback_arg_t dummy, struct redis_list *list, void *ptr) { - struct call_media *med = ptr; - codec_store_add_raw(&med->codecs, rbl_cb_plts_g(s, list, ptr)); - return 0; +static const char *redis_decode_codec_iter(str *value, unsigned int i, helper_arg arg) { + struct codec_store *store = arg.generic; + str *decoded = redis_parser->unescape(value->s, value->len); + struct call_media *media = store->media; + rtp_payload_type *pt = rbl_cb_plts_g(decoded, NULL, media); + g_free(decoded); + if (!pt) + return "invalid payload type"; + codec_store_add_raw(store, pt); + return NULL; +} + +int redis_decode_codec_store(const ng_parser_t *parser, parser_arg list, struct codec_store *store) { + const ng_parser_t *saved = redis_parser; + redis_parser = parser; + const char *err = parser->list_iter(parser, list, redis_decode_codec_iter, NULL, store); + redis_parser = saved; + return err ? -1 : 0; } static int json_medias(call_t *c, struct redis_list *medias, struct redis_list *tags, parser_arg arg) @@ -1708,9 +1736,9 @@ static int json_medias(call_t *c, struct redis_list *medias, struct redis_list * "media_flags")) return -1; - if (redis_hash_get_sdes_params(&med->sdes_in, rh, "sdes_in") < 0) + if (redis_decode_sdes_params(&med->sdes_in, rh, "sdes_in") < 0) return -1; - if (redis_hash_get_sdes_params(&med->sdes_out, rh, "sdes_out") < 0) + if (redis_decode_sdes_params(&med->sdes_out, rh, "sdes_out") < 0) return -1; /* bandwidth data is not critical */ @@ -1718,7 +1746,11 @@ static int json_medias(call_t *c, struct redis_list *medias, struct redis_list * med->sdp_media_bandwidth.rr = (!redis_hash_get_ld(&il, rh, "bandwidth_rr")) ? il : -1; med->sdp_media_bandwidth.rs = (!redis_hash_get_ld(&il, rh, "bandwidth_rs")) ? il : -1; - json_build_list_cb(NULL, c, "payload_types", i, NULL, rbl_cb_plts_r, med, arg); + char payload_key[64]; + snprintf(payload_key, sizeof(payload_key), "payload_types-%u", i); + parser_arg payloads = redis_parser->dict_get_expect(arg, payload_key, BENCODE_LIST); + if (payloads.gen && redis_decode_codec_store(redis_parser, payloads, &med->codecs)) + return -1; /* XXX dtls */ /* link monologue */ @@ -1918,6 +1950,8 @@ static int json_link_streams(call_t *c, struct redis_list *streams, if (json_build_list(&ps->sfds, c, "stream_sfds", i, sfds, arg)) return -1; + for (__auto_type sfd_link = ps->sfds.head; sfd_link; sfd_link = sfd_link->next) + stream_fd_inc(sfd_link->data); if (json_build_list(&q, c, "rtp_sinks", i, streams, arg)) return -1; @@ -2034,6 +2068,11 @@ static int json_link_maps(call_t *c, struct redis_list *maps, if (json_build_list_cb(&em->intf_sfds, c, "map_sfds", em->unique_id, sfds, rbl_cb_intf_sfds, em, arg)) return -1; + for (__auto_type l = em->intf_sfds.head; l; l = l->next) { + struct sfd_intf_list *il = l->data; + for (__auto_type k = il->list.head; k; k = k->next) + stream_fd_inc(k->data); + } } return 0; } @@ -2225,6 +2264,13 @@ static void json_restore_call(struct redis *r, const str *callid, bool foreign) err = "failed to link maps"; if (json_link_maps(c, &maps, &sfds, root)) goto err8; + if (!redis_hash_get_str(&s, &call, "checkpoint-data") + && call_checkpoint_deserialize(c, &s)) { + /* Checkpoints are auxiliary state. A corrupt or unsupported payload must + * disable rollback atomically, not discard an otherwise usable call. */ + call_checkpoint_free_all(c); + ilog(LOG_WARNING, "Ignoring invalid checkpoint data while restoring call"); + } // presence of this key determines whether we were recording at all if (!redis_hash_get_str(&s, &call, "recording_meta_prefix")) { @@ -2263,7 +2309,7 @@ static void json_restore_call(struct redis *r, const str *callid, bool foreign) err4: json_destroy_list(&tags); err3: - json_destroy_hash(&call); + redis_hash_destroy(&call); err2: rwlock_unlock_w(&c->master_lock); err1: @@ -2458,6 +2504,23 @@ int redis_restore(struct redis *r, bool foreign, int db) { #define JSON_SET_SIMPLE_CSTR(a,d) parser->dict_add_str_dup(inner, a, STR_PTR(d)) #define JSON_SET_SIMPLE_STR(a,d) parser->dict_add_str_dup(inner, a, d) +void redis_encode_codec_store(const ng_parser_t *parser, parser_arg list, + const struct codec_store *store) +{ + char tmp[1024]; + for (__auto_type l = store->codec_prefs.head; l; l = l->next) { + rtp_payload_type *pt = l->data; + size_t len = rtpe_snprintf(tmp, sizeof(tmp), "%u/" STR_FORMAT "/%u/" STR_FORMAT + "/%i/%i/" STR_FORMAT "/" STR_FORMAT, + pt->payload_type, STR_FMT(&pt->encoding), pt->clock_rate, + STR_FMT(&pt->encoding_parameters), pt->bitrate, pt->ptime, + STR_FMT(&pt->format_parameters), STR_FMT(&pt->codec_opts)); + char encoded[len * 3 + 1]; + str value = parser->escape(encoded, tmp, len); + parser->list_add_str_dup(list, &value); + } +} + static void json_update_crypto_params(const ng_parser_t *parser, parser_arg inner, const char *key, struct crypto_params *p) { @@ -2476,9 +2539,8 @@ static void json_update_crypto_params(const ng_parser_t *parser, parser_arg inne JSON_SET_NSTRING_LEN("%s-mki", key, p->mki_len, (char *) p->mki); } -static int json_update_sdes_params(const ng_parser_t *parser, parser_arg inner, const char *pref, - unsigned int unique_id, - const char *k, sdes_q *q) +int redis_encode_sdes_params(const ng_parser_t *parser, parser_arg inner, const char *k, + const sdes_q *q) { unsigned int iter = 0; char keybuf[32]; @@ -2501,8 +2563,7 @@ static int json_update_sdes_params(const ng_parser_t *parser, parser_arg inner, return 0; } -static void json_update_dtls_fingerprint(const ng_parser_t *parser, parser_arg inner, const char *pref, - unsigned int unique_id, +void redis_encode_dtls_fingerprint(const ng_parser_t *parser, parser_arg inner, const struct dtls_fingerprint *f) { if (!f->hash_func) @@ -2542,6 +2603,19 @@ static str redis_encode_json(ng_parser_ctx_t *ctx, call_t *c, void **to_free) { JSON_SET_SIMPLE_STR("recording_metadata", &c->metadata); JSON_SET_SIMPLE("block_dtmf","%i", c->block_dtmf); JSON_SET_SIMPLE("call_flags", "%" PRIu64, atomic64_get_na(&c->call_flags)); + void *checkpoint_free = NULL; + str checkpoint_data = call_checkpoint_serialize(c, &checkpoint_free); + if (checkpoint_data.s) { + /* Unlike the bounded fields using JSON_SET_SIMPLE_LEN(), a checkpoint + * grows with the call graph and must not use a VLA on the Redis thread's + * configurable stack. parser->escape() can require up to 3x input. */ + char *encoded = g_malloc_n(checkpoint_data.len + 1, 3); + str encoded_str = parser->escape(encoded, checkpoint_data.s, + checkpoint_data.len); + parser->dict_add_str_dup(inner, "checkpoint-data", &encoded_str); + g_free(encoded); + g_free(checkpoint_free); + } if (c->created_from.len) JSON_SET_SIMPLE_STR("created_from", &c->created_from); @@ -2750,11 +2824,9 @@ static str redis_encode_json(ng_parser_ctx_t *ctx, call_t *c, void **to_free) { if (media->sdp_media_bandwidth.rs >= 0) JSON_SET_SIMPLE("bandwidth_rs","%ld", media->sdp_media_bandwidth.rs); - json_update_sdes_params(parser, inner, "media", media->unique_id, "sdes_in", - &media->sdes_in); - json_update_sdes_params(parser, inner, "media", media->unique_id, "sdes_out", - &media->sdes_out); - json_update_dtls_fingerprint(parser, inner, "media", media->unique_id, &media->fingerprint); + redis_encode_sdes_params(parser, inner, "sdes_in", &media->sdes_in); + redis_encode_sdes_params(parser, inner, "sdes_out", &media->sdes_out); + redis_encode_dtls_fingerprint(parser, inner, &media->fingerprint); } snprintf(tmp, sizeof(tmp), "streams-%u", media->unique_id); @@ -2773,15 +2845,7 @@ static str redis_encode_json(ng_parser_ctx_t *ctx, call_t *c, void **to_free) { snprintf(tmp, sizeof(tmp), "payload_types-%u", media->unique_id); inner = parser->dict_add_list_dup(root, tmp); - for (__auto_type m = media->codecs.codec_prefs.head; m; m = m->next) { - rtp_payload_type *pt = m->data; - JSON_ADD_LIST_STRING("%u/" STR_FORMAT "/%u/" STR_FORMAT "/%i/%i/" - STR_FORMAT "/" STR_FORMAT, - pt->payload_type, STR_FMT(&pt->encoding), - pt->clock_rate, STR_FMT(&pt->encoding_parameters), - pt->bitrate, pt->ptime, STR_FMT(&pt->format_parameters), - STR_FMT(&pt->codec_opts)); - } + redis_encode_codec_store(parser, inner, &media->codecs); // SSRC table dump // XXX needs fixing diff --git a/docs/ng_control_protocol.md b/docs/ng_control_protocol.md index 4eb291b1d..7eb3f1864 100644 --- a/docs/ng_control_protocol.md +++ b/docs/ng_control_protocol.md @@ -845,6 +845,12 @@ Optionally included keys are: dictionary. The response dictionary may also contain the optional key `message` with an explanatory string. No other key is required in the response dictionary. + * `rollback` + + Indicates that the controlling SIP proxy understands the `rollback` + message. If `rollback` is listed, *rtpengine* includes it in a + `supported` list in the response. + * `to-interface` Contains a string identifying the network interface pertaining to the @@ -1390,6 +1396,19 @@ Spaces in each string may be replaced by hyphens. address that has been learned before. If there's a mismatch, the packet will be dropped and not forwarded. +* `track state` + + Enables rollback checkpoints for the selected dialogue. Before applying an + `offer`, *rtpengine* records the media state from the last completed + offer/answer exchange. A successful `answer` commits the exchange and + discards the pending checkpoint. Once enabled, subsequent offers for the + dialogue are checkpointed without repeating the flag. The spelling + `track-state` is equivalent. + + Checkpointing is opt-in because a pending checkpoint retains media and + cryptographic configuration. Calls that do not use this flag do not retain + that state. + * `trickle ICE` Useful for `offer` messages when ICE is advertised to also advertise @@ -1806,8 +1825,12 @@ An example of a complete `offer` request dictionary could be (SDP body abbreviat "ICE": "force", "transport protocol": "RTP/SAVPF", "media address": "2001:d8::6f24:65b", "DTLS": "passive" } -A response message contains only the key `sdp` in addition to `result`, which contains the re-written -SDP body that the SIP proxy should insert into the SIP message. +A response message contains the key `sdp` in addition to `result`, which contains the re-written +SDP body that the SIP proxy should insert into the SIP message. If rollback +checkpointing is enabled, the response also contains `generation`. For an +`offer`, this is the generation of the exchange just opened. An `answer` +returns the generation it committed. If `supports` requested a supported +extension, the response can also contain a `supported` list. Example response: @@ -1850,6 +1873,56 @@ the `direction` key in the `answer` message. The reply message is identical as in the `offer` reply. +## `rollback` Message + +The `rollback` message restores a dialogue to the media state from its last +completed offer/answer exchange without deleting the call. It is intended for +use when an SDP offer has already been applied by *rtpengine* but the remote +endpoint subsequently rejects the signalling transaction. The signalling +element must issue the message explicitly; *rtpengine* does not observe SIP +transaction outcomes. + +The request must contain `call-id`, `from-tag`, and `to-tag`. It may also +contain: + +* `via-branch` + + Selects a particular fork using the same dialogue matching rules as other + NG messages. + +* `generation` + + An optional integer safety check. A pending checkpoint is restored only if + the value matches the generation returned by the corresponding `offer`. + A mismatch is a successful no-op, which prevents a delayed failure from + undoing a newer negotiation. + +The successful response contains `rolled-back`, set to `1` if a pending +checkpoint was restored or `0` if there was no matching pending checkpoint. +It also contains `generation` when checkpointing is enabled, identifying the +committed generation after the operation. This includes no-op responses caused +by no outstanding snapshot or a generation mismatch, as long as the dialogue +has checkpointing enabled. Repeating a successful rollback is therefore safe +and returns `rolled-back: 0` together with the unchanged committed generation. + +Rollback restores addresses and ports, codecs and payload mappings, transport +profile, media direction, and SDES configuration including keys. ICE +credentials are restored and connectivity checks reconstruct candidate-pair +and nomination state. DTLS fingerprint, TLS ID, and setup/role configuration +are restored, but the live OpenSSL association is not serializable and must +perform a new handshake. + +When calls are restored from Redis, checkpoint data is auxiliary: a malformed +or unsupported checkpoint payload is discarded in full while the call itself +is restored without rollback capability. + +Example request and response: + + { "command": "rollback", "call-id": "cfBXzDSZqhYNcXM", + "from-tag": "mS9rSAn0Cr", "to-tag": "yB3KjLa9", "generation": 2 } + + { "result": "ok", "rolled-back": 1, "generation": 1 } + ## `delete` Message The `delete` message must contain at least the keys `call-id` and `from-tag` and may optionally include diff --git a/include/call.h b/include/call.h index 470ff72db..45484bcf6 100644 --- a/include/call.h +++ b/include/call.h @@ -66,7 +66,7 @@ enum message_type { || (opmode == OP_UNSUBSCRIBE || opmode == OP_START_RECORDING) \ || (opmode == OP_STOP_RECORDING || opmode == OP_PAUSE_RECORDING) \ || (opmode == OP_INJECT_START || opmode == OP_INJECT_STOP) \ - || (opmode == OP_OTHER)) + || (opmode == OP_ROLLBACK || opmode == OP_OTHER)) #define IS_OP_DIRECTIONAL(opmode) \ ((opmode == OP_BLOCK_DTMF || opmode == OP_BLOCK_MEDIA) \ @@ -821,6 +821,7 @@ struct call { atomic64 call_flags; unsigned int update_iter; unsigned int media_rec_slots; + struct call_checkpoint *checkpoints; }; @@ -971,6 +972,7 @@ enum thread_looper_action call_timer(void); void __rtp_stats_update(rtp_stats_ht dst, struct codec_store *); bool __init_stream(struct packet_stream *ps); +void call_media_stop(struct call_media *); const rtp_payload_type *__rtp_stats_codec(struct call_media *m); diff --git a/include/call_checkpoint.h b/include/call_checkpoint.h new file mode 100644 index 000000000..8f0295a2a --- /dev/null +++ b/include/call_checkpoint.h @@ -0,0 +1,20 @@ +#ifndef __CALL_CHECKPOINT_H__ +#define __CALL_CHECKPOINT_H__ + +#include +#include +#include "str.h" + +struct call; +struct call_monologue; +struct call_checkpoint; + +uint64_t call_checkpoint_offer(struct call *, struct call_monologue *, struct call_monologue *, bool); +uint64_t call_checkpoint_answer(struct call *, struct call_monologue *, struct call_monologue *, bool *); +int call_checkpoint_rollback(struct call *, struct call_monologue *, struct call_monologue *, + uint64_t, bool, uint64_t *); +void call_checkpoint_free_all(struct call *); +str call_checkpoint_serialize(struct call *, void **); +int call_checkpoint_deserialize(struct call *, const str *); + +#endif diff --git a/include/call_flags.h b/include/call_flags.h index bbf5128a6..6129df446 100644 --- a/include/call_flags.h +++ b/include/call_flags.h @@ -288,6 +288,8 @@ RTPE_NG_FLAGS_STR_CASE_HT_PARAMS t38_no_iaf:1, t38_fec:1, supports_load_limit:1, + supports_rollback:1, + track_state:1, dtls_off:1, sdes_off:1, sdes_unencrypted_srtp:1, diff --git a/include/call_interfaces.h b/include/call_interfaces.h index 732b806c9..be15e484e 100644 --- a/include/call_interfaces.h +++ b/include/call_interfaces.h @@ -31,6 +31,7 @@ str call_query_udp(char **); const char *call_ping_ng(ng_command_ctx_t *ctx); const char *call_offer_ng(ng_command_ctx_t *, const char *); const char *call_answer_ng(ng_command_ctx_t *); +const char *call_rollback_ng(ng_command_ctx_t *); const char *call_delete_ng(ng_command_ctx_t *); const char *call_query_ng(ng_command_ctx_t *); const char *call_list_ng(ng_command_ctx_t *); diff --git a/include/control_ng.h b/include/control_ng.h index 54534be7b..7eaf91a83 100644 --- a/include/control_ng.h +++ b/include/control_ng.h @@ -5,6 +5,7 @@ X(OP_PING, "ping", "ping", "Ping", call_ping_ng) \ XA(OP_OFFER, "offer", "offer", "Offer", call_offer_ng) \ X(OP_ANSWER, "answer", "answer", "Answer", call_answer_ng) \ + X(OP_ROLLBACK, "rollback", "rollback", "Rollback", call_rollback_ng) \ X(OP_DELETE, "delete", "delete", "Delete", call_delete_ng) \ X(OP_QUERY, "query", "query", "Query", call_query_ng) \ X(OP_LIST, "list", "list", "List", call_list_ng) \ diff --git a/include/ice.h b/include/ice.h index 84708826e..3f7d4c080 100644 --- a/include/ice.h +++ b/include/ice.h @@ -159,6 +159,7 @@ void ice_update(struct ice_agent *, struct stream_params *, bool allow_restart); void ice_start(struct ice_agent *); void ice_shutdown(struct ice_agent **); void ice_restart(struct ice_agent *); +void ice_rollback(struct ice_agent *, const str [2], const str [2], const candidate_q *); void ice_candidates_free(candidate_q *); void ice_remote_candidates(candidate_q *, struct ice_agent *); diff --git a/include/redis.h b/include/redis.h index 550008ed1..4071b8eee 100644 --- a/include/redis.h +++ b/include/redis.h @@ -11,6 +11,9 @@ #include "helpers.h" #include "call.h" #include "str.h" +#include "control_ng.h" +#include "crypto.h" +#include "dtls.h" #define REDIS_RESTORE_NUM_THREADS 4 @@ -77,6 +80,16 @@ struct redis_list { void **ptrs; }; +int redis_encode_sdes_params(const ng_parser_t *, parser_arg, const char *, const sdes_q *); +void redis_encode_dtls_fingerprint(const ng_parser_t *, parser_arg, + const struct dtls_fingerprint *); +int redis_decode_sdes_params(sdes_q *, const struct redis_hash *, const char *); +int redis_decode_dtls_fingerprint(struct dtls_fingerprint *, const struct redis_hash *); +int redis_hash_from_parser(struct redis_hash *, const ng_parser_t *, parser_arg); +void redis_hash_destroy(struct redis_hash *); +void redis_encode_codec_store(const ng_parser_t *, parser_arg, const struct codec_store *); +int redis_decode_codec_store(const ng_parser_t *, parser_arg, struct codec_store *); + extern struct redis *rtpe_redis; extern struct redis *rtpe_redis_write; diff --git a/t/Makefile b/t/Makefile index acf9644b4..690e02b32 100644 --- a/t/Makefile +++ b/t/Makefile @@ -81,7 +81,8 @@ include ../lib/common.Makefile daemon-tests-templ-def daemon-tests-templ-def-offer daemon-tests-t38 daemon-tests-evs-dtx \ daemon-tests-transform daemon-tests-http daemon-tests-heuristic daemon-tests-asymmetric \ daemon-tests-dtx-no-shift daemon-tests-rtcp daemon-tests-redis-subscribe daemon-tests-rtp-ext \ - daemon-tests-bundle daemon-tests-dtls \ + daemon-tests-bundle daemon-tests-dtls daemon-tests-rollback \ + daemon-tests-rollback-redis \ daemon-tests-recording daemon-tests-create daemon-tests-alias \ daemon-tests-kernel-api @@ -132,7 +133,7 @@ daemon-tests: daemon-tests-main daemon-tests-jb daemon-tests-pubsub daemon-tests daemon-tests-sdp-orig-replacements daemon-tests-moh daemon-tests-evs-dtx daemon-tests-transform \ daemon-tests-transcode-config daemon-tests-codec-prefs daemon-tests-http daemon-tests-heuristic \ daemon-tests-asymmetric daemon-tests-rtcp daemon-tests-redis-subscribe daemon-tests-rtp-ext \ - daemon-tests-bundle daemon-tests-dtls \ + daemon-tests-bundle daemon-tests-dtls daemon-tests-rollback daemon-tests-rollback-redis \ daemon-tests-recording daemon-tests-create \ daemon-tests-dtx daemon-tests-dtx-cn daemon-tests-dtx-no-shift \ daemon-tests-alias \ @@ -144,6 +145,15 @@ daemon-test-deps: tests-preload.so daemon-tests-main: daemon-test-deps ./auto-test-helper "$@" perl -I../perl auto-daemon-tests.pl +daemon-tests-rollback: daemon-test-deps + ./auto-test-helper "$@" perl -I../perl auto-daemon-tests-rollback.pl + +daemon-tests-rollback-redis: daemon-test-deps + RTPE_REDIS_FORMAT=native ./auto-test-helper "$@-native" \ + perl -I../perl auto-daemon-tests-rollback-redis.pl + RTPE_REDIS_FORMAT=json ./auto-test-helper "$@-json" \ + perl -I../perl auto-daemon-tests-rollback-redis.pl + daemon-tests-jb: daemon-test-deps ./auto-test-helper "$@" perl -I../perl auto-daemon-tests-jb.pl @@ -359,6 +369,7 @@ test-stats: test-stats.o \ $(COMMONOBJS) \ ../daemon/arena.o \ ../daemon/audio_player.o \ + ../daemon/call_checkpoint.o \ ../daemon/call_flags.strhash.o \ ../daemon/call_interfaces.o \ ../daemon/call.o \ @@ -421,6 +432,7 @@ test-transcode: test-transcode.o \ $(COMMONOBJS) \ ../daemon/arena.o \ ../daemon/audio_player.o \ + ../daemon/call_checkpoint.o \ ../daemon/call_flags.strhash.o \ ../daemon/call_interfaces.o \ ../daemon/call.o \ diff --git a/t/auto-daemon-tests-rollback-redis.pl b/t/auto-daemon-tests-rollback-redis.pl new file mode 100644 index 000000000..9fbbffd5d --- /dev/null +++ b/t/auto-daemon-tests-rollback-redis.pl @@ -0,0 +1,315 @@ +#!/usr/bin/perl + +use strict; +use warnings; +use Bencode; +use JSON; +use NGCP::Rtpengine::AutoTest; +use Socket qw(AF_INET SOCK_STREAM sockaddr_in inet_aton); +use Test::More; + +my $redis_format = $ENV{RTPE_REDIS_FORMAT} // 'json'; + +# Fake Redis listener. Keep the last SET payload in this process and return it +# to the next daemon instance through KEYS/GET, modelling a takeover without an +# external redis-server dependency. +my $redis_listener; +socket($redis_listener, AF_INET, SOCK_STREAM, 0) or die; +bind($redis_listener, sockaddr_in(6379, inet_aton('203.0.113.42'))) or die; +listen($redis_listener, 10) or die; + +my ($redis_fd, $saved_call_id, $saved_record); +my $launches = 0; + +sub redis_read_exact { + my ($fd, $len) = @_; + my $buf = ''; + while (length($buf) < $len) { + my $part; + recv($fd, $part, $len - length($buf), 0) or die; + $buf .= $part; + } + return $buf; +} + +sub redis_read_line { + my ($fd) = @_; + my $buf = ''; + while ($buf !~ /\r\n\z/) { + $buf .= redis_read_exact($fd, 1); + } + $buf =~ s/\r\n\z//; + return $buf; +} + +sub redis_command { + my ($fd) = @_; + my $intro = redis_read_line($fd); + $intro =~ /^\*(\d+)\z/ or die "invalid Redis array: $intro"; + my @args; + for (1 .. $1) { + my $bulk = redis_read_line($fd); + $bulk =~ /^\$(\d+)\z/ or die "invalid Redis bulk string: $bulk"; + push @args, redis_read_exact($fd, $1); + redis_read_exact($fd, 2) eq "\r\n" or die "invalid Redis bulk terminator"; + } + return @args; +} + +sub redis_reply_simple { + my ($fd, $reply) = @_; + send($fd, "+$reply\r\n", 0) or die; +} + +sub redis_reply_bulk { + my ($fd, $value) = @_; + send($fd, '$' . length($value) . "\r\n$value\r\n", 0) or die; +} + +sub redis_preamble { + my ($name) = @_; + my $fd; + accept($fd, $redis_listener) or die; + + is_deeply([redis_command($fd)], ['PING'], "$name PING"); + redis_reply_simple($fd, 'PONG'); + is_deeply([redis_command($fd)], ['SELECT', '15'], "$name SELECT"); + redis_reply_simple($fd, 'OK'); + is_deeply([redis_command($fd)], ['INFO'], "$name INFO"); + redis_reply_bulk($fd, "role:master\r\n"); + is_deeply([redis_command($fd)], ['TYPE', 'calls'], "$name TYPE"); + redis_reply_simple($fd, 'none'); + + return $fd; +} + +$NGCP::Rtpengine::AutoTest::launch_cb = sub { + $launches++; + $redis_fd = redis_preamble("launch $launches"); + + is_deeply([redis_command($redis_fd)], ['PING'], "launch $launches restore PING"); + redis_reply_simple($redis_fd, 'PONG'); + is_deeply([redis_command($redis_fd)], ['KEYS', '*'], "launch $launches restore KEYS"); + + if (!defined $saved_record) { + send($redis_fd, "*0\r\n", 0) or die; + return; + } + + send($redis_fd, "*1\r\n\$" . length($saved_call_id) + . "\r\n$saved_call_id\r\n", 0) or die; + + my $restore_fd = redis_preamble("launch $launches restore worker"); + is_deeply([redis_command($restore_fd)], ['GET', $saved_call_id], + "launch $launches restore GET"); + redis_reply_bulk($restore_fd, $saved_record); +}; + +my ($expected_generation, $expected_pending_generation, $expected_pending); + +sub decode_record { + my ($record) = @_; + return $redis_format eq 'json' ? decode_json($record) : Bencode::bdecode($record, 1); +} + +sub encode_record { + my ($record) = @_; + return encode_json($record) if $redis_format eq 'json'; + my $as_strings; + $as_strings = sub { + my ($value) = @_; + return { map { $_ => $as_strings->($value->{$_}) } keys %$value } + if ref($value) eq 'HASH'; + return [ map { $as_strings->($_) } @$value ] if ref($value) eq 'ARRAY'; + my $copy = $value; + return \$copy; + }; + return Bencode::bencode($as_strings->($record)); +} + +sub checkpoint_with_invalid_generation_type { + my ($record) = @_; + my $decoded = decode_record($record); + my $checkpoint_data = $decoded->{json}{'checkpoint-data'}; + $checkpoint_data =~ s/%([0-9a-fA-F]{2})/chr(hex($1))/ge + if $redis_format eq 'json'; + my $checkpoint = decode_json($checkpoint_data); + $checkpoint->{checkpoints}[0]{generation} = 'banana'; + $checkpoint_data = encode_json($checkpoint); + $checkpoint_data =~ s/([^A-Za-z0-9_.~-])/sprintf('%%%02X', ord($1))/ge + if $redis_format eq 'json'; + $decoded->{json}{'checkpoint-data'} = $checkpoint_data; + return encode_record($decoded); +} + +sub inspect_checkpoint { + my ($record) = @_; + my $decoded = decode_record($record); + my $checkpoint_data = $decoded->{json}{'checkpoint-data'}; + ok(defined $checkpoint_data, "$redis_format record contains checkpoint-data"); + $checkpoint_data =~ s/%([0-9a-fA-F]{2})/chr(hex($1))/ge + if $redis_format eq 'json'; + + my $checkpoint = decode_json($checkpoint_data); + is($checkpoint->{version}, 1, "$redis_format checkpoint payload version"); + is(scalar @{$checkpoint->{checkpoints}}, 1, + "$redis_format record contains one dialogue checkpoint"); + my $dialogue = $checkpoint->{checkpoints}[0]; + is($dialogue->{generation}, $expected_generation, + "$redis_format committed generation serialized"); + is($dialogue->{'pending-generation'}, $expected_pending_generation, + "$redis_format pending generation serialized"); + is($dialogue->{pending} ? 1 : 0, $expected_pending, + "$redis_format pending state serialized"); +} + +sub expect_redis_update { + my ($generation, $pending_generation, $pending) = @_; + ($expected_generation, $expected_pending_generation, $expected_pending) + = ($generation, $pending_generation, $pending); + + $NGCP::Rtpengine::req_cb = sub { + is_deeply([redis_command($redis_fd)], ['PING'], "$redis_format update PING"); + redis_reply_simple($redis_fd, 'PONG'); + + my @set = redis_command($redis_fd); + is($set[0], 'SET', "$redis_format update uses SET"); + is($set[3], 'EX', "$redis_format update carries expiry"); + like($set[4], qr/^\d+\z/, "$redis_format update expiry is numeric"); + ($saved_call_id, $saved_record) = @set[1, 2]; + inspect_checkpoint($saved_record); + redis_reply_simple($redis_fd, 'OK'); + }; +} + +sub expect_redis_update_without_checkpoint { + $NGCP::Rtpengine::req_cb = sub { + is_deeply([redis_command($redis_fd)], ['PING'], + "$redis_format invalid-checkpoint update PING"); + redis_reply_simple($redis_fd, 'PONG'); + my @set = redis_command($redis_fd); + is($set[0], 'SET', "$redis_format invalid-checkpoint update uses SET"); + is($set[3], 'EX', "$redis_format invalid-checkpoint update carries expiry"); + my $decoded = decode_record($set[2]); + ok(!exists $decoded->{json}{'checkpoint-data'}, + "$redis_format invalid checkpoint is omitted from the restored call record"); + redis_reply_simple($redis_fd, 'OK'); + }; +} + +sub redis_rtpe_req { + my ($generation, $pending_generation, $pending, @request) = @_; + expect_redis_update($generation, $pending_generation, $pending); + my $response = rtpe_req(@request); + $NGCP::Rtpengine::req_cb = undef; + return $response; +} + +sub sdp { + my ($address, $port, $ufrag, $pwd, $key, $direction) = @_; + return "v=0\r\no=- 2 2 IN IP4 $address\r\ns=rollback-redis-secure\r\n" + . "c=IN IP4 $address\r\nt=0 0\r\nm=audio $port RTP/SAVP 0\r\n" + . "a=rtpmap:0 PCMU/8000\r\na=$direction\r\n" + . "a=ice-ufrag:$ufrag\r\na=ice-pwd:$pwd\r\n" + . "a=candidate:1 1 UDP 2130706431 $address $port typ host\r\n" + . "a=crypto:1 AES_CM_128_HMAC_SHA1_80 inline:$key\r\n"; +} + +sub secure_parameters { + my ($body) = @_; + my @parameters = $body =~ /^(a=(?:ice-ufrag|ice-pwd):.*|a=crypto:1 .*)$/mg; + return \@parameters; +} + +my @daemon_args = (qw(--config-file=none -t -1 -i foo/203.0.113.1 + -n 2233 -c 12346 -f -L 7 -E --redis-num-threads=1), + "--redis=203.0.113.42:6379/15", "--redis-format=$redis_format"); +$NGCP::Rtpengine::AutoTest::port = 2233; +autotest_start(@daemon_args) or die; + +new_call; +my ($call_id, $from_tag, $to_tag) = (cid(), ft(), tt()); +my $via_branch = 'rollback-redis-branch'; +my $response = redis_rtpe_req(0, 1, 1, 'offer', 'tracked Redis offer', { + 'from-tag' => $from_tag, 'via-branch' => $via_branch, flags => ['track-state'], + sdp => sdp('198.51.100.80', 12000, 'oldRedisUfrag', + 'oldRedisPassword012345678', 'MTIzNDU2Nzg5MDEyMzQ1Njc4OTAxMjM0NTY3ODkw', 'sendrecv'), +}); +my $secure_parameters = secure_parameters($response->{sdp}); +is($response->{generation}, 1, 'initial Redis generation opened'); + +redis_rtpe_req(1, 0, 0, 'answer', 'tracked Redis answer', { + 'from-tag' => $from_tag, 'to-tag' => $to_tag, 'via-branch' => $via_branch, + sdp => sdp('198.51.100.81', 12010, 'answerRedisUfrag', + 'answerRedisPassword012345', 'QUJDREVGR0hJSktMTU5PUFFSU1RVVldYWVo5ODc2', 'sendrecv'), +}); +$response = redis_rtpe_req(1, 2, 1, 'offer', 'offer later rejected by the far end', { + 'from-tag' => $from_tag, 'to-tag' => $to_tag, 'via-branch' => $via_branch, + sdp => sdp('198.51.100.82', 12020, 'newRedisUfrag', + 'newRedisPassword012345678', 'YWJjZGVmZ2hpamtsbW5vcHFyc3R1dnd4eXowMTIz', 'sendonly'), +}); +is($response->{generation}, 2, 'pending Redis generation stored'); +my $pending_record = $saved_record; + +NGCP::Rtpengine::AutoTest::shut_rtpe(); +$NGCP::Rtpengine::req_cb = undef; +autotest_start(@daemon_args) or die; + +$response = redis_rtpe_req(1, 0, 0, 'rollback', 'rollback after Redis takeover', { + 'call-id' => $call_id, 'from-tag' => $from_tag, 'to-tag' => $to_tag, + 'via-branch' => $via_branch, generation => 2, +}); +is($response->{'rolled-back'}, 1, 'pending checkpoint survives Redis takeover'); +is($response->{generation}, 1, 'committed generation survives Redis takeover'); + +my $query = rtpe_req('query', 'query rolled-back Redis call', {'call-id' => $call_id}); +ok($query->{tags}{$from_tag}{medias}[0]{streams}[0]{'local port'}, + 'selected media socket survives takeover rollback'); +is($query->{tags}{$from_tag}{medias}[0]{streams}[0]{endpoint}{address}, '198.51.100.80', + 'remote media endpoint survives takeover rollback'); + +$response = redis_rtpe_req(1, 2, 1, 'offer', 'verify restored Redis media state', { + 'call-id' => $call_id, 'from-tag' => $from_tag, 'to-tag' => $to_tag, + 'via-branch' => $via_branch, + sdp => sdp('198.51.100.80', 12000, 'oldRedisUfrag', + 'oldRedisPassword012345678', 'MTIzNDU2Nzg5MDEyMzQ1Njc4OTAxMjM0NTY3ODkw', 'sendrecv'), +}); +is($response->{generation}, 2, 'tracking remains enabled after takeover rollback'); +is_deeply(secure_parameters($response->{sdp}), $secure_parameters, + 'ICE credentials and committed SDES key survive takeover rollback'); + +redis_rtpe_req(2, 0, 0, 'answer', 'commit exchange after takeover rollback', { + 'call-id' => $call_id, 'from-tag' => $from_tag, 'to-tag' => $to_tag, + 'via-branch' => $via_branch, + sdp => sdp('198.51.100.81', 12010, 'answerRedisUfrag', + 'answerRedisPassword012345', 'QUJDREVGR0hJSktMTU5PUFFSU1RVVldYWVo5ODc2', 'sendrecv'), +}); + +NGCP::Rtpengine::AutoTest::shut_rtpe(); +$NGCP::Rtpengine::req_cb = undef; +autotest_start(@daemon_args) or die; +$response = redis_rtpe_req(2, 0, 0, 'rollback', 'rollback after committed state crosses Redis', { + 'call-id' => $call_id, 'from-tag' => $from_tag, 'to-tag' => $to_tag, + 'via-branch' => $via_branch, generation => 3, +}); +is($response->{'rolled-back'}, 0, 'committed checkpoint remains consumed in Redis'); +is($response->{generation}, 2, 'new committed generation persists in Redis'); + +NGCP::Rtpengine::AutoTest::shut_rtpe(); +$NGCP::Rtpengine::req_cb = undef; +$saved_record = checkpoint_with_invalid_generation_type($pending_record); +autotest_start(@daemon_args) or die; +$query = rtpe_req('query', 'query call restored without invalid checkpoint', { + 'call-id' => $call_id, +}); +ok($query->{tags}{$from_tag}, 'invalid checkpoint data does not discard the restored call'); +expect_redis_update_without_checkpoint(); +$response = rtpe_req('rollback', 'invalid checkpoint degrades to no rollback state', { + 'call-id' => $call_id, 'from-tag' => $from_tag, 'to-tag' => $to_tag, + 'via-branch' => $via_branch, generation => 2, +}); +$NGCP::Rtpengine::req_cb = undef; +is($response->{'rolled-back'}, 0, 'type-invalid checkpoint is discarded atomically'); + +NGCP::Rtpengine::AutoTest::shut_rtpe(); +done_testing; diff --git a/t/auto-daemon-tests-rollback.pl b/t/auto-daemon-tests-rollback.pl new file mode 100644 index 000000000..e9d125349 --- /dev/null +++ b/t/auto-daemon-tests-rollback.pl @@ -0,0 +1,412 @@ +#!/usr/bin/perl + +use strict; +use warnings; +use NGCP::Rtpengine::Test; +use NGCP::Rtpengine::AutoTest; +use NGCP::Rtpclient::ICE; +use NGCP::Rtpclient::DTLS; +use IO::Multiplex; +use Socket qw(MSG_DONTWAIT); +use Test::More; + +autotest_start(qw(--config-file=none -t -1 -i 203.0.113.1 + -n 2223 -c 12345 -f -L 7 -E -u 2222)) or die; + +sub sdp { + my ($address, $port, $payload, $direction, $extra_media) = @_; + my $codec = $payload == 8 ? 'PCMA' : 'PCMU'; + my $ret = "v=0\r\no=- 1 1 IN IP4 $address\r\ns=rollback\r\n" + . "c=IN IP4 $address\r\nt=0 0\r\nm=audio $port RTP/AVP $payload\r\n" + . "a=rtpmap:$payload $codec/8000\r\na=$direction\r\n"; + $ret .= "m=video " . ($port + 2) . " RTP/AVP 96\r\n" + . "a=rtpmap:96 VP8/90000\r\na=$direction\r\n" if $extra_media; + return $ret; +} + +sub rollback { + my ($generation, $via_branch) = @_; + my %req = ('call-id' => cid(), 'from-tag' => ft(), 'to-tag' => tt()); + $req{generation} = $generation if defined $generation; + $req{'via-branch'} = $via_branch if defined $via_branch; + return rtpe_req('rollback', 'rollback state', \%req); +} + +sub negotiated_tags { + my ($state) = @_; + # __fill_stream() refreshes ps->last_packet_us on each offer. Query exposes + # that value, truncated to seconds, as both "last packet" and "last user + # packet". It is liveness state rather than negotiated media state. + # Whole-tag comparisons are safe only while no media has flowed: once packets + # arrive, the per-media SSRC lists also contain traffic-derived statistics. + # Such tests must compare only the negotiated fields they intend to restore. + for my $tag (values %{$state->{tags}}) { + for my $media (@{$tag->{medias}}) { + for my $stream (@{$media->{streams}}) { + delete $stream->{'last packet'}; + delete $stream->{'last user packet'}; + } + } + } + return $state->{tags}; +} + +sub secure_sdp { + my ($address, $port, $ufrag, $pwd, $key, $direction) = @_; + return "v=0\r\no=- 2 2 IN IP4 $address\r\ns=rollback-secure\r\n" + . "c=IN IP4 $address\r\nt=0 0\r\nm=audio $port RTP/SAVP 0\r\n" + . "a=rtpmap:0 PCMU/8000\r\na=$direction\r\n" + . "a=ice-ufrag:$ufrag\r\na=ice-pwd:$pwd\r\n" + . "a=candidate:1 1 UDP 2130706431 $address $port typ host\r\n" + . "a=candidate:1 2 UDP 2130706430 $address " . ($port + 1) . " typ host\r\n" + . "a=crypto:1 AES_CM_128_HMAC_SHA1_80 inline:$key\r\n"; +} + +sub dtls_sdp { + my ($address, $port, $fingerprint, $tls_id, $setup) = @_; + return "v=0\r\no=- 3 3 IN IP4 $address\r\ns=rollback-dtls\r\n" + . "c=IN IP4 $address\r\nt=0 0\r\nm=audio $port UDP/TLS/RTP/SAVP 0\r\n" + . "a=rtpmap:0 PCMU/8000\r\na=setup:$setup\r\n" + . "a=fingerprint:sha-256 $fingerprint\r\na=tls-id:$tls_id\r\n"; +} + +sub secure_parameters { + my ($sdp) = @_; + my @parameters = $sdp =~ /^(a=(?:ice-ufrag|ice-pwd):.*|a=crypto:1 .*)$/mg; + return \@parameters; +} + +my ($rollback_dtls, $rollback_dtls_mux, $rollback_dtls_connected, + @rollback_dtls_components); +my $rollback_dtls_output = sub { + my ($component, $data) = @_; + my ($socket, $port) = @{$rollback_dtls_components[$component]}; + snd($socket, $port, $data); +}; + +sub mux_input { + my ($self, $mux, $fh, $input) = @_; + my $peer = $mux->udp_peer($fh); + $rollback_dtls->input($fh, $input, $peer); + for my $component (@$rollback_dtls) { + return unless $component->{_connected}; + } + return if $rollback_dtls_connected; + $rollback_dtls_connected = 1; + pass('DTLS re-handshake succeeds with restored configuration'); + $mux->endloop(); +} + +sub consecutive_offer_rollback { + my ($use_generation, $label) = @_; + new_call; + rtpe_req('offer', "$label initial offer", { + 'from-tag' => ft(), flags => ['track-state'], + sdp => sdp('198.51.100.60', 10000, 0, 'sendrecv'), + }); + rtpe_req('answer', "$label initial answer", { + 'from-tag' => ft(), 'to-tag' => tt(), + sdp => sdp('198.51.100.61', 11000, 0, 'sendrecv'), + }); + my $committed_state = rtpe_req('query', "$label committed state", {}); + + my $first = rtpe_req('offer', "$label first pending offer", { + 'from-tag' => ft(), 'to-tag' => tt(), + sdp => sdp('198.51.100.62', 10020, 8, 'sendonly'), + }); + is($first->{generation}, 2, "$label first offer opens generation 2"); + my $second = rtpe_req('offer', "$label second pending offer", { + 'from-tag' => ft(), 'to-tag' => tt(), + sdp => sdp('198.51.100.63', 10030, 8, 'recvonly', 1), + }); + is($second->{generation}, $first->{generation}, + "$label second offer preserves the pending generation"); + + my $rolled_back = rollback($use_generation ? $first->{generation} : undef); + is($rolled_back->{'rolled-back'}, 1, "$label rollback consumes the original snapshot"); + my $restored_state = rtpe_req('query', "$label restored state", {}); + is_deeply(negotiated_tags($restored_state), negotiated_tags($committed_state), + "$label rollback restores the originally committed state"); +} + +new_call; +my $resp = rtpe_req('offer', 'tracked initial offer', { + 'from-tag' => ft(), flags => ['track-state'], supports => ['rollback'], + sdp => sdp('198.51.100.10', 4000, 0, 'sendrecv'), +}); +is($resp->{generation}, 1, 'initial offer opens generation 1'); +is_deeply($resp->{supported}, ['rollback'], 'rollback capability advertised'); + +$resp = rtpe_req('answer', 'tracked initial answer', { + 'from-tag' => ft(), 'to-tag' => tt(), + sdp => sdp('198.51.100.20', 5000, 0, 'sendrecv'), +}); +is($resp->{generation}, 1, 'initial answer commits generation 1'); +my $committed = rtpe_req('query', 'query committed state', {}); + +$resp = rtpe_req('offer', 'failed renegotiation', { + 'from-tag' => ft(), 'to-tag' => tt(), + sdp => sdp('198.51.100.11', 4010, 8, 'sendonly', 1), +}); +is($resp->{generation}, 2, 'renegotiation opens generation 2'); +$resp = rollback(99); +is($resp->{'rolled-back'}, 0, 'generation mismatch is a no-op'); +is($resp->{generation}, 1, 'generation mismatch returns the committed generation'); +$resp = rollback(2); +is($resp->{'rolled-back'}, 1, 'matching generation rolls back'); +is($resp->{generation}, 1, 'rollback returns committed generation'); +my $restored = rtpe_req('query', 'query restored state', {}); +is_deeply(negotiated_tags($restored), negotiated_tags($committed), + 'query media state restored'); +$resp = rollback(); +is($resp->{'rolled-back'}, 0, 'repeated rollback is a no-op'); + +rtpe_req('offer', 'completed renegotiation offer', { + 'from-tag' => ft(), 'to-tag' => tt(), + sdp => sdp('198.51.100.12', 4020, 8, 'sendrecv'), +}); +$resp = rtpe_req('answer', 'completed renegotiation answer', { + 'from-tag' => ft(), 'to-tag' => tt(), + sdp => sdp('198.51.100.22', 5020, 8, 'sendrecv'), +}); +is($resp->{generation}, 2, 'completed renegotiation commits generation 2'); +$resp = rollback(2); +is($resp->{'rolled-back'}, 0, 'completed exchange cannot be rolled back'); + +consecutive_offer_rollback(0, 'consecutive offers without generation'); +consecutive_offer_rollback(1, 'consecutive offers with first generation'); + +my ($subscription_a, $subscription_b, $subscription_sink) = new_call( + [qw(198.51.100.80 12000)], + [qw(198.51.100.80 12010)], + [qw(198.51.100.80 12020)], +); +my $subscription_offer = rtpe_req('offer', 'subscription rollback initial offer', { + 'from-tag' => ft(), flags => ['track-state'], + sdp => sdp('198.51.100.80', 12000, 0, 'sendrecv'), +}); +my ($subscription_port_a) = $subscription_offer->{sdp} =~ /^m=audio (\d+)/m; +my $subscription_answer = rtpe_req('answer', 'subscription rollback initial answer', { + 'from-tag' => ft(), 'to-tag' => tt(), + sdp => sdp('198.51.100.80', 12010, 0, 'sendrecv'), +}); +my ($subscription_port_b) = $subscription_answer->{sdp} =~ /^m=audio (\d+)/m; +my $subscription = rtpe_req('subscribe request', 'subscription before rollback', { + 'from-tag' => ft(), +}); +my ($subscription_port_sink) = $subscription->{sdp} =~ /^m=audio (\d+)/m; +ok($subscription_port_a && $subscription_port_b && $subscription_port_sink, + 'subscription relay ports captured'); +rtpe_req('subscribe answer', 'subscription before rollback answer', { + 'from-tag' => $subscription->{'from-tag'}, + 'to-tag' => $subscription->{'to-tag'}, + sdp => sdp('198.51.100.80', 12020, 0, 'recvonly'), +}); +snd($subscription_a, $subscription_port_b, + rtp(0, 5000, 8000, 0x7890, "\x55" x 160)); +rcv($subscription_b, $subscription_port_a, + rtpm(0, 5000, 8000, 0x7890, "\x55" x 160)); +rcv($subscription_sink, $subscription_port_sink, + rtpm(0, 5000, 8000, 0x7890, "\x55" x 160)); +my $subscription_committed = rtpe_req('query', 'subscription committed state', {}); +my $subscription_pending = rtpe_req('offer', 'subscription rejected renegotiation', { + 'from-tag' => ft(), 'to-tag' => tt(), + sdp => sdp('198.51.100.81', 12030, 8, 'sendonly'), +}); +my $subscription_rollback = rollback($subscription_pending->{generation}); +is($subscription_rollback->{'rolled-back'}, 1, + 'subscription call rolls back rejected renegotiation'); +my $subscription_restored = rtpe_req('query', 'subscription restored state', {}); +# Do not compare the complete query: restoring the committed remote endpoint is +# itself an endpoint change and therefore triggers call_stream_crypto_reset(), +# which intentionally resets the SSRC's ext_seq along with the crypto context. +for my $tag (keys %{$subscription_committed->{tags}}) { + is_deeply($subscription_restored->{tags}{$tag}{subscriptions}, + $subscription_committed->{tags}{$tag}{subscriptions}, + "rollback preserves subscriptions for $tag"); + is_deeply($subscription_restored->{tags}{$tag}{subscribers}, + $subscription_committed->{tags}{$tag}{subscribers}, + "rollback preserves subscribers for $tag"); +} +is_deeply($subscription_restored->{tags}{ft()}{medias}[0]{streams}[0]{endpoint}, + $subscription_committed->{tags}{ft()}{medias}[0]{streams}[0]{endpoint}, + 'active subscription resolves to the restored source endpoint'); +snd($subscription_a, $subscription_port_b, + rtp(0, 5001, 8160, 0x7890, "\x66" x 160)); +rcv($subscription_b, $subscription_port_a, + rtpm(0, 5001, 8160, 0x7890, "\x66" x 160)); +rcv($subscription_sink, $subscription_port_sink, + rtpm(0, 5001, 8160, 0x7890, "\x66" x 160)); + +new_call; +rtpe_req('offer', 'untracked offer', { + 'from-tag' => ft(), sdp => sdp('198.51.100.30', 6000, 0, 'sendrecv'), +}); +rtpe_req('answer', 'untracked answer', { + 'from-tag' => ft(), 'to-tag' => tt(), + sdp => sdp('198.51.100.40', 7000, 0, 'sendrecv'), +}); +$resp = rollback(); +is($resp->{'rolled-back'}, 0, 'untracked call has no checkpoint'); + +my ($secure_sock) = new_call([qw(198.51.100.50 8000)]); +my $secure_offer = secure_sdp('198.51.100.50', 8000, 'oldUfrag', + 'oldPassword0123456789012', 'MTIzNDU2Nzg5MDEyMzQ1Njc4OTAxMjM0NTY3ODkw', 'sendrecv'); +$resp = rtpe_req('offer', 'tracked ICE and SDES offer', { + 'from-tag' => ft(), 'via-branch' => 'rollback-branch', + flags => ['track-state'], sdp => $secure_offer, +}); +my $secure_parameters = secure_parameters($resp->{sdp}); +rtpe_req('answer', 'tracked ICE and SDES answer', { + 'from-tag' => ft(), 'to-tag' => tt(), + sdp => secure_sdp('198.51.100.51', 9000, 'answerUfrag', + 'answerPassword0123456789', 'QUJDREVGR0hJSktMTU5PUFFSU1RVVldYWVo5ODc2', 'sendrecv'), +}); +rtpe_req('offer', 'ICE restart and SDES rekey', { + 'from-tag' => ft(), 'to-tag' => tt(), 'via-branch' => 'rollback-branch', + sdp => secure_sdp('198.51.100.52', 8010, 'newUfrag', + 'newPassword0123456789012', 'YWJjZGVmZ2hpamtsbW5vcHFyc3R1dnd4eXowMTIz', 'sendonly'), +}); +my $branch_error = rtpe_raw_req({command => 'rollback', 'call-id' => cid(), + 'from-tag' => ft(), 'to-tag' => tt(), 'via-branch' => 'wrong-branch'}); +like($branch_error, qr/Unknown dialogue/, 'incorrect via-branch does not select the dialogue'); +$resp = rollback(2, 'rollback-branch'); +is($resp->{'rolled-back'}, 1, 'ICE restart and SDES rekey roll back'); +$resp = rtpe_req('offer', 'replay pre-restart ICE and SDES offer', { + 'from-tag' => ft(), 'to-tag' => tt(), 'via-branch' => 'rollback-branch', sdp => $secure_offer, +}); +is_deeply(secure_parameters($resp->{sdp}), $secure_parameters, + 'local ICE credentials and SDES key are restored'); +my ($secure_port) = $resp->{sdp} =~ /^m=audio (\d+)/m; +my ($local_ufrag) = $resp->{sdp} =~ /^a=ice-ufrag:(\S+)/m; +my ($local_pwd) = $resp->{sdp} =~ /^a=ice-pwd:(\S+)/m; +my @restored_check = rcv($secure_sock, -1, + qr/^\x00\x01\x00.\x21\x12\xa4\x42(............)/s); +snd($secure_sock, $secure_port, NGCP::Rtpclient::ICE::stun_succ( + $secure_port, $restored_check[2], 'oldPassword0123456789012')); +while (1) { + my $discard = ''; + last unless defined $secure_sock->recv($discard, 65535, MSG_DONTWAIT); +} +my ($stun_packet) = NGCP::Rtpclient::ICE::stun_req(0, 65527, 1, + 'oldUfrag', $local_ufrag, $local_pwd); +snd($secure_sock, $secure_port, $stun_packet); +rcv($secure_sock, -1, qr/^\x01\x01\x00.\x21\x12\xa4\x42/s); +pass('restored ICE credentials authenticate a connectivity check'); +rollback(2, 'rollback-branch'); + +my ($dtls_sock) = new_call([qw(198.51.100.55 9500)]); +$rollback_dtls_mux = IO::Multiplex->new(); +$rollback_dtls_mux->set_callback_object(__PACKAGE__); +$rollback_dtls = NGCP::Rtpclient::DTLS::Group->new($rollback_dtls_mux, + $rollback_dtls_output, [[$dtls_sock]]); +my $original_fingerprint = $rollback_dtls->[0]->fingerprint(); +my $original_dtls_offer = dtls_sdp('198.51.100.55', 9500, $original_fingerprint, + 'rollback-original', 'passive'); +rtpe_req('offer', 'tracked DTLS offer', { + 'from-tag' => ft(), flags => ['track-state'], SDES => 'off', + sdp => $original_dtls_offer, +}); +my $dtls_answer = rtpe_req('answer', 'tracked DTLS answer', { + 'from-tag' => ft(), 'to-tag' => tt(), SDES => 'off', + sdp => dtls_sdp('198.51.100.56', 9510, join(':', ('BB') x 32), + 'rollback-answer', 'active'), +}); +my ($restored_dtls_port) = $dtls_answer->{sdp} =~ /^m=audio (\d+)/m; +ok($restored_dtls_port, 'committed DTLS relay port captured'); +rtpe_req('offer', 'DTLS fingerprint and role change later rejected', { + 'from-tag' => ft(), 'to-tag' => tt(), SDES => 'off', + sdp => dtls_sdp('198.51.100.55', 9500, join(':', ('AA') x 32), + 'rollback-rejected', 'active'), +}); +my $dtls_rollback = rollback(2); +is($dtls_rollback->{'rolled-back'}, 1, 'DTLS configuration rolls back'); +$rollback_dtls_mux->add($dtls_sock); +@rollback_dtls_components = ([$dtls_sock, $restored_dtls_port]); +$rollback_dtls->accept(); +$rollback_dtls_mux->loop(); +rtpe_req('delete', 'delete DTLS rollback call', { + 'from-tag' => ft(), 'to-tag' => tt(), +}); + +new_call; +my $fork_from = ft(); +my $fork_a = 'fork-a-' . tt(); +my $fork_b = 'fork-b-' . tt(); +rtpe_req('offer', 'fork A initial offer', { + 'from-tag' => $fork_from, 'via-branch' => 'fork-a', flags => ['track-state'], + sdp => sdp('198.51.100.60', 10000, 0, 'sendrecv'), +}); +rtpe_req('answer', 'fork A initial answer', { + 'from-tag' => $fork_from, 'to-tag' => $fork_a, 'via-branch' => 'fork-a', + sdp => sdp('198.51.100.61', 10010, 0, 'sendrecv'), +}); +rtpe_req('offer', 'fork B initial offer', { + 'from-tag' => $fork_from, 'via-branch' => 'fork-b', flags => ['track-state'], + sdp => sdp('198.51.100.62', 10020, 0, 'sendrecv'), +}); +rtpe_req('answer', 'fork B initial answer', { + 'from-tag' => $fork_from, 'to-tag' => $fork_b, 'via-branch' => 'fork-b', + sdp => sdp('198.51.100.63', 10030, 0, 'sendrecv'), +}); +my $fork_a_offer = rtpe_req('offer', 'fork A rejected renegotiation', { + 'from-tag' => $fork_from, 'to-tag' => $fork_a, 'via-branch' => 'fork-a', + sdp => sdp('198.51.100.64', 10040, 8, 'sendonly'), +}); +my $fork_b_offer = rtpe_req('offer', 'fork B rejected renegotiation', { + 'from-tag' => $fork_from, 'to-tag' => $fork_b, 'via-branch' => 'fork-b', + sdp => sdp('198.51.100.65', 10050, 8, 'recvonly'), +}); +is($fork_a_offer->{generation}, 2, 'fork A has its own pending generation'); +is($fork_b_offer->{generation}, 2, 'fork B has its own pending generation'); +my $fork_error = rtpe_raw_req({command => 'rollback', 'call-id' => cid(), + 'from-tag' => $fork_from, 'to-tag' => $fork_b, 'via-branch' => 'fork-a'}); +like($fork_error, qr/Unknown dialogue/, 'branch and to-tag must identify the same fork'); +my $fork_a_rollback = rtpe_req('rollback', 'rollback fork A', { + 'from-tag' => $fork_from, 'to-tag' => $fork_a, 'via-branch' => 'fork-a', generation => 2, +}); +is($fork_a_rollback->{'rolled-back'}, 1, 'fork A rolls back independently'); +my $fork_a_repeat = rtpe_req('rollback', 'repeat rollback fork A', { + 'from-tag' => $fork_from, 'to-tag' => $fork_a, 'via-branch' => 'fork-a', generation => 2, +}); +is($fork_a_repeat->{'rolled-back'}, 0, 'fork A checkpoint was consumed'); +my $fork_b_rollback = rtpe_req('rollback', 'rollback fork B', { + 'from-tag' => $fork_from, 'to-tag' => $fork_b, 'via-branch' => 'fork-b', generation => 2, +}); +is($fork_b_rollback->{'rolled-back'}, 1, 'fork B checkpoint remains pending'); + +new_call; +rtpe_req('offer', 'stress initial offer', { + 'from-tag' => ft(), flags => ['track-state'], + sdp => sdp('198.51.100.70', 11000, 0, 'sendrecv'), +}); +rtpe_req('answer', 'stress initial answer', { + 'from-tag' => ft(), 'to-tag' => tt(), + sdp => sdp('198.51.100.71', 11010, 0, 'sendrecv'), +}); +for my $iteration (1 .. 10) { + my $generation = rtpe_req('offer', "stress offer $iteration", { + 'from-tag' => ft(), 'to-tag' => tt(), + sdp => sdp('198.51.100.72', 11020 + $iteration * 2, + $iteration % 2 ? 8 : 0, $iteration % 2 ? 'sendonly' : 'recvonly'), + }); + my $rolled_back = rollback($generation->{generation}); + is($rolled_back->{'rolled-back'}, 1, "stress rollback $iteration consumes checkpoint"); +} +my $pending_delete = rtpe_req('offer', 'leave checkpoint pending for delete', { + 'from-tag' => ft(), 'to-tag' => tt(), + sdp => sdp('198.51.100.73', 11100, 8, 'sendonly'), +}); +is($pending_delete->{generation}, 2, 'delete test leaves generation 2 pending'); +rtpe_req('delete', 'delete call after repeated rollbacks', { + 'from-tag' => ft(), 'to-tag' => tt(), +}); + +my $error = rtpe_raw_req({command => 'rollback', 'call-id' => 'unknown-call', + 'from-tag' => 'from', 'to-tag' => 'to'}); +like($error, qr/Unknown call-id/, 'unknown call is an error'); +$error = rtpe_raw_req({command => 'rollback', 'call-id' => cid(), + 'from-tag' => 'unknown-tag', 'to-tag' => tt()}); +like($error, qr/Unknown dialogue/, 'unknown dialogue is an error'); + +done_testing; diff --git a/t/test-stats.c b/t/test-stats.c index ae99d878a..2c563f4da 100644 --- a/t/test-stats.c +++ b/t/test-stats.c @@ -112,6 +112,13 @@ int main(void) { "answers_ps_max 0 150\n" "answers_ps_avg 0 150\n" "answer_count 0 150\n" + "rollback_time_min 0.000000 150\n" + "rollback_time_max 0.000000 150\n" + "rollback_time_avg 0.000000 150\n" + "rollbacks_ps_min 0 150\n" + "rollbacks_ps_max 0 150\n" + "rollbacks_ps_avg 0 150\n" + "rollback_count 0 150\n" "delete_time_min 0.000000 150\n" "delete_time_max 0.000000 150\n" "delete_time_avg 0.000000 150\n" @@ -596,6 +603,14 @@ int main(void) { "0.000000\n" "avganswerdelay\n" "0.000000\n" + "Min/Max/Avg rollback processing delay\n" + "0.000000/0.000000/0.000000 sec\n" + "minrollbackdelay\n" + "0.000000\n" + "maxrollbackdelay\n" + "0.000000\n" + "avgrollbackdelay\n" + "0.000000\n" "Min/Max/Avg delete processing delay\n" "0.000000/0.000000/0.000000 sec\n" "mindeletedelay\n" @@ -860,6 +875,14 @@ int main(void) { "0\n" "avganswerrequestrate\n" "0\n" + "Min/Max/Avg rollback requests per second\n" + "0/0/0 per sec\n" + "minrollbackrequestrate\n" + "0\n" + "maxrollbackrequestrate\n" + "0\n" + "avgrollbackrequestrate\n" + "0\n" "Min/Max/Avg delete requests per second\n" "0/0/0 per sec\n" "mindeleterequestrate\n" @@ -1263,7 +1286,7 @@ int main(void) { "{\n" "proxies\n" "[\n" - " Proxy | Ping | Offer | Answer | Delete | Query | List | StartRec | StopRec | PauseRec | StartFwd | StopFwd | BlkDTMF | UnblkDTMF | BlkMedia | UnblkMedia | PlayMedia | StopMedia | PlayDTMF | Stats | SlnMedia | UnslnMedia | Pub | SubReq | SubAns | Unsub | InjStart | InjStop | Conn | CLI | Trnsfm | Create | CrtAnsw | Mesh \n" + " Proxy | Ping | Offer | Answer | Rollback | Delete | Query | List | StartRec | StopRec | PauseRec | StartFwd | StopFwd | BlkDTMF | UnblkDTMF | BlkMedia | UnblkMedia | PlayMedia | StopMedia | PlayDTMF | Stats | SlnMedia | UnslnMedia | Pub | SubReq | SubAns | Unsub | InjStart | InjStop | Conn | CLI | Trnsfm | Create | CrtAnsw | Mesh \n" "\n" "]\n" "totalpingcount\n" @@ -1272,6 +1295,8 @@ int main(void) { "0\n" "totalanswercount\n" "0\n" + "totalrollbackcount\n" + "0\n" "totaldeletecount\n" "0\n" "totalquerycount\n" @@ -1365,7 +1390,7 @@ int main(void) { "size\n" "16777208\n" "used\n" - "464\n" + "472\n" "}\n" "]\n" "}\n" @@ -1401,6 +1426,13 @@ int main(void) { "answers_ps_max 0 150\n" "answers_ps_avg 0 150\n" "answer_count 0 150\n" + "rollback_time_min 0.000000 150\n" + "rollback_time_max 0.000000 150\n" + "rollback_time_avg 0.000000 150\n" + "rollbacks_ps_min 0 150\n" + "rollbacks_ps_max 0 150\n" + "rollbacks_ps_avg 0 150\n" + "rollback_count 0 150\n" "delete_time_min 0.000000 150\n" "delete_time_max 0.000000 150\n" "delete_time_avg 0.000000 150\n" @@ -1885,6 +1917,14 @@ int main(void) { "0.000000\n" "avganswerdelay\n" "0.000000\n" + "Min/Max/Avg rollback processing delay\n" + "0.000000/0.000000/0.000000 sec\n" + "minrollbackdelay\n" + "0.000000\n" + "maxrollbackdelay\n" + "0.000000\n" + "avgrollbackdelay\n" + "0.000000\n" "Min/Max/Avg delete processing delay\n" "0.000000/0.000000/0.000000 sec\n" "mindeletedelay\n" @@ -2149,6 +2189,14 @@ int main(void) { "0\n" "avganswerrequestrate\n" "0\n" + "Min/Max/Avg rollback requests per second\n" + "0/0/0 per sec\n" + "minrollbackrequestrate\n" + "0\n" + "maxrollbackrequestrate\n" + "0\n" + "avgrollbackrequestrate\n" + "0\n" "Min/Max/Avg delete requests per second\n" "0/0/0 per sec\n" "mindeleterequestrate\n" @@ -2552,7 +2600,7 @@ int main(void) { "{\n" "proxies\n" "[\n" - " Proxy | Ping | Offer | Answer | Delete | Query | List | StartRec | StopRec | PauseRec | StartFwd | StopFwd | BlkDTMF | UnblkDTMF | BlkMedia | UnblkMedia | PlayMedia | StopMedia | PlayDTMF | Stats | SlnMedia | UnslnMedia | Pub | SubReq | SubAns | Unsub | InjStart | InjStop | Conn | CLI | Trnsfm | Create | CrtAnsw | Mesh \n" + " Proxy | Ping | Offer | Answer | Rollback | Delete | Query | List | StartRec | StopRec | PauseRec | StartFwd | StopFwd | BlkDTMF | UnblkDTMF | BlkMedia | UnblkMedia | PlayMedia | StopMedia | PlayDTMF | Stats | SlnMedia | UnslnMedia | Pub | SubReq | SubAns | Unsub | InjStart | InjStop | Conn | CLI | Trnsfm | Create | CrtAnsw | Mesh \n" "\n" "]\n" "totalpingcount\n" @@ -2561,6 +2609,8 @@ int main(void) { "0\n" "totalanswercount\n" "0\n" + "totalrollbackcount\n" + "0\n" "totaldeletecount\n" "0\n" "totalquerycount\n" @@ -2654,7 +2704,7 @@ int main(void) { "size\n" "16777208\n" "used\n" - "464\n" + "472\n" "}\n" "]\n" "}\n" @@ -2687,6 +2737,13 @@ int main(void) { "answers_ps_max 0 150\n" "answers_ps_avg 0 150\n" "answer_count 1 150\n" + "rollback_time_min 0.000000 150\n" + "rollback_time_max 0.000000 150\n" + "rollback_time_avg 0.000000 150\n" + "rollbacks_ps_min 0 150\n" + "rollbacks_ps_max 0 150\n" + "rollbacks_ps_avg 0 150\n" + "rollback_count 0 150\n" "delete_time_min 0.000000 150\n" "delete_time_max 0.000000 150\n" "delete_time_avg 0.000000 150\n" @@ -3171,6 +3228,14 @@ int main(void) { "3.200000\n" "avganswerdelay\n" "3.200000\n" + "Min/Max/Avg rollback processing delay\n" + "0.000000/0.000000/0.000000 sec\n" + "minrollbackdelay\n" + "0.000000\n" + "maxrollbackdelay\n" + "0.000000\n" + "avgrollbackdelay\n" + "0.000000\n" "Min/Max/Avg delete processing delay\n" "0.000000/0.000000/0.000000 sec\n" "mindeletedelay\n" @@ -3435,6 +3500,14 @@ int main(void) { "0\n" "avganswerrequestrate\n" "0\n" + "Min/Max/Avg rollback requests per second\n" + "0/0/0 per sec\n" + "minrollbackrequestrate\n" + "0\n" + "maxrollbackrequestrate\n" + "0\n" + "avgrollbackrequestrate\n" + "0\n" "Min/Max/Avg delete requests per second\n" "0/0/0 per sec\n" "mindeleterequestrate\n" @@ -3838,7 +3911,7 @@ int main(void) { "{\n" "proxies\n" "[\n" - " Proxy | Ping | Offer | Answer | Delete | Query | List | StartRec | StopRec | PauseRec | StartFwd | StopFwd | BlkDTMF | UnblkDTMF | BlkMedia | UnblkMedia | PlayMedia | StopMedia | PlayDTMF | Stats | SlnMedia | UnslnMedia | Pub | SubReq | SubAns | Unsub | InjStart | InjStop | Conn | CLI | Trnsfm | Create | CrtAnsw | Mesh \n" + " Proxy | Ping | Offer | Answer | Rollback | Delete | Query | List | StartRec | StopRec | PauseRec | StartFwd | StopFwd | BlkDTMF | UnblkDTMF | BlkMedia | UnblkMedia | PlayMedia | StopMedia | PlayDTMF | Stats | SlnMedia | UnslnMedia | Pub | SubReq | SubAns | Unsub | InjStart | InjStop | Conn | CLI | Trnsfm | Create | CrtAnsw | Mesh \n" "\n" "]\n" "totalpingcount\n" @@ -3847,6 +3920,8 @@ int main(void) { "0\n" "totalanswercount\n" "0\n" + "totalrollbackcount\n" + "0\n" "totaldeletecount\n" "0\n" "totalquerycount\n" @@ -3940,7 +4015,7 @@ int main(void) { "size\n" "16777208\n" "used\n" - "464\n" + "472\n" "}\n" "]\n" "}\n" @@ -3992,6 +4067,13 @@ int main(void) { "answers_ps_max 0 157\n" "answers_ps_avg 0 157\n" "answer_count 1 157\n" + "rollback_time_min 0.000000 157\n" + "rollback_time_max 0.000000 157\n" + "rollback_time_avg 0.000000 157\n" + "rollbacks_ps_min 0 157\n" + "rollbacks_ps_max 0 157\n" + "rollbacks_ps_avg 0 157\n" + "rollback_count 0 157\n" "delete_time_min 0.000000 157\n" "delete_time_max 0.000000 157\n" "delete_time_avg 0.000000 157\n" @@ -4476,6 +4558,14 @@ int main(void) { "0.000000\n" "avganswerdelay\n" "0.000000\n" + "Min/Max/Avg rollback processing delay\n" + "0.000000/0.000000/0.000000 sec\n" + "minrollbackdelay\n" + "0.000000\n" + "maxrollbackdelay\n" + "0.000000\n" + "avgrollbackdelay\n" + "0.000000\n" "Min/Max/Avg delete processing delay\n" "0.000000/0.000000/0.000000 sec\n" "mindeletedelay\n" @@ -4740,6 +4830,14 @@ int main(void) { "0\n" "avganswerrequestrate\n" "0\n" + "Min/Max/Avg rollback requests per second\n" + "0/0/0 per sec\n" + "minrollbackrequestrate\n" + "0\n" + "maxrollbackrequestrate\n" + "0\n" + "avgrollbackrequestrate\n" + "0\n" "Min/Max/Avg delete requests per second\n" "0/0/0 per sec\n" "mindeleterequestrate\n" @@ -5143,7 +5241,7 @@ int main(void) { "{\n" "proxies\n" "[\n" - " Proxy | Ping | Offer | Answer | Delete | Query | List | StartRec | StopRec | PauseRec | StartFwd | StopFwd | BlkDTMF | UnblkDTMF | BlkMedia | UnblkMedia | PlayMedia | StopMedia | PlayDTMF | Stats | SlnMedia | UnslnMedia | Pub | SubReq | SubAns | Unsub | InjStart | InjStop | Conn | CLI | Trnsfm | Create | CrtAnsw | Mesh \n" + " Proxy | Ping | Offer | Answer | Rollback | Delete | Query | List | StartRec | StopRec | PauseRec | StartFwd | StopFwd | BlkDTMF | UnblkDTMF | BlkMedia | UnblkMedia | PlayMedia | StopMedia | PlayDTMF | Stats | SlnMedia | UnslnMedia | Pub | SubReq | SubAns | Unsub | InjStart | InjStop | Conn | CLI | Trnsfm | Create | CrtAnsw | Mesh \n" "\n" "]\n" "totalpingcount\n" @@ -5152,6 +5250,8 @@ int main(void) { "0\n" "totalanswercount\n" "0\n" + "totalrollbackcount\n" + "0\n" "totaldeletecount\n" "0\n" "totalquerycount\n" @@ -5245,7 +5345,7 @@ int main(void) { "size\n" "16777208\n" "used\n" - "464\n" + "472\n" "}\n" "]\n" "}\n" @@ -5286,6 +5386,13 @@ int main(void) { "answers_ps_max 0 157\n" "answers_ps_avg 0 157\n" "answer_count 1 157\n" + "rollback_time_min 0.000000 157\n" + "rollback_time_max 0.000000 157\n" + "rollback_time_avg 0.000000 157\n" + "rollbacks_ps_min 0 157\n" + "rollbacks_ps_max 0 157\n" + "rollbacks_ps_avg 0 157\n" + "rollback_count 0 157\n" "delete_time_min 0.000000 157\n" "delete_time_max 0.000000 157\n" "delete_time_avg 0.000000 157\n" @@ -5770,6 +5877,14 @@ int main(void) { "0.000000\n" "avganswerdelay\n" "0.000000\n" + "Min/Max/Avg rollback processing delay\n" + "0.000000/0.000000/0.000000 sec\n" + "minrollbackdelay\n" + "0.000000\n" + "maxrollbackdelay\n" + "0.000000\n" + "avgrollbackdelay\n" + "0.000000\n" "Min/Max/Avg delete processing delay\n" "0.000000/0.000000/0.000000 sec\n" "mindeletedelay\n" @@ -6034,6 +6149,14 @@ int main(void) { "0\n" "avganswerrequestrate\n" "0\n" + "Min/Max/Avg rollback requests per second\n" + "0/0/0 per sec\n" + "minrollbackrequestrate\n" + "0\n" + "maxrollbackrequestrate\n" + "0\n" + "avgrollbackrequestrate\n" + "0\n" "Min/Max/Avg delete requests per second\n" "0/0/0 per sec\n" "mindeleterequestrate\n" @@ -6437,7 +6560,7 @@ int main(void) { "{\n" "proxies\n" "[\n" - " Proxy | Ping | Offer | Answer | Delete | Query | List | StartRec | StopRec | PauseRec | StartFwd | StopFwd | BlkDTMF | UnblkDTMF | BlkMedia | UnblkMedia | PlayMedia | StopMedia | PlayDTMF | Stats | SlnMedia | UnslnMedia | Pub | SubReq | SubAns | Unsub | InjStart | InjStop | Conn | CLI | Trnsfm | Create | CrtAnsw | Mesh \n" + " Proxy | Ping | Offer | Answer | Rollback | Delete | Query | List | StartRec | StopRec | PauseRec | StartFwd | StopFwd | BlkDTMF | UnblkDTMF | BlkMedia | UnblkMedia | PlayMedia | StopMedia | PlayDTMF | Stats | SlnMedia | UnslnMedia | Pub | SubReq | SubAns | Unsub | InjStart | InjStop | Conn | CLI | Trnsfm | Create | CrtAnsw | Mesh \n" "\n" "]\n" "totalpingcount\n" @@ -6446,6 +6569,8 @@ int main(void) { "0\n" "totalanswercount\n" "0\n" + "totalrollbackcount\n" + "0\n" "totaldeletecount\n" "0\n" "totalquerycount\n" @@ -6539,7 +6664,7 @@ int main(void) { "size\n" "16777208\n" "used\n" - "464\n" + "472\n" "}\n" "]\n" "}\n" @@ -6574,6 +6699,13 @@ int main(void) { "answers_ps_max 0 200\n" "answers_ps_avg 0 200\n" "answer_count 1 200\n" + "rollback_time_min 0.000000 200\n" + "rollback_time_max 0.000000 200\n" + "rollback_time_avg 0.000000 200\n" + "rollbacks_ps_min 0 200\n" + "rollbacks_ps_max 0 200\n" + "rollbacks_ps_avg 0 200\n" + "rollback_count 0 200\n" "delete_time_min 0.000000 200\n" "delete_time_max 0.000000 200\n" "delete_time_avg 0.000000 200\n" @@ -7058,6 +7190,14 @@ int main(void) { "0.000000\n" "avganswerdelay\n" "0.000000\n" + "Min/Max/Avg rollback processing delay\n" + "0.000000/0.000000/0.000000 sec\n" + "minrollbackdelay\n" + "0.000000\n" + "maxrollbackdelay\n" + "0.000000\n" + "avgrollbackdelay\n" + "0.000000\n" "Min/Max/Avg delete processing delay\n" "0.000000/0.000000/0.000000 sec\n" "mindeletedelay\n" @@ -7322,6 +7462,14 @@ int main(void) { "0\n" "avganswerrequestrate\n" "0\n" + "Min/Max/Avg rollback requests per second\n" + "0/0/0 per sec\n" + "minrollbackrequestrate\n" + "0\n" + "maxrollbackrequestrate\n" + "0\n" + "avgrollbackrequestrate\n" + "0\n" "Min/Max/Avg delete requests per second\n" "0/0/0 per sec\n" "mindeleterequestrate\n" @@ -7725,7 +7873,7 @@ int main(void) { "{\n" "proxies\n" "[\n" - " Proxy | Ping | Offer | Answer | Delete | Query | List | StartRec | StopRec | PauseRec | StartFwd | StopFwd | BlkDTMF | UnblkDTMF | BlkMedia | UnblkMedia | PlayMedia | StopMedia | PlayDTMF | Stats | SlnMedia | UnslnMedia | Pub | SubReq | SubAns | Unsub | InjStart | InjStop | Conn | CLI | Trnsfm | Create | CrtAnsw | Mesh \n" + " Proxy | Ping | Offer | Answer | Rollback | Delete | Query | List | StartRec | StopRec | PauseRec | StartFwd | StopFwd | BlkDTMF | UnblkDTMF | BlkMedia | UnblkMedia | PlayMedia | StopMedia | PlayDTMF | Stats | SlnMedia | UnslnMedia | Pub | SubReq | SubAns | Unsub | InjStart | InjStop | Conn | CLI | Trnsfm | Create | CrtAnsw | Mesh \n" "\n" "]\n" "totalpingcount\n" @@ -7734,6 +7882,8 @@ int main(void) { "0\n" "totalanswercount\n" "0\n" + "totalrollbackcount\n" + "0\n" "totaldeletecount\n" "0\n" "totalquerycount\n" @@ -7827,7 +7977,7 @@ int main(void) { "size\n" "16777208\n" "used\n" - "464\n" + "472\n" "}\n" "]\n" "}\n" @@ -7865,6 +8015,13 @@ int main(void) { "answers_ps_max 0 200\n" "answers_ps_avg 0 200\n" "answer_count 1 200\n" + "rollback_time_min 0.000000 200\n" + "rollback_time_max 0.000000 200\n" + "rollback_time_avg 0.000000 200\n" + "rollbacks_ps_min 0 200\n" + "rollbacks_ps_max 0 200\n" + "rollbacks_ps_avg 0 200\n" + "rollback_count 0 200\n" "delete_time_min 0.000000 200\n" "delete_time_max 0.000000 200\n" "delete_time_avg 0.000000 200\n" @@ -8349,6 +8506,14 @@ int main(void) { "0.000000\n" "avganswerdelay\n" "0.000000\n" + "Min/Max/Avg rollback processing delay\n" + "0.000000/0.000000/0.000000 sec\n" + "minrollbackdelay\n" + "0.000000\n" + "maxrollbackdelay\n" + "0.000000\n" + "avgrollbackdelay\n" + "0.000000\n" "Min/Max/Avg delete processing delay\n" "0.000000/0.000000/0.000000 sec\n" "mindeletedelay\n" @@ -8613,6 +8778,14 @@ int main(void) { "0\n" "avganswerrequestrate\n" "0\n" + "Min/Max/Avg rollback requests per second\n" + "0/0/0 per sec\n" + "minrollbackrequestrate\n" + "0\n" + "maxrollbackrequestrate\n" + "0\n" + "avgrollbackrequestrate\n" + "0\n" "Min/Max/Avg delete requests per second\n" "0/0/0 per sec\n" "mindeleterequestrate\n" @@ -9016,7 +9189,7 @@ int main(void) { "{\n" "proxies\n" "[\n" - " Proxy | Ping | Offer | Answer | Delete | Query | List | StartRec | StopRec | PauseRec | StartFwd | StopFwd | BlkDTMF | UnblkDTMF | BlkMedia | UnblkMedia | PlayMedia | StopMedia | PlayDTMF | Stats | SlnMedia | UnslnMedia | Pub | SubReq | SubAns | Unsub | InjStart | InjStop | Conn | CLI | Trnsfm | Create | CrtAnsw | Mesh \n" + " Proxy | Ping | Offer | Answer | Rollback | Delete | Query | List | StartRec | StopRec | PauseRec | StartFwd | StopFwd | BlkDTMF | UnblkDTMF | BlkMedia | UnblkMedia | PlayMedia | StopMedia | PlayDTMF | Stats | SlnMedia | UnslnMedia | Pub | SubReq | SubAns | Unsub | InjStart | InjStop | Conn | CLI | Trnsfm | Create | CrtAnsw | Mesh \n" "\n" "]\n" "totalpingcount\n" @@ -9025,6 +9198,8 @@ int main(void) { "0\n" "totalanswercount\n" "0\n" + "totalrollbackcount\n" + "0\n" "totaldeletecount\n" "0\n" "totalquerycount\n" @@ -9118,7 +9293,7 @@ int main(void) { "size\n" "16777208\n" "used\n" - "464\n" + "472\n" "}\n" "]\n" "}\n"