1#![forbid(unsafe_code)]
2
3use super::*;
4use crate::wire::ProtocolLimits;
5
6pub(crate) fn write_tpc_txn_switch_body(
11 writer: &mut TtcWriter,
12 operation: u32,
13 flags: u32,
14 timeout: u32,
15 xid: Option<&[u8]>,
16) {
17 writer.write_ub4(operation);
18 writer.write_u8(0); writer.write_ub4(0); if let Some(global_txn_id) = xid {
21 let mut xid_bytes = global_txn_id.to_vec();
25 xid_bytes.resize(128, 0);
26 writer.write_ub4(SESSIONLESS_FORMAT_ID);
27 writer.write_ub4(u32::try_from(global_txn_id.len()).unwrap_or(0)); writer.write_ub4(0); writer.write_u8(1); writer.write_ub4(u32::try_from(xid_bytes.len()).unwrap_or(0));
31 writer.write_ub4(flags);
32 writer.write_ub4(timeout);
33 writer.write_u8(1); writer.write_u8(1); writer.write_u8(1); writer.write_u8(0); writer.write_ub4(0); writer.write_u8(0); writer.write_ub4(0); writer.write_raw(&xid_bytes);
41 writer.write_ub4(0); } else {
43 writer.write_ub4(0); writer.write_ub4(0); writer.write_ub4(0); writer.write_u8(0); writer.write_ub4(0); writer.write_ub4(flags);
49 writer.write_ub4(timeout);
50 writer.write_u8(1); writer.write_u8(1); writer.write_u8(1); writer.write_u8(0); writer.write_ub4(0); writer.write_u8(0); writer.write_ub4(0); writer.write_ub4(0); }
59}
60
61pub fn build_tpc_txn_switch_payload_with_seq(
66 seq_num: u8,
67 token_num: u64,
68 operation: u32,
69 flags: u32,
70 timeout: u32,
71 xid: Option<&[u8]>,
72) -> Vec<u8> {
73 let mut writer = TtcWriter::new();
74 writer.write_function_code_with_seq(TNS_FUNC_TPC_TXN_SWITCH, seq_num);
75 writer.write_ub8(token_num);
76 write_tpc_txn_switch_body(&mut writer, operation, flags, timeout, xid);
77 writer.into_bytes()
78}
79
80pub fn build_sessionless_piggyback(
87 seq_num: u8,
88 token_num: u64,
89 operation: u32,
90 flags: u32,
91 timeout: u32,
92 xid: Option<&[u8]>,
93) -> Vec<u8> {
94 let mut writer = TtcWriter::new();
95 writer.write_u8(TNS_MSG_TYPE_PIGGYBACK);
96 writer.write_u8(TNS_FUNC_TPC_TXN_SWITCH);
97 writer.write_u8(seq_num);
98 writer.write_ub8(token_num);
99 write_tpc_txn_switch_body(&mut writer, operation, flags, timeout, xid);
100 writer.into_bytes()
101}
102
103pub fn decode_sessionless_txn_state(binary: &[u8]) -> Result<Option<SessionlessTxnState>> {
108 if binary.len() < 2 {
109 return Err(ProtocolError::TtcDecode("short sessionless txn state"));
110 }
111 let state = binary[binary.len() - 2];
112 let sync_version = binary[binary.len() - 1];
113 if sync_version != 1 {
114 return Err(ProtocolError::TtcDecode("unknown transaction sync version"));
115 }
116 if state & TNS_TPC_TXNID_SYNC_UNSET != 0 {
117 Ok(Some(SessionlessTxnState::Unset))
118 } else if state & TNS_TPC_TXNID_SYNC_SET != 0 {
119 Ok(Some(SessionlessTxnState::Set {
120 started_on_server: state & TNS_TPC_TXNID_SYNC_SERVER != 0,
121 }))
122 } else {
123 Ok(None)
124 }
125}
126
127pub fn parse_tpc_txn_switch_response(
132 payload: &[u8],
133 capabilities: ClientCapabilities,
134) -> Result<Option<SessionlessTxnState>> {
135 parse_tpc_txn_switch_response_with_limits(payload, capabilities, ProtocolLimits::DEFAULT)
136}
137
138pub fn parse_tpc_txn_switch_response_with_limits(
139 payload: &[u8],
140 capabilities: ClientCapabilities,
141 limits: ProtocolLimits,
142) -> Result<Option<SessionlessTxnState>> {
143 let mut reader = TtcReader::with_limits(payload, limits)?;
144 let mut state = None;
145 while reader.remaining() > 0 {
146 let message_type = reader.read_u8()?;
147 match message_type {
148 0 => {}
149 TNS_MSG_TYPE_STATUS => {
150 let _call_status = reader.read_ub4()?;
151 let _seq = reader.read_ub2()?;
152 }
153 TNS_MSG_TYPE_PARAMETER => {
154 let _application_value = reader.read_ub4()?;
157 let context_len = reader.read_ub2()?;
158 if context_len > 0 {
159 reader.skip(usize::from(context_len))?;
160 }
161 }
162 TNS_MSG_TYPE_SERVER_SIDE_PIGGYBACK => {
163 if let Some(update) = skip_server_side_piggyback(&mut reader)? {
164 state = Some(update);
165 }
166 }
167 TNS_MSG_TYPE_END_OF_RESPONSE => break,
168 TNS_MSG_TYPE_ERROR => {
169 let info = parse_server_error_info(&mut reader, capabilities.ttc_field_version)?;
170 if info.number != 0 {
171 return Err(ProtocolError::ServerErrorInfo(Box::new(
172 info.into_details(),
173 )));
174 }
175 }
176 _ => break,
177 }
178 }
179 Ok(state)
180}
181
182pub fn build_begin_pipeline_piggyback(seq_num: u8, token_num: u64, pipeline_mode: u8) -> Vec<u8> {
190 let mut writer = TtcWriter::new();
191 writer.write_u8(TNS_MSG_TYPE_PIGGYBACK);
192 writer.write_u8(TNS_FUNC_PIPELINE_BEGIN);
193 writer.write_u8(seq_num);
194 writer.write_ub8(token_num);
195 writer.write_ub2(0); writer.write_u8(0); writer.write_u8(pipeline_mode);
198 writer.into_bytes()
199}
200
201pub fn build_end_pipeline_payload_with_seq(seq_num: u8) -> Vec<u8> {
206 let mut writer = TtcWriter::new();
207 writer.write_function_code_with_seq(TNS_FUNC_PIPELINE_END, seq_num);
208 writer.write_ub8(0); writer.write_ub4(0); writer.into_bytes()
211}
212
213#[derive(Clone, Debug)]
217pub struct TpcXid<'a> {
218 pub format_id: u32,
219 pub global_transaction_id: &'a [u8],
220 pub branch_qualifier: &'a [u8],
221}
222
223fn write_xid_descriptor(writer: &mut TtcWriter, xid: Option<&TpcXid<'_>>) {
230 match xid {
231 Some(xid) => {
232 writer.write_ub4(xid.format_id);
233 writer.write_ub4(u32::try_from(xid.global_transaction_id.len()).unwrap_or(0));
234 writer.write_ub4(u32::try_from(xid.branch_qualifier.len()).unwrap_or(0));
235 writer.write_u8(1); writer.write_ub4(128); }
238 None => {
239 writer.write_ub4(0); writer.write_ub4(0); writer.write_ub4(0); writer.write_u8(0); writer.write_ub4(0); }
245 }
246}
247
248fn write_xid_block_bytes(writer: &mut TtcWriter, xid: &TpcXid<'_>) {
251 let mut xid_bytes = Vec::with_capacity(128);
252 xid_bytes.extend_from_slice(xid.global_transaction_id);
253 xid_bytes.extend_from_slice(xid.branch_qualifier);
254 xid_bytes.resize(128, 0);
255 writer.write_raw(&xid_bytes);
256}
257
258pub fn build_tpc_switch_payload_with_seq(
265 seq_num: u8,
266 operation: u32,
267 flags: u32,
268 timeout: u32,
269 xid: Option<&TpcXid<'_>>,
270 context: Option<&[u8]>,
271) -> Vec<u8> {
272 build_tpc_switch_payload_with_seq_and_version(
279 seq_num,
280 operation,
281 flags,
282 timeout,
283 xid,
284 context,
285 TNS_CCAP_FIELD_VERSION_23_1_EXT_1,
286 )
287}
288
289#[allow(clippy::too_many_arguments)]
296pub fn build_tpc_switch_payload_with_seq_and_version(
297 seq_num: u8,
298 operation: u32,
299 flags: u32,
300 timeout: u32,
301 xid: Option<&TpcXid<'_>>,
302 context: Option<&[u8]>,
303 ttc_field_version: u8,
304) -> Vec<u8> {
305 let mut writer = TtcWriter::new();
306 writer.write_function_header(TNS_FUNC_TPC_TXN_SWITCH, seq_num, ttc_field_version);
307 writer.write_ub4(operation);
308 match context {
309 Some(context) => {
310 writer.write_u8(1); writer.write_ub4(u32::try_from(context.len()).unwrap_or(0));
312 }
313 None => {
314 writer.write_u8(0); writer.write_ub4(0); }
317 }
318 write_xid_descriptor(&mut writer, xid);
319 writer.write_ub4(flags);
320 writer.write_ub4(timeout);
321 writer.write_u8(1); writer.write_u8(1); writer.write_u8(1); writer.write_u8(0); writer.write_ub4(0); writer.write_u8(0); writer.write_ub4(0); if let Some(context) = context {
329 writer.write_raw(context);
330 }
331 if let Some(xid) = xid {
332 write_xid_block_bytes(&mut writer, xid);
333 }
334 writer.write_ub4(0); writer.into_bytes()
336}
337
338pub fn build_tpc_change_state_payload_with_seq(
344 seq_num: u8,
345 operation: u32,
346 requested_state: u32,
347 flags: u32,
348 xid: Option<&TpcXid<'_>>,
349 context: Option<&[u8]>,
350) -> Vec<u8> {
351 build_tpc_change_state_payload_with_seq_and_version(
355 seq_num,
356 operation,
357 requested_state,
358 flags,
359 xid,
360 context,
361 TNS_CCAP_FIELD_VERSION_23_1_EXT_1,
362 )
363}
364
365#[allow(clippy::too_many_arguments)]
369pub fn build_tpc_change_state_payload_with_seq_and_version(
370 seq_num: u8,
371 operation: u32,
372 requested_state: u32,
373 flags: u32,
374 xid: Option<&TpcXid<'_>>,
375 context: Option<&[u8]>,
376 ttc_field_version: u8,
377) -> Vec<u8> {
378 let mut writer = TtcWriter::new();
379 writer.write_function_header(TNS_FUNC_TPC_TXN_CHANGE_STATE, seq_num, ttc_field_version);
380 writer.write_ub4(operation);
381 match context {
382 Some(context) => {
383 writer.write_u8(1); writer.write_ub4(u32::try_from(context.len()).unwrap_or(0));
385 }
386 None => {
387 writer.write_u8(0); writer.write_ub4(0); }
390 }
391 write_xid_descriptor(&mut writer, xid);
392 writer.write_ub4(0); writer.write_ub4(requested_state);
394 writer.write_u8(1); writer.write_ub4(flags);
396 if let Some(context) = context {
397 writer.write_raw(context);
398 }
399 if let Some(xid) = xid {
400 write_xid_block_bytes(&mut writer, xid);
401 }
402 writer.into_bytes()
403}
404
405pub fn parse_tpc_switch_response(
410 payload: &[u8],
411 capabilities: ClientCapabilities,
412) -> Result<TpcSwitchResponse> {
413 parse_tpc_switch_response_with_limits(payload, capabilities, ProtocolLimits::DEFAULT)
414}
415
416pub fn parse_tpc_switch_response_with_limits(
417 payload: &[u8],
418 capabilities: ClientCapabilities,
419 limits: ProtocolLimits,
420) -> Result<TpcSwitchResponse> {
421 let mut reader = TtcReader::with_limits(payload, limits)?;
422 let mut response = TpcSwitchResponse::default();
423 while reader.remaining() > 0 {
424 let message_type = reader.read_u8()?;
425 match message_type {
426 0 => {}
427 TNS_MSG_TYPE_STATUS => {
428 let call_status = reader.read_ub4()?;
429 let _seq = reader.read_ub2()?;
430 response.txn_in_progress = call_status & TNS_EOCS_FLAGS_TXN_IN_PROGRESS != 0;
431 }
432 TNS_MSG_TYPE_PARAMETER => {
433 let _application_value = reader.read_ub4()?;
436 let context_len = reader.read_ub2()?;
437 let context = reader.read_raw(usize::from(context_len))?;
438 response.context = context.to_vec();
439 }
440 TNS_MSG_TYPE_SERVER_SIDE_PIGGYBACK => {
441 if let Some(update) = skip_server_side_piggyback(&mut reader)? {
442 response.sessionless_state = Some(update);
443 }
444 }
445 TNS_MSG_TYPE_END_OF_RESPONSE => break,
446 TNS_MSG_TYPE_ERROR => {
447 let info = parse_server_error_info(&mut reader, capabilities.ttc_field_version)?;
448 if info.number != 0 {
449 return Err(ProtocolError::ServerErrorInfo(Box::new(
454 info.into_details(),
455 )));
456 }
457 response.txn_in_progress = info.call_status & TNS_EOCS_FLAGS_TXN_IN_PROGRESS != 0;
460 }
461 _ => break,
462 }
463 }
464 Ok(response)
465}
466
467pub fn parse_tpc_change_state_response(
472 payload: &[u8],
473 capabilities: ClientCapabilities,
474) -> Result<TpcChangeStateResponse> {
475 parse_tpc_change_state_response_with_limits(payload, capabilities, ProtocolLimits::DEFAULT)
476}
477
478pub fn parse_tpc_change_state_response_with_limits(
479 payload: &[u8],
480 capabilities: ClientCapabilities,
481 limits: ProtocolLimits,
482) -> Result<TpcChangeStateResponse> {
483 let mut reader = TtcReader::with_limits(payload, limits)?;
484 let mut response = TpcChangeStateResponse::default();
485 while reader.remaining() > 0 {
486 let message_type = reader.read_u8()?;
487 match message_type {
488 0 => {}
489 TNS_MSG_TYPE_STATUS => {
490 let call_status = reader.read_ub4()?;
491 let _seq = reader.read_ub2()?;
492 response.txn_in_progress = call_status & TNS_EOCS_FLAGS_TXN_IN_PROGRESS != 0;
493 }
494 TNS_MSG_TYPE_PARAMETER => {
495 response.state = reader.read_ub4()?;
498 }
499 TNS_MSG_TYPE_SERVER_SIDE_PIGGYBACK => {
500 skip_server_side_piggyback(&mut reader)?;
501 }
502 TNS_MSG_TYPE_END_OF_RESPONSE => break,
503 TNS_MSG_TYPE_ERROR => {
504 let info = parse_server_error_info(&mut reader, capabilities.ttc_field_version)?;
505 if info.number != 0 {
506 return Err(ProtocolError::ServerErrorInfo(Box::new(
511 info.into_details(),
512 )));
513 }
514 response.txn_in_progress = info.call_status & TNS_EOCS_FLAGS_TXN_IN_PROGRESS != 0;
517 }
518 _ => break,
519 }
520 }
521 Ok(response)
522}
523
524pub(crate) fn skip_keyword_value_pairs(reader: &mut TtcReader<'_>, num_pairs: u16) -> Result<()> {
525 read_keyword_value_pairs_for_txn_state(reader, num_pairs).map(|_| ())
526}
527
528pub(crate) fn read_keyword_value_pairs_for_txn_state(
533 reader: &mut TtcReader<'_>,
534 num_pairs: u16,
535) -> Result<Option<SessionlessTxnState>> {
536 let mut state = None;
537 for _ in 0..num_pairs {
538 if reader.read_ub2()? > 0 {
539 let _text_value = reader.read_bytes()?;
540 }
541 let mut binary_value = None;
542 if reader.read_ub2()? > 0 {
543 binary_value = reader.read_bytes()?;
544 }
545 let keyword_num = reader.read_ub2()?;
546 if keyword_num == TNS_KEYWORD_NUM_TRANSACTION_ID {
547 if let Some(binary) = binary_value.as_deref() {
548 if let Some(update) = decode_sessionless_txn_state(binary)? {
549 state = Some(update);
550 }
551 }
552 }
553 }
554 Ok(state)
555}
556
557#[cfg(test)]
558mod tpc_tests {
559 use super::*;
560
561 fn xid() -> ([u8; 7], [u8; 8]) {
562 (*b"txn4400", *b"branchId")
563 }
564
565 #[test]
566 fn tpc_begin_payload_encodes_format_branch_and_128_byte_xid() {
567 let (gtid, bqual) = xid();
568 let tpc_xid = TpcXid {
569 format_id: 4400,
570 global_transaction_id: >id,
571 branch_qualifier: &bqual,
572 };
573 let payload = build_tpc_switch_payload_with_seq(
574 4,
575 TNS_TPC_TXN_START,
576 TPC_TXN_FLAGS_NEW,
577 0,
578 Some(&tpc_xid),
579 None,
580 );
581 assert_eq!(&payload[..3], &[3, TNS_FUNC_TPC_TXN_SWITCH, 4]);
583 let body = &payload[4..];
584 assert_eq!(&body[..4], &[1, 1, 0, 0]);
586 assert_eq!(&body[4..7], &[2, 0x11, 0x30]);
588 assert_eq!(&body[7..14], &[1, 7, 1, 8, 1, 1, 0x80]);
591 let block_start = payload.len() - 128 - 1;
594 let block = &payload[block_start..block_start + 128];
595 assert_eq!(&block[..7], b"txn4400");
596 assert_eq!(&block[7..15], b"branchId");
597 assert!(block[15..].iter().all(|&byte| byte == 0));
598 }
599
600 #[test]
601 fn tpc_end_payload_echoes_context() {
602 let context = vec![0xAAu8; 168];
603 let payload =
604 build_tpc_switch_payload_with_seq(7, TNS_TPC_TXN_DETACH, 0, 0, None, Some(&context));
605 let body = &payload[4..];
606 assert_eq!(&body[..5], &[1, 2, 1, 1, 0xA8]);
608 assert!(payload
610 .windows(context.len())
611 .any(|window| window == context.as_slice()));
612 }
613
614 #[test]
615 fn change_state_prepare_payload_shape() {
616 let (gtid, bqual) = xid();
617 let tpc_xid = TpcXid {
618 format_id: 4400,
619 global_transaction_id: >id,
620 branch_qualifier: &bqual,
621 };
622 let payload = build_tpc_change_state_payload_with_seq(
623 8,
624 TNS_TPC_TXN_PREPARE,
625 TNS_TPC_TXN_STATE_PREPARE,
626 0,
627 Some(&tpc_xid),
628 None,
629 );
630 assert_eq!(&payload[..3], &[3, TNS_FUNC_TPC_TXN_CHANGE_STATE, 8]);
631 let body = &payload[4..];
632 assert_eq!(&body[..4], &[1, 3, 0, 0]);
634 }
635
636 #[test]
637 fn switch_response_captures_context_and_txn_bit() {
638 let mut payload = Vec::new();
641 payload.push(TNS_MSG_TYPE_PARAMETER);
642 payload.push(0); payload.extend_from_slice(&[2, 0, 4]); payload.extend_from_slice(&[0xDE, 0xAD, 0xBE, 0xEF]);
645 payload.push(TNS_MSG_TYPE_STATUS);
646 payload.extend_from_slice(&[1, 3]); payload.extend_from_slice(&[0]); payload.push(TNS_MSG_TYPE_END_OF_RESPONSE);
649
650 let response =
651 parse_tpc_switch_response(&payload, ClientCapabilities::default()).expect("decode");
652 assert_eq!(response.context, vec![0xDE, 0xAD, 0xBE, 0xEF]);
653 assert!(response.txn_in_progress);
654 }
655
656 #[test]
657 fn switch_response_end_status_clears_txn_bit() {
658 let mut payload = Vec::new();
660 payload.push(TNS_MSG_TYPE_STATUS);
661 payload.extend_from_slice(&[1, 1]); payload.extend_from_slice(&[0]); payload.push(TNS_MSG_TYPE_END_OF_RESPONSE);
664
665 let response =
666 parse_tpc_switch_response(&payload, ClientCapabilities::default()).expect("decode");
667 assert!(!response.txn_in_progress);
668 }
669
670 #[test]
671 fn change_state_response_reads_out_state() {
672 let mut payload = Vec::new();
674 payload.push(TNS_MSG_TYPE_PARAMETER);
675 payload.extend_from_slice(&[1, 1]); payload.push(TNS_MSG_TYPE_STATUS);
677 payload.extend_from_slice(&[1, 1]); payload.extend_from_slice(&[0]); payload.push(TNS_MSG_TYPE_END_OF_RESPONSE);
680
681 let response = parse_tpc_change_state_response(&payload, ClientCapabilities::default())
682 .expect("decode");
683 assert_eq!(response.state, TNS_TPC_TXN_STATE_REQUIRES_COMMIT);
684 assert!(!response.txn_in_progress);
685 }
686}