Module: Pray::Sync
- Defined in:
- lib/pray/sync.rb
Class Method Summary collapse
- .fetch_bytes(peer_source, artifact) ⇒ Object
- .fetch_json(url) ⇒ Object
- .load_known_peers(root) ⇒ Object
- .load_local_versions(root, package_name) ⇒ Object
- .load_sync_peers(root) ⇒ Object
- .normalize_peer(peer) ⇒ Object
- .registry_version_from_transport(data) ⇒ Object
- .same_identity?(left, right) ⇒ Boolean
- .sync_package(root, peer_source, metadata, package_versions) ⇒ Object
- .synchronize_registry(root, peer_sources) ⇒ Object
- .trim_slash(value) ⇒ Object
- .upsert_known_peer(peers, peer) ⇒ Object
- .write_artifact(root, artifact_path, bytes) ⇒ Object
- .write_known_peers(root, peers) ⇒ Object
- .write_local_index(root, package_versions) ⇒ Object
Class Method Details
.fetch_bytes(peer_source, artifact) ⇒ Object
207 208 209 210 211 212 213 |
# File 'lib/pray/sync.rb', line 207 def fetch_bytes(peer_source, artifact) if artifact.start_with?("http://", "https://") return Registry.http_get(artifact).b end Registry.http_get("#{trim_slash(peer_source)}/#{artifact.delete_prefix("/")}").b end |
.fetch_json(url) ⇒ Object
201 202 203 204 205 |
# File 'lib/pray/sync.rb', line 201 def fetch_json(url) JSON.parse(Registry.http_get(url)) rescue JSON::ParserError => error raise Error.parse("federation response", error.) end |
.load_known_peers(root) ⇒ Object
169 170 171 172 173 174 175 176 |
# File 'lib/pray/sync.rb', line 169 def load_known_peers(root) path = File.join(root, "v1", "peers.json") return [] unless File.file?(path) Array(JSON.parse(File.read(path))).map { |peer| normalize_peer(peer) } rescue JSON::ParserError [] end |
.load_local_versions(root, package_name) ⇒ Object
142 143 144 145 146 |
# File 'lib/pray/sync.rb', line 142 def load_local_versions(root, package_name) path = Publish.(root, package_name) = Publish.(path, package_name) .versions.to_h { |version| [version.version, version] } end |
.load_sync_peers(root) ⇒ Object
75 76 77 78 79 80 81 82 83 84 |
# File 'lib/pray/sync.rb', line 75 def load_sync_peers(root) path = File.join(root, "v1", "peers.json") unless File.file?(path) raise Error.unsupported("no federation peers configured") end Array(JSON.parse(File.read(path))).map { |peer| normalize_peer(peer) } rescue JSON::ParserError => error raise Error.parse("peer list", error.) end |
.normalize_peer(peer) ⇒ Object
193 194 195 196 197 198 199 |
# File 'lib/pray/sync.rb', line 193 def normalize_peer(peer) url = peer["url"].to_s raise Error.parse("peer list", "peer url is required") if url.empty? {"name" => peer["name"].to_s.empty? ? url : peer["name"], "url" => url, "public" => !!peer["public"]} end |
.registry_version_from_transport(data) ⇒ Object
119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 |
# File 'lib/pray/sync.rb', line 119 def registry_version_from_transport(data) signer = data.dig("publisher", "id") || data["signer"] fingerprint = data.dig("publisher", "key_fingerprint") || data["signer_fingerprint"] signature = data.dig("signature", "value") || data["signature"] RegistryPackageVersion.new( version: data["version"].to_s, artifact: data["artifact"].to_s, artifact_hash: data["artifact_hash"].to_s, tree_hash: data["tree_hash"], yanked: data.fetch("yanked", false), targets: Array(data["targets"]), exports: Array(data["exports"]), signer: signer, signer_fingerprint: fingerprint, published_at: data["published_at"], signature: signature.is_a?(Hash) ? signature["value"] : signature ) end |
.same_identity?(left, right) ⇒ Boolean
138 139 140 |
# File 'lib/pray/sync.rb', line 138 def same_identity?(left, right) left.artifact_hash == right.artifact_hash && left.tree_hash == right.tree_hash end |
.sync_package(root, peer_source, metadata, package_versions) ⇒ Object
86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 |
# File 'lib/pray/sync.rb', line 86 def sync_package(root, peer_source, , package_versions) name = ["name"] package_versions[name] ||= load_local_versions(root, name) Array(["versions"]).each do |version_data| version = registry_version_from_transport(version_data) existing = package_versions[name][version.version] if existing if same_identity?(existing, version) package_versions[name][version.version] = existing next end raise Error.integrity( "conflicting metadata for package #{name} version #{version.version}" ) end artifact_hash = version.artifact_hash unless artifact_hash raise Error.integrity("federation package #{name} #{version.version} is missing an artifact hash") end bytes = fetch_bytes(peer_source, version.artifact) computed = Hashing.sha256_prefixed(bytes) if computed != artifact_hash raise Error.integrity("artifact hash mismatch for #{name} #{version.version}") end write_artifact(root, version.artifact, bytes) package_versions[name][version.version] = version end end |
.synchronize_registry(root, peer_sources) ⇒ Object
11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 |
# File 'lib/pray/sync.rb', line 11 def synchronize_registry(root, peer_sources) root = File.(root) if peer_sources.empty? raise Error.unsupported("no federation peers configured") end peer_sources.each do |peer| if peer.start_with?("pray+ssh://", "ssh+pray://") raise Error.unsupported("pray_ssh sync peers are not implemented yet in pray-cli Ruby") end end known_peers = load_known_peers(root) pending = peer_sources.dup discovered = peer_sources.to_set peer_sources.each { |url| upsert_known_peer(known_peers, {"name" => url, "url" => url, "public" => false}) } package_versions = {} peer_count = 0 while (peer_source = pending.shift) next unless discovered.delete?(peer_source) peer_count += 1 discovery = fetch_json("#{trim_slash(peer_source)}/.well-known/pray-federation.json") unless discovery["spec"] == "pray-federation-v1" raise Error.resolution("peer #{peer_source} does not speak the pray federation protocol") end Array(discovery["peers"]).each do |peer| peer = normalize_peer(peer) next if peer["url"] == peer_source upsert_known_peer(known_peers, peer) next if pending.include?(peer["url"]) || discovered.include?(peer["url"]) pending << peer["url"] discovered << peer["url"] end index = fetch_json("#{trim_slash(peer_source)}/v1/sync/index") unless index["spec"] == "prayfile-distribution-1" raise Error.resolution( "peer #{peer_source} returned unsupported registry index spec: #{index["spec"]}" ) end Array(index["packages"]).each do |summary| name = summary["name"] = fetch_json("#{trim_slash(peer_source)}/v1/sync/package/#{name}") unless ["name"] == name raise Error.resolution( "peer #{peer_source} returned mismatched package metadata for #{name}" ) end sync_package(root, peer_source, , package_versions) end end write_known_peers(root, known_peers) write_local_index(root, package_versions) {peers: peer_count, packages: package_versions.length, known_peers: known_peers.length} end |
.trim_slash(value) ⇒ Object
215 216 217 |
# File 'lib/pray/sync.rb', line 215 def trim_slash(value) value.to_s.sub(%r{/+\z}, "") end |
.upsert_known_peer(peers, peer) ⇒ Object
184 185 186 187 188 189 190 191 |
# File 'lib/pray/sync.rb', line 184 def upsert_known_peer(peers, peer) existing = peers.find { |entry| entry["url"] == peer["url"] } if existing existing.merge!(peer) else peers << peer end end |
.write_artifact(root, artifact_path, bytes) ⇒ Object
162 163 164 165 166 167 |
# File 'lib/pray/sync.rb', line 162 def write_artifact(root, artifact_path, bytes) relative = PathSafety.sanitize_relative_path(artifact_path) path = File.join(root, relative) FileUtils.mkdir_p(File.dirname(path)) File.binwrite(path, bytes) end |
.write_known_peers(root, peers) ⇒ Object
178 179 180 181 182 |
# File 'lib/pray/sync.rb', line 178 def write_known_peers(root, peers) path = File.join(root, "v1", "peers.json") FileUtils.mkdir_p(File.dirname(path)) File.write(path, JSON.pretty_generate(peers)) end |
.write_local_index(root, package_versions) ⇒ Object
148 149 150 151 152 153 154 155 156 157 158 159 160 |
# File 'lib/pray/sync.rb', line 148 def write_local_index(root, package_versions) package_versions.each do |name, versions| = RegistryPackageMetadata.new(name: name, versions: versions.values) Publish.( Publish.(root, name), ) end index = Publish.load_registry_index(root) names = index.packages.to_set package_versions.each_key { |name| names << name } index.packages = names.sort Publish.write_registry_index(root, index) end |