diff --git a/.gitignore b/.gitignore index 599be4e..0840e68 100644 --- a/.gitignore +++ b/.gitignore @@ -1,4 +1,5 @@ *.beam *.ez /build +**/build erl_crash.dump diff --git a/README.md b/README.md index 800a53e..3b67cff 100644 --- a/README.md +++ b/README.md @@ -1,6 +1,6 @@ # Factos -Prototype context-first event-sourcing helpers for Gleam and KurrentDB. +Prototype context-first event-sourcing helpers for Gleam. This package deliberately does not model Event Sourcing as aggregates. A command capability reads the facts relevant to one decision, folds a temporary decision @@ -79,7 +79,7 @@ pub fn registration_decider() { `Decider` is inspired by FModel: it is a small pure domain component made from three things only: initial state, a decision function, and an evolution function. -It does not know about KurrentDB, projections, subscriptions, HTTP, retries, or +It does not know about SQLite, KurrentDB, projections, subscriptions, HTTP, retries, or repositories. You can test a decider without any storage: @@ -130,39 +130,70 @@ factos.tag("account:abc123") This is intentionally opinionated. Tags duplicate selected payload information, but they make consistency and query needs visible at the event-store boundary. +## Backends + +Factos is split into a small core package and backend packages: + +1. `factos` contains the store-independent domain primitives: `Decider`, `View`, `Query`, `Tag`, `Recorded`, and `Context`. +2. `backends/factos_sqlight` provides module `factos/sqlight` for SQLite via the `sqlight` package. +3. `backends/factos_kurrentdb_erlang` provides module `factos/kurrentdb` for KurrentDB on Erlang. + +Backends own their storage codecs and storage errors. The core package does not depend on either SQLite or KurrentDB. + ## Codecs -The app owns encoding and decoding. The codec returns both the domain event and -the event-store metadata needed by the context API. +The app owns encoding and decoding. Backend codecs return both the domain event +and the event-store metadata needed by the context API. ```gleam -pub fn codec() -> factos.EventCodec(Event, DecodeError) { - factos.EventCodec(encode: encode, decode: decode_event) +import factos/sqlight as factos_sqlight +import sqlight + +pub fn codec() -> factos_sqlight.EventCodec(Event, DecodeError) { + factos_sqlight.EventCodec(encode: encode, decode: decode_event) } ``` -`encode` returns a `factos.Proposed` with the domain event, event type, tags, and -the KurrentDB append message. `decode` returns a `factos.Decoded` with the domain -event, event type, and tags read from the stored event. +For SQLite, `encode` returns a `factos_sqlight.Proposed` with an id, domain event, +event type, tags, and encoded bytes. `decode` returns a `factos.Decoded` with the +domain event, event type, and tags read from the stored event. The library does not force one JSON shape for tags. In a real app, store tags in custom metadata or payload fields consistently, and decode them at the boundary. +## SQLite Backend + +```gleam +import factos/sqlight as factos_sqlight + +use connection <- sqlight.with_connection("file:events.sqlite3") +let assert Ok(Nil) = factos_sqlight.migrate(connection) + +factos_sqlight.dispatch_context( + connection, + stream: "facts", + query: username_context("renata"), + decider: registration_decider(), + codec: codec(), + command: RegisterUser("renata"), +) +``` + +The SQLite backend stores an append-only `factos_events` table and uses `BEGIN IMMEDIATE` while dispatching commands. It can enforce `FailIfEventsMatch(query, after)` inside the same SQLite transaction. + ## Reading A Context ```gleam -factos.read_context( +factos_sqlight.read_context( connection, query: username_context("renata"), decider: registration_decider(), codec: codec(), - timeout: 5000, ) ``` -`read_context` reads from `$all`, applies a server-side event-type filter when -possible, decodes events, filters them by the full query, folds state, and -returns a `factos.Context`. +Backend `read_context` functions load events, decode them, filter them by the +full query, fold state with the decider, and return a `factos.Context`. The returned context includes: @@ -182,25 +213,23 @@ factos.FailIfEventsMatch(query, after: position) That means: append these new facts only if no facts matching the command context appeared after the position used for the decision. -This prototype models that condition, but the current KurrentDB-backed context -dispatch returns `UnsupportedAppendCondition` for it because regular KurrentDB -append checks stream revisions, not arbitrary event-type/tag queries. +The SQLite backend enforces this condition transactionally. The KurrentDB backend +models the same condition, but returns `UnsupportedAppendCondition` for it because +regular KurrentDB append checks stream revisions, not arbitrary event-type/tag +queries. ```gleam -factos.dispatch_context( +factos_sqlight.dispatch_context( connection, stream: "facts", query: username_context("renata"), decider: registration_decider(), codec: codec(), command: RegisterUser("renata"), - timeout: 5000, ) ``` -Use this shape for stores that support DCB-style atomic append conditions. With -the current KurrentDB operation set, it is still useful as the honest API shape, -but it cannot complete successfully for `FailIfEventsMatch` yet. +Use this shape for stores that support DCB-style atomic append conditions. ## Stream Consistency @@ -208,13 +237,12 @@ KurrentDB can safely protect a single stream with expected revision checks. This is still useful when a stream is the right consistency boundary. ```gleam -factos.dispatch_stream( +factos_sqlight.dispatch_stream( connection, stream: "user-renata", decider: registration_decider(), codec: codec(), command: RegisterUser("renata"), - timeout: 5000, ) ``` @@ -251,9 +279,9 @@ Factos intentionally stops at pure projection computation. It does not provide a materialized-view repository abstraction yet; persistence and delivery choices belong outside the domain component. -## KurrentDB Tradeoffs +## KurrentDB Backend Tradeoffs -KurrentDB support available through the current dependency: +KurrentDB support available through `factos_kurrentdb_erlang`: 1. Read a single stream. 2. Append to a stream with `NoStream`, `Revision(n)`, `StreamExists`, or `Any`. diff --git a/backends/factos_kurrentdb_erlang/certs/ca.crt b/backends/factos_kurrentdb_erlang/certs/ca.crt new file mode 100644 index 0000000..3969515 --- /dev/null +++ b/backends/factos_kurrentdb_erlang/certs/ca.crt @@ -0,0 +1,19 @@ +-----BEGIN CERTIFICATE----- +MIIDFzCCAf+gAwIBAgIUS6TmEGEyiTmKvTvYz7hXvTcgPiUwDQYJKoZIhvcNAQEL +BQAwGzEZMBcGA1UEAwwQa3VycmVudGRiLWRldi1jYTAeFw0yNjA2MjYwMjUxMDBa +Fw0zNjA2MjMwMjUxMDBaMBsxGTAXBgNVBAMMEGt1cnJlbnRkYi1kZXYtY2EwggEi +MA0GCSqGSIb3DQEBAQUAA4IBDwAwggEKAoIBAQCK4BycVSZAe/Osd3o9GvG4Ncxr +mB2RZlm85YyVB6e4Mdln2cFyfkp94JbLQDT3fY30u8Q4ep4xf3Nn3l6MuTubeOQV +hyAffTpvxrWE+5Oa4LJzovhyaWiM83dgVg0MmOhh+3LU7k672uTV+Uca0G14gjIf +JFUW+0Zl3alr+70OT1f00c/xtICY4tEQu9VTGenM+9eTuhrrUpjz2hLAsLTvEhP4 +kdpPl/moKp7FYOiu+96bdc8yGdm073pLeR80ZAc+Xhwq6RI2HDWCVukX7yXy1X1F +8RS4rjWqRQKaW3DSuwhYiWhcFRdYDv9VlLD18I4HccFvNbvQA0fd5bS0U9ahAgMB +AAGjUzBRMB0GA1UdDgQWBBQ9qKjS1qDSiWW8d4P04+60R+NIkDAfBgNVHSMEGDAW +gBQ9qKjS1qDSiWW8d4P04+60R+NIkDAPBgNVHRMBAf8EBTADAQH/MA0GCSqGSIb3 +DQEBCwUAA4IBAQAtsEdwhDm/4EXTbqar26Qr6Cuo3P8zMxf+mVoQplLAUQa/loJX +GxqbsgHhX4awaxLbSStSZrCtQNw6H0rTwiQRDi9UbA154+7DgoWWEr+fwt7eHzSb +c2ZOPI225u+HfJTC9imVW8F+npvNgFaoT36OvBU0ktvvl6dchULNW5nfGqswub/p +VkGdU4AFGm/yW0hLOh7yqIeLrtlN/iBLh0fKe9LG2VZ6YJYFj0fWD30KU5ZvwU7E +c1zSxfUX93+3/+Ft/sRqFtphG5hXugApwsxqEtM6HYD7Byotee9wjbukkyiNjZDg +Logwi8U+XYfio89HE508iM/DIzNVa/l6PSOX +-----END CERTIFICATE----- diff --git a/backends/factos_kurrentdb_erlang/certs/ca.key b/backends/factos_kurrentdb_erlang/certs/ca.key new file mode 100644 index 0000000..84fa594 --- /dev/null +++ b/backends/factos_kurrentdb_erlang/certs/ca.key @@ -0,0 +1,28 @@ +-----BEGIN PRIVATE KEY----- +MIIEvQIBADANBgkqhkiG9w0BAQEFAASCBKcwggSjAgEAAoIBAQCK4BycVSZAe/Os +d3o9GvG4NcxrmB2RZlm85YyVB6e4Mdln2cFyfkp94JbLQDT3fY30u8Q4ep4xf3Nn +3l6MuTubeOQVhyAffTpvxrWE+5Oa4LJzovhyaWiM83dgVg0MmOhh+3LU7k672uTV ++Uca0G14gjIfJFUW+0Zl3alr+70OT1f00c/xtICY4tEQu9VTGenM+9eTuhrrUpjz +2hLAsLTvEhP4kdpPl/moKp7FYOiu+96bdc8yGdm073pLeR80ZAc+Xhwq6RI2HDWC +VukX7yXy1X1F8RS4rjWqRQKaW3DSuwhYiWhcFRdYDv9VlLD18I4HccFvNbvQA0fd +5bS0U9ahAgMBAAECggEAArq9gzbS9vPctd4+CD0LJMrFZRS2+Y5qe3mDQCNXsPmF +V2q+qCV6aTOQoSdmlxnoECgf1tiVmv1RV0h2ASPrm45WZMQsbeQCIdPk6cuAQtwx +U6+ffI+s7O7E0Q9V57JKaHEWyF+z61Ilut0gsDKqISMFcUpfdAF9qGdBMwCuTj16 +cJEaYacpFdDk/u3h9Csuu0YZiXxpHx2VmSAFchQOY9pYMNifOGYfr6ZqU73EoQ0e +/bXOqENV3x9OSZjC3ccinbtsuU6K9Z6xJOVZhOz9VRE9/o2W8R+UasVmnNT/4ZbB +r09OYU1wizbwfpFLDWs3E2qFYQVhF1AZtdMtPgEqkQKBgQC9yv4qVXUfUse1n0oY +CX8zBov8O4lXknq64ALwmEoNjynNwEm2DXpulP9/07NrV9XhPRledNasUp2mKL4c ++WDG5qjdoM2zE3Fx1QVMMq+p2QZaFmflps+qJ/BCOmEOF7S0cvOH+GhjwEsXrU1Q +c8tsNX2uhhZlsIKkvcikw6sVUQKBgQC7Ug1kviNiOPro/wOuWKtw6VNvp9lkU+up +nbXnUizMrs0Mrn/m8LjGAum8vmgXFUnH46MPQHZj6VvL4dRebfEMuCRz3qiUNMKi +3bpFR6Bl0Tn297h8pqpNratkOKdBD/yr2sFm4gorJX0I1Z+JAsgB25LNc8W9qFCf +XTOLc2qYUQKBgQCDrTuL6YB5+/fdFafVZ3ld0HP8yt2t6U3HK7Y+cJooMCSDwJ4j +ddR0tmFRsXIwzl7wh3B7bTqnkiYYavoDpi0zskKEiZVNYfb6UB390Mi5YX4bsKHi +3koDtvPlLxW5Lk9MRtiZhIoAcyBmS/FxGPWQnMgW9qbBZKYvYBC955diEQKBgArA +sQghagKPZsfNK7bsXBsFKcb1CaOataJs7S40J2IwfpDFy43EL7ceH7C39V2t2Shi +Rs/vUVx23tAbTIeHJBko0N7d3ytyw+F5fOHRNMHjesJUggCVyJzg5T/BiMhRVJ3A +1u1C+HZ1lnHVYW0J/dUtd4XXqXgzmz0qqnTM0UehAoGAQ2J8WhZz2lffopenz1sy +WqRFyMPgO8Xij/C2liErmVO6IycWUxTVDKdj78g61+c1JLdvLspJV49Dr7gWUWOd +Br4hj9nMVJavmrQJdFh5mwL6gpeUtrIH78XHozkYPNoaM/hdvJTiVLP7cvctAFhZ +yGx2LWMKcfpUQzyBGdmw0Ss= +-----END PRIVATE KEY----- diff --git a/backends/factos_kurrentdb_erlang/certs/ca.srl b/backends/factos_kurrentdb_erlang/certs/ca.srl new file mode 100644 index 0000000..d946817 --- /dev/null +++ b/backends/factos_kurrentdb_erlang/certs/ca.srl @@ -0,0 +1 @@ +4CDBD69624CF67183CB3BBA528EC27499110E3F5 diff --git a/backends/factos_kurrentdb_erlang/certs/node.conf b/backends/factos_kurrentdb_erlang/certs/node.conf new file mode 100644 index 0000000..b60c6b9 --- /dev/null +++ b/backends/factos_kurrentdb_erlang/certs/node.conf @@ -0,0 +1,16 @@ +[req] +distinguished_name = req_distinguished_name +req_extensions = v3_req +prompt = no + +[req_distinguished_name] +CN = localhost + +[v3_req] +keyUsage = keyEncipherment, dataEncipherment, digitalSignature +extendedKeyUsage = serverAuth +subjectAltName = @alt_names + +[alt_names] +DNS.1 = localhost +IP.1 = 127.0.0.1 diff --git a/backends/factos_kurrentdb_erlang/certs/node.crt b/backends/factos_kurrentdb_erlang/certs/node.crt new file mode 100644 index 0000000..3efe93a --- /dev/null +++ b/backends/factos_kurrentdb_erlang/certs/node.crt @@ -0,0 +1,20 @@ +-----BEGIN CERTIFICATE----- +MIIDPjCCAiagAwIBAgIUTNvWliTPZxg8s7ulKOwnSZEQ4/UwDQYJKoZIhvcNAQEL +BQAwGzEZMBcGA1UEAwwQa3VycmVudGRiLWRldi1jYTAeFw0yNjA2MjYwMjUxMDBa +Fw0zNjA2MjMwMjUxMDBaMBQxEjAQBgNVBAMMCWxvY2FsaG9zdDCCASIwDQYJKoZI +hvcNAQEBBQADggEPADCCAQoCggEBAOYTrWSFzIDQBwFEpjPNLrUOPh6GK9C82tjC +WQJCsIaoX/Q9UqHyxy+hN2XCmU4CC7nmGh6AF5RWRrEzJ2LxFazFf2lXuHgsXagS +0QYrAYgOhhlrF6v0E0aIq97bgKDCKLDSIabMzNwDS9rLjXj6EZfUFqEUIUozRCHC +CZWfk5uYGPURdJwU7BH6v1ZlI6RXRDLUAm4OuZZfndNvKg1up1FDw3KSJrNpSGuD +ZlKSMTVMA2gp8RLzRinyttFSfcL8EuxfByNO/cArT8KWcf0XgJHPyy0DrYzsj38q +GySIENQrCRBwLDp9KLDN83Mj2gHFV3U74WpRxMyYP5EfWiRaDtUCAwEAAaOBgDB+ +MAsGA1UdDwQEAwIEsDATBgNVHSUEDDAKBggrBgEFBQcDATAaBgNVHREEEzARggls +b2NhbGhvc3SHBH8AAAEwHQYDVR0OBBYEFKu0c/YP5lN1OUHT8MOejJEoYyfgMB8G +A1UdIwQYMBaAFD2oqNLWoNKJZbx3g/Tj7rRH40iQMA0GCSqGSIb3DQEBCwUAA4IB +AQBEMF5AU7nTikUoq/AhnVqWRsZFu6kuTWaR0gRwK9VujHyU7SDVN9IpkQMtK0Xi +mJVbLZI62ynEJHHXPvuIQ6hVb3IJtWAgKZffhKDQpbaNwzwWzmvHTOPRpIgMbIdT +VA6irNunbCvmegLql0Sqkg37i/mQ3+gmkr02MEeqCfWUDfNodz6looxh1RCIfaIC +zBRtFZtjMUj0nRURqGljZEZUUMrJayDr6DWKqMfHKyoDlvVLbicJufmrJSvpVnJJ +KmxrnKRQet7O8pnkHyiuqrnSTIKXYr1UG2fr9eKgD8H1AIEuagx+Uc2M5P7IgtYX +HYBNv9cJRMh0c4j97LBJYtHo +-----END CERTIFICATE----- diff --git a/backends/factos_kurrentdb_erlang/certs/node.csr b/backends/factos_kurrentdb_erlang/certs/node.csr new file mode 100644 index 0000000..97b6f9b --- /dev/null +++ b/backends/factos_kurrentdb_erlang/certs/node.csr @@ -0,0 +1,17 @@ +-----BEGIN CERTIFICATE REQUEST----- +MIICqDCCAZACAQAwFDESMBAGA1UEAwwJbG9jYWxob3N0MIIBIjANBgkqhkiG9w0B +AQEFAAOCAQ8AMIIBCgKCAQEA5hOtZIXMgNAHAUSmM80utQ4+HoYr0Lza2MJZAkKw +hqhf9D1SofLHL6E3ZcKZTgILueYaHoAXlFZGsTMnYvEVrMV/aVe4eCxdqBLRBisB +iA6GGWsXq/QTRoir3tuAoMIosNIhpszM3ANL2suNePoRl9QWoRQhSjNEIcIJlZ+T +m5gY9RF0nBTsEfq/VmUjpFdEMtQCbg65ll+d028qDW6nUUPDcpIms2lIa4NmUpIx +NUwDaCnxEvNGKfK20VJ9wvwS7F8HI079wCtPwpZx/ReAkc/LLQOtjOyPfyobJIgQ +1CsJEHAsOn0osM3zcyPaAcVXdTvhalHEzJg/kR9aJFoO1QIDAQABoE8wTQYJKoZI +hvcNAQkOMUAwPjALBgNVHQ8EBAMCBLAwEwYDVR0lBAwwCgYIKwYBBQUHAwEwGgYD +VR0RBBMwEYIJbG9jYWxob3N0hwR/AAABMA0GCSqGSIb3DQEBCwUAA4IBAQCNx33M +Y8qlqrBAUyqnb0/9wNq9ixwx5m2mOokw6SqHU2YGr1QYQ/5VCU4TCahbslsGHQWz +Y94D5ZW+UCoexkcAb1P84T64QL8Zm8hBO0/rBpAWREaGJpGGsISWus+1RBa1B8OI +Phpm2JmVJ+LqIfHuQJccYYnrlkjqLuN5jmjs65s4WSyUDYoYpoBxD3AL9sIROANx +hJ/tmqyHBSPkN2YWpqk3EHwg2TSMZyPqjMrlUaDFKEYbWx18ksHQbe0bhSHRIZRZ +vtZCLNxek98FXuQaXm2auHORglBAq3O2TyOArhYQ7Y6cNC35HD2XtUFZyGz0f1zd +2KcU7s1pwouesu3Y +-----END CERTIFICATE REQUEST----- diff --git a/backends/factos_kurrentdb_erlang/certs/node.key b/backends/factos_kurrentdb_erlang/certs/node.key new file mode 100644 index 0000000..4f9924b --- /dev/null +++ b/backends/factos_kurrentdb_erlang/certs/node.key @@ -0,0 +1,28 @@ +-----BEGIN PRIVATE KEY----- +MIIEvgIBADANBgkqhkiG9w0BAQEFAASCBKgwggSkAgEAAoIBAQDmE61khcyA0AcB +RKYzzS61Dj4ehivQvNrYwlkCQrCGqF/0PVKh8scvoTdlwplOAgu55hoegBeUVkax +Mydi8RWsxX9pV7h4LF2oEtEGKwGIDoYZaxer9BNGiKve24Cgwiiw0iGmzMzcA0va +y414+hGX1BahFCFKM0QhwgmVn5ObmBj1EXScFOwR+r9WZSOkV0Qy1AJuDrmWX53T +byoNbqdRQ8NykiazaUhrg2ZSkjE1TANoKfES80Yp8rbRUn3C/BLsXwcjTv3AK0/C +lnH9F4CRz8stA62M7I9/KhskiBDUKwkQcCw6fSiwzfNzI9oBxVd1O+FqUcTMmD+R +H1okWg7VAgMBAAECggEAB9qrnGch5lLTrmQmzVVnjwg7sCSR6dgMm4I08ioPJyWn +0ulmAP/N828MOlXUkHBq8I9tnFVwmKCCXMm7gjnrLMD4OsMjGb0f/F0aFB0TOg8O +3l7Eydq07r87KMoy/6npJDIkMnLC2o7lP8SboYHd6GI13I1YnpUV8hYS6C/wpMrR +mdBU6XThExWOrFtcOrfbX7DImpjoAamrPDkGjaeuJczTm0vZTXl8w+ThAswS+zjo +fxIqetk47qMqjr7kkanOu24bNAxnAiQKu6Kz0623prQ6iXqjIph7CDbwLtuQZ6mP +s8V6MWXGJY5gZyoAZ0FJ5CsqhcdyTvvTbJf3nVr7+QKBgQD5YtfcKuLhYrxaVDwL +YCVacgooOnqV+jcqZ06qCKm5MHhQWCYrXuvalYOfHBt2ugkHz3No5arZTG2qbXBF +pnqhxmCLlUptPcMJXQHPCQAiginiUx70Dhk0IYRK3TQivbW1ybx/g1TIKso6ou30 +dTwfBQlFESt4cse7WPV/xu7aSQKBgQDsLbzm/cq9bCVy5HaNLCCCgbFq0jq7VMsP +6Z9lTtyyOm3cEz8k27y1OQwPrny0NhqJZs37xjgPDUgJdCqdbefHkflCz23NLc/f +rdHNOA4s8+KHItH5Z4bR3EOBLR/krcpx/XompqB9qhWxBxa8hsOeSFBITvgeBCaw +oD8cx/AwLQKBgQDJ1OU+mrbkEjS+Jk4yJq4UdRcjV7C+kLL07ocLtdcmucOlwrGh +iED5tue/bdAMVqPYXlzZGIcdNm3K8Kdct0+ofhTE4x5JKyMeANfl5zLkutOLCBqV +CpP7TOT0cfIv67mUVqDn0jJbjcX9jr9miTsPH9RQwYSdBsf/KBAIScglgQKBgCrg +mtzs0nPVQG89XvB+RGCtHwKfrB36ZOs8pL2Ftbd9uBguPlZ4tifIdZIbQXSOJf8v +9NFyyRaieKOOvXXbUCsBK1mfwvVvDcA0FFTHintKw6N5BNncm7NZ4799675edtR/ +CkAeHCD0Uf/To6MSbE0+H6UhARah9kw2q36UJdz5AoGBALMMNzhW5bUEJKtFYrOF +nF3I8RZA1+pO3nHLpPQML37E9NQfrpWi5a+OXzlL5YbRq55KY5HYBqwMmrJGUfT9 +rAVdQ0sbrgjdOl6Tm1Ey6fipPImfi0Y7d5IC6jSMmtoZvfml4cpo5xrNfRzl7IUm +jrPQ9bpMzWg8RRsCo7pKvVbE +-----END PRIVATE KEY----- diff --git a/backends/factos_kurrentdb_erlang/certs/trusted/ca.crt b/backends/factos_kurrentdb_erlang/certs/trusted/ca.crt new file mode 100644 index 0000000..3969515 --- /dev/null +++ b/backends/factos_kurrentdb_erlang/certs/trusted/ca.crt @@ -0,0 +1,19 @@ +-----BEGIN CERTIFICATE----- +MIIDFzCCAf+gAwIBAgIUS6TmEGEyiTmKvTvYz7hXvTcgPiUwDQYJKoZIhvcNAQEL +BQAwGzEZMBcGA1UEAwwQa3VycmVudGRiLWRldi1jYTAeFw0yNjA2MjYwMjUxMDBa +Fw0zNjA2MjMwMjUxMDBaMBsxGTAXBgNVBAMMEGt1cnJlbnRkYi1kZXYtY2EwggEi +MA0GCSqGSIb3DQEBAQUAA4IBDwAwggEKAoIBAQCK4BycVSZAe/Osd3o9GvG4Ncxr +mB2RZlm85YyVB6e4Mdln2cFyfkp94JbLQDT3fY30u8Q4ep4xf3Nn3l6MuTubeOQV +hyAffTpvxrWE+5Oa4LJzovhyaWiM83dgVg0MmOhh+3LU7k672uTV+Uca0G14gjIf +JFUW+0Zl3alr+70OT1f00c/xtICY4tEQu9VTGenM+9eTuhrrUpjz2hLAsLTvEhP4 +kdpPl/moKp7FYOiu+96bdc8yGdm073pLeR80ZAc+Xhwq6RI2HDWCVukX7yXy1X1F +8RS4rjWqRQKaW3DSuwhYiWhcFRdYDv9VlLD18I4HccFvNbvQA0fd5bS0U9ahAgMB +AAGjUzBRMB0GA1UdDgQWBBQ9qKjS1qDSiWW8d4P04+60R+NIkDAfBgNVHSMEGDAW +gBQ9qKjS1qDSiWW8d4P04+60R+NIkDAPBgNVHRMBAf8EBTADAQH/MA0GCSqGSIb3 +DQEBCwUAA4IBAQAtsEdwhDm/4EXTbqar26Qr6Cuo3P8zMxf+mVoQplLAUQa/loJX +GxqbsgHhX4awaxLbSStSZrCtQNw6H0rTwiQRDi9UbA154+7DgoWWEr+fwt7eHzSb +c2ZOPI225u+HfJTC9imVW8F+npvNgFaoT36OvBU0ktvvl6dchULNW5nfGqswub/p +VkGdU4AFGm/yW0hLOh7yqIeLrtlN/iBLh0fKe9LG2VZ6YJYFj0fWD30KU5ZvwU7E +c1zSxfUX93+3/+Ft/sRqFtphG5hXugApwsxqEtM6HYD7Byotee9wjbukkyiNjZDg +Logwi8U+XYfio89HE508iM/DIzNVa/l6PSOX +-----END CERTIFICATE----- diff --git a/backends/factos_kurrentdb_erlang/compose.yml b/backends/factos_kurrentdb_erlang/compose.yml new file mode 100644 index 0000000..1fca35d --- /dev/null +++ b/backends/factos_kurrentdb_erlang/compose.yml @@ -0,0 +1,28 @@ +services: + kurrentdb: + image: docker.kurrent.io/kurrent-latest/kurrentdb:latest + environment: + - KURRENTDB_CLUSTER_SIZE=1 + - KURRENTDB_RUN_PROJECTIONS=All + - KURRENTDB_START_STANDARD_PROJECTIONS=true + - KURRENTDB_NODE_PORT=2113 + - KURRENTDB_INSECURE=false + - KURRENTDB_ALLOW_ANONYMOUS_ENDPOINT_ACCESS=true + - KURRENTDB_ALLOW_ANONYMOUS_STREAM_ACCESS=true + - KURRENTDB_ENABLE_ATOM_PUB_OVER_HTTP=true + - KURRENTDB_CERTIFICATE_FILE=/certs/node.crt + - KURRENTDB_CERTIFICATE_PRIVATE_KEY_FILE=/certs/node.key + - KURRENTDB_TRUSTED_ROOT_CERTIFICATES_PATH=/certs/trusted + volumes: + - ./certs:/certs:ro + ports: + - "2113:2113" + healthcheck: + test: + [ + "CMD-SHELL", + "curl --fail --cacert /certs/ca.crt https://localhost:2113/health/live || exit 1", + ] + interval: 2s + timeout: 5s + retries: 30 diff --git a/backends/factos_kurrentdb_erlang/gleam.toml b/backends/factos_kurrentdb_erlang/gleam.toml new file mode 100644 index 0000000..b94ce90 --- /dev/null +++ b/backends/factos_kurrentdb_erlang/gleam.toml @@ -0,0 +1,14 @@ +name = "factos_kurrentdb_erlang" +version = "1.0.0" + +[dependencies] +factos = { path = "../.." } +gleam_stdlib = ">= 1.0.0 and < 2.0.0" +kurrentdb = { path = "../../../kurrentdb" } +kurrentdb_erlang = { path = "../../../kurrentdb/backends/kurrentdb_erlang" } +youid = ">= 1.6.0 and < 2.0.0" + +[dev_dependencies] +gleeunit = ">= 1.0.0 and < 2.0.0" +gleam_json = ">= 3.0.0 and < 4.0.0" +global_value = ">= 1.0.0 and < 2.0.0" diff --git a/backends/factos_kurrentdb_erlang/manifest.toml b/backends/factos_kurrentdb_erlang/manifest.toml new file mode 100644 index 0000000..92fa693 --- /dev/null +++ b/backends/factos_kurrentdb_erlang/manifest.toml @@ -0,0 +1,43 @@ +# Do not manually edit this file, it is managed by Gleam. +# +# This file locks the dependency versions used, to make your build +# deterministic and to prevent unexpected versions from being included +# in your application. +# +# You should check this file into your source control repository. + +packages = [ + { name = "certifi", version = "2.17.0", build_tools = ["rebar3"], requirements = [], otp_app = "certifi", source = "hex", outer_checksum = "8122798A17F0293C80DAADA25D0F81C7F4D708C73FEF782C7C9B1950E26E4D21" }, + { name = "factos", version = "1.0.0", build_tools = ["gleam"], requirements = ["gleam_stdlib"], source = "local", path = "../.." }, + { name = "gleam_crypto", version = "1.6.0", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "gleam_crypto", source = "hex", outer_checksum = "2DE9E4EF53CF6FEE049D4F765731F7178F7A11AEFAE00EEE63BF7536B354AD3F" }, + { name = "gleam_erlang", version = "1.3.0", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "gleam_erlang", source = "hex", outer_checksum = "1124AD3AA21143E5AF0FC5CF3D9529F6DB8CA03E43A55711B60B6B7B3874375C" }, + { name = "gleam_hackney", version = "1.4.0", build_tools = ["gleam"], requirements = ["gleam_http", "gleam_stdlib", "hackney"], source = "local", path = "../../../gleam_hackney" }, + { name = "gleam_http", version = "4.3.0", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "gleam_http", source = "hex", outer_checksum = "82EA6A717C842456188C190AFB372665EA56CE13D8559BF3B1DD9E40F619EE0C" }, + { name = "gleam_json", version = "3.1.0", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "gleam_json", source = "hex", outer_checksum = "44FDAA8847BE8FC48CA7A1C089706BD54BADCC4C45B237A992EDDF9F2CDB2836" }, + { name = "gleam_otp", version = "1.2.0", build_tools = ["gleam"], requirements = ["gleam_erlang", "gleam_stdlib"], otp_app = "gleam_otp", source = "hex", outer_checksum = "BA6A294E295E428EC1562DC1C11EA7530DCB981E8359134BEABC8493B7B2258E" }, + { name = "gleam_stdlib", version = "1.0.3", build_tools = ["gleam"], requirements = [], otp_app = "gleam_stdlib", source = "hex", outer_checksum = "1F543AFBA5D33DA493E6087F4E4C4F20D899411343512686C98A8ABB2963CF22" }, + { name = "gleam_time", version = "1.8.0", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "gleam_time", source = "hex", outer_checksum = "533D8723774D61AD4998324F5DD1DABDCDBFABAFB9E87CB5D03C6955448FC97D" }, + { name = "gleeunit", version = "1.11.0", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "gleeunit", source = "hex", outer_checksum = "EC31ABA74256AEA531EDF8169931D775BBB384FED0A8A1BDC4DD9354E3E21826" }, + { name = "global_value", version = "1.0.0", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "global_value", source = "hex", outer_checksum = "23F74C91A7B819C43ABCCBF49DAD5BB8799D81F2A3736BA9A534BD47F309FF4F" }, + { name = "h2", version = "0.10.2", build_tools = ["rebar3"], requirements = [], otp_app = "h2", source = "hex", outer_checksum = "497A899F338B42E6A0B292524E635B0CE6F9379FA39395C8E38D06351CD9B9CF" }, + { name = "hackney", version = "4.4.5", build_tools = ["rebar3"], requirements = ["certifi", "h2", "idna", "mimerl", "parse_trans", "quic", "ssl_verify_fun", "webtransport"], otp_app = "hackney", source = "hex", outer_checksum = "6D72BEF4E135E94C522C271E11FBB6933EFB0006EF235A3933807D0BE73B71EC" }, + { name = "idna", version = "7.1.0", build_tools = ["rebar3"], requirements = [], otp_app = "idna", source = "hex", outer_checksum = "6AE959A025BF36DF61A8CAB8508D9654891B5426A84C44D82DEAFFD6DDF8C71F" }, + { name = "kurrentdb", version = "1.0.0", build_tools = ["gleam"], requirements = ["gleam_http", "gleam_json", "gleam_stdlib", "youid"], source = "local", path = "../../../kurrentdb" }, + { name = "kurrentdb_erlang", version = "1.0.0", build_tools = ["gleam"], requirements = ["gleam_erlang", "gleam_hackney", "gleam_http", "gleam_otp", "gleam_stdlib", "kurrentdb", "youid"], source = "local", path = "../../../kurrentdb/backends/kurrentdb_erlang" }, + { name = "mimerl", version = "1.5.0", build_tools = ["rebar3"], requirements = [], otp_app = "mimerl", source = "hex", outer_checksum = "DB648CE065BAE14EA84CA8B5DD123F42F49417CEF693541110BF6F9E9BE9ECC4" }, + { name = "parse_trans", version = "3.4.2", build_tools = ["rebar3"], requirements = [], otp_app = "parse_trans", source = "hex", outer_checksum = "4C25347DE3B7C35732D32E69AB43D1CEEE0BEAE3F3B3ADE1B59CBD3DD224D9CA" }, + { name = "quic", version = "1.6.5", build_tools = ["rebar3"], requirements = [], otp_app = "quic", source = "hex", outer_checksum = "DE1A88972C33201A50D1A17C8C4A14528BD1D8F25EF705897D680AE312D0AA78" }, + { name = "ssl_verify_fun", version = "1.1.7", build_tools = ["mix", "rebar3", "make"], requirements = [], otp_app = "ssl_verify_fun", source = "hex", outer_checksum = "FE4C190E8F37401D30167C8C405EDA19469F34577987C76DDE613E838BBC67F8" }, + { name = "webtransport", version = "0.4.1", build_tools = ["rebar3"], requirements = ["h2", "quic"], otp_app = "webtransport", source = "hex", outer_checksum = "006E4E52A8F03B69201D4637C85B424A4DDCCC1909F32C5742FED8D495A75174" }, + { name = "youid", version = "1.6.0", build_tools = ["gleam"], requirements = ["gleam_crypto", "gleam_stdlib", "gleam_time"], otp_app = "youid", source = "hex", outer_checksum = "7A3ABA44B1B38BC2BDCB5474C5317AA372BE58DFBC649815EE08B03526DDA18D" }, +] + +[requirements] +factos = { path = "../.." } +gleam_json = { version = ">= 3.0.0 and < 4.0.0" } +gleam_stdlib = { version = ">= 1.0.0 and < 2.0.0" } +gleeunit = { version = ">= 1.0.0 and < 2.0.0" } +global_value = { version = ">= 1.0.0 and < 2.0.0" } +kurrentdb = { path = "../../../kurrentdb" } +kurrentdb_erlang = { path = "../../../kurrentdb/backends/kurrentdb_erlang" } +youid = { version = ">= 1.6.0 and < 2.0.0" } diff --git a/backends/factos_kurrentdb_erlang/scripts/generate-dev-certs.sh b/backends/factos_kurrentdb_erlang/scripts/generate-dev-certs.sh new file mode 100644 index 0000000..515028f --- /dev/null +++ b/backends/factos_kurrentdb_erlang/scripts/generate-dev-certs.sh @@ -0,0 +1,53 @@ +#!/usr/bin/env sh +set -eu + +script_dir=$(CDPATH= cd -- "$(dirname -- "$0")" && pwd) +backend_dir=$(CDPATH= cd -- "$script_dir/.." && pwd) +cert_dir="$backend_dir/certs" +trusted_dir="$cert_dir/trusted" + +mkdir -p "$cert_dir" +mkdir -p "$trusted_dir" + +openssl req -x509 -newkey rsa:2048 -days 3650 -nodes \ + -keyout "$cert_dir/ca.key" \ + -out "$cert_dir/ca.crt" \ + -subj "/CN=kurrentdb-dev-ca" + +cat >"$cert_dir/node.conf" < Proposed(event), + decode: fn(read_stream.RecordedEvent) -> + Result(factos.Decoded(event), decode_error), + ) +} + +pub type Error(domain_error, decode_error) { + DomainError(domain_error) + DecodeError(decode_error) + ReadError(kurrentdb_erlang.Error(read_stream.ResponseError)) + AppendError(kurrentdb_erlang.Error(append_to_stream.ResponseError)) + ReadTimedOut + UnsupportedAppendCondition(factos.AppendCondition) +} + +pub fn read_context( + connection: kurrentdb_erlang.Connection, + query query: factos.Query, + decider decider: factos.Decider(command, state, event, domain_error), + codec codec: EventCodec(event, decode_error), + timeout timeout: Int, +) -> Result(factos.Context(event, state), Error(domain_error, decode_error)) { + let factos.Decider(initial, _, evolve) = decider + + read_context_events(connection, query, initial, evolve, codec, timeout) +} + +pub fn dispatch_context( + connection: kurrentdb_erlang.Connection, + stream stream_name: String, + query query: factos.Query, + decider decider: factos.Decider(command, state, event, domain_error), + codec codec: EventCodec(event, decode_error), + command command: command, + timeout timeout: Int, +) -> Result(append_to_stream.Append, Error(domain_error, decode_error)) { + use context <- result.try(read_context( + connection, + query: query, + decider: decider, + codec: codec, + timeout: timeout, + )) + use pair <- result.try( + factos.decide_context(context, command, decider) + |> result.map_error(DomainError), + ) + let #(context, events) = pair + + append_with_condition( + connection, + stream_name, + events, + codec, + context.append_condition, + timeout, + ) +} + +pub fn load_stream( + connection: kurrentdb_erlang.Connection, + stream stream_name: String, + decider decider: factos.Decider(command, state, event, domain_error), + codec codec: EventCodec(event, decode_error), + timeout timeout: Int, +) -> Result( + factos.LoadedStream(event, state), + Error(domain_error, decode_error), +) { + let factos.Decider(initial, _, evolve) = decider + + load_stream_events(connection, stream_name, initial, evolve, codec, timeout) +} + +pub fn dispatch_stream( + connection: kurrentdb_erlang.Connection, + stream stream_name: String, + decider decider: factos.Decider(command, state, event, domain_error), + codec codec: EventCodec(event, decode_error), + command command: command, + timeout timeout: Int, +) -> Result(append_to_stream.Append, Error(domain_error, decode_error)) { + let factos.Decider(initial, decide, evolve) = decider + + use loaded <- result.try(load_stream_events( + connection, + stream_name, + initial, + evolve, + codec, + timeout, + )) + + use events <- result.try( + decide(loaded.state, command) + |> result.map_error(DomainError), + ) + + append_stream_events( + connection, + stream_name, + events, + codec, + loaded.revision, + timeout, + ) +} + +fn append_with_condition( + connection: kurrentdb_erlang.Connection, + stream_name: String, + events: List(event), + codec: EventCodec(event, decode_error), + condition: factos.AppendCondition, + timeout: Int, +) -> Result(append_to_stream.Append, Error(domain_error, decode_error)) { + case condition { + factos.NoAppendCondition -> + append_to_stream_with_config( + connection, + stream_name, + events, + codec, + append_to_stream.configure() |> append_to_stream.any, + timeout, + ) + factos.FailIfEventsMatch(_, _) -> + Error(UnsupportedAppendCondition(condition)) + } +} + +fn read_context_events( + connection: kurrentdb_erlang.Connection, + query: factos.Query, + initial: state, + evolve: fn(state, event) -> state, + codec: EventCodec(event, decode_error), + timeout: Int, +) -> Result(factos.Context(event, state), Error(domain_error, decode_error)) { + let stream = + kurrentdb_erlang.read_all( + connection, + config: read_all.configure() + |> read_all.filter(query_to_read_all_filter(query)), + ) + + receive_context( + stream, + query, + initial, + evolve, + codec, + [], + factos.NoPosition, + timeout, + ) +} + +fn load_stream_events( + connection: kurrentdb_erlang.Connection, + stream_name: String, + initial: state, + evolve: fn(state, event) -> state, + codec: EventCodec(event, decode_error), + timeout: Int, +) -> Result( + factos.LoadedStream(event, state), + Error(domain_error, decode_error), +) { + let stream = + kurrentdb_erlang.read_stream( + connection, + stream: stream_name, + config: read_stream.configure(), + ) + + receive_stream( + stream, + stream_name, + initial, + evolve, + codec, + [], + factos.NoEvents, + timeout, + ) +} + +fn append_stream_events( + connection: kurrentdb_erlang.Connection, + stream_name: String, + events: List(event), + codec: EventCodec(event, decode_error), + expected: factos.Revision, + timeout: Int, +) -> Result(append_to_stream.Append, Error(domain_error, decode_error)) { + append_to_stream_with_config( + connection, + stream_name, + events, + codec, + append_config(expected), + timeout, + ) +} + +fn append_to_stream_with_config( + connection: kurrentdb_erlang.Connection, + stream_name: String, + events: List(event), + codec: EventCodec(event, decode_error), + config: append_to_stream.Configuration, + timeout: Int, +) -> Result(append_to_stream.Append, Error(domain_error, decode_error)) { + case events { + [] -> + Ok(append_to_stream.Append( + current_revision: -1, + position: append_to_stream.NoPositionReturned, + )) + [_, ..] -> { + let EventCodec(encode, _) = codec + let task = + kurrentdb_erlang.append_to_stream( + connection, + stream: stream_name, + events: list.map(events, fn(event) { + let Proposed(message: message, ..) = encode(event) + message + }), + config: config, + ) + + kurrentdb_erlang.await(task, within: timeout) + |> result.map_error(AppendError) + } + } +} + +fn receive_context( + stream: kurrentdb_erlang.Stream, + query: factos.Query, + state: state, + evolve: fn(state, event) -> state, + codec: EventCodec(event, decode_error), + events: List(factos.Recorded(event)), + position: factos.SequencePosition, + timeout: Int, +) -> Result(factos.Context(event, state), Error(domain_error, decode_error)) { + case kurrentdb_erlang.receive(stream, within: timeout) { + Error(Nil) -> { + kurrentdb_erlang.close(stream) + Error(ReadTimedOut) + } + Ok(kurrentdb_erlang.StreamFinished) -> { + kurrentdb_erlang.close(stream) + Ok(factos.Context( + query:, + state:, + events: list.reverse(events), + position:, + append_condition: factos.FailIfEventsMatch(query, position), + )) + } + Ok(kurrentdb_erlang.StreamFailed(error)) -> { + kurrentdb_erlang.close(stream) + Error(ReadError(error)) + } + Ok(kurrentdb_erlang.ReadEvent(event)) -> + receive_context_event( + stream, + query, + state, + evolve, + codec, + events, + position, + timeout, + event, + ) + Ok(kurrentdb_erlang.ReadMessage(read_stream.ReadEvent(event))) -> + receive_context_event( + stream, + query, + state, + evolve, + codec, + events, + position, + timeout, + event, + ) + Ok(kurrentdb_erlang.ReadMessage(read_stream.LastAllStreamPosition(read_stream.Position( + commit_position: commit_position, + prepare_position: _, + )))) -> + receive_context( + stream, + query, + state, + evolve, + codec, + events, + factos.highest_position( + position, + factos.SequencePosition(commit_position), + ), + timeout, + ) + Ok(kurrentdb_erlang.ReadMessage(_)) -> + receive_context( + stream, + query, + state, + evolve, + codec, + events, + position, + timeout, + ) + } +} + +fn receive_context_event( + stream: kurrentdb_erlang.Stream, + query: factos.Query, + state: state, + evolve: fn(state, event) -> state, + codec: EventCodec(event, decode_error), + events: List(factos.Recorded(event)), + position: factos.SequencePosition, + timeout: Int, + read_event: read_stream.ReadEvent, +) -> Result(factos.Context(event, state), Error(domain_error, decode_error)) { + use recorded <- result.try(decode_recorded(read_event, codec)) + let next_position = factos.highest_position(position, recorded.position) + + case factos.matches_query(recorded, query) { + True -> + receive_context( + stream, + query, + evolve(state, recorded.event), + evolve, + codec, + [recorded, ..events], + next_position, + timeout, + ) + False -> + receive_context( + stream, + query, + state, + evolve, + codec, + events, + next_position, + timeout, + ) + } +} + +fn receive_stream( + stream: kurrentdb_erlang.Stream, + stream_name: String, + state: state, + evolve: fn(state, event) -> state, + codec: EventCodec(event, decode_error), + events: List(factos.Recorded(event)), + revision: factos.Revision, + timeout: Int, +) -> Result( + factos.LoadedStream(event, state), + Error(domain_error, decode_error), +) { + case kurrentdb_erlang.receive(stream, within: timeout) { + Error(Nil) -> { + kurrentdb_erlang.close(stream) + Error(ReadTimedOut) + } + Ok(kurrentdb_erlang.StreamFinished) -> { + kurrentdb_erlang.close(stream) + Ok(factos.LoadedStream( + stream: stream_name, + state:, + events: list.reverse(events), + revision:, + )) + } + Ok(kurrentdb_erlang.StreamFailed(error)) -> { + kurrentdb_erlang.close(stream) + case error { + kurrentdb_erlang.OperationError(read_stream.ReadStreamNotFound(_)) -> + Ok(factos.LoadedStream( + stream: stream_name, + state:, + events: [], + revision: factos.NoEvents, + )) + _ -> Error(ReadError(error)) + } + } + Ok(kurrentdb_erlang.ReadEvent(event)) -> + receive_stream_event( + stream, + stream_name, + state, + evolve, + codec, + events, + timeout, + event, + ) + Ok(kurrentdb_erlang.ReadMessage(read_stream.ReadEvent(event))) -> + receive_stream_event( + stream, + stream_name, + state, + evolve, + codec, + events, + timeout, + event, + ) + Ok(kurrentdb_erlang.ReadMessage(_)) -> + receive_stream( + stream, + stream_name, + state, + evolve, + codec, + events, + revision, + timeout, + ) + } +} + +fn receive_stream_event( + stream: kurrentdb_erlang.Stream, + stream_name: String, + state: state, + evolve: fn(state, event) -> state, + codec: EventCodec(event, decode_error), + events: List(factos.Recorded(event)), + timeout: Int, + read_event: read_stream.ReadEvent, +) -> Result( + factos.LoadedStream(event, state), + Error(domain_error, decode_error), +) { + use recorded <- result.try(decode_recorded(read_event, codec)) + + receive_stream( + stream, + stream_name, + evolve(state, recorded.event), + evolve, + codec, + [recorded, ..events], + factos.CurrentRevision(recorded.revision), + timeout, + ) +} + +fn decode_recorded( + read_event: read_stream.ReadEvent, + codec: EventCodec(event, decode_error), +) -> Result(factos.Recorded(event), Error(domain_error, decode_error)) { + let recorded_event = case read_event { + read_stream.Recorded(event) -> event + read_stream.Resolved(event: event, ..) -> event + } + + let EventCodec(_, decode) = codec + use decoded <- result.try( + decode(recorded_event) + |> result.map_error(DecodeError), + ) + + let factos.Decoded(event, type_, tags) = decoded + Ok(factos.Recorded( + id: uuid.to_string(recorded_event.id), + stream: recorded_event.stream, + revision: recorded_event.revision, + position: factos.SequencePosition(recorded_event.commit_position), + type_: type_, + tags: tags, + event: event, + )) +} + +fn query_to_read_all_filter(query: factos.Query) -> read_all.Filter { + let types = query_event_type_names(query) + + case types { + [] -> read_all.NoFilter + [_, ..] -> read_all.EventTypePrefix(types, window: read_all.FilterMax(1000)) + } +} + +fn query_event_type_names(query: factos.Query) -> List(String) { + case query { + factos.AllEvents -> [] + factos.Query(items) -> collect_type_names(items, []) + } +} + +fn collect_type_names( + items: List(factos.QueryItem), + names: List(String), +) -> List(String) { + case items { + [] -> list.reverse(names) + [factos.QueryItem(types, _), ..rest] -> + collect_type_names(rest, prepend_type_names(types, names)) + } +} + +fn prepend_type_names( + types: List(factos.EventType), + names: List(String), +) -> List(String) { + case types { + [] -> names + [type_, ..rest] -> + prepend_type_names(rest, [factos.event_type_name(type_), ..names]) + } +} + +fn append_config(revision: factos.Revision) -> append_to_stream.Configuration { + case revision { + factos.NoEvents -> + append_to_stream.configure() |> append_to_stream.no_stream + factos.CurrentRevision(revision) -> + append_to_stream.configure() + |> append_to_stream.expected_revision(revision) + } +} diff --git a/backends/factos_kurrentdb_erlang/test/factos_kurrentdb_erlang_test.gleam b/backends/factos_kurrentdb_erlang/test/factos_kurrentdb_erlang_test.gleam new file mode 100644 index 0000000..e8c339f --- /dev/null +++ b/backends/factos_kurrentdb_erlang/test/factos_kurrentdb_erlang_test.gleam @@ -0,0 +1,240 @@ +import factos +import factos/kurrentdb as factos_kurrentdb +import gleam/bit_array +import gleam/int +import gleam/list +import gleam/option +import gleam/result +import gleeunit +import global_value +import kurrentdb +import kurrentdb/operation/append_to_stream +import kurrentdb/operation/read_stream +import kurrentdb_erlang +import youid/uuid + +pub fn main() -> Nil { + gleeunit.main() +} + +const connection_string = "kurrentdb://admin:changeit@localhost:2113?tls=true" + +const timeout = 10_000 + +type CounterCommand { + Increment +} + +type CounterEvent { + Incremented(value: Int) +} + +type CounterState { + CounterState(total: Int) +} + +type DecodeError { + UnknownEvent + InvalidData +} + +pub fn backend_package_compiles_test() { + let module_name = "factos/kurrentdb" + assert module_name == "factos/kurrentdb" +} + +pub fn dispatch_stream_handles_many_events_integration_test() { + let stream_name = unique_name("counter-stream") + let event_type = unique_name("FactosCounterIncremented") + + let assert Ok(append_to_stream.Append(current_revision: 99, position: _)) = + dispatch_counter_stream_many(stream_name, event_type, 100) + + let assert Ok(loaded) = + factos_kurrentdb.load_stream( + connection(), + stream: stream_name, + decider: counter_decider(), + codec: counter_codec(event_type), + timeout: timeout, + ) + + assert loaded.state == CounterState(100) + assert loaded.revision == factos.CurrentRevision(99) + assert list.length(loaded.events) == 100 +} + +pub fn read_context_handles_many_streams_integration_test() { + let event_type = unique_name("FactosCounterContextIncremented") + let query = + factos.query([ + factos.query_item(types: [factos.event_type(event_type)], tags: [ + factos.tag("counter:load"), + ]), + ]) + + let assert Ok(append_to_stream.Append(current_revision: 0, position: _)) = + dispatch_counter_context_streams_many(event_type, 50) + + let assert Ok(context) = + factos_kurrentdb.read_context( + connection(), + query: query, + decider: counter_decider(), + codec: counter_codec(event_type), + timeout: timeout, + ) + + assert context.state == CounterState(50) + assert list.length(context.events) == 50 + assert context.position != factos.NoPosition +} + +fn connection() -> kurrentdb_erlang.Connection { + global_value.create_with_unique_name( + "factos_kurrentdb_erlang.global.data", + fn() { + let assert Ok(client) = + kurrentdb.from_connection_string(connection_string) + + let assert Ok(connection) = + kurrentdb_erlang.new(client) + |> kurrentdb_erlang.verify_ca_certificate_file("certs/ca.crt") + |> kurrentdb_erlang.start(option.None) + + connection + }, + ) +} + +fn dispatch_counter_stream_many( + stream_name: String, + event_type: String, + remaining: Int, +) -> Result(append_to_stream.Append, factos_kurrentdb.Error(Nil, DecodeError)) { + let result = + factos_kurrentdb.dispatch_stream( + connection(), + stream: stream_name, + decider: counter_decider(), + codec: counter_codec(event_type), + command: Increment, + timeout: timeout, + ) + + case remaining, result { + 1, _ -> result + _, Ok(_) -> + dispatch_counter_stream_many(stream_name, event_type, remaining - 1) + _, Error(error) -> Error(error) + } +} + +fn dispatch_counter_context_streams_many( + event_type: String, + remaining: Int, +) -> Result(append_to_stream.Append, factos_kurrentdb.Error(Nil, DecodeError)) { + let result = + factos_kurrentdb.dispatch_stream( + connection(), + stream: unique_name("counter-context"), + decider: counter_decider(), + codec: counter_codec(event_type), + command: Increment, + timeout: timeout, + ) + + case remaining, result { + 1, _ -> result + _, Ok(_) -> dispatch_counter_context_streams_many(event_type, remaining - 1) + _, Error(error) -> Error(error) + } +} + +fn counter_decider() -> factos.Decider( + CounterCommand, + CounterState, + CounterEvent, + Nil, +) { + factos.decider( + initial: CounterState(0), + decide: counter_decide, + evolve: counter_evolve, + ) +} + +fn counter_decide( + state: CounterState, + command: CounterCommand, +) -> Result(List(CounterEvent), Nil) { + let CounterState(total) = state + case command { + Increment -> Ok([Incremented(total + 1)]) + } +} + +fn counter_evolve(state: CounterState, event: CounterEvent) -> CounterState { + let CounterState(total) = state + case event { + Incremented(_) -> CounterState(total + 1) + } +} + +fn counter_codec( + event_type: String, +) -> factos_kurrentdb.EventCodec(CounterEvent, DecodeError) { + factos_kurrentdb.EventCodec( + encode: encode_counter_event(_, event_type), + decode: decode_counter_event(_, event_type), + ) +} + +fn encode_counter_event( + event: CounterEvent, + event_type: String, +) -> factos_kurrentdb.Proposed(CounterEvent) { + case event { + Incremented(value) -> + factos_kurrentdb.Proposed( + event: event, + type_: factos.event_type(event_type), + tags: [factos.tag("counter:load")], + message: append_to_stream.binary_event( + uuid: uuid.v4(), + type_: event_type, + data: bit_array.from_string(int.to_string(value)), + ), + ) + } +} + +fn decode_counter_event( + stored: read_stream.RecordedEvent, + event_type: String, +) -> Result(factos.Decoded(CounterEvent), DecodeError) { + case list.key_find(stored.metadata, "type") { + Ok(type_name) if type_name == event_type -> { + use text <- result.try( + bit_array.to_string(stored.data) + |> result.map_error(fn(_) { InvalidData }), + ) + use value <- result.try( + int.parse(text) + |> result.map_error(fn(_) { InvalidData }), + ) + Ok( + factos.Decoded( + event: Incremented(value), + type_: factos.event_type(event_type), + tags: [factos.tag("counter:load")], + ), + ) + } + Ok(_) | Error(_) -> Error(UnknownEvent) + } +} + +fn unique_name(prefix: String) -> String { + prefix <> "-" <> uuid.to_string(uuid.v4()) +} diff --git a/backends/factos_sqlight/gleam.toml b/backends/factos_sqlight/gleam.toml new file mode 100644 index 0000000..cea193f --- /dev/null +++ b/backends/factos_sqlight/gleam.toml @@ -0,0 +1,10 @@ +name = "factos_sqlight" +version = "1.0.0" + +[dependencies] +factos = { path = "../.." } +gleam_stdlib = ">= 1.0.0 and < 2.0.0" +sqlight = ">= 1.1.0 and < 2.0.0" + +[dev_dependencies] +gleeunit = ">= 1.0.0 and < 2.0.0" diff --git a/backends/factos_sqlight/manifest.toml b/backends/factos_sqlight/manifest.toml new file mode 100644 index 0000000..a41f74c --- /dev/null +++ b/backends/factos_sqlight/manifest.toml @@ -0,0 +1,21 @@ +# Do not manually edit this file, it is managed by Gleam. +# +# This file locks the dependency versions used, to make your build +# deterministic and to prevent unexpected versions from being included +# in your application. +# +# You should check this file into your source control repository. + +packages = [ + { name = "esqlite", version = "0.9.0", build_tools = ["rebar3"], requirements = [], otp_app = "esqlite", source = "hex", outer_checksum = "CCF72258A4EE152EC7AD92AA9A03552EB6CA1B06B65C93AD5B6E55C302E05855" }, + { name = "factos", version = "1.0.0", build_tools = ["gleam"], requirements = ["gleam_stdlib"], source = "local", path = "../.." }, + { name = "gleam_stdlib", version = "1.0.3", build_tools = ["gleam"], requirements = [], otp_app = "gleam_stdlib", source = "hex", outer_checksum = "1F543AFBA5D33DA493E6087F4E4C4F20D899411343512686C98A8ABB2963CF22" }, + { name = "gleeunit", version = "1.11.0", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "gleeunit", source = "hex", outer_checksum = "EC31ABA74256AEA531EDF8169931D775BBB384FED0A8A1BDC4DD9354E3E21826" }, + { name = "sqlight", version = "1.1.0", build_tools = ["gleam"], requirements = ["esqlite", "gleam_stdlib"], otp_app = "sqlight", source = "hex", outer_checksum = "ECA1A4B45C35EB9EFCEEB7FAAC7BF5D8B2C777A7C1FC8A9C12CB67D54CED42E7" }, +] + +[requirements] +factos = { path = "../.." } +gleam_stdlib = { version = ">= 1.0.0 and < 2.0.0" } +gleeunit = { version = ">= 1.0.0 and < 2.0.0" } +sqlight = { version = ">= 1.1.0 and < 2.0.0" } diff --git a/backends/factos_sqlight/src/factos/sqlight.gleam b/backends/factos_sqlight/src/factos/sqlight.gleam new file mode 100644 index 0000000..722677d --- /dev/null +++ b/backends/factos_sqlight/src/factos/sqlight.gleam @@ -0,0 +1,536 @@ +//// SQLite backend for Factos using the `sqlight` package. + +import factos +import gleam/dynamic/decode +import gleam/list +import gleam/result +import gleam/string +import sqlight + +pub type Proposed(event) { + Proposed( + id: String, + event: event, + type_: factos.EventType, + tags: List(factos.Tag), + data: BitArray, + ) +} + +pub type StoredEvent { + StoredEvent( + position: Int, + id: String, + stream: String, + revision: Int, + type_: factos.EventType, + tags: List(factos.Tag), + data: BitArray, + ) +} + +pub type EventCodec(event, decode_error) { + EventCodec( + encode: fn(event) -> Proposed(event), + decode: fn(StoredEvent) -> Result(factos.Decoded(event), decode_error), + ) +} + +pub type Append { + Append(current_revision: Int, position: factos.SequencePosition) +} + +pub type Error(domain_error, decode_error) { + DomainError(domain_error) + DecodeError(decode_error) + StoreError(sqlight.Error) + AppendConditionFailed(factos.AppendCondition) +} + +pub fn migrate(connection: sqlight.Connection) -> Result(Nil, sqlight.Error) { + sqlight.exec( + " + create table if not exists factos_events ( + position integer primary key autoincrement, + id text not null, + stream text not null, + revision integer not null, + type text not null, + tags text not null, + data blob not null, + unique(stream, revision) + ); + create index if not exists factos_events_stream_revision + on factos_events(stream, revision); + create index if not exists factos_events_position + on factos_events(position); + ", + on: connection, + ) +} + +pub fn read_context( + connection: sqlight.Connection, + query query: factos.Query, + decider decider: factos.Decider(command, state, event, domain_error), + codec codec: EventCodec(event, decode_error), +) -> Result(factos.Context(event, state), Error(domain_error, decode_error)) { + let factos.Decider(initial, _, evolve) = decider + + use events <- result.try(read_matching_events(connection, query, codec)) + let position = highest_recorded_position(events) + + Ok(factos.Context( + query:, + state: factos.evolve_recorded( + initial: initial, + events: events, + evolve: evolve, + ), + events: events, + position: position, + append_condition: factos.FailIfEventsMatch(query, position), + )) +} + +pub fn dispatch_context( + connection: sqlight.Connection, + stream stream_name: String, + query query: factos.Query, + decider decider: factos.Decider(command, state, event, domain_error), + codec codec: EventCodec(event, decode_error), + command command: command, +) -> Result(Append, Error(domain_error, decode_error)) { + use _ <- result.try( + sqlight.exec("begin immediate", on: connection) + |> result.map_error(StoreError), + ) + + let result = { + use context <- result.try(read_context( + connection, + query: query, + decider: decider, + codec: codec, + )) + use pair <- result.try( + factos.decide_context(context, command, decider) + |> result.map_error(DomainError), + ) + let #(context, events) = pair + + append_with_condition( + connection, + stream_name, + events, + codec, + context.append_condition, + ) + } + + finish_transaction(connection, result) +} + +pub fn load_stream( + connection: sqlight.Connection, + stream stream_name: String, + decider decider: factos.Decider(command, state, event, domain_error), + codec codec: EventCodec(event, decode_error), +) -> Result( + factos.LoadedStream(event, state), + Error(domain_error, decode_error), +) { + let factos.Decider(initial, _, evolve) = decider + use events <- result.try(read_stream_events(connection, stream_name, codec)) + + Ok(factos.LoadedStream( + stream: stream_name, + state: factos.evolve_recorded( + initial: initial, + events: events, + evolve: evolve, + ), + events: events, + revision: stream_revision(events), + )) +} + +pub fn dispatch_stream( + connection: sqlight.Connection, + stream stream_name: String, + decider decider: factos.Decider(command, state, event, domain_error), + codec codec: EventCodec(event, decode_error), + command command: command, +) -> Result(Append, Error(domain_error, decode_error)) { + use _ <- result.try( + sqlight.exec("begin immediate", on: connection) + |> result.map_error(StoreError), + ) + + let result = { + use loaded <- result.try(load_stream( + connection, + stream: stream_name, + decider: decider, + codec: codec, + )) + let factos.Decider(_, decide, _) = decider + use events <- result.try( + decide(loaded.state, command) + |> result.map_error(DomainError), + ) + + append_stream_events( + connection, + stream_name, + events, + codec, + loaded.revision, + ) + } + + finish_transaction(connection, result) +} + +fn finish_transaction( + connection: sqlight.Connection, + result: Result(Append, Error(domain_error, decode_error)), +) -> Result(Append, Error(domain_error, decode_error)) { + case result { + Ok(append) -> + case sqlight.exec("commit", on: connection) { + Ok(Nil) -> Ok(append) + Error(error) -> Error(StoreError(error)) + } + Error(error) -> { + let _ = sqlight.exec("rollback", on: connection) + Error(error) + } + } +} + +fn append_with_condition( + connection: sqlight.Connection, + stream_name: String, + events: List(event), + codec: EventCodec(event, decode_error), + condition: factos.AppendCondition, +) -> Result(Append, Error(domain_error, decode_error)) { + case condition { + factos.NoAppendCondition -> + append_current_stream(connection, stream_name, events, codec) + factos.FailIfEventsMatch(query, after) -> + case has_matching_events_after(connection, query, after) { + Error(error) -> Error(StoreError(error)) + Ok(True) -> Error(AppendConditionFailed(condition)) + Ok(False) -> + append_current_stream(connection, stream_name, events, codec) + } + } +} + +fn append_current_stream( + connection: sqlight.Connection, + stream_name: String, + events: List(event), + codec: EventCodec(event, decode_error), +) -> Result(Append, Error(domain_error, decode_error)) { + use revision <- result.try( + current_revision(connection, stream_name) + |> result.map_error(StoreError), + ) + append_stream_events( + connection, + stream_name, + events, + codec, + factos.CurrentRevision(revision), + ) +} + +fn append_stream_events( + connection: sqlight.Connection, + stream_name: String, + events: List(event), + codec: EventCodec(event, decode_error), + expected: factos.Revision, +) -> Result(Append, Error(domain_error, decode_error)) { + case events { + [] -> + Ok(Append( + current_revision: revision_to_int(expected), + position: factos.NoPosition, + )) + [_, ..] -> { + use current <- result.try( + current_revision(connection, stream_name) + |> result.map_error(StoreError), + ) + case expected_matches(expected, current) { + False -> Error(AppendConditionFailed(factos.NoAppendCondition)) + True -> + insert_events( + connection, + stream_name, + events, + codec, + current + 1, + factos.NoPosition, + ) + } + } + } +} + +fn insert_events( + connection: sqlight.Connection, + stream_name: String, + events: List(event), + codec: EventCodec(event, decode_error), + revision: Int, + position: factos.SequencePosition, +) -> Result(Append, Error(domain_error, decode_error)) { + case events { + [] -> Ok(Append(current_revision: revision - 1, position: position)) + [event, ..rest] -> { + let EventCodec(encode, _) = codec + let Proposed(id, _, type_, tags, data) = encode(event) + use positions <- result.try( + sqlight.query( + " + insert into factos_events (id, stream, revision, type, tags, data) + values (?, ?, ?, ?, ?, ?) + returning position + ", + on: connection, + with: [ + sqlight.text(id), + sqlight.text(stream_name), + sqlight.int(revision), + sqlight.text(factos.event_type_name(type_)), + sqlight.text(tags_to_text(tags)), + sqlight.blob(data), + ], + expecting: int_field_decoder(), + ) + |> result.map_error(StoreError), + ) + let position = case positions { + [position, ..] -> factos.SequencePosition(position) + [] -> position + } + insert_events( + connection, + stream_name, + rest, + codec, + revision + 1, + position, + ) + } + } +} + +fn read_matching_events( + connection: sqlight.Connection, + query: factos.Query, + codec: EventCodec(event, decode_error), +) -> Result(List(factos.Recorded(event)), Error(domain_error, decode_error)) { + use rows <- result.try( + sqlight.query( + "select position, id, stream, revision, type, tags, data from factos_events order by position", + on: connection, + with: [], + expecting: stored_event_decoder(), + ) + |> result.map_error(StoreError), + ) + decode_rows(rows, codec) + |> result.map(list.filter(_, factos.matches_query(_, query))) +} + +fn read_stream_events( + connection: sqlight.Connection, + stream_name: String, + codec: EventCodec(event, decode_error), +) -> Result(List(factos.Recorded(event)), Error(domain_error, decode_error)) { + use rows <- result.try( + sqlight.query( + "select position, id, stream, revision, type, tags, data from factos_events where stream = ? order by revision", + on: connection, + with: [sqlight.text(stream_name)], + expecting: stored_event_decoder(), + ) + |> result.map_error(StoreError), + ) + decode_rows(rows, codec) +} + +fn decode_rows( + rows: List(StoredEvent), + codec: EventCodec(event, decode_error), +) -> Result(List(factos.Recorded(event)), Error(domain_error, decode_error)) { + case rows { + [] -> Ok([]) + [row, ..rest] -> { + use recorded <- result.try(decode_row(row, codec)) + use rest <- result.try(decode_rows(rest, codec)) + Ok([recorded, ..rest]) + } + } +} + +fn decode_row( + row: StoredEvent, + codec: EventCodec(event, decode_error), +) -> Result(factos.Recorded(event), Error(domain_error, decode_error)) { + let EventCodec(_, decode_event) = codec + use decoded <- result.try(decode_event(row) |> result.map_error(DecodeError)) + let factos.Decoded(event, type_, tags) = decoded + let StoredEvent(position, id, stream, revision, _, _, _) = row + + Ok(factos.Recorded( + id: id, + stream: stream, + revision: revision, + position: factos.SequencePosition(position), + type_: type_, + tags: tags, + event: event, + )) +} + +fn stored_event_decoder() -> decode.Decoder(StoredEvent) { + use position <- decode.field(0, decode.int) + use id <- decode.field(1, decode.string) + use stream <- decode.field(2, decode.string) + use revision <- decode.field(3, decode.int) + use type_name <- decode.field(4, decode.string) + use tags <- decode.field(5, decode.string) + use data <- decode.field(6, decode.bit_array) + decode.success(StoredEvent( + position: position, + id: id, + stream: stream, + revision: revision, + type_: factos.event_type(type_name), + tags: tags_from_text(tags), + data: data, + )) +} + +fn current_revision( + connection: sqlight.Connection, + stream_name: String, +) -> Result(Int, sqlight.Error) { + sqlight.query( + "select coalesce(max(revision), -1) from factos_events where stream = ?", + on: connection, + with: [sqlight.text(stream_name)], + expecting: int_field_decoder(), + ) + |> result.map(fn(rows) { + case rows { + [revision, ..] -> revision + [] -> -1 + } + }) +} + +fn has_matching_events_after( + connection: sqlight.Connection, + query: factos.Query, + after: factos.SequencePosition, +) -> Result(Bool, sqlight.Error) { + let after_position = case after { + factos.NoPosition -> -1 + factos.SequencePosition(position) -> position + } + sqlight.query( + "select type, tags from factos_events where position > ?", + on: connection, + with: [sqlight.int(after_position)], + expecting: query_match_decoder(), + ) + |> result.map( + list.any(_, fn(pair) { + let #(type_, tags) = pair + matches_query_parts(type_, tags, query) + }), + ) +} + +fn query_match_decoder() -> decode.Decoder( + #(factos.EventType, List(factos.Tag)), +) { + use type_name <- decode.field(0, decode.string) + use tags <- decode.field(1, decode.string) + decode.success(#(factos.event_type(type_name), tags_from_text(tags))) +} + +fn int_field_decoder() -> decode.Decoder(Int) { + use value <- decode.field(0, decode.int) + decode.success(value) +} + +fn matches_query_parts( + type_: factos.EventType, + tags: List(factos.Tag), + query: factos.Query, +) -> Bool { + factos.matches_query( + factos.Recorded( + id: "", + stream: "", + revision: 0, + position: factos.NoPosition, + type_: type_, + tags: tags, + event: Nil, + ), + query, + ) +} + +fn stream_revision(events: List(factos.Recorded(event))) -> factos.Revision { + case list.reverse(events) { + [] -> factos.NoEvents + [event, ..] -> factos.CurrentRevision(event.revision) + } +} + +fn highest_recorded_position( + events: List(factos.Recorded(event)), +) -> factos.SequencePosition { + case list.reverse(events) { + [] -> factos.NoPosition + [event, ..] -> event.position + } +} + +fn expected_matches(expected: factos.Revision, current: Int) -> Bool { + case expected { + factos.NoEvents -> current == -1 + factos.CurrentRevision(revision) -> current == revision + } +} + +fn revision_to_int(revision: factos.Revision) -> Int { + case revision { + factos.NoEvents -> -1 + factos.CurrentRevision(revision) -> revision + } +} + +fn tags_to_text(tags: List(factos.Tag)) -> String { + tags + |> list.map(factos.tag_value) + |> string.join(with: "\n") +} + +fn tags_from_text(tags: String) -> List(factos.Tag) { + case string.is_empty(tags) { + True -> [] + False -> tags |> string.split(on: "\n") |> list.map(factos.tag) + } +} diff --git a/backends/factos_sqlight/test/factos_sqlight_test.gleam b/backends/factos_sqlight/test/factos_sqlight_test.gleam new file mode 100644 index 0000000..7198c09 --- /dev/null +++ b/backends/factos_sqlight/test/factos_sqlight_test.gleam @@ -0,0 +1,315 @@ +import factos +import factos/sqlight as factos_sqlight +import gleam/bit_array +import gleam/int +import gleam/list +import gleam/result +import gleeunit +import sqlight + +pub fn main() -> Nil { + gleeunit.main() +} + +type Command { + RegisterUser(username: String) +} + +type Event { + UserRegistered(username: String) +} + +type State { + Available + Taken +} + +type DomainError { + AlreadyTaken +} + +type DecodeError { + UnknownEvent + InvalidData +} + +type CounterCommand { + Increment +} + +type CounterEvent { + Incremented(value: Int) +} + +type CounterState { + CounterState(total: Int) +} + +pub fn dispatch_stream_persists_events_test() { + use connection <- sqlight.with_connection(":memory:") + let assert Ok(Nil) = factos_sqlight.migrate(connection) + + let assert Ok(factos_sqlight.Append(current_revision: 0, position: _)) = + factos_sqlight.dispatch_stream( + connection, + stream: "user-renata", + decider: decider(), + codec: codec(), + command: RegisterUser("renata"), + ) + + let assert Ok(loaded) = + factos_sqlight.load_stream( + connection, + stream: "user-renata", + decider: decider(), + codec: codec(), + ) + + assert loaded.state == Taken + assert loaded.revision == factos.CurrentRevision(0) +} + +pub fn dispatch_stream_handles_many_events_test() { + use connection <- sqlight.with_connection(":memory:") + let assert Ok(Nil) = factos_sqlight.migrate(connection) + + let assert Ok(factos_sqlight.Append( + current_revision: 249, + position: factos.SequencePosition(_), + )) = dispatch_counter_stream_many(connection, 250) + + let assert Ok(loaded) = + factos_sqlight.load_stream( + connection, + stream: "counter-load", + decider: counter_decider(), + codec: counter_codec(), + ) + + assert loaded.state == CounterState(250) + assert loaded.revision == factos.CurrentRevision(249) + assert list.length(loaded.events) == 250 +} + +pub fn dispatch_context_handles_many_streams_test() { + use connection <- sqlight.with_connection(":memory:") + let assert Ok(Nil) = factos_sqlight.migrate(connection) + + let query = + factos.query([ + factos.query_item(types: [factos.event_type("Incremented")], tags: [ + factos.tag("counter:load"), + ]), + ]) + + let assert Ok(factos_sqlight.Append( + current_revision: 0, + position: factos.SequencePosition(_), + )) = dispatch_counter_context_many(connection, query, 100) + + let assert Ok(context) = + factos_sqlight.read_context( + connection, + query: query, + decider: counter_decider(), + codec: counter_codec(), + ) + + assert context.state == CounterState(100) + assert list.length(context.events) == 100 + assert context.position != factos.NoPosition +} + +fn decider() -> factos.Decider(Command, State, Event, DomainError) { + factos.decider(initial: Available, decide: decide, evolve: evolve) +} + +fn decide(state: State, command: Command) -> Result(List(Event), DomainError) { + case state, command { + Available, RegisterUser(username) -> Ok([UserRegistered(username)]) + Taken, RegisterUser(_) -> Error(AlreadyTaken) + } +} + +fn evolve(_state: State, _event: Event) -> State { + Taken +} + +fn codec() -> factos_sqlight.EventCodec(Event, DecodeError) { + factos_sqlight.EventCodec(encode: encode, decode: decode_event) +} + +fn encode(event: Event) -> factos_sqlight.Proposed(Event) { + case event { + UserRegistered(username) -> + factos_sqlight.Proposed( + id: "event-" <> username, + event: event, + type_: factos.event_type("UserRegistered"), + tags: [factos.tag("username:" <> username)], + data: bit_array.from_string(username), + ) + } +} + +fn decode_event( + stored: factos_sqlight.StoredEvent, +) -> Result(factos.Decoded(Event), DecodeError) { + case factos.event_type_name(stored.type_) { + "UserRegistered" -> { + use username <- result.try( + bit_array.to_string(stored.data) + |> result.map_error(fn(_) { InvalidData }), + ) + Ok(factos.Decoded( + event: UserRegistered(username), + type_: stored.type_, + tags: stored.tags, + )) + } + _ -> Error(UnknownEvent) + } +} + +fn dispatch_counter_stream_many( + connection: sqlight.Connection, + remaining: Int, +) -> Result(factos_sqlight.Append, factos_sqlight.Error(Nil, DecodeError)) { + case remaining { + 0 -> + factos_sqlight.dispatch_stream( + connection, + stream: "counter-load", + decider: counter_decider(), + codec: counter_codec(), + command: Increment, + ) + _ -> { + let result = + factos_sqlight.dispatch_stream( + connection, + stream: "counter-load", + decider: counter_decider(), + codec: counter_codec(), + command: Increment, + ) + case remaining, result { + 1, _ -> result + _, Ok(_) -> dispatch_counter_stream_many(connection, remaining - 1) + _, Error(error) -> Error(error) + } + } + } +} + +fn dispatch_counter_context_many( + connection: sqlight.Connection, + query: factos.Query, + remaining: Int, +) -> Result(factos_sqlight.Append, factos_sqlight.Error(Nil, DecodeError)) { + case remaining { + 0 -> + factos_sqlight.dispatch_context( + connection, + stream: "counter-context-0", + query: query, + decider: counter_decider(), + codec: counter_codec(), + command: Increment, + ) + _ -> { + let stream_name = "counter-context-" <> int.to_string(remaining) + let result = + factos_sqlight.dispatch_context( + connection, + stream: stream_name, + query: query, + decider: counter_decider(), + codec: counter_codec(), + command: Increment, + ) + case remaining, result { + 1, _ -> result + _, Ok(_) -> + dispatch_counter_context_many(connection, query, remaining - 1) + _, Error(error) -> Error(error) + } + } + } +} + +fn counter_decider() -> factos.Decider( + CounterCommand, + CounterState, + CounterEvent, + Nil, +) { + factos.decider( + initial: CounterState(0), + decide: counter_decide, + evolve: counter_evolve, + ) +} + +fn counter_decide( + state: CounterState, + command: CounterCommand, +) -> Result(List(CounterEvent), Nil) { + let CounterState(total) = state + case command { + Increment -> Ok([Incremented(total + 1)]) + } +} + +fn counter_evolve(state: CounterState, event: CounterEvent) -> CounterState { + let CounterState(total) = state + case event { + Incremented(_) -> CounterState(total + 1) + } +} + +fn counter_codec() -> factos_sqlight.EventCodec(CounterEvent, DecodeError) { + factos_sqlight.EventCodec( + encode: encode_counter_event, + decode: decode_counter_event, + ) +} + +fn encode_counter_event( + event: CounterEvent, +) -> factos_sqlight.Proposed(CounterEvent) { + case event { + Incremented(value) -> + factos_sqlight.Proposed( + id: "counter-event-" <> int.to_string(value), + event: event, + type_: factos.event_type("Incremented"), + tags: [factos.tag("counter:load")], + data: bit_array.from_string(int.to_string(value)), + ) + } +} + +fn decode_counter_event( + stored: factos_sqlight.StoredEvent, +) -> Result(factos.Decoded(CounterEvent), DecodeError) { + case factos.event_type_name(stored.type_) { + "Incremented" -> { + use text <- result.try( + bit_array.to_string(stored.data) + |> result.map_error(fn(_) { InvalidData }), + ) + use value <- result.try( + int.parse(text) + |> result.map_error(fn(_) { InvalidData }), + ) + Ok(factos.Decoded( + event: Incremented(value), + type_: stored.type_, + tags: stored.tags, + )) + } + _ -> Error(UnknownEvent) + } +} diff --git a/gleam.toml b/gleam.toml index 6e91a47..4cec0a8 100644 --- a/gleam.toml +++ b/gleam.toml @@ -14,10 +14,6 @@ version = "1.0.0" [dependencies] gleam_stdlib = ">= 1.0.0 and < 2.0.0" -gleam_json = ">= 3.1.0 and < 4.0.0" -kurrentdb = { path = "../kurrentdb" } -kurrentdb_erlang = { path = "../kurrentdb/backends/kurrentdb_erlang" } -youid = ">= 1.6.0 and < 2.0.0" [dev_dependencies] gleeunit = ">= 1.0.0 and < 2.0.0" diff --git a/manifest.toml b/manifest.toml index e92a8cc..1fc09fd 100644 --- a/manifest.toml +++ b/manifest.toml @@ -7,33 +7,10 @@ # You should check this file into your source control repository. packages = [ - { name = "certifi", version = "2.17.0", build_tools = ["rebar3"], requirements = [], otp_app = "certifi", source = "hex", outer_checksum = "8122798A17F0293C80DAADA25D0F81C7F4D708C73FEF782C7C9B1950E26E4D21" }, - { name = "gleam_crypto", version = "1.6.0", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "gleam_crypto", source = "hex", outer_checksum = "2DE9E4EF53CF6FEE049D4F765731F7178F7A11AEFAE00EEE63BF7536B354AD3F" }, - { name = "gleam_erlang", version = "1.3.0", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "gleam_erlang", source = "hex", outer_checksum = "1124AD3AA21143E5AF0FC5CF3D9529F6DB8CA03E43A55711B60B6B7B3874375C" }, - { name = "gleam_hackney", version = "1.4.0", build_tools = ["gleam"], requirements = ["gleam_http", "gleam_stdlib", "hackney"], source = "local", path = "../gleam_hackney" }, - { name = "gleam_http", version = "4.3.0", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "gleam_http", source = "hex", outer_checksum = "82EA6A717C842456188C190AFB372665EA56CE13D8559BF3B1DD9E40F619EE0C" }, - { name = "gleam_json", version = "3.1.0", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "gleam_json", source = "hex", outer_checksum = "44FDAA8847BE8FC48CA7A1C089706BD54BADCC4C45B237A992EDDF9F2CDB2836" }, - { name = "gleam_otp", version = "1.2.0", build_tools = ["gleam"], requirements = ["gleam_erlang", "gleam_stdlib"], otp_app = "gleam_otp", source = "hex", outer_checksum = "BA6A294E295E428EC1562DC1C11EA7530DCB981E8359134BEABC8493B7B2258E" }, { name = "gleam_stdlib", version = "1.0.3", build_tools = ["gleam"], requirements = [], otp_app = "gleam_stdlib", source = "hex", outer_checksum = "1F543AFBA5D33DA493E6087F4E4C4F20D899411343512686C98A8ABB2963CF22" }, - { name = "gleam_time", version = "1.8.0", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "gleam_time", source = "hex", outer_checksum = "533D8723774D61AD4998324F5DD1DABDCDBFABAFB9E87CB5D03C6955448FC97D" }, { name = "gleeunit", version = "1.11.0", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "gleeunit", source = "hex", outer_checksum = "EC31ABA74256AEA531EDF8169931D775BBB384FED0A8A1BDC4DD9354E3E21826" }, - { name = "h2", version = "0.10.2", build_tools = ["rebar3"], requirements = [], otp_app = "h2", source = "hex", outer_checksum = "497A899F338B42E6A0B292524E635B0CE6F9379FA39395C8E38D06351CD9B9CF" }, - { name = "hackney", version = "4.4.5", build_tools = ["rebar3"], requirements = ["certifi", "h2", "idna", "mimerl", "parse_trans", "quic", "ssl_verify_fun", "webtransport"], otp_app = "hackney", source = "hex", outer_checksum = "6D72BEF4E135E94C522C271E11FBB6933EFB0006EF235A3933807D0BE73B71EC" }, - { name = "idna", version = "7.1.0", build_tools = ["rebar3"], requirements = [], otp_app = "idna", source = "hex", outer_checksum = "6AE959A025BF36DF61A8CAB8508D9654891B5426A84C44D82DEAFFD6DDF8C71F" }, - { name = "kurrentdb", version = "1.0.0", build_tools = ["gleam"], requirements = ["gleam_http", "gleam_json", "gleam_stdlib", "youid"], source = "local", path = "../kurrentdb" }, - { name = "kurrentdb_erlang", version = "1.0.0", build_tools = ["gleam"], requirements = ["gleam_erlang", "gleam_hackney", "gleam_http", "gleam_otp", "gleam_stdlib", "kurrentdb", "youid"], source = "local", path = "../kurrentdb/backends/kurrentdb_erlang" }, - { name = "mimerl", version = "1.5.0", build_tools = ["rebar3"], requirements = [], otp_app = "mimerl", source = "hex", outer_checksum = "DB648CE065BAE14EA84CA8B5DD123F42F49417CEF693541110BF6F9E9BE9ECC4" }, - { name = "parse_trans", version = "3.4.2", build_tools = ["rebar3"], requirements = [], otp_app = "parse_trans", source = "hex", outer_checksum = "4C25347DE3B7C35732D32E69AB43D1CEEE0BEAE3F3B3ADE1B59CBD3DD224D9CA" }, - { name = "quic", version = "1.6.5", build_tools = ["rebar3"], requirements = [], otp_app = "quic", source = "hex", outer_checksum = "DE1A88972C33201A50D1A17C8C4A14528BD1D8F25EF705897D680AE312D0AA78" }, - { name = "ssl_verify_fun", version = "1.1.7", build_tools = ["mix", "rebar3", "make"], requirements = [], otp_app = "ssl_verify_fun", source = "hex", outer_checksum = "FE4C190E8F37401D30167C8C405EDA19469F34577987C76DDE613E838BBC67F8" }, - { name = "webtransport", version = "0.4.1", build_tools = ["rebar3"], requirements = ["h2", "quic"], otp_app = "webtransport", source = "hex", outer_checksum = "006E4E52A8F03B69201D4637C85B424A4DDCCC1909F32C5742FED8D495A75174" }, - { name = "youid", version = "1.6.0", build_tools = ["gleam"], requirements = ["gleam_crypto", "gleam_stdlib", "gleam_time"], otp_app = "youid", source = "hex", outer_checksum = "7A3ABA44B1B38BC2BDCB5474C5317AA372BE58DFBC649815EE08B03526DDA18D" }, ] [requirements] -gleam_json = { version = ">= 3.1.0 and < 4.0.0" } gleam_stdlib = { version = ">= 1.0.0 and < 2.0.0" } gleeunit = { version = ">= 1.0.0 and < 2.0.0" } -kurrentdb = { path = "../kurrentdb" } -kurrentdb_erlang = { path = "../kurrentdb/backends/kurrentdb_erlang" } -youid = { version = ">= 1.6.0 and < 2.0.0" } diff --git a/src/factos.gleam b/src/factos.gleam index 837ef8f..3171453 100644 --- a/src/factos.gleam +++ b/src/factos.gleam @@ -1,23 +1,13 @@ -//// Context-first event-sourcing helpers for Gleam and KurrentDB. +//// Store-independent event-sourcing domain primitives. //// -//// Event Sourcing does not require aggregates. A command capability should read -//// the facts relevant to its decision, fold a temporary decision model, decide -//// which new facts to record, and then record those facts only when the relevant -//// context is still stable. This module models that API directly. -//// -//// KurrentDB's regular append API can guarantee stream revision consistency. It -//// cannot, through the operations used here, atomically guarantee a DCB-style -//// "fail if events matching this query were appended after position X" check. -//// The API keeps those concepts separate instead of pretending they are the same. +//// Factos keeps the domain model in the application. The core module models +//// facts, command contexts, pure decision components, and pure views. Concrete +//// storage concerns live in backend packages such as `factos_sqlight` and +//// `factos_kurrentdb_erlang`. import gleam/list import gleam/option.{type Option, None, Some} import gleam/result -import kurrentdb/operation/append_to_stream -import kurrentdb/operation/read_all -import kurrentdb/operation/read_stream -import kurrentdb_erlang -import youid/uuid.{type Uuid} pub type EventType { EventType(String) @@ -38,7 +28,7 @@ pub type QueryItem { pub type SequencePosition { NoPosition - SequencePosition(commit_position: Int, prepare_position: Int) + SequencePosition(Int) } pub type AppendCondition { @@ -67,32 +57,12 @@ pub type Decoded(event) { Decoded(event: event, type_: EventType, tags: List(Tag)) } -pub type Proposed(event) { - Proposed( - event: event, - type_: EventType, - tags: List(Tag), - message: append_to_stream.Event, - ) -} - -pub type EventCodec(event, decode_error) { - EventCodec( - encode: fn(event) -> Proposed(event), - decode: fn(read_stream.RecordedEvent) -> - Result(Decoded(event), decode_error), - ) -} - pub type Recorded(event) { Recorded( - id: Uuid, + id: String, stream: String, revision: Int, position: SequencePosition, - metadata: List(#(String, String)), - custom_metadata: BitArray, - data: BitArray, type_: EventType, tags: List(Tag), event: event, @@ -118,23 +88,24 @@ pub type LoadedStream(event, state) { ) } -pub type Error(domain_error, decode_error) { - DomainError(domain_error) - DecodeError(decode_error) - ReadError(kurrentdb_erlang.Error(read_stream.ResponseError)) - AppendError(kurrentdb_erlang.Error(append_to_stream.ResponseError)) - ReadTimedOut - UnsupportedAppendCondition(AppendCondition) -} - pub fn event_type(name: String) -> EventType { EventType(name) } +pub fn event_type_name(event_type: EventType) -> String { + let EventType(name) = event_type + name +} + pub fn tag(value: String) -> Tag { Tag(value) } +pub fn tag_value(tag: Tag) -> String { + let Tag(value) = tag + value +} + pub fn query(items: List(QueryItem)) -> Query { case items { [] -> AllEvents @@ -188,6 +159,15 @@ pub fn compute_state( Ok(fold_events(state, events, evolve)) } +pub fn evolve_recorded( + initial initial: state, + events events: List(Recorded(event)), + evolve evolve: fn(state, event) -> state, +) -> state { + use state, recorded <- list.fold(events, initial) + evolve(state, recorded.event) +} + pub fn project( view view: View(state, event), events events: List(event), @@ -217,213 +197,16 @@ pub fn merge_views( #(first_evolve(first_state, event), second_evolve(second_state, event)) } -pub fn read_context( - connection: kurrentdb_erlang.Connection, - query query: Query, - decider decider: Decider(command, state, event, domain_error), - codec codec: EventCodec(event, decode_error), - timeout timeout: Int, -) -> Result(Context(event, state), Error(domain_error, decode_error)) { - let Decider(initial, _, evolve) = decider - - read_context_events(connection, query, initial, evolve, codec, timeout) -} - pub fn decide_context( context: Context(event, state), command: command, decider: Decider(command, state, event, domain_error), -) -> Result( - #(Context(event, state), List(event)), - Error(domain_error, decode_error), -) { +) -> Result(#(Context(event, state), List(event)), domain_error) { let Decider(_, decide, _) = decider - use events <- result.try( - decide(context.state, command) - |> result.map_error(DomainError), - ) - + use events <- result.try(decide(context.state, command)) Ok(#(context, events)) } -pub fn dispatch_context( - connection: kurrentdb_erlang.Connection, - stream stream_name: String, - query query: Query, - decider decider: Decider(command, state, event, domain_error), - codec codec: EventCodec(event, decode_error), - command command: command, - timeout timeout: Int, -) -> Result(append_to_stream.Append, Error(domain_error, decode_error)) { - use context <- result.try(read_context( - connection, - query: query, - decider: decider, - codec: codec, - timeout: timeout, - )) - use pair <- result.try(decide_context(context, command, decider)) - let #(context, events) = pair - - append_with_condition( - connection, - stream_name, - events, - codec, - context.append_condition, - timeout, - ) -} - -fn append_with_condition( - connection: kurrentdb_erlang.Connection, - stream_name: String, - events: List(event), - codec: EventCodec(event, decode_error), - condition: AppendCondition, - timeout: Int, -) -> Result(append_to_stream.Append, Error(domain_error, decode_error)) { - case condition { - NoAppendCondition -> - append_to_stream_with_config( - connection, - stream_name, - events, - codec, - append_to_stream.configure() |> append_to_stream.any, - timeout, - ) - FailIfEventsMatch(_, _) -> Error(UnsupportedAppendCondition(condition)) - } -} - -fn read_context_events( - connection: kurrentdb_erlang.Connection, - query: Query, - initial: state, - evolve: fn(state, event) -> state, - codec: EventCodec(event, decode_error), - timeout: Int, -) -> Result(Context(event, state), Error(domain_error, decode_error)) { - let stream = - kurrentdb_erlang.read_all( - connection, - config: read_all.configure() - |> read_all.filter(query_to_read_all_filter(query)), - ) - - receive_context( - stream, - query, - initial, - evolve, - codec, - [], - NoPosition, - timeout, - ) -} - -fn load_stream_events( - connection: kurrentdb_erlang.Connection, - stream_name: String, - initial: state, - evolve: fn(state, event) -> state, - codec: EventCodec(event, decode_error), - timeout: Int, -) -> Result(LoadedStream(event, state), Error(domain_error, decode_error)) { - let stream = - kurrentdb_erlang.read_stream( - connection, - stream: stream_name, - config: read_stream.configure(), - ) - - receive_stream( - stream, - stream_name, - initial, - evolve, - codec, - [], - NoEvents, - timeout, - ) -} - -pub fn load_stream( - connection: kurrentdb_erlang.Connection, - stream stream_name: String, - decider decider: Decider(command, state, event, domain_error), - codec codec: EventCodec(event, decode_error), - timeout timeout: Int, -) -> Result(LoadedStream(event, state), Error(domain_error, decode_error)) { - let Decider(initial, _, evolve) = decider - - load_stream_events(connection, stream_name, initial, evolve, codec, timeout) -} - -pub fn dispatch_stream( - connection: kurrentdb_erlang.Connection, - stream stream_name: String, - decider decider: Decider(command, state, event, domain_error), - codec codec: EventCodec(event, decode_error), - command command: command, - timeout timeout: Int, -) -> Result(append_to_stream.Append, Error(domain_error, decode_error)) { - let Decider(initial, decide, evolve) = decider - - use loaded <- result.try(load_stream_events( - connection, - stream_name, - initial, - evolve, - codec, - timeout, - )) - - use events <- result.try( - decide(loaded.state, command) - |> result.map_error(DomainError), - ) - - append_stream_events( - connection, - stream_name, - events, - codec, - loaded.revision, - timeout, - ) -} - -fn append_stream_events( - connection: kurrentdb_erlang.Connection, - stream_name: String, - events: List(event), - codec: EventCodec(event, decode_error), - expected: Revision, - timeout: Int, -) -> Result(append_to_stream.Append, Error(domain_error, decode_error)) { - append_to_stream_with_config( - connection, - stream_name, - events, - codec, - append_config(expected), - timeout, - ) -} - -fn fold_events( - initial: state, - events: List(event), - evolve: fn(state, event) -> state, -) -> state { - use state, event <- list.fold(events, initial) - evolve(state, event) -} - pub fn matches_query(recorded: Recorded(event), query: Query) -> Bool { case query { AllEvents -> True @@ -438,336 +221,21 @@ pub fn highest_position( case left, right { NoPosition, position -> position position, NoPosition -> position - SequencePosition(left_commit, left_prepare), - SequencePosition(right_commit, right_prepare) - -> - case - left_commit > right_commit - || { left_commit == right_commit && left_prepare >= right_prepare } - { - True -> left - False -> right + SequencePosition(left), SequencePosition(right) -> + case left >= right { + True -> SequencePosition(left) + False -> SequencePosition(right) } } } -fn append_to_stream_with_config( - connection: kurrentdb_erlang.Connection, - stream_name: String, +fn fold_events( + initial: state, events: List(event), - codec: EventCodec(event, decode_error), - config: append_to_stream.Configuration, - timeout: Int, -) -> Result(append_to_stream.Append, Error(domain_error, decode_error)) { - case events { - [] -> - Ok(append_to_stream.Append( - current_revision: -1, - position: append_to_stream.NoPositionReturned, - )) - [_, ..] -> { - let EventCodec(encode, _) = codec - let task = - kurrentdb_erlang.append_to_stream( - connection, - stream: stream_name, - events: list.map(events, fn(event) { - let Proposed(message: message, ..) = encode(event) - message - }), - config: config, - ) - - kurrentdb_erlang.await(task, within: timeout) - |> result.map_error(AppendError) - } - } -} - -fn receive_context( - stream: kurrentdb_erlang.Stream, - query: Query, - state: state, - evolve: fn(state, event) -> state, - codec: EventCodec(event, decode_error), - events: List(Recorded(event)), - position: SequencePosition, - timeout: Int, -) -> Result(Context(event, state), Error(domain_error, decode_error)) { - case kurrentdb_erlang.receive(stream, within: timeout) { - Error(Nil) -> { - kurrentdb_erlang.close(stream) - Error(ReadTimedOut) - } - Ok(kurrentdb_erlang.StreamFinished) -> { - kurrentdb_erlang.close(stream) - Ok(Context( - query:, - state:, - events: list.reverse(events), - position:, - append_condition: FailIfEventsMatch(query, position), - )) - } - Ok(kurrentdb_erlang.StreamFailed(error)) -> { - kurrentdb_erlang.close(stream) - Error(ReadError(error)) - } - Ok(kurrentdb_erlang.ReadEvent(event)) -> - receive_context_event( - stream, - query, - state, - evolve, - codec, - events, - position, - timeout, - event, - ) - Ok(kurrentdb_erlang.ReadMessage(read_stream.ReadEvent(event))) -> - receive_context_event( - stream, - query, - state, - evolve, - codec, - events, - position, - timeout, - event, - ) - Ok(kurrentdb_erlang.ReadMessage(read_stream.LastAllStreamPosition(read_stream.Position( - commit_position: commit_position, - prepare_position: prepare_position, - )))) -> - receive_context( - stream, - query, - state, - evolve, - codec, - events, - highest_position( - position, - SequencePosition(commit_position, prepare_position), - ), - timeout, - ) - Ok(kurrentdb_erlang.ReadMessage(_)) -> - receive_context( - stream, - query, - state, - evolve, - codec, - events, - position, - timeout, - ) - } -} - -fn receive_context_event( - stream: kurrentdb_erlang.Stream, - query: Query, - state: state, - evolve: fn(state, event) -> state, - codec: EventCodec(event, decode_error), - events: List(Recorded(event)), - position: SequencePosition, - timeout: Int, - read_event: read_stream.ReadEvent, -) -> Result(Context(event, state), Error(domain_error, decode_error)) { - use recorded <- result.try(decode_recorded(read_event, codec)) - let next_position = highest_position(position, recorded.position) - - case matches_query(recorded, query) { - True -> - receive_context( - stream, - query, - evolve(state, recorded.event), - evolve, - codec, - [recorded, ..events], - next_position, - timeout, - ) - False -> - receive_context( - stream, - query, - state, - evolve, - codec, - events, - next_position, - timeout, - ) - } -} - -fn receive_stream( - stream: kurrentdb_erlang.Stream, - stream_name: String, - state: state, evolve: fn(state, event) -> state, - codec: EventCodec(event, decode_error), - events: List(Recorded(event)), - revision: Revision, - timeout: Int, -) -> Result(LoadedStream(event, state), Error(domain_error, decode_error)) { - case kurrentdb_erlang.receive(stream, within: timeout) { - Error(Nil) -> { - kurrentdb_erlang.close(stream) - Error(ReadTimedOut) - } - Ok(kurrentdb_erlang.StreamFinished) -> { - kurrentdb_erlang.close(stream) - Ok(LoadedStream( - stream: stream_name, - state:, - events: list.reverse(events), - revision:, - )) - } - Ok(kurrentdb_erlang.StreamFailed(error)) -> { - kurrentdb_erlang.close(stream) - case error { - kurrentdb_erlang.OperationError(read_stream.ReadStreamNotFound(_)) -> - Ok(LoadedStream( - stream: stream_name, - state:, - events: [], - revision: NoEvents, - )) - _ -> Error(ReadError(error)) - } - } - Ok(kurrentdb_erlang.ReadEvent(event)) -> - receive_stream_event( - stream, - stream_name, - state, - evolve, - codec, - events, - timeout, - event, - ) - Ok(kurrentdb_erlang.ReadMessage(read_stream.ReadEvent(event))) -> - receive_stream_event( - stream, - stream_name, - state, - evolve, - codec, - events, - timeout, - event, - ) - Ok(kurrentdb_erlang.ReadMessage(_)) -> - receive_stream( - stream, - stream_name, - state, - evolve, - codec, - events, - revision, - timeout, - ) - } -} - -fn receive_stream_event( - stream: kurrentdb_erlang.Stream, - stream_name: String, - state: state, - evolve: fn(state, event) -> state, - codec: EventCodec(event, decode_error), - events: List(Recorded(event)), - timeout: Int, - read_event: read_stream.ReadEvent, -) -> Result(LoadedStream(event, state), Error(domain_error, decode_error)) { - use recorded <- result.try(decode_recorded(read_event, codec)) - - receive_stream( - stream, - stream_name, - evolve(state, recorded.event), - evolve, - codec, - [recorded, ..events], - CurrentRevision(recorded.revision), - timeout, - ) -} - -fn decode_recorded( - read_event: read_stream.ReadEvent, - codec: EventCodec(event, decode_error), -) -> Result(Recorded(event), Error(domain_error, decode_error)) { - let EventCodec(_, decode) = codec - use decoded <- result.try( - decode(read_event.event) - |> result.map_error(DecodeError), - ) - - let Decoded(event, type_, tags) = decoded - Ok(Recorded( - id: read_event.event.id, - stream: read_event.event.stream, - revision: read_event.event.revision, - position: SequencePosition( - read_event.event.commit_position, - read_event.event.prepare_position, - ), - metadata: read_event.event.metadata, - custom_metadata: read_event.event.custom_metadata, - data: read_event.event.data, - type_: type_, - tags: tags, - event: event, - )) -} - -fn query_to_read_all_filter(query: Query) -> read_all.Filter { - let types = query_event_type_names(query) - - case types { - [] -> read_all.NoFilter - [_, ..] -> read_all.EventTypePrefix(types, window: read_all.FilterMax(1000)) - } -} - -fn query_event_type_names(query: Query) -> List(String) { - case query { - AllEvents -> [] - Query(items) -> collect_type_names(items, []) - } -} - -fn collect_type_names( - items: List(QueryItem), - names: List(String), -) -> List(String) { - case items { - [] -> list.reverse(names) - [QueryItem(types, _), ..rest] -> - collect_type_names(rest, prepend_type_names(types, names)) - } -} - -fn prepend_type_names( - types: List(EventType), - names: List(String), -) -> List(String) { - case types { - [] -> names - [EventType(name), ..rest] -> prepend_type_names(rest, [name, ..names]) - } +) -> state { + use state, event <- list.fold(events, initial) + evolve(state, event) } fn matches_item(recorded: Recorded(event), item: QueryItem) -> Bool { @@ -795,12 +263,3 @@ fn matches_tags(event_tags: List(Tag), required_tags: List(Tag)) -> Bool { } } } - -fn append_config(revision: Revision) -> append_to_stream.Configuration { - case revision { - NoEvents -> append_to_stream.configure() |> append_to_stream.no_stream - CurrentRevision(revision) -> - append_to_stream.configure() - |> append_to_stream.expected_revision(revision) - } -} diff --git a/test/factos_test.gleam b/test/factos_test.gleam index 0226962..1b99973 100644 --- a/test/factos_test.gleam +++ b/test/factos_test.gleam @@ -2,7 +2,6 @@ import factos import gleam/int import gleam/option.{None, Some} import gleeunit -import youid/uuid pub fn main() -> Nil { gleeunit.main() @@ -106,8 +105,8 @@ pub fn empty_query_matches_all_events_test() { } pub fn highest_position_keeps_later_position_test() { - let early = factos.SequencePosition(commit_position: 10, prepare_position: 10) - let later = factos.SequencePosition(commit_position: 11, prepare_position: 0) + let early = factos.SequencePosition(10) + let later = factos.SequencePosition(11) assert factos.highest_position(early, later) == later assert factos.highest_position(factos.NoPosition, early) == early @@ -234,16 +233,10 @@ fn recorded( } factos.Recorded( - id: uuid.v7(), + id: "event-" <> int.to_string(revision), stream: "user-1", revision: revision, - position: factos.SequencePosition( - commit_position: revision, - prepare_position: revision, - ), - metadata: [], - custom_metadata: <<>>, - data: <<>>, + position: factos.SequencePosition(revision), type_: type_, tags: tags, event: event,