diff --git a/ChipDnaClient.csproj b/ChipDnaClient.csproj
index 78e719f..2cf4680 100644
--- a/ChipDnaClient.csproj
+++ b/ChipDnaClient.csproj
@@ -247,6 +247,7 @@
+
diff --git a/Client.cs b/Client.cs
index 2499ab6..485e047 100644
--- a/Client.cs
+++ b/Client.cs
@@ -25,6 +25,10 @@ namespace Creditcall.ChipDna.Client
public TaskCompletionSource> TransactionCompletionSource { get; set; }
+ internal readonly TransactionProgressTap Progress = new TransactionProgressTap();
+
+ public void SuppressUncertainProgress() { Progress.Suppress(); }
+
private volatile bool voiceReferralRequired;
private volatile bool voiceReferralDone;
private readonly bool saveReceipt;
@@ -241,11 +245,13 @@ namespace Creditcall.ChipDna.Client
parameters.Add(GetExtraParams("StartTransaction"));
Console.WriteLine(parameters.ToString());
+ Progress.Begin(reference);
Response response = chipDnaClientLib.StartTransaction(parameters);
string errors;
if (response.GetValue(ParameterKeys.Errors, out errors) && !string.IsNullOrEmpty(errors))
{
+ Progress.StartFailed(errors);
var result = new Dictionary
{
{ ParameterKeys.Errors, errors },
@@ -1301,6 +1307,7 @@ namespace Creditcall.ChipDna.Client
private void ChipDnaClientLibOnTransactionFinished(object sender, EventParameters transactionFinishedEventArgs)
{
+ Progress.End();
var result = new Dictionary();
var transactionFinishedStringBuilder = new StringBuilder();
transactionFinishedStringBuilder.Append("Transaction Finished Event Parameters: ");
@@ -1371,6 +1378,7 @@ namespace Creditcall.ChipDna.Client
private void ChipDnaClientLibOnSignatureVerificationRequested(object sender, EventParameters signatureVerificationRequestedEventArgs)
{
+ Progress.Report("SIGNATURE", signatureVerificationRequestedEventArgs);
var transactionFinishedStringBuilder = new StringBuilder();
transactionFinishedStringBuilder.Append("Signature Verification Requested Event Parameters: ");
string reference = null;
@@ -1426,17 +1434,20 @@ namespace Creditcall.ChipDna.Client
private void ChipDnaClientLibOnTransactionPause(object sender, EventParameters transactionPauseEventArgs)
{
+ Progress.Report("PAUSE", transactionPauseEventArgs);
Console.WriteLine("Transaction Pause Event Parameters: {0}{1} *Waiting For Continue Command --> Press 'L' To Continue*",
transactionPauseEventArgs, "\n");
}
private void ChipDnaClientLibOnTransactionUpdate(object sender, EventParameters transactionUpdateEventArgs)
{
+ Progress.Report("UPDATE", transactionUpdateEventArgs);
Console.WriteLine("Transaction Update Event Parameters: {0}", transactionUpdateEventArgs);
}
private void ChipDnaClientLibOnCardNotification(object sender, EventParameters cardNotificationEventArgs)
{
+ Progress.Report("CARD_STATUS", cardNotificationEventArgs);
Console.WriteLine("Card Notification Event Parameters: {0}", cardNotificationEventArgs);
}
@@ -1507,6 +1518,7 @@ namespace Creditcall.ChipDna.Client
private void ChipDnaClientLibOnErrorEvent(object sender, ClientHelperErrorEventArgs errorEventArgs)
{
+ Progress.SuppressIfActive();
var result = new Dictionary
{
{ ParameterKeys.Errors, errorEventArgs.RaisedException.Message}
diff --git a/ClientApp.cs b/ClientApp.cs
index fc4b145..b941a70 100644
--- a/ClientApp.cs
+++ b/ClientApp.cs
@@ -72,6 +72,7 @@ namespace Creditcall.ChipDna.Client
private static void StartHttpServer()
{
HttpListener listener = new HttpListener();
+ listener.Prefixes.Add("http://127.0.0.1:18181/start-transaction-stream/");
listener.Prefixes.Add("http://127.0.0.1:18181/start-transaction/");
listener.Prefixes.Add("http://127.0.0.1:18181/confirm-transaction/");
listener.Prefixes.Add("http://127.0.0.1:18181/transaction-information/");
@@ -80,6 +81,7 @@ namespace Creditcall.ChipDna.Client
listener.Start();
Console.WriteLine("Listening on:");
+ Console.WriteLine(" http://127.0.0.1:18181/start-transaction-stream/");
Console.WriteLine(" http://127.0.0.1:18181/start-transaction/");
Console.WriteLine(" http://127.0.0.1:18181/confirm-transaction/");
Console.WriteLine(" http://127.0.0.1:18181/transaction-information/");
@@ -101,7 +103,11 @@ namespace Creditcall.ChipDna.Client
continue;
}
- if (rawUrl.StartsWith("/start-transaction", StringComparison.OrdinalIgnoreCase))
+ if (rawUrl.StartsWith("/start-transaction-stream", StringComparison.OrdinalIgnoreCase))
+ {
+ HandleStartTransactionStream(context);
+ }
+ else if (rawUrl.StartsWith("/start-transaction", StringComparison.OrdinalIgnoreCase))
{
HandleStartTransaction(context);
@@ -257,47 +263,7 @@ namespace Creditcall.ChipDna.Client
Console.WriteLine($"Received transaction request: Amount={payload.Amount}, Type={payload.TransactionType}");
- client.TransactionCompletionSource = new TaskCompletionSource>();
-
- // start transaction (non-blocking)
- client.PerformStartTransaction(payload.Amount, payload.TransactionType);
-
- // Wait for completion with timeout
- var task = client.TransactionCompletionSource.Task;
- if (!task.Wait(TimeSpan.FromSeconds(transactionTimeoutSeconds)))
- {
- // timeout
- var errorMap = new Dictionary
- {
- { ParameterKeys.TransactionResult, "TIMEOUT" },
- { ParameterKeys.Errors, "Transaction timed out waiting for device response" }
- };
-
- var wrappedTimeout = new SerializableKeyValueList(errorMap);
- context.Response.ContentType = "application/xml";
- new XmlSerializer(typeof(SerializableKeyValueList)).Serialize(context.Response.OutputStream, wrappedTimeout);
- return;
- }
-
- // Task completed — handle possible exception
- Dictionary transactionResult;
- try
- {
- transactionResult = task.Result; // safe now, Wait returned
- }
- catch (AggregateException agg)
- {
- var first = agg.Flatten().InnerExceptions[0];
- var errMap = new Dictionary
- {
- { ParameterKeys.TransactionResult, "ERROR" },
- { ParameterKeys.Errors, $"Transaction task faulted: {first.Message}" }
- };
- var wrapped = new SerializableKeyValueList(errMap);
- context.Response.ContentType = "application/xml";
- new XmlSerializer(typeof(SerializableKeyValueList)).Serialize(context.Response.OutputStream, wrapped);
- return;
- }
+ var transactionResult = ExecuteTransaction(payload);
var responseSerializer = new XmlSerializer(typeof(SerializableKeyValueList));
context.Response.ContentType = "application/xml";
@@ -316,6 +282,54 @@ namespace Creditcall.ChipDna.Client
}
+ private static Dictionary ExecuteTransaction(TransactionPayload payload)
+ {
+ client.TransactionCompletionSource = new TaskCompletionSource>();
+ return TransactionExecution.Run(
+ () => client.PerformStartTransaction(payload.Amount, payload.TransactionType),
+ () => client.TransactionCompletionSource.Task,
+ client.SuppressUncertainProgress,
+ TimeSpan.FromSeconds(transactionTimeoutSeconds));
+ }
+
+ private static void HandleStartTransactionStream(HttpListenerContext context)
+ {
+ context.Response.ContentType = "application/x-ndjson; charset=utf-8";
+ context.Response.SendChunked = true;
+ var writer = new TransactionStreamWriter(context.Response.OutputStream);
+ var delivery = Task.Run(() => writer.Deliver());
+ Action progress = writer.Progress;
+ try
+ {
+ TransactionPayload payload;
+ using (var reader = new StreamReader(context.Request.InputStream))
+ using (var json = new Newtonsoft.Json.JsonTextReader(reader))
+ {
+ payload = new Newtonsoft.Json.JsonSerializer().Deserialize(json);
+ }
+ if (payload == null)
+ throw new InvalidDataException("Missing transaction payload");
+
+ var result = client.Progress.Observe(progress, () => ExecuteTransaction(payload));
+ writer.Complete(new { type = "result", result = TransactionStreamFields.Required(result) });
+ }
+ catch (Exception)
+ {
+ writer.Complete(new { type = "error", error = "Transaction stream failed" });
+ }
+ finally
+ {
+ // Retain serial response completion, as with the XML endpoint.
+ // Hardlink's existing per-call timeout bounds its connection;
+ // no new delivery timer may truncate an approved final result.
+ try { delivery.GetAwaiter().GetResult(); }
+ finally
+ {
+ try { context.Response.Close(); } catch (HttpListenerException) { }
+ }
+ }
+ }
+
private static string GetAbsolutePath(string fileName)
{
var path = Path.GetDirectoryName(Assembly.GetExecutingAssembly().Location);
diff --git a/Tests/StreamingTests.cs b/Tests/StreamingTests.cs
new file mode 100644
index 0000000..e3c489d
--- /dev/null
+++ b/Tests/StreamingTests.cs
@@ -0,0 +1,531 @@
+using System;
+using System.Collections.Generic;
+using System.IO;
+using System.Linq;
+using System.Reflection;
+using System.Text;
+using System.Threading.Tasks;
+using Newtonsoft.Json.Linq;
+
+namespace Creditcall.ChipDna.Client
+{
+ internal static class StreamingTests
+ {
+ private static int Main()
+ {
+ var tests = new Action[] {
+ SharedExecutionReturnsAuthoritativeResultOnce,
+ TimeoutDoesNotCancelCompletion,
+ StartAndTaskErrorsKeepExceptionBehavior,
+ SubscriptionIsRemovedOnSuccessAndException,
+ OwnershipAcceptsReferenceLessProgressAndSuppressesPermanently,
+ RepeatedReferenceSuppressesProgress,
+ NativeMiuraSequenceKeepsFinalAuthoritative,
+ NativeMiuraRecoveryCyclesKeepFinalAuthoritative,
+ NativeNotificationUsesNotificationKey,
+ NativeOwnershipAndRemovalBoundaries,
+ PreDispatchFailureReleasesOnlyUnusedReference,
+ UnprovenStartFailuresSuppressObservation,
+ AsyncErrorSuppressesOnlyActiveObservation,
+ TimeoutUnsubscribesAndKeepsFuturePaymentsEnabled,
+ QueueCapacityCannotDropFinal,
+ ConcurrentCallbacksKeepValidLinesAndOneFinal,
+ DisconnectDoesNotInterruptExecution,
+ FinalFieldsExcludeSensitiveParameters
+ };
+ int failed = 0;
+ foreach (var test in tests)
+ {
+ try { test(); Console.WriteLine("PASS " + test.Method.Name); }
+ catch (Exception ex) { failed++; Console.WriteLine("FAIL " + test.Method.Name + ": " + ex); }
+ }
+ Console.WriteLine($"{tests.Length - failed}/{tests.Length} tests passed");
+ return failed == 0 ? 0 : 1;
+ }
+
+ private static void Check(bool condition, string message)
+ {
+ if (!condition) throw new Exception(message);
+ }
+
+ private static EventParameters Native(string key, string value, string reference = null)
+ {
+ var fields = new Dictionary();
+ if (key != null) fields[key] = value;
+ if (reference != null) fields["REFERENCE"] = reference;
+ // SDK event construction is internal; no ClientHelper or terminal is started.
+ return (EventParameters)Activator.CreateInstance(typeof(EventParameters),
+ BindingFlags.Instance | BindingFlags.NonPublic, null, new object[] { fields }, null);
+ }
+
+ private static void CheckUnsubscribed(TransactionProgressTap tap)
+ {
+ var listener = typeof(TransactionProgressTap).GetField("listener", BindingFlags.Instance | BindingFlags.NonPublic);
+ Check(listener.GetValue(tap) == null, "progress subscription survived execution");
+ }
+
+ private static void SharedExecutionReturnsAuthoritativeResultOnce()
+ {
+ foreach (var outcome in new[] { "Approved", "Declined", "Cancelled" })
+ {
+ var expected = new Dictionary { { "TRANSACTION_RESULT", outcome }, { "REFERENCE", "existing-ref" } };
+ var source = new TaskCompletionSource>();
+ int starts = 0, suppressions = 0;
+ var result = TransactionExecution.Run(() => { starts++; source.SetResult(expected); },
+ () => source.Task, () => suppressions++, TimeSpan.FromSeconds(1));
+ Check(ReferenceEquals(expected, result) && starts == 1 && suppressions == 0,
+ "shared execution changed result or start count for " + outcome);
+ }
+ }
+
+ private static void TimeoutDoesNotCancelCompletion()
+ {
+ var source = new TaskCompletionSource>();
+ int starts = 0, suppressions = 0;
+ var result = TransactionExecution.Run(() => starts++, () => source.Task, () => suppressions++, TimeSpan.Zero);
+ Check(result["TRANSACTION_RESULT"] == "TIMEOUT" &&
+ result["ERRORS"] == "Transaction timed out waiting for device response", "timeout fields changed");
+ Check(starts == 1 && suppressions == 1 && !source.Task.IsCompleted, "timeout cancelled or restarted the transaction");
+ source.SetResult(new Dictionary { { "TRANSACTION_RESULT", "Approved" } });
+ Check(source.Task.Result["TRANSACTION_RESULT"] == "Approved", "late SDK completion was disabled");
+ }
+
+ private static void StartAndTaskErrorsKeepExceptionBehavior()
+ {
+ int suppressions = 0;
+ try
+ {
+ TransactionExecution.Run(() => { throw new InvalidOperationException("start"); },
+ () => Task.FromResult(new Dictionary()), () => suppressions++, TimeSpan.Zero);
+ throw new Exception("expected original start exception");
+ }
+ catch (InvalidOperationException) { }
+ var fault = new TaskCompletionSource>();
+ fault.SetException(new IOException("SDK fault"));
+ try
+ {
+ TransactionExecution.Run(() => { }, () => fault.Task, () => suppressions++, TimeSpan.FromSeconds(1));
+ throw new Exception("expected original Task.Wait AggregateException");
+ }
+ catch (AggregateException) { }
+ Check(suppressions == 2, "ambiguous exception did not suppress observation");
+ }
+
+ private static void SubscriptionIsRemovedOnSuccessAndException()
+ {
+ var tap = new TransactionProgressTap();
+ int first = 0, second = 0;
+ tap.Observe((s,v) => first++, () => {
+ tap.Begin("first");
+ tap.Report("UPDATE", Native("UPDATE", "CardRequested", "first"));
+ tap.End();
+ return new Dictionary();
+ });
+ CheckUnsubscribed(tap);
+ tap.Begin("after-first");
+ tap.Report("UPDATE", Native("UPDATE", "CardRequested", "after-first"));
+ tap.End();
+ try
+ {
+ tap.Observe((s,v) => second++, () => {
+ tap.Begin("second");
+ tap.Report("UPDATE", Native("UPDATE", "CardRequested", "second"));
+ throw new IOException("test");
+ });
+ }
+ catch (IOException) { }
+ tap.Report("UPDATE", Native("UPDATE", "CardRequested", "second"));
+ Check(first == 1 && second == 1, $"subscriptions leaked: {first}/{second}");
+ CheckUnsubscribed(tap);
+ }
+
+ private static void OwnershipAcceptsReferenceLessProgressAndSuppressesPermanently()
+ {
+ var tap = new TransactionProgressTap();
+ int callbacks = 0, payments = 0;
+ tap.Observe((s,v) => callbacks++, () => {
+ payments++; tap.Begin("first");
+ tap.Report("UPDATE", Native("UPDATE", "CardRequested", null));
+ tap.Report("UPDATE", Native("UPDATE", "CardRequested", "other"));
+ tap.Report("UPDATE", Native("UPDATE", "CardRequested", "first"));
+ tap.Suppress();
+ tap.Report("UPDATE", Native("UPDATE", "CardRequested", "first"));
+ tap.End();
+ return new Dictionary();
+ });
+ tap.Observe((s,v) => callbacks++, () => {
+ payments++; tap.Begin("later");
+ tap.Report("UPDATE", Native("UPDATE", "CardRequested", "first"));
+ tap.Report("UPDATE", Native("UPDATE", "CardRequested", "later"));
+ tap.End();
+ return new Dictionary();
+ });
+ Check(callbacks == 2 && payments == 2, "active reference-less progress was dropped, suppression leaked, or later payments were disabled");
+ }
+
+ private static void RepeatedReferenceSuppressesProgress()
+ {
+ var tap = new TransactionProgressTap();
+ int callbacks = 0;
+ tap.Observe((s,v) => callbacks++, () => {
+ tap.Begin("same-second-reference");
+ tap.End();
+ tap.Begin("same-second-reference");
+ tap.Report("UPDATE", Native("UPDATE", "CardRequested", "same-second-reference"));
+ return new Dictionary();
+ });
+ Check(callbacks == 0, "reused reference attributed an ambiguous event");
+ }
+
+ private static void NativeMiuraSequenceKeepsFinalAuthoritative()
+ {
+ foreach (var outcome in new[] { "Approved", "Declined" })
+ using (var output = new MemoryStream())
+ {
+ var tap = new TransactionProgressTap();
+ var writer = new TransactionStreamWriter(output);
+ var expected = new Dictionary { { "TRANSACTION_RESULT", outcome }, { "REFERENCE", "miura-ref" } };
+ int starts = 0;
+ var result = tap.Observe(writer.Progress, () => TransactionExecution.Run(() => {
+ starts++;
+ tap.Begin("miura-ref");
+ tap.Report("UPDATE", Native("UPDATE", "CardRequested"));
+ tap.Report("UPDATE", Native("UPDATE", "CardRequested"));
+ tap.Report("CARD_STATUS", Native("NOTIFICATION", "Inserted"));
+ foreach (var update in new[] { "TransactionStarted", "PinEntryStarted", "OnlineAuthRequested", "OnlineAuthCompleted" })
+ tap.Report("UPDATE", Native("UPDATE", update));
+ tap.End(); // TransactionFinished closes observation before final processing.
+ tap.Report("UPDATE", Native("UPDATE", "CardRemovalRequested"));
+ tap.Report("CARD_STATUS", Native("NOTIFICATION", "Removed"));
+ }, () => Task.FromResult(expected), tap.Suppress, TimeSpan.FromSeconds(1)));
+ writer.Complete(new { type = "result", result });
+ writer.Deliver();
+ var frames = Frames(output);
+ var values = frames.Take(frames.Length - 1).Select(f => f["value"].Value());
+ Check(values.SequenceEqual(new[] { "CardRequested", "CardRequested", "Inserted", "TransactionStarted", "PinEntryStarted", "OnlineAuthRequested", "OnlineAuthCompleted" }),
+ "SDK 3.17 Miura sequence lost or changed normal reference-less progress");
+ Check(starts == 1 && ReferenceEquals(result, expected) && frames.Length == 8 &&
+ frames.Last()["result"]["TRANSACTION_RESULT"].Value() == outcome,
+ "progress changed start count or authoritative " + outcome + " result");
+ CheckUnsubscribed(tap);
+ }
+ }
+
+ private static void NativeMiuraRecoveryCyclesKeepFinalAuthoritative()
+ {
+ foreach (var cycles in new[] { 1, 3 })
+ foreach (var outcome in new[] { "Approved", "Declined" })
+ using (var output = new MemoryStream())
+ {
+ var tap = new TransactionProgressTap();
+ var writer = new TransactionStreamWriter(output);
+ var expected = new Dictionary { { "TRANSACTION_RESULT", outcome } };
+ int starts = 0;
+ var result = tap.Observe(writer.Progress, () => TransactionExecution.Run(() => {
+ starts++;
+ tap.Begin("recovery");
+ for (int i = 0; i < cycles; i++)
+ {
+ tap.Report("UPDATE", Native("UPDATE", "CardRequested"));
+ tap.Report("CARD_STATUS", Native("NOTIFICATION", "Inserted"));
+ tap.Report("UPDATE", Native("UPDATE", "CardRemovalRequested"));
+ tap.Report("CARD_STATUS", Native("NOTIFICATION", "Removed"));
+ }
+ tap.Report("UPDATE", Native("UPDATE", "CardRequested"));
+ tap.Report("CARD_STATUS", Native("NOTIFICATION", "Inserted"));
+ foreach (var update in new[] { "TransactionStarted", "PinEntryStarted", "OnlineAuthRequested", "OnlineAuthCompleted" })
+ tap.Report("UPDATE", Native("UPDATE", update));
+ tap.End();
+ tap.Report("UPDATE", Native("UPDATE", "CardRemovalRequested"));
+ }, () => Task.FromResult(expected), tap.Suppress, TimeSpan.FromSeconds(1)));
+ writer.Complete(new { type = "result", result });
+ writer.Deliver();
+ var frames = Frames(output);
+ var expectedProgress = Enumerable.Range(0, cycles)
+ .SelectMany(i => new[] { "UPDATE:CardRequested", "CARD_STATUS:Inserted", "UPDATE:CardRemovalRequested" })
+ .Concat(new[] { "UPDATE:CardRequested", "CARD_STATUS:Inserted", "UPDATE:TransactionStarted",
+ "UPDATE:PinEntryStarted", "UPDATE:OnlineAuthRequested", "UPDATE:OnlineAuthCompleted" });
+ Check(frames.Take(frames.Length - 1).All(f => f["type"].Value() == "status") &&
+ frames.Take(frames.Length - 1).Select(f => f["source"] + ":" + f["value"]).SequenceEqual(expectedProgress),
+ "Miura recovery prompts were dropped, reordered, deduplicated, or included Removed");
+ Check(starts == 1 && ReferenceEquals(result, expected) &&
+ frames.Count(f => f["type"].Value() == "result") == 1 &&
+ frames.Last()["result"]["TRANSACTION_RESULT"].Value() == outcome,
+ "recovery progress changed transaction execution or the authoritative final result");
+ CheckUnsubscribed(tap);
+ }
+ }
+
+ private static void NativeNotificationUsesNotificationKey()
+ {
+ var tap = new TransactionProgressTap();
+ var delivered = new List();
+ tap.Observe((source, value) => delivered.Add(source + ":" + value), () => {
+ tap.Begin("native-key");
+ tap.Report("CARD_STATUS", Native("CARD_STATUS", "Inserted"));
+ foreach (var value in new[] { "Inserted", "Tapped", "Swiped" })
+ tap.Report("CARD_STATUS", Native("NOTIFICATION", value));
+ tap.Report("CARD_STATUS", Native("NOTIFICATION", "Removed", "native-key"));
+ tap.Report("CARD_STATUS", Native("NOTIFICATION", "Unknown"));
+ tap.Report("UPDATE", Native("PAN", "4761730000001133"));
+ tap.Report("UPDATE", null);
+ tap.End();
+ return new Dictionary();
+ });
+ Check(delivered.SequenceEqual(new[] { "CARD_STATUS:Inserted", "CARD_STATUS:Tapped", "CARD_STATUS:Swiped" }),
+ "native NOTIFICATION extraction used CARD_STATUS, exposed card data, or changed the wire source");
+ }
+
+ private static void NativeOwnershipAndRemovalBoundaries()
+ {
+ var tap = new TransactionProgressTap();
+ var delivered = new List();
+ var removals = new[] { "CardRemovalRequested", "CardRemovalEnforced" };
+ tap.Observe((source, value) => delivered.Add(source + ":" + value), () => {
+ foreach (var removal in removals)
+ foreach (var reference in new[] { null, "", "first", "other" })
+ tap.Report("UPDATE", Native("UPDATE", removal, reference));
+ tap.Begin("first");
+ tap.Report("UPDATE", Native("UPDATE", "CardRequested", ""));
+ tap.Report("UPDATE", Native("UPDATE", "PinEntryStarted", "other"));
+ tap.Report("SIGNATURE", Native(null, null));
+ tap.Report("PAUSE", Native("PAUSE_STATE", "PostCardDetails"));
+ foreach (var removal in removals)
+ foreach (var reference in new[] { null, "", "first", "other" })
+ {
+ tap.Report("UPDATE", Native("UPDATE", removal, reference));
+ tap.Report("CARD_STATUS", Native("NOTIFICATION", "Removed", reference));
+ }
+ tap.End();
+ foreach (var removal in removals)
+ foreach (var reference in new[] { null, "", "first", "other" })
+ tap.Report("UPDATE", Native("UPDATE", removal, reference));
+ tap.Report("UPDATE", Native("UPDATE", "CardRequested"));
+ tap.Begin("next");
+ // Accepted best-effort limitation: unidentified progress, including removal,
+ // is allowed in this new window even if it came from an older queue.
+ foreach (var removal in removals)
+ {
+ tap.Report("UPDATE", Native("UPDATE", removal));
+ tap.Report("UPDATE", Native("UPDATE", removal, "first"));
+ }
+ tap.Report("CARD_STATUS", Native("NOTIFICATION", "Removed"));
+ tap.Report("UPDATE", Native("UPDATE", "PinEntryStarted", "first"));
+ tap.Report("UPDATE", Native("UPDATE", "CardRequested"));
+ tap.Report("UPDATE", Native("UPDATE", "Approved"));
+ tap.End();
+ return new Dictionary();
+ });
+ Check(delivered.SequenceEqual(new[] { "UPDATE:CardRequested", "SIGNATURE:Requested", "PAUSE:Paused",
+ "UPDATE:CardRemovalRequested", "UPDATE:CardRemovalRequested", "UPDATE:CardRemovalRequested",
+ "UPDATE:CardRemovalEnforced", "UPDATE:CardRemovalEnforced", "UPDATE:CardRemovalEnforced",
+ "UPDATE:CardRemovalRequested", "UPDATE:CardRemovalEnforced", "UPDATE:CardRequested" }),
+ "active-window, explicit-reference, or removal ownership changed");
+ }
+
+ private static void PreDispatchFailureReleasesOnlyUnusedReference()
+ {
+ var tap = new TransactionProgressTap();
+ int starts = 0, callbacks = 0;
+ foreach (var rejected in new[] { true, false })
+ {
+ var expected = rejected
+ ? new Dictionary { { "ERRORS", "ClientNotConnectedToServer" }, { "REFERENCE", "same-second" } }
+ : new Dictionary { { "TRANSACTION_RESULT", "Approved" }, { "REFERENCE", "same-second" } };
+ var result = tap.Observe((source, value) => callbacks++, () => TransactionExecution.Run(() => {
+ starts++;
+ tap.Begin("same-second");
+ if (rejected) tap.StartFailed("ClientNotConnectedToServer");
+ else tap.Report("UPDATE", Native("UPDATE", "CardRequested"));
+ if (!rejected) tap.End();
+ tap.Report("UPDATE", Native("UPDATE", "PinEntryStarted"));
+ }, () => Task.FromResult(expected), tap.Suppress, TimeSpan.FromSeconds(1)));
+ Check(ReferenceEquals(result, expected), "pre-dispatch handling changed the financial result");
+ CheckUnsubscribed(tap);
+ }
+ Check(starts == 2 && callbacks == 1, "proven pre-dispatch rejection prevented a safe same-reference retry");
+ }
+
+ private static void UnprovenStartFailuresSuppressObservation()
+ {
+ foreach (var error in new[] { "AmountTooSmall", "MissingParameter", "NoPinPadsAvailable", "TransactionInProgress",
+ "ClientNotConnectedToServer,AmountTooSmall", " ClientNotConnectedToServer", "500005", "Unknown" })
+ {
+ var tap = new TransactionProgressTap();
+ int starts = 0, callbacks = 0;
+ foreach (var first in new[] { true, false })
+ {
+ var expected = first ? new Dictionary { { "ERRORS", error } }
+ : new Dictionary { { "TRANSACTION_RESULT", "Approved" } };
+ var result = tap.Observe((source, value) => callbacks++, () => TransactionExecution.Run(() => {
+ starts++;
+ tap.Begin(first ? "failed" : "later");
+ if (first) tap.StartFailed(error);
+ // A subsequent safe rejection must never clear an existing latch.
+ tap.StartFailed("ClientNotConnectedToServer");
+ tap.Report("UPDATE", Native("UPDATE", "CardRequested"));
+ foreach (var removal in new[] { "CardRemovalRequested", "CardRemovalEnforced" })
+ {
+ tap.Report("UPDATE", Native("UPDATE", removal));
+ tap.Report("UPDATE", Native("UPDATE", removal, "later"));
+ }
+ tap.Report("UPDATE", Native("UPDATE", "CardRequested", "later"));
+ tap.End();
+ }, () => Task.FromResult(expected), tap.Suppress, TimeSpan.FromSeconds(1)));
+ Check(ReferenceEquals(result, expected), "observational handling changed result for " + error);
+ }
+ Check(starts == 2 && callbacks == 0, "unproven error recovered progress or disabled payment: " + error);
+ }
+ }
+
+ private static void AsyncErrorSuppressesOnlyActiveObservation()
+ {
+ var tap = new TransactionProgressTap();
+ int callbacks = 0, starts = 0;
+ tap.SuppressIfActive(); // Idle SDK errors have no active transaction to make ambiguous.
+ foreach (var reference in new[] { "first", "later" })
+ {
+ var expected = new Dictionary { { "TRANSACTION_RESULT", "Approved" } };
+ var result = tap.Observe((source, value) => callbacks++, () => TransactionExecution.Run(() => {
+ starts++;
+ tap.Begin(reference);
+ tap.Report("UPDATE", Native("UPDATE", "CardRequested"));
+ tap.SuppressIfActive();
+ foreach (var removal in new[] { "CardRemovalRequested", "CardRemovalEnforced" })
+ {
+ tap.Report("UPDATE", Native("UPDATE", removal));
+ tap.Report("UPDATE", Native("UPDATE", removal, reference));
+ }
+ tap.Report("UPDATE", Native("UPDATE", "PinEntryStarted"));
+ tap.End();
+ }, () => Task.FromResult(expected), tap.Suppress, TimeSpan.FromSeconds(1)));
+ Check(ReferenceEquals(result, expected), "asynchronous SDK error altered financial execution");
+ }
+ Check(starts == 2 && callbacks == 1, "asynchronous error failed to latch suppression or suppressed while idle");
+ }
+
+ private static void TimeoutUnsubscribesAndKeepsFuturePaymentsEnabled()
+ {
+ var tap = new TransactionProgressTap();
+ var completion = new TaskCompletionSource>();
+ int starts = 0, callbacks = 0;
+ var timeout = tap.Observe((source, value) => callbacks++, () => TransactionExecution.Run(() => {
+ starts++;
+ tap.Begin("timeout");
+ tap.Report("UPDATE", Native("UPDATE", "CardRequested"));
+ }, () => completion.Task, tap.Suppress, TimeSpan.Zero));
+ Check(timeout["TRANSACTION_RESULT"] == "TIMEOUT" && !completion.Task.IsCompleted, "timeout changed financial completion");
+ CheckUnsubscribed(tap);
+ tap.Report("UPDATE", Native("UPDATE", "PinEntryStarted"));
+ completion.SetResult(new Dictionary { { "TRANSACTION_RESULT", "Approved" } });
+ tap.End();
+ var expected = new Dictionary { { "TRANSACTION_RESULT", "Declined" } };
+ var later = tap.Observe((source, value) => callbacks++, () => TransactionExecution.Run(() => {
+ starts++;
+ tap.Begin("later");
+ tap.Report("UPDATE", Native("UPDATE", "CardRequested"));
+ foreach (var removal in new[] { "CardRemovalRequested", "CardRemovalEnforced" })
+ {
+ tap.Report("UPDATE", Native("UPDATE", removal));
+ tap.Report("UPDATE", Native("UPDATE", removal, "later"));
+ }
+ tap.Report("UPDATE", Native("UPDATE", "PinEntryStarted", "later"));
+ tap.End();
+ }, () => Task.FromResult(expected), tap.Suppress, TimeSpan.FromSeconds(1)));
+ Check(starts == 2 && callbacks == 1 && ReferenceEquals(later, expected), "timeout suppression leaked progress or blocked future payments");
+ CheckUnsubscribed(tap);
+ }
+
+ private static JObject[] Frames(MemoryStream stream)
+ {
+ var text = Encoding.UTF8.GetString(stream.ToArray());
+ Check(text.Length == 0 || text[0] != '\uFEFF', "stream contained a BOM");
+ return text.Split(new[] { "\r\n", "\n" }, StringSplitOptions.RemoveEmptyEntries).Select(JObject.Parse).ToArray();
+ }
+
+ private static void QueueCapacityCannotDropFinal()
+ {
+ using (var output = new MemoryStream())
+ {
+ var writer = new TransactionStreamWriter(output);
+ for (int i = 0; i < 1000; i++) writer.Progress("UPDATE", "CardRequested");
+ writer.Complete(new { type = "result", result = new { TRANSACTION_RESULT = "Approved" } });
+ writer.Complete(new { type = "error", error = "duplicate" });
+ writer.Progress("UPDATE", "PinEntryStarted");
+ writer.Deliver();
+ writer.Deliver();
+ var frames = Frames(output);
+ Check(frames.Length == 129 && frames.Last()["type"].Value() == "result", "full queue dropped or duplicated final");
+ Check(frames.Take(128).All(f => f["value"].Value() == "CardRequested"), "progress arrived after final");
+ }
+ }
+
+ private static void ConcurrentCallbacksKeepValidLinesAndOneFinal()
+ {
+ using (var output = new MemoryStream())
+ {
+ var writer = new TransactionStreamWriter(output);
+ var delivery = Task.Run(() => writer.Deliver());
+ var tap = new TransactionProgressTap();
+ tap.Observe(writer.Progress, () => {
+ tap.Begin("concurrent");
+ Parallel.For(0, 5000, i => tap.Report("UPDATE", Native("UPDATE", "PinEntryStarted")));
+ tap.End();
+ return new Dictionary();
+ });
+ writer.Complete(new { type = "result", result = new { TRANSACTION_RESULT = "Declined" } });
+ Check(delivery.Wait(TimeSpan.FromSeconds(5)), "writer did not finish");
+ var frames = Frames(output);
+ Check(frames.Count(f => f["type"].Value() == "result") == 1 &&
+ frames.Last()["type"].Value() == "result", "concurrent callbacks corrupted terminal framing");
+ Check(frames.Where(f => f["type"].Value() == "status").All(f =>
+ f["source"].Value() == "UPDATE" && f["value"].Value() == "PinEntryStarted"),
+ "concurrent callbacks corrupted status fields");
+ }
+ }
+
+ private sealed class DisconnectedStream : MemoryStream
+ {
+ public override void Write(byte[] buffer, int offset, int count) { throw new IOException("disconnected"); }
+ }
+
+ private static void DisconnectDoesNotInterruptExecution()
+ {
+ using (var output = new DisconnectedStream())
+ {
+ var writer = new TransactionStreamWriter(output);
+ writer.Progress("UPDATE", "CardRequested");
+ writer.Deliver();
+ var expected = new Dictionary { { "TRANSACTION_RESULT", "Approved" } };
+ int starts = 0;
+ var result = TransactionExecution.Run(() => starts++, () => Task.FromResult(expected),
+ () => { throw new Exception("disconnect suppressed financial execution"); }, TimeSpan.FromSeconds(1));
+ writer.Complete(new { type = "result", result });
+ writer.Progress("UPDATE", "CardRequested");
+ Check(starts == 1 && ReferenceEquals(result, expected), "disconnect changed final transaction result");
+ }
+ }
+
+ private static void FinalFieldsExcludeSensitiveParameters()
+ {
+ var fields = TransactionStreamFields.Required(new Dictionary {
+ { "TRANSACTION_RESULT", "Approved" }, { "PAN", "1234567890123456" },
+ { "PAN_MASKED", "1234567890123456" }, { "CVV", "123" }, { "PIN", "1234" },
+ { "TRACK_DATA", "track-secret" }, { "RAW", "raw-secret" }, { "CARD_HASH", "hash" },
+ { "CARD_REFERENCE", "ref" }, { "RECEIPT_DATA_CARDHOLDER", "receipt" }
+ });
+ Check(!fields.ContainsKey("PAN") && !fields.ContainsKey("CVV") && !fields.ContainsKey("PIN") &&
+ !fields.ContainsKey("TRACK_DATA") && !fields.ContainsKey("RAW") && fields["PAN_MASKED"] == "",
+ "sensitive fields crossed structured boundary");
+ Check(fields["CARD_HASH"] == "hash" && fields["CARD_REFERENCE"] == "ref" && fields["RECEIPT_DATA_CARDHOLDER"] == "receipt", "required final fields were lost");
+ Check(TransactionStreamFields.MaskedCardNumber("123456******1234") == "123456******1234" &&
+ TransactionStreamFields.MaskedCardNumber("476173******1133") == "476173******1133" &&
+ TransactionStreamFields.MaskedCardNumber("************1133") == "************1133" &&
+ TransactionStreamFields.MaskedCardNumber("1234567890123456****") == "" &&
+ TransactionStreamFields.MaskedCardNumber("34373631373330303030303031313333") == "", "valid masking rejected or encoding/full digits treated as masking");
+ Check(!TransactionStreamFields.IsProgress("PAN", "CardRequested") &&
+ !TransactionStreamFields.IsProgress("UPDATE", "1234567890123456") &&
+ !TransactionStreamFields.IsProgress("CARD_STATUS", "Removed"), "unsafe progress permitted");
+ }
+ }
+}
diff --git a/Tests/StreamingTests.csproj b/Tests/StreamingTests.csproj
new file mode 100644
index 0000000..6be403b
--- /dev/null
+++ b/Tests/StreamingTests.csproj
@@ -0,0 +1,23 @@
+
+
+
+ Exe
+ v4.7.2
+ ChipDNAClient.StreamingTests
+ bin\Debug\
+ latest
+
+
+
+
+
+ ..\CompiledLibs\Creditcall.ChipDna.ClientLib.dll
+
+
+ ..\packages\Newtonsoft.Json.11.0.2\lib\net45\Newtonsoft.Json.dll
+
+ TransactionStreaming.cs
+
+
+
+
diff --git a/TransactionStreaming.cs b/TransactionStreaming.cs
new file mode 100644
index 0000000..1363233
--- /dev/null
+++ b/TransactionStreaming.cs
@@ -0,0 +1,269 @@
+using System;
+using System.Collections.Generic;
+using System.Collections.Concurrent;
+using System.IO;
+using System.Threading;
+using System.Threading.Tasks;
+using Newtonsoft.Json;
+
+namespace Creditcall.ChipDna.Client
+{
+ // This is the original start/wait mechanism shared by the XML and NDJSON adapters.
+ internal static class TransactionExecution
+ {
+ internal static Dictionary Run(Action start,
+ Func>> completion, Action suppress, TimeSpan timeout)
+ {
+ try
+ {
+ start();
+ var task = completion();
+ // Task.Wait intentionally retains its existing exception behavior.
+ if (!task.Wait(timeout))
+ {
+ suppress();
+ return new Dictionary
+ {
+ { ParameterKeys.TransactionResult, "TIMEOUT" },
+ { ParameterKeys.Errors, "Transaction timed out waiting for device response" }
+ };
+ }
+ try { return task.Result; }
+ catch (AggregateException agg)
+ {
+ var first = agg.Flatten().InnerExceptions[0];
+ return new Dictionary
+ {
+ { ParameterKeys.TransactionResult, "ERROR" },
+ { ParameterKeys.Errors, $"Transaction task faulted: {first.Message}" }
+ };
+ }
+ }
+ catch
+ {
+ suppress();
+ throw;
+ }
+ }
+ }
+
+ internal sealed class TransactionProgressTap
+ {
+ private readonly object gate = new object();
+ private readonly HashSet references = new HashSet(StringComparer.Ordinal);
+ private Action listener;
+ private string reference;
+ private bool suppressed;
+
+ internal Dictionary Observe(Action callback, Func> execute)
+ {
+ lock (gate) { listener += callback; }
+ try { return execute(); }
+ finally { lock (gate) { listener -= callback; } }
+ }
+ internal void Begin(string value)
+ {
+ lock (gate)
+ {
+ if (suppressed) return;
+ // The existing reference generation can repeat within one second.
+ // A repeated reference cannot prove callback ownership.
+ if (reference != null || string.IsNullOrEmpty(value) || !references.Add(value))
+ suppressed = true;
+ reference = value;
+ }
+ }
+ internal void End() { lock (gate) { reference = null; } }
+ internal void Suppress() { lock (gate) { suppressed = true; reference = null; } }
+ internal void SuppressIfActive()
+ {
+ lock (gate)
+ {
+ if (reference == null) return;
+ suppressed = true;
+ reference = null;
+ }
+ }
+
+ internal void StartFailed(string errors)
+ {
+ lock (gate)
+ {
+ if (suppressed) return;
+ // SDK 3.17 StartCommand returns this exact error before SendRequest.
+ // No other error category proves that dispatch never happened.
+ if (errors == ClientErrorCodes.ClientNotConnectedToServer.ToString())
+ {
+ if (reference != null) references.Remove(reference);
+ reference = null;
+ return;
+ }
+ suppressed = true;
+ reference = null;
+ }
+ }
+
+ internal void Report(string source, EventParameters parameters)
+ {
+ lock (gate)
+ {
+ if (suppressed || reference == null || parameters == null) return;
+ string eventReference;
+ parameters.GetValue(ParameterKeys.Reference, out eventReference);
+ bool hasReference = !string.IsNullOrEmpty(eventReference);
+ if (hasReference && eventReference != reference) return;
+
+ string value;
+ switch (source)
+ {
+ case "UPDATE":
+ if (!parameters.GetValue(ParameterKeys.Update, out value)) return;
+ break;
+ case "CARD_STATUS":
+ if (!parameters.GetValue(ParameterKeys.Notification, out value)) return;
+ break;
+ case "SIGNATURE": value = "Requested"; break;
+ case "PAUSE": value = "Paused"; break;
+ default: return;
+ }
+ if (!TransactionStreamFields.IsProgress(source, value)) return;
+
+ // Normal SDK 3.17 progress has no reference. The serial active window
+ // provides best-effort UI ownership, not a financial decision: a queued
+ // callback, including removal, may briefly show stale progress in a later window.
+ // The listener only attempts a bounded, nonblocking enqueue.
+ listener?.Invoke(source, value);
+ }
+ }
+ }
+
+ internal sealed class TransactionStreamWriter
+ {
+ private readonly object gate = new object();
+ private readonly ConcurrentQueue progress = new ConcurrentQueue();
+ private readonly AutoResetEvent available = new AutoResetEvent(false);
+ private readonly Stream output;
+ private string terminal;
+ private bool completed;
+ private bool disconnected;
+ private bool delivering;
+
+ internal TransactionStreamWriter(Stream output) { this.output = output; }
+
+ internal void Progress(string source, string value)
+ {
+ if (!TransactionStreamFields.IsProgress(source, value)) return;
+ lock (gate)
+ {
+ if (completed || disconnected || progress.Count >= 128) return;
+ progress.Enqueue(JsonConvert.SerializeObject(new { type = "status", source, value }));
+ available.Set();
+ }
+ }
+
+ internal void Complete(object frame)
+ {
+ lock (gate)
+ {
+ if (completed || disconnected) return;
+ completed = true;
+ terminal = JsonConvert.SerializeObject(frame);
+ available.Set();
+ }
+ }
+
+ internal void Deliver()
+ {
+ lock (gate) { if (delivering) return; delivering = true; }
+ try
+ {
+ // One writer owns all progress and terminal writes.
+ using (var writer = new StreamWriter(output, new System.Text.UTF8Encoding(false), 1024, true))
+ {
+ while (true)
+ {
+ string line;
+ bool isTerminal = false;
+ lock (gate)
+ {
+ if (!progress.TryDequeue(out line) && completed)
+ {
+ line = terminal;
+ isTerminal = true;
+ }
+ }
+ if (line == null) { available.WaitOne(); continue; }
+ writer.WriteLine(line);
+ writer.Flush();
+ if (isTerminal) return;
+ }
+ }
+ }
+ catch (Exception ex) when (ex is IOException || ex is System.Net.HttpListenerException || ex is ObjectDisposedException)
+ {
+ lock (gate) { disconnected = true; }
+ }
+ finally
+ {
+ lock (gate)
+ {
+ disconnected = true;
+ available.Dispose();
+ }
+ }
+ }
+ }
+
+ internal static class TransactionStreamFields
+ {
+ internal static bool IsProgress(string source, string value)
+ {
+ switch (source)
+ {
+ case "UPDATE":
+ switch (value)
+ {
+ case "CardRequested": case "ProvideCard": case "InsertOrSwipeCard":
+ case "InsertOrContactless": case "InsertCard": case "SwipeCard":
+ case "CardRemovalRequested": case "CardRemovalEnforced":
+ case "PinEntryStarted": case "PinEntryPrompt": case "PinEntryInProgress":
+ case "PinEntryFailed": case "PinEntrySuccessful": case "TransactionStarted":
+ case "OnlineAuthRequested": case "OnlineAuthorizeRequest": case "OnlineAuthCompleted":
+ return true;
+ }
+ return false;
+ case "CARD_STATUS": return value == "Inserted" || value == "Tapped" || value == "Swiped";
+ case "SIGNATURE": return value == "Requested";
+ case "PAUSE": return value == "Paused";
+ default: return false;
+ }
+ }
+
+ internal static Dictionary Required(Dictionary fields)
+ {
+ var result = new Dictionary();
+ // Receipts remain local to hardlink; no raw parameter bag reaches the kiosk.
+ foreach (var key in new[] { "TRANSACTION_RESULT", "ERRORS", "ERROR_DESCRIPTION",
+ "REFERENCE", "PAN_MASKED", "CARD_SCHEME_ID", "EXPIRY_DATE", "CARD_HASH",
+ "CARD_REFERENCE", "RECEIPT_DATA_CARDHOLDER" })
+ {
+ string value;
+ if (fields.TryGetValue(key, out value))
+ result[key] = key == "PAN_MASKED" ? MaskedCardNumber(value) : value;
+ }
+ return result;
+ }
+
+ internal static string MaskedCardNumber(string value)
+ {
+ int digits = 0, masks = 0;
+ foreach (char ch in value ?? "")
+ {
+ if (ch >= '0' && ch <= '9') digits++;
+ else if (ch == '*' || ch == 'x' || ch == 'X') masks++;
+ else if (ch != ' ' && ch != '-') return "";
+ }
+ return masks >= 4 && digits <= 10 ? value : "";
+ }
+ }
+}