Module: Pray::Sync

Defined in:
lib/pray/sync.rb

Class Method Summary collapse

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.message)
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.message)
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

Returns:

  • (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.expand_path(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