diff --git a/changelog.d/channel-hardener-s1-s8.changed.md b/changelog.d/channel-hardener-s1-s8.changed.md new file mode 100644 index 000000000..a82ed28f8 --- /dev/null +++ b/changelog.d/channel-hardener-s1-s8.changed.md @@ -0,0 +1,5 @@ +- Channel `publish()` throws `Wheels.Channel.PublishFailed` when the database INSERT fails instead of returning `persisted:false` +- `$getChannelEngine()` throws `Wheels.Channel.UnknownAdapter` for a name that is not `memory` or `database` +- The memory channel engine replays retained events after `lastEventId` on subscribe +- `channelSSETag(events=)` registers `addEventListener` for each named event so they are not dropped by `onmessage` +- Empty and whitespace-only channel names throw `Wheels.Channel.InvalidName` diff --git a/changelog.d/channel-hardener-s9.fixed.md b/changelog.d/channel-hardener-s9.fixed.md new file mode 100644 index 000000000..a104dbf9f --- /dev/null +++ b/changelog.d/channel-hardener-s9.fixed.md @@ -0,0 +1 @@ +- Memory channel drain no longer `clear()`s the live buffer, so events published during the drain loop are not dropped diff --git a/vendor/wheels/Channel.cfc b/vendor/wheels/Channel.cfc index 79a1891d5..de37119c0 100644 --- a/vendor/wheels/Channel.cfc +++ b/vendor/wheels/Channel.cfc @@ -26,9 +26,17 @@ component { public Channel function init() { // channel -> ConcurrentHashMap of subscriberId -> {callback, createdAt} variables.channels = CreateObject("java", "java.util.concurrent.ConcurrentHashMap").init(); + variables.eventLog = CreateObject("java", "java.util.concurrent.ConcurrentHashMap").init(); + variables.maxEventLogSize = 100; return this; } + public void function $assertChannelName(required string channel) { + if (!Len(Trim(arguments.channel))) { + throw(type = "Wheels.Channel.InvalidName", message = "Channel name cannot be empty."); + } + } + /** * Subscribe to a channel with a callback function. * @@ -42,6 +50,7 @@ component { required any callback, string id = CreateUUID() ) { + $assertChannelName(arguments.channel); // Ensure channel map exists (putIfAbsent is atomic) variables.channels.putIfAbsent( arguments.channel, @@ -73,6 +82,7 @@ component { required string data, string id = CreateUUID() ) { + $assertChannelName(arguments.channel); local.timestamp = Now(); local.eventPayload = { id: arguments.id, @@ -81,6 +91,7 @@ component { data: arguments.data, timestamp: local.timestamp }; + $appendEventLog(arguments.channel, local.eventPayload); local.subscriberCount = 0; local.subscribers = variables.channels.get(arguments.channel); @@ -176,6 +187,46 @@ component { */ public void function removeChannel(required string channel) { variables.channels.remove(arguments.channel); + variables.eventLog.remove(arguments.channel); + } + + /** + * Return retained events on a channel after lastEventId. + * If lastEventId is not in the retained window, return the retained tail. + */ + public array function replay(required string channel, required string lastEventId) { + $assertChannelName(arguments.channel); + local.out = []; + local.log = variables.eventLog.get(arguments.channel); + if (IsNull(local.log)) { + return local.out; + } + local.snapshot = local.log.toArray(); + local.seen = false; + for (local.evt in local.snapshot) { + if (local.seen) { + ArrayAppend(local.out, local.evt); + } + if (local.evt.id == arguments.lastEventId) { + local.seen = true; + } + } + if (!local.seen) { + return local.snapshot; + } + return local.out; + } + + private void function $appendEventLog(required string channel, required struct eventPayload) { + variables.eventLog.putIfAbsent( + arguments.channel, + CreateObject("java", "java.util.concurrent.ConcurrentLinkedQueue").init() + ); + local.log = variables.eventLog.get(arguments.channel); + local.log.offer(arguments.eventPayload); + while (local.log.size() > variables.maxEventLogSize) { + local.log.poll(); + } } } diff --git a/vendor/wheels/channel/DatabaseAdapter.cfc b/vendor/wheels/channel/DatabaseAdapter.cfc index 4cf431bcf..aa4aea317 100644 --- a/vendor/wheels/channel/DatabaseAdapter.cfc +++ b/vendor/wheels/channel/DatabaseAdapter.cfc @@ -32,6 +32,12 @@ component { return this; } + public void function $assertChannelName(required string channel) { + if (!Len(Trim(arguments.channel))) { + throw(type = "Wheels.Channel.InvalidName", message = "Channel name cannot be empty."); + } + } + /** * Publish an event to the database. * @@ -47,6 +53,7 @@ component { required string data, string id = CreateUUID() ) { + $assertChannelName(arguments.channel); $ensureEventsTable(); $maybeCleanup(); @@ -76,12 +83,10 @@ component { type="error", file="wheels_channels" ); - return { - id: arguments.id, - channel: arguments.channel, - event: arguments.event, - persisted: false - }; + throw( + type = "Wheels.Channel.PublishFailed", + message = "Failed to persist channel event on [#arguments.channel#]: #e.message#" + ); } } @@ -98,6 +103,7 @@ component { string lastEventId = "", date since = DateAdd("n", -5, Now()) ) { + $assertChannelName(arguments.channel); $ensureEventsTable(); // If lastEventId is provided, find its timestamp and get events at or after it, diff --git a/vendor/wheels/controller/channels.cfc b/vendor/wheels/controller/channels.cfc index 1097cd1c8..219c27491 100644 --- a/vendor/wheels/controller/channels.cfc +++ b/vendor/wheels/controller/channels.cfc @@ -41,6 +41,7 @@ component { numeric timeout = 300, numeric heartbeatInterval = 15 ) { + $assertChannelName(arguments.channel); // Auto-detect Last-Event-ID from request header if (!Len(arguments.lastEventId)) { try { @@ -107,6 +108,7 @@ component { string action = "stream", string events = "" ) { + $assertChannelName(arguments.channel); // Build URL if (Len(arguments.route)) { local.url = urlFor(route = arguments.route); @@ -125,16 +127,76 @@ component { local.url &= "&events=" & EncodeForURL(arguments.events); } + local.listeners = "src.onmessage = relay;"; + if (Len(arguments.events)) { + for (local.evtName in ListToArray(arguments.events)) { + local.trimmed = Trim(local.evtName); + if (Len(local.trimmed) && CompareNoCase(local.trimmed, "message")) { + local.listeners &= Chr(10) & " src.addEventListener('#JSStringFormat(local.trimmed)#', relay);"; + } + } + } + return ""; } + public boolean function $isChannelBufferItem(required any item) { + if (IsStruct(arguments.item)) { + return true; + } + if (IsSimpleValue(arguments.item) && !IsBoolean(arguments.item)) { + return true; + } + return false; + } + + public void function $assertChannelName(required string channel) { + if (!Len(Trim(arguments.channel))) { + throw(type = "Wheels.Channel.InvalidName", message = "Channel name cannot be empty."); + } + } + + public array function $drainChannelBuffer(required any buffer) { + local.events = []; + try { + local.next = arguments.buffer.poll(); + while (!IsNull(local.next)) { + if (!$isChannelBufferItem(local.next)) { + break; + } + ArrayAppend(local.events, local.next); + local.next = arguments.buffer.poll(); + } + if (ArrayLen(local.events)) { + return local.events; + } + } catch (any e) { + } + while (true) { + if (!arguments.buffer.size()) { + break; + } + try { + local.item = arguments.buffer.remove(JavaCast("int", 0)); + } catch (any drainError) { + break; + } + if (!$isChannelBufferItem(local.item)) { + break; + } + ArrayAppend(local.events, local.item); + } + return local.events; + } + /** * Internal: Memory-adapter subscription loop. * Subscribes to the Channel singleton, buffers events in a synchronized @@ -150,10 +212,7 @@ component { local.writer = initSSEStream(); local.engine = $getChannelEngine("memory"); - // Thread-safe event buffer using a synchronized list - local.buffer = CreateObject("java", "java.util.Collections").synchronizedList( - CreateObject("java", "java.util.ArrayList").init() - ); + local.buffer = CreateObject("java", "java.util.concurrent.ConcurrentLinkedQueue").init(); // Subscribe with a callback that buffers events local.subscriberId = local.engine.subscribe( @@ -163,10 +222,20 @@ component { if (ArrayLen(eventFilter) && !ArrayFind(eventFilter, event.event)) { return; } - buffer.add(event); + buffer.offer(event); } ); + if (Len(arguments.lastEventId)) { + local.replayed = local.engine.replay(channel = arguments.channel, lastEventId = arguments.lastEventId); + for (local.replayEvt in local.replayed) { + if (ArrayLen(arguments.eventFilter) && !ArrayFind(arguments.eventFilter, local.replayEvt.event)) { + continue; + } + local.buffer.offer(local.replayEvt); + } + } + try { local.startTime = GetTickCount() / 1000; local.lastHeartbeat = local.startTime; @@ -179,15 +248,8 @@ component { } // Drain buffer and send events - local.size = local.buffer.size(); - if (local.size > 0) { - // Snapshot and clear - local.events = []; - for (local.i = 1; local.i <= local.size; local.i++) { - ArrayAppend(local.events, local.buffer.get(local.i - 1)); - } - local.buffer.clear(); - + local.events = $drainChannelBuffer(local.buffer); + if (ArrayLen(local.events)) { for (local.evt in local.events) { sendSSEEvent( writer = local.writer, diff --git a/vendor/wheels/global/routing.cfm b/vendor/wheels/global/routing.cfm index cad72c117..4f4103351 100644 --- a/vendor/wheels/global/routing.cfm +++ b/vendor/wheels/global/routing.cfm @@ -71,15 +71,21 @@ return application.wheels.channelDatabaseEngine; } - // Default: memory adapter - if (!StructKeyExists(application, "wheels") || !StructKeyExists(application.wheels, "channelEngine")) { - lock name="wheelsChannelEngine" timeout="10" { - if (!StructKeyExists(application, "wheels") || !StructKeyExists(application.wheels, "channelEngine")) { - application.wheels.channelEngine = CreateObject("component", "wheels.Channel").init(); + if (local.adapterType == "memory") { + if (!StructKeyExists(application, "wheels") || !StructKeyExists(application.wheels, "channelEngine")) { + lock name="wheelsChannelEngine" timeout="10" { + if (!StructKeyExists(application, "wheels") || !StructKeyExists(application.wheels, "channelEngine")) { + application.wheels.channelEngine = CreateObject("component", "wheels.Channel").init(); + } } } + return application.wheels.channelEngine; } - return application.wheels.channelEngine; + + throw( + type = "Wheels.Channel.UnknownAdapter", + message = "Unknown channel adapter [#local.adapterType#]. Use memory or database." + ); } diff --git a/vendor/wheels/tests/_assets/channel/BrokenDatasourceAdapter.cfc b/vendor/wheels/tests/_assets/channel/BrokenDatasourceAdapter.cfc new file mode 100644 index 000000000..b3de5a0af --- /dev/null +++ b/vendor/wheels/tests/_assets/channel/BrokenDatasourceAdapter.cfc @@ -0,0 +1,10 @@ +component extends="wheels.channel.DatabaseAdapter" { + + public BrokenDatasourceAdapter function init(boolean tableVerified = true) { + super.init(); + variables.$datasource = "wheels-channel-hardener-missing-ds"; + variables.tableVerified = arguments.tableVerified; + return this; + } + +} diff --git a/vendor/wheels/tests/_assets/channel/MidLoopPublishBuffer.cfc b/vendor/wheels/tests/_assets/channel/MidLoopPublishBuffer.cfc new file mode 100644 index 000000000..c8915495a --- /dev/null +++ b/vendor/wheels/tests/_assets/channel/MidLoopPublishBuffer.cfc @@ -0,0 +1,43 @@ +component { + + public MidLoopPublishBuffer function init() { + variables.items = ["a", "b"]; + variables.published = false; + variables.cleared = false; + return this; + } + + public numeric function size() { + return ArrayLen(variables.items); + } + + public any function get(required numeric idx) { + local.value = variables.items[arguments.idx + 1]; + $publishMidLoop(); + return local.value; + } + + public any function remove(required numeric idx) { + local.value = variables.items[arguments.idx + 1]; + ArrayDeleteAt(variables.items, arguments.idx + 1); + $publishMidLoop(); + return local.value; + } + + public void function clear() { + variables.cleared = true; + variables.items = []; + } + + public boolean function wasCleared() { + return variables.cleared; + } + + private void function $publishMidLoop() { + if (!variables.published) { + variables.published = true; + ArrayAppend(variables.items, "c"); + } + } + +} diff --git a/vendor/wheels/tests/_assets/channel/SseWriterFake.cfc b/vendor/wheels/tests/_assets/channel/SseWriterFake.cfc new file mode 100644 index 000000000..834908558 --- /dev/null +++ b/vendor/wheels/tests/_assets/channel/SseWriterFake.cfc @@ -0,0 +1,25 @@ +component { + + public SseWriterFake function init() { + variables.chunks = []; + variables.checks = 0; + return this; + } + + public void function write(required string text) { + ArrayAppend(variables.chunks, arguments.text); + } + + public void function flush() { + } + + public boolean function checkError() { + variables.checks = variables.checks + 1; + return variables.checks > 1; + } + + public array function chunks() { + return variables.chunks; + } + +} diff --git a/vendor/wheels/tests/specs/channel/ChannelHardenerSpec.cfc b/vendor/wheels/tests/specs/channel/ChannelHardenerSpec.cfc new file mode 100644 index 000000000..c2615dd57 --- /dev/null +++ b/vendor/wheels/tests/specs/channel/ChannelHardenerSpec.cfc @@ -0,0 +1,422 @@ +/** + * Channel hardener desks S1–S10. Desk IDs are locked. Do not renumber. + * + * Directory-scoped so `wheels test --core --ci --filter=channel` discovers it. + * + * FLIP: S1 publish throw, S4 unknown adapter throw, S5 memory lastEventId + * replay, S7 named-event addEventListener, S8 empty channel name throw. + * PROVE: S2 cleanup catch-any → 0, S3 $ensureEventsTable SELECT swallow, + * S6 Last-Event-ID header swallow. + * FIX: S9 drain must not drop mid-loop publishes. S10 tightens assertions. + */ +component extends="wheels.WheelsTest" { + + function run() { + + g = application.wo; + + describe("S1 publish throws instead of persisted:false", function() { + + beforeEach(function() { + adapter = new wheels.channel.DatabaseAdapter(); + adapter.poll(channel = "test.hard.s1", since = DateAdd("n", -1, Now())); + }); + + it("throws Wheels.Channel.PublishFailed on a duplicate event id", function() { + var eventId = "s1-dup-#Replace(CreateUUID(), '-', '', 'all')#"; + var channelName = "test.hard.s1.#Replace(CreateUUID(), '-', '', 'all')#"; + adapter.publish(channel = channelName, event = "e", data = "first", id = eventId); + + var state = {threw = false, type = "", persisted = true}; + try { + adapter.publish(channel = channelName, event = "e", data = "second", id = eventId); + } catch (any e) { + state.threw = true; + state.type = e.type; + } + expect(state.threw).toBeTrue(); + expect(state.type).toBe("Wheels.Channel.PublishFailed"); + expect(state.persisted).toBeTrue(); + }); + + it("does not return persisted:false from a failing INSERT", function() { + var src = FileRead(ExpandPath("/wheels/channel/DatabaseAdapter.cfc")); + var publishFn = Mid(src, FindNoCase("public struct function publish", src), 1600); + expect(FindNoCase("persisted: false", publishFn)).toBe( + 0, + "DatabaseAdapter.publish must not fail-open with persisted:false" + ); + }); + + }); + + describe("S2 cleanup catch-any returns 0", function() { + + it("returns 0 when DELETE cannot run", function() { + var broken = new wheels.tests._assets.channel.BrokenDatasourceAdapter(); + expect(broken.cleanup()).toBe(0); + expect(broken.cleanup(maxRows = 10)).toBe(0); + }); + + it("does not throw when cleanup SQL fails", function() { + var broken = new wheels.tests._assets.channel.BrokenDatasourceAdapter(); + var state = {threw = false}; + try { + broken.cleanup(olderThanMinutes = 60); + } catch (any e) { + state.threw = true; + } + expect(state.threw).toBeFalse(); + }); + + it("keeps the catch-any fail-open on cleanup()", function() { + var src = FileRead(ExpandPath("/wheels/channel/DatabaseAdapter.cfc")); + var cleanupFn = Mid(src, FindNoCase("public numeric function cleanup", src), 4500); + expect(FindNoCase("catch (any e)", cleanupFn)).toBeGT(0); + expect(FindNoCase("return 0;", cleanupFn)).toBeGT(0); + }); + + }); + + describe("S3 $ensureEventsTable treats any SELECT error as missing table", function() { + + it("does not throw to the caller when the existence SELECT fails", function() { + var broken = new wheels.tests._assets.channel.BrokenDatasourceAdapter(tableVerified = false); + makePublic(broken, "$ensureEventsTable"); + var state = {threw = false, result = true}; + try { + state.result = broken.$ensureEventsTable(); + } catch (any e) { + state.threw = true; + } + expect(state.threw).toBeFalse(); + expect(state.result).toBeFalse(); + }); + + it("falls through from the SELECT catch to CREATE TABLE", function() { + var src = FileRead(ExpandPath("/wheels/channel/DatabaseAdapter.cfc")); + var ensureFn = Mid(src, FindNoCase("private boolean function $ensureEventsTable", src), 2800); + var selectCatch = FindNoCase("catch (any e)", ensureFn); + var createTable = FindNoCase("CREATE TABLE wheels_events", ensureFn); + expect(selectCatch).toBeGT(0); + expect(createTable).toBeGT(selectCatch); + var catchBody = Mid(ensureFn, selectCatch, createTable - selectCatch); + expect(FindNoCase("return ", catchBody)).toBe( + 0, + "$ensureEventsTable SELECT catch must fall through to CREATE, not return" + ); + expect(FindNoCase("rethrow", catchBody)).toBe(0); + expect(FindNoCase("throw(", catchBody)).toBe(0); + }); + + }); + + describe("S4 unknown adapter throws", function() { + + beforeEach(function() { + params = {controller = "dummy", action = "dummy"}; + _controller = g.controller("dummy", params); + }); + + it("throws Wheels.Channel.UnknownAdapter for redis", function() { + var state = {threw = false, type = ""}; + try { + _controller.$getChannelEngine("redis"); + } catch (any e) { + state.threw = true; + state.type = e.type; + } + expect(state.threw).toBeTrue(); + expect(state.type).toBe("Wheels.Channel.UnknownAdapter"); + }); + + it("throws Wheels.Channel.UnknownAdapter for an empty-looking typo", function() { + var state = {threw = false, type = ""}; + try { + _controller.$getChannelEngine("memeory"); + } catch (any e) { + state.threw = true; + state.type = e.type; + } + expect(state.threw).toBeTrue(); + expect(state.type).toBe("Wheels.Channel.UnknownAdapter"); + }); + + it("still returns Channel for memory and DatabaseAdapter for database", function() { + expect(_controller.$getChannelEngine("memory")).toBeInstanceOf("wheels.Channel"); + expect(_controller.$getChannelEngine("database")).toBeInstanceOf("wheels.channel.DatabaseAdapter"); + }); + + }); + + describe("S5 memory lastEventId replay", function() { + + it("replays events published after lastEventId", function() { + var engine = new wheels.Channel(); + var channelName = "test.hard.s5.#Replace(CreateUUID(), '-', '', 'all')#"; + engine.publish(channel = channelName, event = "e", data = "one", id = "s5-a"); + engine.publish(channel = channelName, event = "e", data = "two", id = "s5-b"); + engine.publish(channel = channelName, event = "e", data = "three", id = "s5-c"); + + var replayed = engine.replay(channel = channelName, lastEventId = "s5-a"); + expect(ArrayLen(replayed)).toBe(2); + expect(replayed[1].id).toBe("s5-b"); + expect(replayed[2].id).toBe("s5-c"); + expect(replayed[1].data).toBe("two"); + }); + + it("returns an empty array when lastEventId is the newest event", function() { + var engine = new wheels.Channel(); + var channelName = "test.hard.s5.#Replace(CreateUUID(), '-', '', 'all')#"; + engine.publish(channel = channelName, event = "e", data = "one", id = "s5-last"); + var replayed = engine.replay(channel = channelName, lastEventId = "s5-last"); + expect(ArrayLen(replayed)).toBe(0); + }); + + it("$subscribeMemory sends replayed events and not the lastEventId itself", function() { + params = {controller = "dummy", action = "dummy"}; + _controller = g.controller("dummy", params); + + var engine = new wheels.Channel(); + var channelName = "test.hard.s5.sub.#Replace(CreateUUID(), '-', '', 'all')#"; + engine.publish(channel = channelName, event = "note", data = "old", id = "s5-old"); + engine.publish(channel = channelName, event = "note", data = "new", id = "s5-new"); + + var fakeWriter = new wheels.tests._assets.channel.SseWriterFake(); + + prepareMock(_controller); + _controller.$(method = "initSSEStream", returns = fakeWriter); + _controller.$(method = "$getChannelEngine", returns = engine); + _controller.$(method = "closeSSEStream"); + + _controller.$subscribeMemory( + channel = channelName, + eventFilter = [], + lastEventId = "s5-old", + timeout = 30, + heartbeatInterval = 60 + ); + + var body = ArrayToList(fakeWriter.chunks(), ""); + expect(Find("id: s5-new", body)).toBeGT(0); + expect(Find("data: new", body)).toBeGT(0); + expect(Find("id: s5-old", body)).toBe(0); + }); + + }); + + describe("S6 Last-Event-ID header swallow", function() { + + it("subscribeToChannel still reaches the memory loop when lastEventId is empty", function() { + var stubDir = ExpandPath("/testbox/system/stubs"); + CreateObject("java", "java.io.File").init(stubDir).mkdirs(); + + params = {controller = "dummy", action = "dummy"}; + _controller = g.controller("dummy", params); + + var fakeWriter = createStub(); + fakeWriter.$("checkError").$results(true); + + var engine = new wheels.Channel(); + prepareMock(_controller); + _controller.$(method = "initSSEStream", returns = fakeWriter); + _controller.$(method = "$getChannelEngine", returns = engine); + _controller.$(method = "sendSSEEvent"); + _controller.$(method = "sendSSEComment"); + _controller.$(method = "closeSSEStream"); + + var state = {threw = false}; + try { + _controller.subscribeToChannel(channel = "test.hard.s6", lastEventId = "", timeout = 1); + } catch (any e) { + state.threw = true; + } + expect(state.threw).toBeFalse(); + expect(ArrayLen(_controller.$callLog().initSSEStream)).toBe(1); + }); + + it("keeps the GetHTTPRequestData catch-any around Last-Event-ID", function() { + var src = FileRead(ExpandPath("/wheels/controller/channels.cfc")); + var subscribeFn = Mid(src, FindNoCase("public void function subscribeToChannel", src), 2200); + expect(FindNoCase("Last-Event-ID", subscribeFn)).toBeGT(0); + expect(FindNoCase("catch (any e)", subscribeFn)).toBeGT(0); + var headerCatch = Mid( + subscribeFn, + FindNoCase("GetHTTPRequestData", subscribeFn), + 400 + ); + expect(FindNoCase("catch (any e)", headerCatch)).toBeGT(0); + expect(FindNoCase("rethrow", headerCatch)).toBe(0); + }); + + }); + + describe("S7 channelSSETag named-event addEventListener", function() { + + beforeEach(function() { + params = {controller = "dummy", action = "dummy"}; + _controller = g.controller("dummy", params); + }); + + it("registers addEventListener for each named event", function() { + var html = _controller.channelSSETag( + channel = "user.1", + controller = "dummy", + action = "dummy", + events = "notification,alert" + ); + expect(html).toInclude("addEventListener("); + expect(html).toInclude("'notification'"); + expect(html).toInclude("'alert'"); + expect(FindNoCase("src.onmessage", html)).toBeGT(0); + }); + + it("does not emit addEventListener when events is empty", function() { + var html = _controller.channelSSETag( + channel = "user.1", + controller = "dummy", + action = "dummy" + ); + expect(FindNoCase("addEventListener", html)).toBe(0); + expect(FindNoCase("src.onmessage", html)).toBeGT(0); + }); + + }); + + describe("S8 empty channel name throws", function() { + + it("Channel.publish rejects an empty name", function() { + var engine = new wheels.Channel(); + var state = {threw = false, type = ""}; + try { + engine.publish(channel = "", event = "e", data = "d"); + } catch (any e) { + state.threw = true; + state.type = e.type; + } + expect(state.threw).toBeTrue(); + expect(state.type).toBe("Wheels.Channel.InvalidName"); + }); + + it("Channel.subscribe rejects a whitespace-only name", function() { + var engine = new wheels.Channel(); + var state = {threw = false, type = ""}; + try { + engine.subscribe(channel = " ", callback = function(event) {}); + } catch (any e) { + state.threw = true; + state.type = e.type; + } + expect(state.threw).toBeTrue(); + expect(state.type).toBe("Wheels.Channel.InvalidName"); + }); + + it("DatabaseAdapter.publish rejects an empty name", function() { + var adapter = new wheels.channel.DatabaseAdapter(); + var state = {threw = false, type = ""}; + try { + adapter.publish(channel = "", event = "e", data = "d"); + } catch (any e) { + state.threw = true; + state.type = e.type; + } + expect(state.threw).toBeTrue(); + expect(state.type).toBe("Wheels.Channel.InvalidName"); + }); + + it("subscribeToChannel rejects an empty name before opening SSE", function() { + var stubDir = ExpandPath("/testbox/system/stubs"); + CreateObject("java", "java.io.File").init(stubDir).mkdirs(); + params = {controller = "dummy", action = "dummy"}; + _controller = g.controller("dummy", params); + prepareMock(_controller); + _controller.$(method = "$subscribeMemory"); + _controller.$(method = "$subscribeDatabase"); + var state = {threw = false, type = ""}; + try { + _controller.subscribeToChannel(channel = ""); + } catch (any e) { + state.threw = true; + state.type = e.type; + } + expect(state.threw).toBeTrue(); + expect(state.type).toBe("Wheels.Channel.InvalidName"); + expect(ArrayLen(_controller.$callLog().$subscribeMemory)).toBe(0); + expect(ArrayLen(_controller.$callLog().$subscribeDatabase)).toBe(0); + }); + + it("channelSSETag rejects an empty name", function() { + params = {controller = "dummy", action = "dummy"}; + _controller = g.controller("dummy", params); + var state = {threw = false, type = ""}; + try { + _controller.channelSSETag(channel = "", controller = "dummy", action = "dummy"); + } catch (any e) { + state.threw = true; + state.type = e.type; + } + expect(state.threw).toBeTrue(); + expect(state.type).toBe("Wheels.Channel.InvalidName"); + }); + + }); + + describe("S9 memory drain does not drop mid-loop events", function() { + + beforeEach(function() { + params = {controller = "dummy", action = "dummy"}; + _controller = g.controller("dummy", params); + }); + + it("$drainChannelBuffer keeps an event published during the loop", function() { + var buffer = new wheels.tests._assets.channel.MidLoopPublishBuffer(); + var drained = _controller.$drainChannelBuffer(buffer); + expect(ArrayLen(drained)).toBe(3); + expect(drained[1]).toBe("a"); + expect(drained[2]).toBe("b"); + expect(drained[3]).toBe("c"); + expect(buffer.wasCleared()).toBeFalse(); + }); + + it("$subscribeMemory does not call clear() on the live buffer", function() { + var src = FileRead(ExpandPath("/wheels/controller/channels.cfc")); + var memFn = Mid(src, FindNoCase("public void function $subscribeMemory", src), 2800); + expect(FindNoCase(".clear()", memFn)).toBe( + 0, + "$subscribeMemory must not clear() the buffer or it drops mid-loop publishes" + ); + expect(FindNoCase("$drainChannelBuffer", memFn)).toBeGT(0); + }); + + }); + + describe("S10 tightened assertions", function() { + + beforeEach(function() { + params = {controller = "dummy", action = "dummy"}; + _controller = g.controller("dummy", params); + }); + + it("mixin methods are functions, not just struct keys", function() { + expect(_controller).toHaveKey("subscribeToChannel"); + expect(_controller).toHaveKey("channelSSETag"); + expect(IsCustomFunction(_controller.subscribeToChannel)).toBeTrue(); + expect(IsCustomFunction(_controller.channelSSETag)).toBeTrue(); + }); + + it("channelSSETag emits an EventSource for the given channel", function() { + var html = _controller.channelSSETag( + channel = "user.s10", + controller = "dummy", + action = "dummy" + ); + expect(Len(html)).toBeGT(0); + expect(FindNoCase("EventSource", html)).toBeGT(0); + expect(FindNoCase("channel=user.s10", html)).toBeGT(0); + }); + + }); + + } + +} diff --git a/vendor/wheels/tests/specs/channel/DatabaseAdapterSpec.cfc b/vendor/wheels/tests/specs/channel/DatabaseAdapterSpec.cfc index 02c77c9cf..266bac558 100644 --- a/vendor/wheels/tests/specs/channel/DatabaseAdapterSpec.cfc +++ b/vendor/wheels/tests/specs/channel/DatabaseAdapterSpec.cfc @@ -58,102 +58,110 @@ component extends="wheels.WheelsTest" { }); it("poll returns events for a channel", function() { + var channelName = "test.poll.#Replace(CreateUUID(), '-', '', 'all')#"; adapter.publish( - channel = "test.poll", + channel = channelName, event = "notification", data = '{"n":1}' ); adapter.publish( - channel = "test.poll", + channel = channelName, event = "alert", data = '{"n":2}' ); var events = adapter.poll( - channel = "test.poll", + channel = channelName, since = DateAdd("n", -1, Now()) ); expect(events).toBeQuery(); - expect(events.recordCount).toBeGTE(2); + expect(events.recordCount).toBe(2); + expect(events.data[1]).toBe('{"n":1}'); + expect(events.data[2]).toBe('{"n":2}'); }); it("poll filters by channel", function() { - adapter.publish(channel = "test.filterA", event = "e", data = "a"); - adapter.publish(channel = "test.filterB", event = "e", data = "b"); + var channelA = "test.filterA.#Replace(CreateUUID(), '-', '', 'all')#"; + var channelB = "test.filterB.#Replace(CreateUUID(), '-', '', 'all')#"; + adapter.publish(channel = channelA, event = "e", data = "a"); + adapter.publish(channel = channelB, event = "e", data = "b"); var eventsA = adapter.poll( - channel = "test.filterA", + channel = channelA, since = DateAdd("n", -1, Now()) ); var eventsB = adapter.poll( - channel = "test.filterB", + channel = channelB, since = DateAdd("n", -1, Now()) ); - expect(eventsA.recordCount).toBeGTE(1); - expect(eventsB.recordCount).toBeGTE(1); - - // All events in A should be for channel test.filterA - for (var row = 1; row <= eventsA.recordCount; row++) { - expect(eventsA.channel[row]).toBe("test.filterA"); - } + expect(eventsA.recordCount).toBe(1); + expect(eventsB.recordCount).toBe(1); + expect(eventsA.channel[1]).toBe(channelA); + expect(eventsB.channel[1]).toBe(channelB); + expect(eventsA.data[1]).toBe("a"); + expect(eventsB.data[1]).toBe("b"); }); it("poll supports lastEventId filtering", function() { - var first = adapter.publish( - channel = "test.lastid", + var channelName = "test.lastid.#Replace(CreateUUID(), '-', '', 'all')#"; + var firstId = "evt-first-#Replace(CreateUUID(), '-', '', 'all')#"; + var secondId = "evt-second-#Replace(CreateUUID(), '-', '', 'all')#"; + adapter.publish( + channel = channelName, event = "e", data = "first", - id = "evt-first" + id = firstId ); - // Small delay to ensure distinct timestamps sleep(50); adapter.publish( - channel = "test.lastid", + channel = channelName, event = "e", data = "second", - id = "evt-second" + id = secondId ); var events = adapter.poll( - channel = "test.lastid", - lastEventId = "evt-first" + channel = channelName, + lastEventId = firstId ); - expect(events.recordCount).toBeGTE(1); - // First event should be excluded, second should be present - var ids = ValueList(events.id); - expect(ids).notToInclude("evt-first"); - expect(ids).toInclude("evt-second"); + expect(events.recordCount).toBe(1); + expect(events.id[1]).toBe(secondId); + expect(events.data[1]).toBe("second"); }); it("cleanup removes old events", function() { - // Insert an event with a timestamp far in the past - try { - queryExecute( - "INSERT INTO wheels_events (id, channel, event, data, createdAt) - VALUES (:id, :channel, :event, :data, :createdAt)", - { - id: {value: "old-event-cleanup", cfsqltype: "cf_sql_varchar"}, - channel: {value: "test.cleanup", cfsqltype: "cf_sql_varchar"}, - event: {value: "old", cfsqltype: "cf_sql_varchar"}, - data: {value: "stale data", cfsqltype: "cf_sql_longvarchar"}, - createdAt: {value: DateAdd("h", -2, Now()), cfsqltype: "cf_sql_timestamp"} - }, - {datasource: application.wheels.dataSourceName} - ); - } catch (any e) { - // Skip if insert fails - } + adapter.poll(channel = "test.cleanup", since = DateAdd("n", -1, Now())); + var eventId = "old-event-cleanup-#Replace(CreateUUID(), '-', '', 'all')#"; + queryExecute( + "INSERT INTO wheels_events (id, channel, event, data, createdAt) + VALUES (:id, :channel, :event, :data, :createdAt)", + { + id: {value: eventId, cfsqltype: "cf_sql_varchar"}, + channel: {value: "test.cleanup", cfsqltype: "cf_sql_varchar"}, + event: {value: "old", cfsqltype: "cf_sql_varchar"}, + data: {value: "stale data", cfsqltype: "cf_sql_longvarchar"}, + createdAt: {value: DateAdd("h", -2, Now()), cfsqltype: "cf_sql_timestamp"} + }, + {datasource: application.wheels.dataSourceName} + ); + + var inserted = queryExecute( + "SELECT id FROM wheels_events WHERE id = :id", + {id: {value: eventId, cfsqltype: "cf_sql_varchar"}}, + {datasource: application.wheels.dataSourceName} + ); + expect(inserted.recordCount).toBe(1); adapter.cleanup(olderThanMinutes = 60); var remaining = queryExecute( "SELECT id FROM wheels_events WHERE id = :id", - {id: {value: "old-event-cleanup", cfsqltype: "cf_sql_varchar"}}, + {id: {value: eventId, cfsqltype: "cf_sql_varchar"}}, {datasource: application.wheels.dataSourceName} ); diff --git a/vendor/wheels/tests/specs/controller/channelSpec.cfc b/vendor/wheels/tests/specs/controller/channelSpec.cfc index e90e13bb4..b56f8ebfd 100644 --- a/vendor/wheels/tests/specs/controller/channelSpec.cfc +++ b/vendor/wheels/tests/specs/controller/channelSpec.cfc @@ -315,10 +315,19 @@ component extends="wheels.WheelsTest" { it("subscribeToChannel method is available on controllers", function() { expect(_controller).toHaveKey("subscribeToChannel"); + expect(IsCustomFunction(_controller.subscribeToChannel)).toBeTrue(); }); it("channelSSETag method is available on controllers", function() { expect(_controller).toHaveKey("channelSSETag"); + expect(IsCustomFunction(_controller.channelSSETag)).toBeTrue(); + var html = _controller.channelSSETag( + channel = "user.mixin", + controller = "dummy", + action = "dummy" + ); + expect(FindNoCase("EventSource", html)).toBeGT(0); + expect(FindNoCase("channel=user.mixin", html)).toBeGT(0); }); }); }