Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 8 additions & 1 deletion .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,14 @@ on:

jobs:
style:
uses: PlaceOS/.github/.github/workflows/crystal-style.yml@main
runs-on: ubuntu-latest
container: crystallang/crystal:latest
steps:
- uses: actions/checkout@v7.0.1
- name: Format
run: crystal tool format --check
- name: Lint
uses: crystal-ameba/github-action@v1.0.0

test:
uses: PlaceOS/.github/.github/workflows/containerised-test.yml@main
8 changes: 8 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,14 @@ We use [RethinkDB](https://rethinkdb.com) to unify our database and event bus, g
| `PG_LOCK_TIMEOUT` | Timeout on retrying Advisory lock in seconds | 5 |
| `PG_DATABASE_URL` | Or provide a Database DSN | |

## Control-system telemetry notifications

`ControlSystem` uses pg-orm's model-level changefeed policy to ignore updates confined to `signage_last_seen` and `playlist_item_id`. The SQL trigger skips creating a CDC event for those updates, while the timestamp and current item are still persisted. Configuration updates, including updates that also change telemetry, continue to notify; inserts and deletes are unchanged.

The policy is installed when the control-system changefeed is registered. It applies to all writers of these two fields, and no-op updates on `sys` are also silent. Consumers that need current telemetry should query PostgreSQL rather than rely on configuration changefeeds.

Deploy EventBus 1.1.0 or newer to every service that installs CDC triggers before enabling this models version. Older installers can restore the combined trigger and produce unwanted or duplicate update events. pg-orm 2.4.0 passes the model declaration to EventBus; no database migration or core-side filter is required. For rollback, use EventBus's `replace_cdc_update_policy` with the expected current columns; merely removing the model declaration preserves the installed policy.

## Testing

```shell
Expand Down
2 changes: 0 additions & 2 deletions shard.override.yml
Original file line number Diff line number Diff line change
Expand Up @@ -4,8 +4,6 @@ dependencies:
commit: 2c156ab8ce688d676ca912b25ccdc631b265952c
retriable:
github: Sija/retriable.cr
pg-orm:
github: spider-gazelle/pg-orm
sanitize:
github: straight-shoota/sanitize
branch: master
1 change: 1 addition & 0 deletions shard.yml
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,7 @@ dependencies:
# ORM for Postgres built on active-model
pg-orm:
github: spider-gazelle/pg-orm
version: ">= 2.4.0"

# PlaceOS calendar for Tenant model
place_calendar:
Expand Down
93 changes: 93 additions & 0 deletions spec/control_system_changefeed_spec.cr
Original file line number Diff line number Diff line change
@@ -0,0 +1,93 @@
require "./helper"

module PlaceOS::Model
private def self.receive_signage_notification(notifications : Channel(String)) : String
select
when message = notifications.receive
message
when timeout(5.seconds)
raise "Timed out waiting for signage CDC notification"
end
end

private def self.signage_cdc_actions(id : String) : Array(String)
PgORM::Database.connection do |db|
db.query_all("SELECT event_action FROM public.eventbus_cdc_events WHERE event_table = 'sys' AND row_id = $1 ORDER BY id", args: [id], as: String)
end
end

describe "ControlSystem signage changefeed" do
it "persists heartbeats without CDC rows or notifications and retains configuration events" do
system = Generator.control_system
item = Generator.item.save!
notifications = Channel(String).new(32)
barrier = "signage-barrier-#{RANDOM.hex(8)}"
listener = PG::ListenConnection.new(ENV["PG_DATABASE_URL"], ["cdc_events", "signage_spec_barrier"]) do |notification|
payload = notification.payload
if payload == barrier || (JSON.parse(payload)["table"].as_s == "sys" && JSON.parse(payload)["id"].as_s == system.id)
notifications.send(payload)
end
end
feed = ControlSystem.changes
system.save!
id = system.id.as(String)
JSON.parse(receive_signage_notification(notifications))["action"].as_s.should eq("insert")
signage_cdc_actions(id).should eq(["insert"])

# Repeated items and both forms of clearing must persist without CDC work.
{item.id, item.id, nil, item.id, ""}.each do |item_id|
previous = system.signage_last_seen
system.update_last_seen_time(item_id)
system.reload!
system.signage_last_seen.should be > previous
system.playlist_item_id.should eq(item_id.presence)
signage_cdc_actions(id).should eq(["insert"])
end
# A committed notification to the same listener proves prior writes emitted none.
PgORM::Database.connection { |db| db.exec("SELECT pg_notify('signage_spec_barrier', $1)", args: [barrier]) }
receive_signage_notification(notifications).should eq(barrier)

system.name = "renamed-#{RANDOM.hex(8)}"
system.save!
JSON.parse(receive_signage_notification(notifications))["action"].as_s.should eq("update")

PgORM::Database.connection do |db|
db.exec("UPDATE sys SET description = 'mixed update', signage_last_seen = now(), playlist_item_id = $2 WHERE id = $1", args: [id, item.id])
end
JSON.parse(receive_signage_notification(notifications))["action"].as_s.should eq("update")
system.reload!
system.description.should eq("mixed update")
system.playlist_item_id.should eq(item.id)
system.destroy
JSON.parse(receive_signage_notification(notifications))["action"].as_s.should eq("delete")
signage_cdc_actions(id).should eq(["insert", "update", "update", "delete"])
ensure
listener.try &.close
feed.try &.stop
system.try &.delete
item.try &.delete
end

it "keeps other models' update notifications" do
zone = Generator.zone.save!
notifications = Channel(String).new(4)
listener = PG::ListenConnection.new(ENV["PG_DATABASE_URL"], ["cdc_events"]) do |notification|
payload = JSON.parse(notification.payload)
if payload["table"].as_s == "zone" && payload["id"].as_s == zone.id
notifications.send(notification.payload)
end
end
feed = Zone.changes(zone.id)
zone.name = "updated-#{RANDOM.hex(8)}"
zone.save!
JSON.parse(receive_signage_notification(notifications))["action"].as_s.should eq("update")
PgORM::Database.connection do |db|
db.query_all("SELECT event_action FROM public.eventbus_cdc_events WHERE event_table = 'zone' AND row_id = $1 AND event_action = 'update'", args: [zone.id], as: String).should eq(["update"])
end
ensure
listener.try &.close
feed.try &.stop
zone.try &.delete
end
end
end
2 changes: 1 addition & 1 deletion spec/group_permissions_spec.cr
Original file line number Diff line number Diff line change
Expand Up @@ -94,7 +94,7 @@ module PlaceOS::Model
sql = Group.accessible_zones_sql(
authority.id.not_nil!, subsystems || [subsystem], user.id.not_nil!, required,
)
return nil if sql.nil?
return if sql.nil?
::PgORM::Database.connection do |db|
db.query_all("SELECT * FROM #{sql} AS q(zone_id)", &.read(String))
end.sort!
Expand Down
2 changes: 1 addition & 1 deletion spec/shortener_spec.cr
Original file line number Diff line number Diff line change
Expand Up @@ -51,7 +51,7 @@ module PlaceOS::Model
Array.new(16) { Generator.shortener(user: user, authority: authority).save! }
end

shorteners.each { |short| short.persisted?.should be_true }
shorteners.each(&.persisted?.should(be_true))
Shortener.count.should eq 16

ids = shorteners.map(&.id.as(String))
Expand Down
8 changes: 4 additions & 4 deletions spec/storage_spec.cr
Original file line number Diff line number Diff line change
Expand Up @@ -48,23 +48,23 @@ module PlaceOS::Model
it "should return default storage when flagged" do
authority = Generator.authority.save!

def1 = Generator.storage.save!
def1 = Generator.storage(bucket: "default-bucket-1").save!
def1.authority_id.should be_nil
def2 = Generator.storage.save!
def2 = Generator.storage(bucket: "default-bucket-2").save!
def2.authority_id.should be_nil
def1.reload!
def2.reload!

def1.is_default.should be_false
def2.is_default.should be_true

s2 = Generator.storage
s2 = Generator.storage(bucket: "authority-bucket-1")
s2.authority_id = authority.id
s2.is_default = false
s2.save!
ret_store = Storage.storage_or_default(authority.id).id.should eq s2.id

s3 = Generator.storage
s3 = Generator.storage(bucket: "authority-bucket-2")
s3.authority_id = authority.id
s3.is_default = true
s3.save!
Expand Down
10 changes: 4 additions & 6 deletions spec/user_spec.cr
Original file line number Diff line number Diff line change
Expand Up @@ -142,12 +142,10 @@ module PlaceOS::Model
Array.new(4) { Generator.user(admin: true).save! }
.map { |u|
future do
begin
u.destroy
rescue e : Model::Error
e.message.should eq "At least one admin must remain"
errors << e
end
u.destroy
rescue e : Model::Error
e.message.should eq "At least one admin must remain"
errors << e
end
}.each &.get
end
Expand Down
3 changes: 3 additions & 0 deletions src/placeos-models/control_system.cr
Original file line number Diff line number Diff line change
Expand Up @@ -73,6 +73,9 @@ module PlaceOS::Model
attribute signage_last_seen : Time = -> { 5.hours.ago }, converter: PlaceOS::Model::Timestamps::EpochConverter, type: "integer", format: "Int64", mass_assignment: false
belongs_to Playlist::Item, foreign_key: "playlist_item_id"

# Signage telemetry persists without notifying services to reload running drivers.
changefeed_ignore_updates :signage_last_seen, :playlist_item_id

attribute space_config : Hash(String, JSON::Any) = {} of String => JSON::Any

def update_last_seen_time(item_id : String? = nil)
Expand Down
10 changes: 5 additions & 5 deletions src/placeos-models/group.cr
Original file line number Diff line number Diff line change
Expand Up @@ -163,10 +163,10 @@ module PlaceOS::Model
required : Permissions,
) : String?
direct_memberships = user_direct_memberships(authority_id, user_id)
return nil if direct_memberships.empty?
return if direct_memberships.empty?

authority_groups = Group.where(authority_id: authority_id).to_a
return nil if authority_groups.empty?
return if authority_groups.empty?

parent_of = {} of UUID => UUID
all_group_ids = [] of UUID
Expand All @@ -179,7 +179,7 @@ module PlaceOS::Model
end
member_group_ids << gid if g.subsystems.any? { |sub| subsystems.includes?(sub) }
end
return nil if member_group_ids.empty?
return if member_group_ids.empty?

effective_memberships = walk_up_memberships(direct_memberships, all_group_ids, parent_of)
group_depths = compute_group_depths(all_group_ids, parent_of)
Expand All @@ -189,7 +189,7 @@ module PlaceOS::Model
GroupZone.where(group_id: member_group_ids).each do |gz|
(rows_by_group[gz.group_id] ||= [] of GroupZone) << gz
end
return nil if rows_by_group.empty?
return if rows_by_group.empty?

satisfies = ->(perms : Permissions) do
perms.manage? || (perms & required) != Permissions::None
Expand Down Expand Up @@ -224,7 +224,7 @@ module PlaceOS::Model
end
end
end
return nil if anchor_hits.empty? && walk_seeds.empty?
return if anchor_hits.empty? && walk_seeds.empty?

escape = ->(value : String) { "'#{value.gsub("'", "''")}'" }
anchors_sql = anchor_hits.join(", ") { |zid| "(#{escape.call(zid)})" }
Expand Down
12 changes: 5 additions & 7 deletions src/placeos-models/shortener.cr
Original file line number Diff line number Diff line change
Expand Up @@ -71,13 +71,11 @@ module PlaceOS::Model
def save!(**options)
attempts = 0
loop do
begin
return super(**options)
rescue error : ::PgORM::Error::RecordNotSaved
attempts += 1
id_collision = new_record? && error.message.try(&.includes?("shortener_pkey"))
raise error unless id_collision && attempts < CREATE_ID_ATTEMPTS
end
return super(**options)
rescue error : ::PgORM::Error::RecordNotSaved
attempts += 1
id_collision = new_record? && error.message.try(&.includes?("shortener_pkey"))
raise error unless id_collision && attempts < CREATE_ID_ATTEMPTS
end
end

Expand Down
8 changes: 3 additions & 5 deletions src/placeos-models/utilities/settings_helper.cr
Original file line number Diff line number Diff line change
Expand Up @@ -32,11 +32,9 @@ module PlaceOS::Model::Utilities
settings
.each_with_object({} of YAML::Any => YAML::Any) do |setting, acc|
# Parse and merge into accumulated settings hash
begin
acc.merge!(setting.any)
rescue error
Log.warn(exception: error) { "failed to merge all settings: #{setting.inspect}" }
end
acc.merge!(setting.any)
rescue error
Log.warn(exception: error) { "failed to merge all settings: #{setting.inspect}" }
end
end

Expand Down
32 changes: 32 additions & 0 deletions tasks/146-signage-heartbeat.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,32 @@
# Suppress signage heartbeat changefeed events

Parent issue: https://github.com/PlaceOS/local/issues/146
Agent/task: codex-146-models-20260916-root
Branch: ai/146-signage-heartbeat
Start: 2026-09-16

## Contract

ControlSystem declares `changefeed_ignore_updates :signage_last_seen, :playlist_item_id` using pg-orm 2.4.0 and EventBus 1.1.0. `update_last_seen_time` remains one SQL UPDATE. Telemetry persists (including item changes/clearing) without sys CDC events; updates changing configuration still notify, even if telemetry changes too. INSERT/DELETE and unrelated models retain their notifications. The policy applies to all writers once the changefeed is registered; all trigger-installing services must upgrade before enabling it.

## Checklist

- [x] Verify pg-orm v2.4.0 release, local instructions and uncontested claim.
- [x] Reproduce original bug with real PostgreSQL CDC: first heartbeat produces [insert, update] instead of [insert].
- [x] Require pg-orm >=2.4.0 and remove unversioned override; local ignored lockfile resolves pg-orm2.4.0/EventBus1.1.0.
- [x] Add model-level declaration and explain telemetry intent; single UPDATE unchanged, no migration/core filter.
- [x] Focused regression passes: 2 examples, 0 failures/errors; covers persistence, repeated/nil/empty heartbeat, no CDC rows/notifications, configuration/mixed updates, create/delete and Zone.
- [ ] One agent runs specs at a time; full suite, formatting, lint and GitHub CI must pass.
- [ ] Independent review, squash merge, close parent issue and release claim.

## Review

Independent review approved declaration, dependency requirement, test coverage, rollout docs and behavior-preserving lint cleanup. EventBus/pg-orm stages are already merged/released. No migration is required. Rollout must update all CDC installers before enabling the model policy; SQL-backed telemetry reads remain available. Suppression includes telemetry-only writes outside this method. Record test and CI results in PR and parent issue before completion.

### Verification adjustment

The three-file targeted compile exited 137 before specs ran on the shared 8GB Docker VM. Retry with `--threads 1 --no-debug`; do not stop unrelated user containers or prune Docker. If full local compilation still exceeds available memory, run the full suite on GitHub CI's isolated runner after the focused regression passes.

Focused green: `./test --threads 1 --no-debug spec/control_system_changefeed_spec.cr` passes 2 examples (161ms runtime). Full local suite intentionally deferred to CI after broader three-file compiles hit OOM137 twice. `./bin/ameba`: 180 files, zero failures; formatting/diff checks clean. No local suite remains running.

Full CI first run: heartbeat coverage passed both images. Stable: 683 examples, one Zone callback timeout. Unstable: same timeout plus random Storage bucket collision. Zone coverage now asserts SQL CDC rows and PostgreSQL notifications directly; Storage default fixtures use distinct buckets. Models style CI now uses checkout v7.0.1 and Ameba action v1.0.0, replacing the shared legacy action pinned to Crystal 1.14 while formatting used latest Crystal. Full matrix must pass before merge.
Loading